From e48eeef49b3b13dba6cf508bd13116715b767086 Mon Sep 17 00:00:00 2001 From: butubb <1422726308@qq.com> Date: Sun, 16 Aug 2026 08:12:05 +0800 Subject: [PATCH] =?UTF-8?q?feat:=20=E4=BB=BB=E5=8A=A1=E6=8A=A2=E5=8D=A0+?= =?UTF-8?q?=E5=BD=92=E8=BF=98=EF=BC=88=E6=8A=A2=E5=8D=A0=E4=BB=BB=E5=8A=A1?= =?UTF-8?q?=E7=BB=93=E6=9D=9F=E5=90=8E=E8=87=AA=E5=8A=A8=E9=87=8D=E5=90=AF?= =?UTF-8?q?=E8=A2=AB=E6=8A=A2=E5=8D=A0=E4=BB=BB=E5=8A=A1=EF=BC=89+=20?= =?UTF-8?q?=E4=BF=AE=E5=A4=8D=E6=8A=A2=E5=8D=A0=E6=AD=BB=E9=94=81=EF=BC=88?= =?UTF-8?q?stop=5Fdevice=20=E7=A7=BB=E5=87=BA=E9=94=81=E5=9D=97=EF=BC=89+?= =?UTF-8?q?=20wait=20=E6=AD=A5=E9=AA=A4=E5=8F=AF=E4=B8=AD=E6=96=AD=20+=20?= =?UTF-8?q?=E5=BE=AA=E7=8E=AF=E5=9D=97=E6=96=B0=E5=A2=9E=E4=B8=80=E7=9B=B4?= =?UTF-8?q?=E5=BE=AA=E7=8E=AF=E6=A8=A1=E5=BC=8F=EF=BC=88=E9=85=8D=E5=90=88?= =?UTF-8?q?=E5=AE=9A=E6=97=B6=E5=90=AF=E5=8A=A8+=E5=81=9C=E6=AD=A2?= =?UTF-8?q?=E5=AE=9E=E7=8E=B0=E5=85=A8=E5=A4=A9=E5=80=99=E5=85=BB=E5=8F=B7?= =?UTF-8?q?=EF=BC=89?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- core/task_manager.py | 48 ++++++++++++++++++++++++++++++++++++++---- static/admin/editor.js | 10 ++++++--- static/admin/tasks.js | 4 ++++ tasks/generic/task.py | 20 ++++++++++++++++-- 4 files changed, 73 insertions(+), 9 deletions(-) diff --git a/core/task_manager.py b/core/task_manager.py index 86b9b01..c1c7014 100644 --- a/core/task_manager.py +++ b/core/task_manager.py @@ -151,7 +151,14 @@ class TaskJob: g = manager.groups.get(self.target.get("group_name")) serials = list(g.serials) if g else [] else: - # all:返回所有空闲设备(天然只含在线设备) + # all:默认返回所有空闲设备(天然只含在线设备); + # 抢占模式返回全部在线就绪设备(含被占用,执行时抢占) + if self.params.get("preempt"): + try: + return [d["serial"] for d in manager.stf.list_all_devices() + if d.get("present") and d.get("ready")] + except Exception: + return [] try: return [d["serial"] for d in manager.stf.list_free_devices()] except Exception: @@ -688,13 +695,40 @@ class TaskManager: if serial in self._stop_requested: _log.info(f"{serial} 用户已请求停止,取消重试 (job={job.name})") return - # 同一 serial 同时只能一个 worker + # 同一 serial 同时只能一个 worker;不同任务可配置抢占 + preempt = False + preempted_job = None # 被抢占的原任务 id(抢占结束后归还) with self._lock: if serial in self._running: - _log.warning(f"{serial} 已有任务在跑,跳过 (job={job.name})") + cur = self._running[serial] + if cur.get("job_id") == job.id: + # 同一任务重复触发:跳过(原行为) + _log.warning(f"{serial} 已有任务在跑,跳过 (job={job.name})") + return + if not job.params.get("preempt"): + # 不同任务且本任务未开启抢占:跳过 + _log.warning(f"{serial} 已有任务在跑,跳过 (job={job.name})") + return + preempt = True + preempted_job = cur.get("job_id") + if preempt: + # 抢占:锁外停止该设备上的其他任务(stop_device 内部拿同一把锁, + # 在锁内调用会死锁!),等其释放后接管 + _log.warning(f"{serial} 抢占:停止任务 {preempted_job} 后执行 {job.name}") + self.stop_device(serial) + deadline = time.time() + 30 + while time.time() < deadline: + with self._lock: + if serial not in self._running: + break + time.sleep(0.5) + else: + _log.warning(f"{serial} 抢占超时(旧任务 30s 未退出),跳过 (job={job.name})") return + with self._lock: self._running[serial] = {"job_id": job.id, "started_at": time.time(), - "attempt": attempt, "task_type": job.task_type} + "attempt": attempt, "task_type": job.task_type, + "preempted_job": preempted_job} _update_status(serial, task_job=job.name, attempt=attempt, max_attempts=max_attempts) @@ -761,6 +795,12 @@ class TaskManager: # 清除停止标志:整个重试循环结束(成功/失败/停止)后允许下次任务 with self._lock: self._stop_requested.discard(serial) + # 归还:本任务是抢占任务,结束后自动重新启动被抢占的原任务 + if preempted_job: + j = self.jobs.get(preempted_job) + if j and j.enabled: + _log.info(f"{serial} 抢占任务 {job.name} 结束,归还设备给任务 {j.name}") + threading.Thread(target=self._run_job, args=(j,), daemon=True).start() # ---- 运行控制 ---- def stop_device(self, serial): diff --git a/static/admin/editor.js b/static/admin/editor.js index c56a88d..eb46d33 100644 --- a/static/admin/editor.js +++ b/static/admin/editor.js @@ -341,13 +341,14 @@ var _stepEditor={ h+='
'; h+='