Files
auto_control/tasks/generic/task.py
T
butubb 97cec6e211 feat(通知): 系统级 Webhook 通知子系统(事件目录 + 可插拔适配器 + 防刷屏)
平台此前出了问题只能靠人盯页面。现在各组件统一走 `notifier.notify(事件, **字段)`,
推到企业微信 / 自建服务;**所有可通知点都登记进事件目录,默认全关,用户按 webhook 勾选**。

## 架构(`core/notifier.py` + `core/notify_events.py`)
业务线程调 `notify()` → 只做内存操作(读配置快照/匹配订阅/入队)→ 返回;
后台 1 个 dispatcher(聚合 + 每 hook 限流 + 折叠摘要)+ 3 个 sender(真实 HTTP、退避重试)
负责真正发出去。硬约束:**notify 零 DB、零 HTTP、零阻塞、异常不冒泡**——所以任务线程里
可以直接调(不用 app_context、不用 try/except),但**必须放在所有 `with self._lock` 之外**。

- **事件目录 36 条**(任务批次/单设备/设备/Worker/业务/安装/系统/AI),支持 `task.*` 通配订阅;
  语义分工避免重复告警:`worker.*` 是单次尝试级,`task.device.success/failed` 是唯一权威结论。
- **适配器可插拔**:`wecom`(markdown,按 4096 **字节**截断、超限不截半个汉字)+
  `json`(模板占位符,替换值按 JSON 转义,保存前干跑校验);钉钉/飞书留了插槽(前端置灰)。
- **防打爆四层**:聚合窗口(默认 30s,同批次合并成一条并带样本)→ 令牌桶限流(默认 18/分,
  对齐企微硬限 20)→ 被限流的**折叠成摘要不丢弃** → 有界队列背压。取舍:失败通知最多延迟
  一个窗口,换来群不被刷屏。

## 安全与存储
- 配置只落 `app_meta.notify_webhooks` 一个键(**不建表** → 不涉及备份覆盖红线)。
- URL 本身就是凭据(企微 `?key=`)→ 接口回显/发送记录/日志一律 `mask_url()/scrub()`;
  编辑时留空即不修改;secret 永不回显。DATA_MODEL 的明文凭据告警补上了这一条。
- 发送记录:内存环形缓冲 200 条(重启清空)+ 独立 `logs/notify.log`。

## 接入点(每个都放在锁外、不改 return 顺序)
task_manager(批次开始/结束用新增的 `_BatchTracker` 统一在 finally 计数、单设备成功/失败/
离线/重试/停止/抢占/归还/cron 停止)、device_worker 心跳看门狗、generic 任务选择器连续失效、
apk 安装开始/完成、设备上下线(**状态沿检测**,只报新变化)、备份导出/恢复、经验巡检、
用户登录、服务启停。

## 前端
系统 Tab 新增「通知 / Webhook」子分栏:多条 webhook 列表(URL 打码)+ 编辑弹窗(格式/URL/
密钥/事件勾选树带 ★建议/聚合/限流/自定义模板/预览)+ 发送测试 + 发送记录。

## 自测
- 进程内逻辑 10 组断言全绿:聚合合并、限流+折叠、无配置/全局关静默丢弃、未知事件、
  内部异常不外泄、JSON 转义(标题含引号换行仍合法)、URL/异常消息脱敏、配置校验。
- 端到端(假 webhook 接收端)17 项断言全绿:真实事件投递(user.login / task.batch.no_device)、
  企微请求体形状、**HTTP 200 + errcode 93000 判为失败**、500 重试 3 次、记录里 URL 打码。
- 韧性:webhook 指向黑洞地址时登录耗时 100~114ms(基线 107~133ms,**异步隔离生效**);
  配置写成坏 JSON 服务照常启动、通知静默不发、日志有 error(服务端实测后已复原)。
- 页面:系统 → 通知 面板/弹窗/36 个事件复选框/预览全部正常,无 JS 报错。
- 自测数据已清理(webhook、自建任务、写坏又复原的配置键)。

