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+='
需先用"点击元素"步骤定位到输入框,本步骤只负责输入文字
'; }else if(step.type==='loop'){ - const lm=p.loop_mode==='time'?'time':'rounds'; + const lm=p.loop_mode==='time'?'time':(p.loop_mode==='forever'?'forever':'rounds'); h+='
'+ + ''+ + '
'+ '
'+ '
'; - h+='
将下方子步骤(拖拽到本卡片内)循环执行指定轮次或时长;按时间时配合"运行时长上限"看预计结束时间
'; + h+='
将下方子步骤循环执行指定轮次/时长;一直循环无限跑直到任务被外部停止——配合调度"定时启动+停止"(如 9:00 启动 / 18:00 停止)实现全天候养号,或配合任务"运行时长上限"自动停止
'; }else if(step.type==='group'){ h+='
动作组:按顺序执行下方子步骤一次。把多个步骤打包成一个可复用的整体。
'; h+='
把本组步骤保存到左侧自定义动作库,方便复用
'; @@ -1057,6 +1058,9 @@ async function saveTask(jobId){ // 离线自动跳过(serial/group 模式):任务触发时跳过 STF 池中不在线的设备 const skipEl=document.getElementById('f-skip_offline'); params.skip_offline = skipEl?skipEl.checked:true; + // 抢占设备:触发时停止目标设备上的其他任务,接管执行 + const preemptEl=document.getElementById('f-preempt'); + params.preempt = preemptEl?preemptEl.checked:false; const targetMode=document.getElementById('f-target_mode').value; const target={mode:targetMode}; diff --git a/static/admin/tasks.js b/static/admin/tasks.js index f840f44..f652d07 100644 --- a/static/admin/tasks.js +++ b/static/admin/tasks.js @@ -108,6 +108,10 @@ function openTaskModal(jobId){ ''+ ''+ ''+ + '
'+ + ''+ + ''+ + '
'+ ''+ '
'+ '
'+ diff --git a/tasks/generic/task.py b/tasks/generic/task.py index 4b5c278..4d5b097 100644 --- a/tasks/generic/task.py +++ b/tasks/generic/task.py @@ -502,14 +502,30 @@ class GenericStepsWorker(BaseWorker): def _exec_wait(self, d, params, depth=0): min_s = float(params.get("min", 1.0)) max_s = float(params.get("max", 3.0)) - random_sleep(min_s, max_s) + # 分片 sleep:每 0.5s 检查停止/超时信号,抢占任务能及时接管(纯 sleep 无法中断) + total = random.uniform(min_s, max_s) + end = time.time() + total + while time.time() < end: + if self.stopped() or self.is_time_up(): + return + time.sleep(min(0.5, end - time.time())) def _exec_loop(self, d, params, depth=0): children = params.get("children", []) if not children: return mode = params.get("loop_mode", "rounds") - if mode == "time": + if mode == "forever": + # 一直循环:无限跑直到任务被外部停止(定时停止 cron / 运行时长上限 / 手动停止)。 + # 配合"定时启动+停止"调度(如 9:00 启动 / 18:00 停止)实现全天候运行。 + i = 0 + while not self.stopped() and not self.is_time_up(): + i += 1 + self.set_action(f"循环第 {i} 轮(持续)") + if i % 10 == 0 or i <= 2: + _log.info(f"[{self.serial}] loop 第 {i} 轮(持续运行)") + self._exec_steps(d, children, depth + 1) + elif mode == "time": # 按时间循环:跑满 loop_duration 秒,同时受全局 stopped/max_duration 约束 duration = float(params.get("loop_duration", 0) or 0) if duration <= 0: