feat: 通用步骤支持概率触发、循环按时间、预计结束时间上报
- 步骤新增 probability 参数:<100 时按百分比概率决定是否执行(实现"偶尔点赞") - loop 支持 loop_mode=rounds/time,按时间循环跑满 loop_duration 秒 - BaseWorker._start_timer 上报 end_time(开始+max_duration),task_manager 透出
This commit is contained in:
@@ -246,6 +246,9 @@ class BaseWorker(threading.Thread):
|
|||||||
def _start_timer(self):
|
def _start_timer(self):
|
||||||
"""子类在 run_task 开头调用,启动运行时长计时。"""
|
"""子类在 run_task 开头调用,启动运行时长计时。"""
|
||||||
self._start_time = time.time()
|
self._start_time = time.time()
|
||||||
|
if self.max_duration > 0:
|
||||||
|
# 预计结束时间(前端监控页展示),到点自动停止
|
||||||
|
_update_status(self.serial, end_time=self._start_time + self.max_duration)
|
||||||
|
|
||||||
def is_time_up(self):
|
def is_time_up(self):
|
||||||
"""是否已达最大运行时长。max_duration=0 时永远返回 False。"""
|
"""是否已达最大运行时长。max_duration=0 时永远返回 False。"""
|
||||||
|
|||||||
@@ -735,6 +735,7 @@ class TaskManager:
|
|||||||
"running_job": r.get("job_id", ""),
|
"running_job": r.get("job_id", ""),
|
||||||
"task_job": w.get("task_job", ""),
|
"task_job": w.get("task_job", ""),
|
||||||
"attempt": r.get("attempt", 0),
|
"attempt": r.get("attempt", 0),
|
||||||
|
"end_time": w.get("end_time", 0),
|
||||||
})
|
})
|
||||||
return result
|
return result
|
||||||
|
|
||||||
|
|||||||
+29
-7
@@ -21,6 +21,7 @@ worker 按 steps 顺序执行,支持 loop 步骤循环、停止信号、进度
|
|||||||
group - 动作组(含 children 步骤列表,按序执行一次,可折叠复用)
|
group - 动作组(含 children 步骤列表,按序执行一次,可折叠复用)
|
||||||
"""
|
"""
|
||||||
import random
|
import random
|
||||||
|
import time
|
||||||
|
|
||||||
from tasks.base import BaseTask, register_task
|
from tasks.base import BaseTask, register_task
|
||||||
from core.device_worker import BaseWorker, _update_status
|
from core.device_worker import BaseWorker, _update_status
|
||||||
@@ -43,7 +44,7 @@ STEP_TYPES = [
|
|||||||
{"type": "wait", "label": "等待", "icon": "⏱",
|
{"type": "wait", "label": "等待", "icon": "⏱",
|
||||||
"params": {"min": 1.0, "max": 3.0}},
|
"params": {"min": 1.0, "max": 3.0}},
|
||||||
{"type": "loop", "label": "循环块", "icon": "↻",
|
{"type": "loop", "label": "循环块", "icon": "↻",
|
||||||
"params": {"max_iterations": 10, "children": []}},
|
"params": {"loop_mode": "rounds", "max_iterations": 10, "loop_duration": 600, "children": []}},
|
||||||
{"type": "group", "label": "动作组", "icon": "📦",
|
{"type": "group", "label": "动作组", "icon": "📦",
|
||||||
"params": {"children": []}},
|
"params": {"children": []}},
|
||||||
]
|
]
|
||||||
@@ -112,6 +113,11 @@ class GenericStepsWorker(BaseWorker):
|
|||||||
stype = step.get("type", "")
|
stype = step.get("type", "")
|
||||||
label = step.get("label", stype)
|
label = step.get("label", stype)
|
||||||
params = step.get("params", {})
|
params = step.get("params", {})
|
||||||
|
# 概率触发:probability=100 必执行,<100 按百分比概率决定本次是否执行
|
||||||
|
prob = float(params.get("probability", 100))
|
||||||
|
if prob < 100 and random.random() * 100 > prob:
|
||||||
|
_log.info(f"[{self.serial}] 步骤 '{label}' 概率 {prob}% 未触发,跳过")
|
||||||
|
return
|
||||||
self.set_action(f"执行: {label}")
|
self.set_action(f"执行: {label}")
|
||||||
_log.info(f"[{self.serial}] 步骤: {label}({stype}) params={params}")
|
_log.info(f"[{self.serial}] 步骤: {label}({stype}) params={params}")
|
||||||
|
|
||||||
@@ -196,16 +202,32 @@ class GenericStepsWorker(BaseWorker):
|
|||||||
random_sleep(min_s, max_s)
|
random_sleep(min_s, max_s)
|
||||||
|
|
||||||
def _exec_loop(self, d, params, depth=0):
|
def _exec_loop(self, d, params, depth=0):
|
||||||
max_iter = int(params.get("max_iterations", 10))
|
|
||||||
children = params.get("children", [])
|
children = params.get("children", [])
|
||||||
if not children:
|
if not children:
|
||||||
return
|
return
|
||||||
for i in range(max_iter):
|
mode = params.get("loop_mode", "rounds")
|
||||||
if self.stopped() or self.is_time_up():
|
if mode == "time":
|
||||||
|
# 按时间循环:跑满 loop_duration 秒,同时受全局 stopped/max_duration 约束
|
||||||
|
duration = float(params.get("loop_duration", 0) or 0)
|
||||||
|
if duration <= 0:
|
||||||
|
_log.warning(f"[{self.serial}] 按时间循环未设置时长,跳过")
|
||||||
return
|
return
|
||||||
self.set_action(f"循环第 {i+1}/{max_iter} 轮")
|
start = time.time()
|
||||||
_log.info(f"[{self.serial}] loop 第 {i+1}/{max_iter} 轮")
|
i = 0
|
||||||
self._exec_steps(d, children, depth + 1)
|
while not self.stopped() and not self.is_time_up() and (time.time() - start) < duration:
|
||||||
|
i += 1
|
||||||
|
self.set_action(f"循环第 {i} 轮(按时间)")
|
||||||
|
_log.info(f"[{self.serial}] loop 第 {i} 轮(按时间, 已{(time.time()-start):.0f}s/{duration:.0f}s)")
|
||||||
|
self._exec_steps(d, children, depth + 1)
|
||||||
|
else:
|
||||||
|
# 按轮次循环
|
||||||
|
max_iter = int(params.get("max_iterations", 10))
|
||||||
|
for i in range(max_iter):
|
||||||
|
if self.stopped() or self.is_time_up():
|
||||||
|
return
|
||||||
|
self.set_action(f"循环第 {i+1}/{max_iter} 轮")
|
||||||
|
_log.info(f"[{self.serial}] loop 第 {i+1}/{max_iter} 轮")
|
||||||
|
self._exec_steps(d, children, depth + 1)
|
||||||
|
|
||||||
def _exec_group(self, d, params, depth=0):
|
def _exec_group(self, d, params, depth=0):
|
||||||
"""动作组:按序执行 children 一次(类似 loop 但不循环)。
|
"""动作组:按序执行 children 一次(类似 loop 但不循环)。
|
||||||
|
|||||||
Reference in New Issue
Block a user