From a9720e8b815bf96f174dd6286505765bb5ef02e6 Mon Sep 17 00:00:00 2001 From: butubb <1422726308@qq.com> Date: Wed, 19 Aug 2026 15:46:18 +0800 Subject: [PATCH] =?UTF-8?q?fix:=20=E6=8A=A2=E5=8D=A0=E5=BD=92=E8=BF=98?= =?UTF-8?q?=E4=B8=A2=E5=A4=B1=E2=80=94=E2=80=94preempted=5Fjob=20=E5=9C=A8?= =?UTF-8?q?=E9=87=8D=E8=AF=95=E5=BE=AA=E7=8E=AF=E5=86=85=E8=A2=AB=E9=87=8D?= =?UTF-8?q?=E7=BD=AE=E4=B8=BA=20None=EF=BC=8C=E5=85=B3=E9=94=AE=E5=AD=97?= =?UTF-8?q?=E4=BB=BB=E5=8A=A1=E4=B8=80=E6=97=A6=E9=87=8D=E8=AF=95=EF=BC=88?= =?UTF-8?q?attempt=E2=89=A52=EF=BC=89=E8=A2=AB=E6=8A=A2=E5=8D=A0=E7=9A=84?= =?UTF-8?q?=E5=85=BB=E5=8F=B7=E4=BB=BB=E5=8A=A1=E6=B0=B8=E4=B8=8D=E6=81=A2?= =?UTF-8?q?=E5=A4=8D=EF=BC=9B=E6=8F=90=E5=8D=87=E5=88=B0=E5=BE=AA=E7=8E=AF?= =?UTF-8?q?=E5=A4=96+=E5=B7=B2=E6=9C=89=E5=80=BC=E4=B8=8D=E8=A6=86?= =?UTF-8?q?=E7=9B=96+=E5=BD=92=E8=BF=98=E6=97=B6=E8=A2=AB=E6=8A=A2?= =?UTF-8?q?=E5=8D=A0=E4=BB=BB=E5=8A=A1=E5=81=9C=E7=94=A8/=E5=88=A0?= =?UTF-8?q?=E9=99=A4=E4=B9=9F=E6=89=93=E6=97=A5=E5=BF=97?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- core/task_manager.py | 11 +++++++++-- 1 file changed, 9 insertions(+), 2 deletions(-) diff --git a/core/task_manager.py b/core/task_manager.py index c92ed65..bfccbf2 100644 --- a/core/task_manager.py +++ b/core/task_manager.py @@ -636,6 +636,9 @@ class TaskManager: # 错峰启动:在各自线程内等待,分摊批量启动的连接/占用压力 if start_delay > 0: time.sleep(start_delay) + # 被抢占的原任务 id(抢占结束后归还)——必须在循环外: + # 重试时若重置为 None,finally 归还逻辑会丢失信息,被抢占任务永不恢复 + preempted_job = None try: for attempt in range(1, max_attempts + 1): # 用户已请求停止 → 不再启动新 attempt @@ -645,7 +648,6 @@ class TaskManager: return # 同一 serial 同时只能一个 worker;不同任务可配置抢占 preempt = False - preempted_job = None # 被抢占的原任务 id(抢占结束后归还) with self._lock: if serial in self._running: cur = self._running[serial] @@ -658,7 +660,8 @@ class TaskManager: _log.warning(f"{serial} 已有任务在跑,跳过 (job={job.name})") return preempt = True - preempted_job = cur.get("job_id") + if not preempted_job: + preempted_job = cur.get("job_id") if preempt: # 抢占:锁外停止该设备上的其他任务(stop_device 内部拿同一把锁, # 在锁内调用会死锁!),等其释放后接管 @@ -749,6 +752,10 @@ class TaskManager: if j and j.enabled: _log.info(f"{serial} 抢占任务 {job.name} 结束,归还设备给任务 {j.name}") threading.Thread(target=self._run_job, args=(j,), daemon=True).start() + else: + # 打日志便于排查:被抢占任务已删除/停用时不会归还,但要知道原因 + _log.info(f"{serial} 抢占任务 {job.name} 结束," + f"被抢占任务 {preempted_job} {'已停用,不归还' if j else '已不存在,不归还'}") # ---- 运行控制 ---- def stop_device(self, serial):