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):