文档:新增 doc/NOTIFY.md(事件表/配置/格式约束/防刷屏/加事件三步骤/排障)并登记进 doc/README;
API.md §2.11;DATA_MODEL 的 app_meta 键表与明文凭据告警;ARCHITECTURE 线程表/分层/扩展点;
根 README 功能索引与日志表。
2026-09-15 13:55:14 +08:00

710 lines
34 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""通用步骤任务定义。
前端步骤编辑器编排步骤链,保存为 params.steps(JSON 数组)。
worker 按 steps 顺序执行,支持 loop 步骤循环、停止信号、进度上报。
步骤 schema:
{
"id": "step_1", # 唯一 id(前端生成)
"type": "open_app", # 步骤类型
"label": "打开抖音", # 用户可读名称
"params": { ... } # 步骤参数(按 type 不同)
}
支持的步骤类型:
open_app - 启动 app
stop_app - 强制结束 app(am force-stop,下次打开冷启动)
screen_on - 亮屏(息屏时唤醒并滑动解锁)
screen_off - 息屏
keep_screen - 保持亮屏/恢复自动息屏(svc power stayon,设备充电时常亮)
key_event - 按键(返回/Home/回车/菜单等,d.press)
swipe - 滑动(方向/时长)
swipe_until - 滑动直到元素出现(最多 N 次,可选中找到后点击)
click - 点击元素(选择器)
click_xy - 点击坐标(屏幕百分比,无选择器时兜底)
long_click - 长按元素(选择器 + 时长)
wait_el - 等待元素出现(条件等待,替代固定时长 wait)
input_text - 输入文字(随机候选/指定文字,可先清空)
wait - 等待时长
loop - 循环块(含 children 步骤列表 + max_iterations)
group - 动作组(含 children 步骤列表,按序执行一次,可折叠复用)
"""
import random
import re
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.logger import get_logger
from core import notifier
_log = get_logger("task.generic")
# 历史抓取器生成的"伪序号"选择器://*[@resource-id="x"][4]
# XPath 里这是"父节点内排第 4",不是"第 4 个匹配"——同属性多实例时 [2..n] 全部失配。
_LEGACY_IDX_XPATH = re.compile(r'^(//\*\[@[^\]]+\])\[(\d+)\](.*)$')
# ---- 「保持亮屏」用的系统设置 ----
# 系统息屏超时(毫秒):保持亮屏时顶到最大,恢复时写回原值
_SCREEN_TIMEOUT_MAX = "2147483647"
_SCREEN_TIMEOUT_DEFAULT = "600000" # 兜底恢复值:10 分钟(原值记不住时用)
_SCREEN_TIMEOUT_BACKUP = {} # serial -> 原始 screen_off_timeout
def _norm_legacy_xpath(sel_val):
"""把历史 `//*[@attr=…][k]` 纠正为 `(//*[@attr=…])[k]`(只改整体前缀,保留后续子路径)。
窄范围:只匹配"属性XPath + 数字谓词"开头的形态;结构路径里
`.../FrameLayout[2]` 的兄弟序号是有意为之,不受影响。
"""
if not isinstance(sel_val, str):
return sel_val
m = _LEGACY_IDX_XPATH.match(sel_val.strip())
if not m:
return sel_val
return f"({m.group(1)})[{m.group(2)}]{m.group(3)}"
# ================== 步骤类型定义(前端操作库 + 后端执行共用)==================
STEP_TYPES = [
{"type": "open_app", "label": "打开App", "icon": "▶",
"params": {"package": "", "wait_home": False, "home_feature": ""}},
{"type": "stop_app", "label": "结束App", "icon": "⏹",
"params": {"package": ""}},
{"type": "screen_on", "label": "亮屏", "icon": "💡",
"params": {}},
{"type": "screen_off", "label": "息屏", "icon": "🌙",
"params": {}},
{"type": "keep_screen", "label": "保持亮屏", "icon": "🔆",
"params": {"mode": "on"}},
{"type": "key_event", "label": "按键", "icon": "⌨️",
"params": {"key": "back"}},
{"type": "swipe", "label": "滑动", "icon": "↕",
"params": {"direction": "up", "duration_min": 0.25, "duration_max": 0.50}},
{"type": "swipe_until", "label": "滑动直到元素", "icon": "🔍",
"params": {"selector_type": "xpath", "selector_value": "", "direction": "up",
"max_swipes": 8, "click_when_found": True}},
{"type": "click", "label": "点击元素", "icon": "✦",
"params": {"selector_type": "description", "selector_value": "", "wait_timeout": 2}},
{"type": "click_xy", "label": "点击坐标", "icon": "🎯",
"params": {"x": 50, "y": 50}},
{"type": "long_click", "label": "长按元素", "icon": "👆",
"params": {"selector_type": "xpath", "selector_value": "", "duration": 1.0, "wait_timeout": 2}},
{"type": "wait_el", "label": "等待元素", "icon": "⏳",
"params": {"selector_type": "xpath", "selector_value": "", "timeout": 10}},
{"type": "input_text", "label": "输入文字", "icon": "⌨",
"params": {"mode": "random", "texts": "你好\n有趣\n支持", "fixed_text": "", "clear_first": True}},
{"type": "clipboard", "label": "剪贴板注入", "icon": "📋",
"params": {"text": "", "paste": True}},
{"type": "wait", "label": "等待", "icon": "⏱",
"params": {"min": 1.0, "max": 3.0}},
{"type": "loop", "label": "循环块", "icon": "↻",
"params": {"loop_mode": "rounds", "max_iterations": 10, "loop_duration": 600, "children": []}},
{"type": "group", "label": "动作组", "icon": "📦",
"params": {"children": []}},
{"type": "if_el", "label": "条件判断", "icon": "❓",
"params": {"selector_type": "xpath", "selector_value": "", "timeout": 3,
"then": [], "else": []}},
]
DEFAULT_PARAMS = {
"max_duration": 0, # 最大运行时长(秒),0=不限时
# 注意:**没有默认 steps**。步骤由前端步骤编辑器编排产出,新建任务时从零开始拖;
# 缺 steps 的任务执行时会明确报错(见 run_task),不会静默空跑。
}
class GenericStepsWorker(BaseWorker):
"""通用步骤 worker:按 params.steps 顺序执行步骤链。
复用 BaseWorker 的设备生命周期(STF 占用/释放、u2 连接、异常、状态上报、stop)。
只实现 run_task,按步骤类型分发到 _exec_<type> 方法。
"""
def __init__(self, serial, params=None):
super().__init__(serial, params)
p = {**DEFAULT_PARAMS, **(self.params or {})}
self.max_duration = int(p.get("max_duration", 0))
self.steps = p.get("steps", [])
self._action_counts = {} # 步骤执行计数 {step_label: count}
self._miss_counts = {} # 选择器健康:selector -> 连续未命中次数
# 某 click 选择器连续未找到元素的次数达到该值,判定可能失效(App 改版)
_MAX_CONSECUTIVE_MISS = 10
def _track_selector_health(self, sel_val, found):
"""记录某选择器的连续未命中次数,达到阈值上报 last_warning 并重置计数。
用于发现"选择器失效但任务仍显示成功"的静默空转问题(如 App 改版后元素找不到了)。
"""
if found:
self._miss_counts.pop(sel_val, None)
return
n = self._miss_counts.get(sel_val, 0) + 1
self._miss_counts[sel_val] = n
if n >= self._MAX_CONSECUTIVE_MISS:
_log.error(f"[{self.serial}] 选择器连续 {n} 次未找到元素,可能已失效(App 改版?):{sel_val}")
_update_status(self.serial, last_warning=f"选择器连续{n}次未命中")
self._miss_counts.pop(sel_val, None) # 已上报,重置避免刷屏
# 通知:这是"任务显示成功但什么都没做"的隐蔽故障,值得人看一眼
# (重置计数后要再攒满 10 次才会再报,天然不刷屏)
notifier.notify("task.selector.invalid", serial=self.serial,
selector=str(sel_val)[:200], miss_count=n)
def run_task(self, d):
"""按 steps 顺序执行。loop 步骤递归执行其 children。"""
if not self.steps:
# 任务没配步骤(DEFAULT_PARAMS 不含默认步骤,步骤只能由编辑器产出):
# 明确报错而不是"跑完 0 步"当成功——设备表/任务页会显示该错误
raise ValueError("通用步骤任务没有可执行步骤:请在「任务」页编辑该任务并添加步骤")
self._start_timer()
self.set_action("开始执行步骤链")
self.set_progress(done=0, total=0, unit="操作",
action_counts=dict(self._action_counts))
try:
self._exec_steps(d, self.steps)
finally:
elapsed = self.elapsed()
summary = f"完成步骤链,运行 {elapsed//60}m{elapsed%60}s"
_update_status(self.serial, current_action=summary)
def _exec_steps(self, d, steps, depth=0):
"""执行步骤列表。depth 防止无限嵌套。"""
if depth > 5:
_log.warning(f"[{self.serial}] 步骤嵌套深度超限(>5),跳过")
return
for step in steps:
if self.stopped() or self.is_time_up():
return
self._exec_one(d, step, depth)
def _exec_one(self, d, step, depth=0):
"""执行单个步骤。"""
stype = step.get("type", "")
label = step.get("label", stype)
params = step.get("params", {})
# 兼容历史选择器://*[@attr=…][k] → (//*[@attr=…])[k](XPath 位置谓词语义)
if params.get("selector_type") == "xpath" and params.get("selector_value"):
fixed = _norm_legacy_xpath(params["selector_value"])
if fixed != params["selector_value"]:
params = dict(params, selector_value=fixed)
_log.info(f"[{self.serial}] 选择器已纠正为: {fixed}")
# 概率触发: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}")
_log.info(f"[{self.serial}] 步骤: {label}({stype}) params={params}")
ret = None
handler = getattr(self, f"_exec_{stype}", None)
if handler:
try:
ret = handler(d, params, depth)
except Exception as e:
_log.error(f"[{self.serial}] 步骤 {label}({stype}) 异常: {e}")
else:
_log.warning(f"[{self.serial}] 未知步骤类型: {stype}")
# 容器步骤(loop/group/if_el)不计入操作计数,避免监控显示噪音
if stype not in ("loop", "group", "if_el"):
self._action_counts[label] = self._action_counts.get(label, 0) + 1
elapsed = self.elapsed()
# done=累计执行的操作次数(无固定总数),前端显示"已执行 N 次操作"
self.set_progress(done=sum(self._action_counts.values()),
total=0, unit="操作",
action_counts=dict(self._action_counts), elapsed=elapsed)
return ret # 命中结果(True/False/None),"测试此步骤"用
# ================== 步骤执行器 ==================
def _exec_screen_on(self, d, params, depth=0):
"""亮屏:息屏时唤醒并滑动解锁。任务执行前用(息屏时 u2 无法操作)。"""
try:
d.unlock()
_log.info(f"[{self.serial}] 亮屏")
except Exception as e:
_log.warning(f"[{self.serial}] 亮屏异常: {e}")
def _exec_screen_off(self, d, params, depth=0):
"""息屏:任务结束后用,避免长时间亮屏(烧屏/发热)。"""
try:
d.screen_off()
_log.info(f"[{self.serial}] 息屏")
except Exception as e:
_log.warning(f"[{self.serial}] 息屏异常: {e}")
def _exec_keep_screen(self, d, params, depth=0):
"""保持亮屏 / 恢复自动息屏。
⚠️ 只用 `svc power stayon true` **不管用**:它管的是「**充电时**屏幕常亮」
(= stay_on_while_plugged_in)。设备走 WiFi 跑任务、没插充电器时它完全不起作用
——以前"保持亮屏"看着像没生效就是这个原因(2026-09-14 用户反馈)。
真正管用的是把系统**息屏超时**顶到最大(`screen_off_timeout`):不充电也不会
自动黑屏。进入时记下原值,`mode=off` 写回;进程重启丢了记录就写回常规的 10 分钟。
"""
mode = params.get("mode", "on")
try:
if mode == "off":
d.shell("svc power stayon false")
orig = _SCREEN_TIMEOUT_BACKUP.pop(self.serial, _SCREEN_TIMEOUT_DEFAULT)
d.shell(f"settings put system screen_off_timeout {orig}")
_log.info(f"[{self.serial}] 恢复自动息屏(超时 {orig}ms)")
else:
cur = ""
try:
cur = (d.shell("settings get system screen_off_timeout") or "").strip()
except Exception:
pass
if self.serial not in _SCREEN_TIMEOUT_BACKUP:
# 没设过时 `settings get` 会返回 "null"(走系统默认)——别把 null 写回去
_SCREEN_TIMEOUT_BACKUP[self.serial] = \
cur if cur.isdigit() else _SCREEN_TIMEOUT_DEFAULT
d.shell("svc power stayon true") # 充电场景(顺带)
d.shell(f"settings put system screen_off_timeout {_SCREEN_TIMEOUT_MAX}")
d.shell("input keyevent 224") # 立刻唤醒(KEYCODE_WAKEUP,单向)
_log.info(f"[{self.serial}] 保持亮屏(息屏超时→最大,原值 "
f"{_SCREEN_TIMEOUT_BACKUP.get(self.serial)}ms)")
except Exception as e:
_log.warning(f"[{self.serial}] keep_screen({mode}) 异常: {e}")
def _exec_stop_app(self, d, params, depth=0):
"""强制结束 App(am force-stop)。清后台不留进程,下次打开为冷启动。"""
package = params.get("package", "")
if not package:
_log.warning(f"[{self.serial}] stop_app 缺少 package 参数")
return
try:
d.app_stop(package)
_log.info(f"[{self.serial}] 已强制结束 {package}")
except Exception as e:
_log.warning(f"[{self.serial}] 结束 {package} 异常: {e}")
def _exec_open_app(self, d, params, depth=0):
package = params.get("package", "")
if not package:
_log.warning("open_app 缺少 package 参数")
return
wait_home = params.get("wait_home", False)
home_feature = params.get("home_feature", "")
d.app_start(package, wait=True)
if wait_home:
home_check = (lambda d: d(descriptionContains=home_feature).exists) if home_feature else (lambda d: True)
if not wait_for_app_home(d, package, home_check, timeout=40):
self.set_action(f"{package} 首页加载超时,继续")
def _exec_key_event(self, d, params, depth=0):
"""按键:返回/Home/回车等(d.press)。退出评论、返回上一页必备。"""
key = params.get("key", "back")
try:
d.press(key)
_log.info(f"[{self.serial}] 按键: {key}")
except Exception as e:
_log.warning(f"[{self.serial}] 按键 {key} 异常: {e}")
def _exec_swipe_until(self, d, params, depth=0):
"""滑动直到元素出现(最多 max_swipes 次),可选找到后点击。
养号核心动作:信息流刷到目标按钮(点赞/评论)再操作。
返回 True=找到, False=未找到, None=配置缺选择器。
"""
sel_type = params.get("selector_type", "xpath")
sel_val = params.get("selector_value", "")
direction = params.get("direction", "up")
max_swipes = max(1, int(params.get("max_swipes", 8)))
click_when_found = params.get("click_when_found", True)
if not sel_val:
_log.warning("swipe_until 缺少选择器")
return None
info = d.info
w, h = info["displayWidth"], info["displayHeight"]
cx = int(w * 0.5)
if direction == "down":
sy, ey = int(h * 0.2), int(h * 0.8)
else:
sy, ey = int(h * 0.8), int(h * 0.2)
for i in range(1, max_swipes + 1):
if self.stopped() or self.is_time_up():
return None
found = False
try:
if sel_type == "xpath":
el = d.xpath(sel_val)
found = el.wait(timeout=1)
else:
el = d(**{sel_type: sel_val})
found = el.exists(timeout=1)
except Exception:
found = False
if found:
self.set_action(f"滑动找到元素(第{i}次)")
if click_when_found:
try:
el.click()
_log.info(f"[{self.serial}] 滑动{i}次后找到并点击: {sel_val}")
except Exception as e:
_log.warning(f"[{self.serial}] 找到但点击失败: {e}")
else:
_log.info(f"[{self.serial}] 滑动{i}次后找到: {sel_val}")
self._track_selector_health(sel_val, True)
return True
d.swipe(cx, sy, cx, ey, 0.3)
time.sleep(0.5)
self._track_selector_health(sel_val, False)
_log.warning(f"[{self.serial}] 滑动{max_swipes}次未找到元素: {sel_val}")
return False
def _exec_click_xy(self, d, params, depth=0):
"""点击坐标(屏幕百分比 0-100,50/50=屏幕中心)。无选择器时兜底。"""
x_pct = float(params.get("x", 50))
y_pct = float(params.get("y", 50))
info = d.info
w, h = info["displayWidth"], info["displayHeight"]
x = int(w * min(max(x_pct, 0), 100) / 100)
y = int(h * min(max(y_pct, 0), 100) / 100)
d.click(x, y)
_log.info(f"[{self.serial}] 点击坐标 ({x},{y})")
def _exec_long_click(self, d, params, depth=0):
"""长按元素(选择器 + 时长秒)。复制链接/呼出菜单用。"""
sel_type = params.get("selector_type", "xpath")
sel_val = params.get("selector_value", "")
duration = float(params.get("duration", 1.0))
timeout = float(params.get("wait_timeout", 2))
if not sel_val:
_log.warning("long_click 缺少选择器")
return None
found = False
try:
if sel_type == "xpath":
el = d.xpath(sel_val)
if el.wait(timeout=timeout):
el.long_click(duration=duration)
found = True
else:
el = d(**{sel_type: sel_val})
if el.exists(timeout=timeout):
el.long_click(duration=duration)
found = True
self._track_selector_health(sel_val, found)
_log.info(f"[{self.serial}] 长按 {'命中' if found else '未找到'}: {sel_val}")
except Exception as e:
_log.warning(f"[{self.serial}] 长按异常: {e}")
return found
def _exec_wait_el(self, d, params, depth=0):
"""等待元素出现(条件等待)。返回 True=出现, False=超时, None=缺选择器。"""
sel_type = params.get("selector_type", "xpath")
sel_val = params.get("selector_value", "")
timeout = float(params.get("timeout", 10))
if not sel_val:
_log.warning("wait_el 缺少选择器")
return None
try:
if sel_type == "xpath":
found = d.xpath(sel_val).wait(timeout=timeout)
else:
found = d(**{sel_type: sel_val}).exists(timeout=timeout)
self._track_selector_health(sel_val, found)
_log.info(f"[{self.serial}] 等待元素 {'出现' if found else '超时未出现'}: {sel_val}")
return found
except Exception as e:
_log.warning(f"[{self.serial}] 等待元素异常: {e}")
return False
def _exec_if_el(self, d, params, depth=0):
"""条件判断:找到元素 → 执行 then 分支;超时未找到 → 执行 else 分支。
selector_type=ocr 时用屏幕 OCR 匹配文字(图片/画布/WebView 里的文字也能找到),
命中且 ocr_click=True 时自动点击该文字中心。
分支里可放任意步骤(含循环/嵌套条件判断),实现"有弹窗就关掉、没弹窗就继续"这类逻辑。
返回 True=命中 / False=未命中 / None=缺选择器或 OCR 不可用。
"""
sel_type = params.get("selector_type", "xpath")
sel_val = params.get("selector_value", "")
timeout = float(params.get("timeout", 3))
if not sel_val:
_log.warning(f"[{self.serial}] if_el 缺少选择器,跳过")
return None
found = False
if sel_type == "ocr":
# OCR 模式:截图 → 识别文字 → 关键词匹配(UI 树里没有的文字也能找到)
try:
from core.ocr import available as _ocr_available, find_on_screen
ok, hint = _ocr_available()
if not ok:
_log.error(f"[{self.serial}] OCR 条件判断不可用: {hint}")
return None
img = d.screenshot()
if img is None:
_log.warning(f"[{self.serial}] OCR 截图失败")
else:
found, center, matched = find_on_screen(img, sel_val)
if found:
if params.get("ocr_click") and center:
d.click(*center)
_log.info(f"[{self.serial}] OCR 命中 '{matched}',已点击 ({center[0]},{center[1]})")
else:
_log.info(f"[{self.serial}] OCR 命中 '{matched}'")
else:
_log.info(f"[{self.serial}] OCR 未命中: '{sel_val}'")
except Exception as e:
_log.warning(f"[{self.serial}] OCR 条件判断异常: {e}")
found = False
else:
try:
if sel_type == "xpath":
found = d.xpath(sel_val).wait(timeout=timeout)
else:
found = d(**{sel_type: sel_val}).exists(timeout=timeout)
except Exception as e:
_log.warning(f"[{self.serial}] 条件判断异常: {e}")
found = False
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)} 步")
self._exec_steps(d, branch, depth + 1)
return found
def _exec_swipe(self, d, params, depth=0):
direction = params.get("direction", "up")
dur_min = float(params.get("duration_min", 0.25))
dur_max = float(params.get("duration_max", 0.50))
dur = random.uniform(dur_min, dur_max)
info = d.info
w, h = info["displayWidth"], info["displayHeight"]
cx = int(w * 0.5)
if direction == "up":
d.swipe(cx, int(h * 0.8), cx, int(h * 0.2), dur)
elif direction == "down":
d.swipe(cx, int(h * 0.2), cx, int(h * 0.8), dur)
elif direction == "left":
d.swipe(int(w * 0.8), int(h * 0.5), int(w * 0.2), int(h * 0.5), dur)
elif direction == "right":
d.swipe(int(w * 0.2), int(h * 0.5), int(w * 0.8), int(h * 0.5), dur)
def _exec_click(self, d, params, depth=0):
sel_type = params.get("selector_type", "xpath")
sel_val = params.get("selector_value", "")
timeout = float(params.get("wait_timeout", 2))
if not sel_val:
_log.warning(f"[{self.serial}] click 跳过:selector_value 为空,请在步骤编辑器中填写或抓取元素")
return
# 兼容旧配置:selector_type 不是 xpath 但值是 XPath(以 // 开头)
if sel_type != "xpath" and sel_val.strip().startswith("//"):
sel_type = "xpath"
found = False
try:
if sel_type == "xpath":
# u2 的 d.xpath() 返回 XPathSelector,exists 是属性需用 wait()
el = d.xpath(sel_val)
if el.wait(timeout=timeout):
el.click()
_log.info(f"[{self.serial}] click 命中: xpath={sel_val}")
found = True
else:
_log.warning(f"[{self.serial}] click 未找到元素: xpath={sel_val} (等待{timeout}s)")
else:
# UiSelector 风格:d(description=xxx) 等,exists() 支持 timeout 参数
el = d(**{sel_type: sel_val})
if el.exists(timeout=timeout):
el.click()
_log.info(f"[{self.serial}] click 命中: {sel_type}={sel_val}")
found = True
else:
_log.warning(f"[{self.serial}] click 未找到元素: {sel_type}={sel_val} (等待{timeout}s)")
self._track_selector_health(sel_val, found)
except Exception as e:
_log.warning(f"[{self.serial}] click 异常: {e}")
return found # 命中结果,"测试此步骤"用
def _exec_clipboard(self, d, params, depth=0):
"""剪贴板注入(ClipInject 通道,Android 10+ 可用),可选立即粘贴。
批量发链接场景:先聚焦评论区输入框 → 注入+粘贴链接1 → 换链接 → 注入+粘贴链接2……
粘贴模式:ClipInject 写入剪贴板后调 atx-agent pasteClipboard 粘贴到当前
焦点输入框;动作失败时兜底 u2 set_text 直接设置聚焦输入框文本(不依赖
剪贴板,与"输入文字"兜底同机制)。
"""
text = str(params.get("text", "") or "")
if not text:
_log.warning(f"[{self.serial}] clipboard 内容为空,跳过")
return False
from core.clipboard_helper import inject_clipboard
ok, msg = inject_clipboard(self.serial, text, d=d)
if not ok:
_log.warning(f"[{self.serial}] 剪贴板注入失败: {msg}")
return False
_log.info(f"[{self.serial}] 剪贴板注入: {text[:60]}({msg})")
if not params.get("paste"):
return True
try:
d.jsonrpc.pasteClipboard()
_log.info(f"[{self.serial}] 已粘贴: {text[:60]}")
return True
except Exception as e:
try:
d(focused=True).set_text(text)
_log.info(f"[{self.serial}] 粘贴兜底 set_text: {text[:60]}")
return True
except Exception as e2:
_log.warning(f"[{self.serial}] 粘贴失败: {e} / {e2}")
return True # 剪贴板已写入,粘贴失败不致命
def _exec_wait(self, d, params, depth=0):
min_s = float(params.get("min", 1.0))
max_s = float(params.get("max", 3.0))
# 分片 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 == "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:
_log.warning(f"[{self.serial}] 按时间循环未设置时长,跳过")
return
start = time.time()
i = 0
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):
"""动作组:按序执行 children 一次(类似 loop 但不循环)。
用于把一系列步骤打包成一个可复用的动作。
"""
children = params.get("children", [])
if not children:
return
self._exec_steps(d, children, depth + 1)
def _exec_input_text(self, d, params, depth=0):
"""在当前焦点输入框输入文字。
需先用"点击元素"步骤定位到输入框,本步骤只负责输入。
mode=random: 从 texts(换行分隔)随机选一条
mode=fixed: 输入 fixed_text
"""
mode = params.get("mode", "random")
if mode == "fixed":
text = params.get("fixed_text", "")
else:
raw = params.get("texts", "")
candidates = [t.strip() for t in raw.splitlines() if t.strip()]
if not candidates:
_log.warning("input_text 无候选文字")
return
text = random.choice(candidates)
if not text:
return
clear_first = params.get("clear_first", True)
try:
# 优先用 u2 的 send_keys(对当前焦点元素输入);clear=True 先清空再输入
d.send_keys(text, clear=clear_first)
_log.info(f"[{self.serial}] input_text: {text}" + ("(已先清空)" if clear_first else ""))
except Exception as e:
# 兜底:找 EditText 设置文本
try:
edit = d(className="android.widget.EditText")
if edit.exists:
if clear_first:
edit.clear_text()
edit.set_text(text)
_log.info(f"[{self.serial}] input_text(set_text): {text}")
else:
_log.warning(f"[{self.serial}] input_text 未找到输入框: {e}")
except Exception as e2:
_log.warning(f"[{self.serial}] input_text 失败: {e2}")
@register_task
class GenericStepsTask(BaseTask):
"""通用步骤任务:前端编辑器编排步骤链,worker 按顺序执行。"""
task_type = "generic_steps"
name = "通用步骤"
description = "通过步骤编辑器编排任务流程,支持循环/点击/滑动/点赞/评论等"
default_params = dict(DEFAULT_PARAMS)
@classmethod
def list_action_types(cls):
"""返回支持的步骤类型列表(前端操作库用)。"""
return STEP_TYPES
@classmethod
def get_action_class(cls, action_type):
return None
def create_worker(self, serial, params):
merged = {**DEFAULT_PARAMS, **(params or {})}
return GenericStepsWorker(serial, params=merged)
def test_step(serial, step):
"""在设备上单步试执行(编辑器"测试此步骤"按钮用)。
直接 adb connect + u2 连接后执行单个步骤,验证选择器是否命中。
返回 (ok, msg, result):result 为执行器返回的命中结果(True/False/None)。
与运行中的任务互不干扰(只读连接,不占用/不释放 STF)。
"""
import uiautomator2 as u2
from concurrent.futures import ThreadPoolExecutor # 与 device_worker 一致(本机 stdlib threading 无此属性)
from core.adb_helper import adb_connect
from core.device_worker import _U2_CONNECT_TIMEOUT
try:
if not adb_connect(serial):
return False, f"adb connect {serial} 失败", None
except Exception as e:
return False, f"adb connect 异常: {e}", None
try:
with ThreadPoolExecutor(max_workers=1) as pool:
fut = pool.submit(u2.connect, serial)
d = fut.result(timeout=_U2_CONNECT_TIMEOUT)
except Exception as e:
return False, f"u2 连接失败: {e}", None
if d is None:
return False, "u2 连接超时", None
try:
w = GenericStepsWorker(serial, {})
result = w._exec_one(d, step)
return True, f"步骤已执行({step.get('label') or step.get('type')})", result
except Exception as e:
return False, f"执行异常: {e}", None