Merge branch 'feat/task-patrol'——任务级公共巡检(含与拟人滑动的合并:import、能力表两处冲突已解)
This commit is contained in:
+167
-6
@@ -26,6 +26,8 @@ worker 按 steps 顺序执行,支持 loop 步骤循环、停止信号、进度
|
||||
wait_el - 等待元素出现(条件等待,替代固定时长 wait)
|
||||
input_text - 输入文字(随机候选/指定文字,可先清空)
|
||||
wait - 等待时长
|
||||
notify - 发自定义通知(标题/正文自己写,推送 webhook)
|
||||
stop_self - 停止本设备的任务(其它设备不受影响)
|
||||
loop - 循环块(含 children 步骤列表 + max_iterations)
|
||||
group - 动作组(含 children 步骤列表,按序执行一次,可折叠复用)
|
||||
"""
|
||||
@@ -35,9 +37,10 @@ import time
|
||||
|
||||
from tasks.base import BaseTask, register_task
|
||||
from core.device_worker import BaseWorker, _update_status
|
||||
from core.u2_helper import ensure_app_running, wait_for_app_home, random_sleep
|
||||
from core.u2_helper import (ensure_app_running, wait_for_app_home, random_sleep,
|
||||
screen_is_on, current_package)
|
||||
from core.logger import get_logger
|
||||
from core import notifier, step_log, humanize
|
||||
from core import notifier, step_log, humanize, patrol
|
||||
|
||||
_log = get_logger("task.generic")
|
||||
|
||||
@@ -110,6 +113,10 @@ STEP_TYPES = [
|
||||
{"type": "if_el", "label": "条件判断", "icon": "❓",
|
||||
"params": {"selector_type": "xpath", "selector_value": "", "timeout": 3,
|
||||
"then": [], "else": []}},
|
||||
{"type": "notify", "label": "发通知", "icon": "🔔",
|
||||
"params": {"title": "", "message": "", "level": "info"}},
|
||||
{"type": "stop_self", "label": "停止本设备", "icon": "⛔",
|
||||
"params": {"reason": ""}},
|
||||
]
|
||||
|
||||
|
||||
@@ -139,6 +146,17 @@ class GenericStepsWorker(BaseWorker):
|
||||
self._path = []
|
||||
self._step_rows = 0 # 本次运行已记录的明细条数(封顶见 _record_step)
|
||||
self._step_capped = False
|
||||
# 公共巡检(任务级配置,独立于步骤画布;见 core/patrol.py):
|
||||
# _wid 内部标识,用来记"下次到点 / 上次命中"两个节奏
|
||||
self._watchers = []
|
||||
for i, w in enumerate(p.get("watchers") or []):
|
||||
if not isinstance(w, dict) or not w.get("enabled", True):
|
||||
continue
|
||||
w = dict(w)
|
||||
w["_wid"] = f"w{i}"
|
||||
self._watchers.append(w)
|
||||
self._watch_next = {} # wid -> 下次到点(缺省 0 = 任务一开始就先查一次)
|
||||
self._watch_last = {} # wid -> 上次命中(冷却用)
|
||||
|
||||
# 某 click 选择器连续未找到元素的次数达到该值,判定可能失效(App 改版)
|
||||
_MAX_CONSECUTIVE_MISS = 10
|
||||
@@ -188,6 +206,7 @@ class GenericStepsWorker(BaseWorker):
|
||||
for i, step in enumerate(steps):
|
||||
if self.stopped() or self.is_time_up():
|
||||
return
|
||||
self._maybe_patrol(d) # 穿插:每步之前看一眼有没有巡检到点了
|
||||
# 压入本步在兄弟里的序号(1 起):容器步骤执行 children 时会在其后
|
||||
# 继续追加,得到 "2.1" 这样的嵌套路径
|
||||
self._path.append(str(i + 1))
|
||||
@@ -282,6 +301,71 @@ class GenericStepsWorker(BaseWorker):
|
||||
except Exception:
|
||||
pass # 记录失败绝不能影响任务执行
|
||||
|
||||
# ================== 公共巡检(任务级,穿插执行) ==================
|
||||
def _maybe_patrol(self, d):
|
||||
"""看一眼有没有巡检到点了;到点就跑一次检查(配置见 core/patrol.py)。
|
||||
|
||||
穿插式:不起线程、不抢屏幕,所以巡检精度受主流程步长影响——某一步卡
|
||||
30 秒,巡检最多晚 30 秒(这是刻意取舍,换来的是绝不会和主流程打架)。
|
||||
"""
|
||||
if not self._watchers or self.stopped():
|
||||
return
|
||||
now = time.time()
|
||||
for w in self._watchers:
|
||||
wid = w["_wid"]
|
||||
if now < self._watch_next.get(wid, 0):
|
||||
continue
|
||||
# 先排下一次再执行:检查本身可能慢(元素检查要 dump UI 树),
|
||||
# 否则慢检查会导致"每分钟"变成"每检查一次"
|
||||
self._watch_next[wid] = now + max(5, int(w.get("interval") or 60))
|
||||
try:
|
||||
self._run_patrol(d, w, now)
|
||||
except Exception as e:
|
||||
_log.warning(f"[{self.serial}] 巡检[{w.get('name') or wid}]异常: {e}")
|
||||
|
||||
def _run_patrol(self, d, w, now):
|
||||
"""执行一次巡检:命中 → 过冷却 → 做动作 + 发通知 + 记一条明细。"""
|
||||
wid, name = w["_wid"], (w.get("name") or patrol.check_label(w.get("check")))
|
||||
hit, detail = patrol.evaluate(d, w)
|
||||
if not hit:
|
||||
return
|
||||
# 命中后冷却:条件持续成立(比如一直熄屏)时不要每分钟都动作/发通知
|
||||
cooldown = max(0, int(w.get("cooldown") or 0))
|
||||
if now - self._watch_last.get(wid, 0) < cooldown:
|
||||
_log.info(f"[{self.serial}] 巡检[{name}]命中({detail})但在冷却期内,不重复动作")
|
||||
return
|
||||
self._watch_last[wid] = now
|
||||
# 文案在**动作之前**渲染:动作会改变设备状态(点亮后 {screen} 就成了"亮屏"),
|
||||
# 而用户要看的是"发现时"的样子——"发现熄屏,已点亮"而不是"发现亮屏,已点亮"
|
||||
title = self._render_vars(d, w.get("title") or "") or f"巡检命中:{name}"
|
||||
message = self._render_vars(d, w.get("message") or "") or detail
|
||||
done = patrol.act(d, w, self)
|
||||
_log.info(f"[{self.serial}] 巡检[{name}]命中:{detail}"
|
||||
f"{(' → ' + done) if done else '(只通知)'}")
|
||||
self.set_action(f"巡检[{name}]:{detail}{(',' + done) if done else ''}")
|
||||
summary = f"{detail}{(' → ' + done) if done else ''}"
|
||||
try: # 明细:只记命中,不记"没事发生"
|
||||
step_log.record(run_id=self.ctx.get("run_id", ""), serial=self.serial,
|
||||
device_name=self.ctx.get("device_name", ""),
|
||||
job_id=self.ctx.get("job_id", ""),
|
||||
job_name=self.ctx.get("job_name", ""),
|
||||
step_label=f"巡检[{name}]", step_type="patrol",
|
||||
result="hit", detail=summary)
|
||||
except Exception:
|
||||
pass
|
||||
if w.get("notify", True):
|
||||
notifier.notify("task.patrol.hit",
|
||||
patrol_name=name,
|
||||
check_label=patrol.check_label(w.get("check")),
|
||||
action_label=patrol.action_label(w.get("action")),
|
||||
detail=detail,
|
||||
job_id=self.ctx.get("job_id", ""),
|
||||
job_name=self.ctx.get("job_name", ""),
|
||||
serial=self.serial,
|
||||
device_name=self.ctx.get("device_name", ""),
|
||||
title=title, message=message,
|
||||
level=w.get("level", "warning"))
|
||||
|
||||
# ================== 步骤执行器 ==================
|
||||
def _exec_screen_on(self, d, params, depth=0):
|
||||
"""亮屏:息屏时唤醒并滑动解锁。任务执行前用(息屏时 u2 无法操作)。"""
|
||||
@@ -494,7 +578,23 @@ class GenericStepsWorker(BaseWorker):
|
||||
_log.warning(f"[{self.serial}] if_el 缺少选择器,跳过")
|
||||
return None
|
||||
found = False
|
||||
if sel_type == "ocr":
|
||||
# 非元素类条件(屏幕状态/前台 App):不参与"选择器健康"统计——它们不是选择器,
|
||||
# 连续未命中不该被记成"选择器失效"
|
||||
dynamic = sel_type in ("screen", "foreground")
|
||||
detail = ""
|
||||
if sel_type == "screen":
|
||||
# 屏幕亮/熄:selector_value 填 on / off(也认 熄屏/灭屏 这类中文)
|
||||
want = str(sel_val).strip().lower()
|
||||
want_off = want in ("off", "0", "false", "no", "熄屏", "灭屏", "黑屏")
|
||||
state = screen_is_on(d)
|
||||
found = (state is False) if want_off else (state is True)
|
||||
detail = (f"屏幕{'熄屏' if state is False else '亮屏' if state else '状态未知'}"
|
||||
f"(判断{'熄屏' if want_off else '亮屏'})")
|
||||
elif sel_type == "foreground":
|
||||
cur = current_package(d)
|
||||
found = cur == str(sel_val).strip()
|
||||
detail = f"前台='{cur or '未知'}'(期望 {sel_val})"
|
||||
elif sel_type == "ocr":
|
||||
# OCR 模式:截图 → 识别文字 → 关键词匹配(UI 树里没有的文字也能找到)
|
||||
try:
|
||||
from core.ocr import available as _ocr_available, find_on_screen
|
||||
@@ -527,10 +627,13 @@ class GenericStepsWorker(BaseWorker):
|
||||
except Exception as e:
|
||||
_log.warning(f"[{self.serial}] 条件判断异常: {e}")
|
||||
found = False
|
||||
self._track_selector_health(sel_val, found)
|
||||
if dynamic:
|
||||
_log.info(f"[{self.serial}] 条件判断 {detail} → {'命中' if found else '未命中'}")
|
||||
else:
|
||||
self._track_selector_health(sel_val, found)
|
||||
branch = params.get("then" if found else "else") or []
|
||||
self.set_action(f"条件{'命中' if found else '未命中'} → 执行{'找到' if found else '未找到'}分支({len(branch)}步)")
|
||||
_log.info(f"[{self.serial}] 条件判断 {'命中' if found else '未命中'}: {sel_val},执行{'then' if found else 'else'}分支 {len(branch)} 步")
|
||||
_log.info(f"[{self.serial}] 条件判断 {'命中' if found else '未命中'}: {detail or sel_val},执行{'then' if found else 'else'}分支 {len(branch)} 步")
|
||||
self._exec_steps(d, branch, depth + 1)
|
||||
return found
|
||||
|
||||
@@ -552,6 +655,61 @@ class GenericStepsWorker(BaseWorker):
|
||||
f"({info['from']}→{info['to']}"
|
||||
f"{',弧线' if info.get('human') else ',直线'})")
|
||||
|
||||
def _exec_notify(self, d, params, depth=0):
|
||||
"""发一条自定义通知(标题/正文由任务自己写,支持占位符)。
|
||||
|
||||
事件是 `task.notify.custom`:**谁能收到取决于 webhook 的事件订阅**——
|
||||
这个步骤只负责"把消息交出去",不关心发给谁(与平台其它通知同一套口径)。
|
||||
"""
|
||||
title = self._render_vars(d, params.get("title") or "")
|
||||
msg = self._render_vars(d, params.get("message") or "")
|
||||
if not (title or msg):
|
||||
_log.warning(f"[{self.serial}] notify 步骤既没标题也没内容,跳过")
|
||||
return None
|
||||
notifier.notify("task.notify.custom",
|
||||
job_id=self.ctx.get("job_id", ""),
|
||||
job_name=self.ctx.get("job_name", ""),
|
||||
serial=self.serial,
|
||||
device_name=self.ctx.get("device_name", ""),
|
||||
title=title or "任务通知", message=msg,
|
||||
level=params.get("level", "info"))
|
||||
_log.info(f"[{self.serial}] 已发自定义通知: {title} {msg}")
|
||||
return True
|
||||
|
||||
def _render_vars(self, d, text):
|
||||
"""替换通知文本里的占位符:{device} {serial} {job} {time} {app} {screen}。
|
||||
|
||||
`{app}` / `{screen}` 要查设备(0.3~0.7s),**只在文本里真用到时才查**,
|
||||
否则每条通知都白付一次设备往返。
|
||||
"""
|
||||
if not text or "{" not in text:
|
||||
return text
|
||||
vals = {"device": self.ctx.get("device_name") or self.serial,
|
||||
"serial": self.serial,
|
||||
"job": self.ctx.get("job_name", ""),
|
||||
"time": time.strftime("%Y-%m-%d %H:%M:%S")}
|
||||
if "{app}" in text:
|
||||
vals["app"] = current_package(d) or "未知"
|
||||
if "{screen}" in text:
|
||||
st = screen_is_on(d)
|
||||
vals["screen"] = "未知" if st is None else ("亮屏" if st else "熄屏")
|
||||
for k, v in vals.items():
|
||||
text = text.replace("{" + k + "}", str(v))
|
||||
return text
|
||||
|
||||
def _exec_stop_self(self, d, params, depth=0):
|
||||
"""停止本设备的任务(其它设备的任务不受影响)。
|
||||
|
||||
实现只是置 worker 的停止位:`_exec_steps` 每步前检查 `stopped()`,所以
|
||||
后续步骤不会再跑;TaskManager 侧把它记成"被停止"而不是失败
|
||||
(见 `core/task_manager.py` 的 `_run_with_retry` 自停分支)。
|
||||
"""
|
||||
reason = (params.get("reason") or "").strip()
|
||||
_log.info(f"[{self.serial}] 任务内主动停止本设备({reason or '未填原因'})")
|
||||
self.set_action(f"停止本设备任务{(':' + reason) if reason else ''}")
|
||||
self.stop()
|
||||
return True
|
||||
|
||||
def _exec_click(self, d, params, depth=0):
|
||||
sel_type = params.get("selector_type", "xpath")
|
||||
sel_val = params.get("selector_value", "")
|
||||
@@ -632,7 +790,10 @@ class GenericStepsWorker(BaseWorker):
|
||||
while time.time() < end:
|
||||
if self.stopped() or self.is_time_up():
|
||||
return
|
||||
time.sleep(min(0.5, end - time.time()))
|
||||
self._maybe_patrol(d) # 长等待里也看一眼,别让巡检被 wait 拖住
|
||||
# max(0, …):巡检本身要花时间(屏幕 0.3s、前台 0.7s、元素检查更久),
|
||||
# 扣掉这段时间后可能已经过点了,负数 sleep 会直接抛异常
|
||||
time.sleep(max(0.0, min(0.5, end - time.time())))
|
||||
|
||||
def _exec_loop(self, d, params, depth=0):
|
||||
children = params.get("children", [])
|
||||
|
||||
Reference in New Issue
Block a user