feat: 任务抢占+归还(抢占任务结束后自动重启被抢占任务)+ 修复抢占死锁(stop_device 移出锁块)+ wait 步骤可中断 + 循环块新增一直循环模式(配合定时启动+停止实现全天候养号)
This commit is contained in:
+44
-4
@@ -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):
|
||||
|
||||
Reference in New Issue
Block a user