1046 lines
38 KiB
Markdown
1046 lines
38 KiB
Markdown
# 任务开发指南
|
||
|
||
面向 `platform-tools`(STF + uiautomator2 + Flask 单页应用,多设备并发任务执行框架)的新开发者。描述架构、核心概念,并给出从 0 到 1 新增一个 app 任务所需的全部模板与规范。看完本文即可上手开发新任务。
|
||
|
||
---
|
||
|
||
## 1. 架构总览
|
||
|
||
### 1.1 分层设计
|
||
|
||
平台按"配置 / 核心 / 任务 / 前端 / 数据 / 日志 / 工具"分层,职责清晰、互不交叉:
|
||
|
||
| 层 | 路径 | 职责 |
|
||
| --- | --- | --- |
|
||
| 配置层 | `config.py` | 项目根配置:STF 服务地址 / adb 路径 / web 端口等基础设施。**不放任务参数**(任务参数属于 `tasks/`) |
|
||
| 核心层 | `core/` | 框架运行时:`logger` 日志、`stf_client` STF API 封装、`adb_helper` adb 操作、`device_worker` Worker 基类、`task_manager` 调度器、`u2_helper` uiautomator2 通用操作、`actions` 全局 Action 基类 |
|
||
| 任务层 | `tasks/` | 每个 app 一个子包,自包含 `task.py` + `actions/`,互不依赖 |
|
||
| 前端层 | `templates/admin/monitor.html` | 单页应用(纯 HTML+CSS+JS,无框架):设备监控 / 任务管理 / 分组 / 日志 / 用户 5 个 Tab |
|
||
| 数据层 | `data/` | SQLite 持久化:`users.db`(用户 + 设备分组 + 任务计划) |
|
||
| 日志层 | `logs/` | 四类日志:`core.log` / `task.log` / `web.log` / `action.log`,10MB 滚动保留 5 份 |
|
||
| 文档层 | `doc/` | 项目文档 |
|
||
| 工具层 | `bin/adb/` | adb 可执行文件 |
|
||
| 脚本层 | `scripts/` | 实用脚本 |
|
||
|
||
### 1.2 目录树
|
||
|
||
```
|
||
platform-tools/
|
||
├── config.py # 根配置(STF/adb/web,不放任务参数)
|
||
├── web_server.py # Flask 入口(JSON API + 登录页 + 单页应用)
|
||
├── main.py # 根目录入口
|
||
├── core/ # 核心程序层
|
||
│ ├── __init__.py
|
||
│ ├── logger.py # 日志器(分文件、10MB 滚动)
|
||
│ ├── stf_client.py # STF API 封装
|
||
│ ├── adb_helper.py # adb 操作(全局锁串行化)
|
||
│ ├── device_worker.py # BaseWorker 基类 + STFDevice + 看门狗
|
||
│ ├── task_manager.py # TaskManager 调度器 + 前台 App 扫描器
|
||
│ ├── u2_helper.py # uiautomator2 通用操作(ensure_app_running/wait_for_app_home)
|
||
│ ├── models.py # SQLAlchemy 模型(User/DeviceGroup/TaskJob)
|
||
│ └── actions/
|
||
│ ├── __init__.py # create_action_registry / register_action / should_trigger
|
||
│ └── base.py # BaseAction 全局基类
|
||
├── tasks/ # 任务定义层
|
||
│ ├── __init__.py # 全局 _TASK_TYPES 注册表
|
||
│ ├── base.py # BaseTask 基类
|
||
│ └── douyin/ # 抖音养号(示例)
|
||
│ ├── __init__.py
|
||
│ ├── task.py # DEFAULT_PARAMS + Worker + Task + @register_task
|
||
│ └── actions/
|
||
│ ├── __init__.py # 先 from .base import ACTIONS,再 from . import like, comment
|
||
│ ├── base.py # ACTIONS = create_action_registry()
|
||
│ ├── like.py # @register_action(ACTIONS) LikeAction
|
||
│ └── comment.py # @register_action(ACTIONS) CommentAction
|
||
├── templates/admin/
|
||
│ ├── monitor.html # 单页应用(5 Tab,纯前端渲染)
|
||
│ └── login.html # 登录页
|
||
├── data/ # 持久化数据
|
||
│ └── users.db # SQLite(用户/分组/任务)
|
||
├── logs/ # 日志(10MB 滚动保留 5 份)
|
||
├── doc/ # 文档
|
||
├── bin/adb/ # adb 工具
|
||
└── scripts/ # 实用脚本
|
||
```
|
||
|
||
### 1.3 数据流
|
||
|
||
```
|
||
┌──────────────┐ 创建 Job ┌─────────────┐ 分发 ┌──────────────┐
|
||
│ 单页应用前端 │ ───────────► │ TaskManager │ ──────► │ Worker(设备) │
|
||
│ (monitor.html│ └─────────────┘ └──────────────┘
|
||
│ fetch + DOM)│ ▲ │
|
||
└──────────────┘ │ 心跳/状态 │ u2 操作
|
||
│ │ ▼
|
||
│ JSON API │ ┌────────────────┐
|
||
▼ │ │ STF Device / adb│
|
||
┌──────────────┐ ┌──────────────┐ └────────────────┘
|
||
│ web_server │ │ 看门狗监控 │
|
||
│ (Flask API) │ └──────────────┘
|
||
└──────────────┘
|
||
│
|
||
▼
|
||
┌──────────────┐
|
||
│ data/users.db│ SQLite 持久化(用户/分组/任务)
|
||
└──────────────┘
|
||
```
|
||
|
||
### 1.4 关键设计决策
|
||
|
||
- **Flask + Flask-Login**:已移除 Flask-Admin(自定义场景下过于受限),改用纯 Flask + 单页应用
|
||
- **单页应用**:`web_server.py` 只提供 JSON API + 登录页,`monitor.html` 纯前端渲染(fetch + DOM 操作),无服务端模板依赖
|
||
- **SQLite 持久化**:替代旧 JSON 文件,支持用户/分组/任务的关系存储
|
||
- **多线程模式**:Flask 启用 `threaded=True` 处理并发请求
|
||
- **设备状态缓存**:`get_status` 带 5 秒缓存,worker 运行状态实时组装
|
||
|
||
---
|
||
|
||
## 2. 核心概念
|
||
|
||
### 2.1 TaskType — 任务类型
|
||
|
||
一个 `TaskType` 描述"做什么"(例如抖音养号、快手养号),由 `Task` 子类 + `Worker` 子类 + `DEFAULT_PARAMS` 组成。每个 `TaskType` 注册到全局 `_TASK_TYPES` 字典(在 `tasks/__init__.py`),key 为任务类型字符串,value 为 `Task` 类。
|
||
|
||
```python
|
||
# tasks/base.py
|
||
_TASK_TYPES = {}
|
||
|
||
def register_task(task_cls):
|
||
"""任务类型注册装饰器(无需传 name,用 task_cls.task_type)"""
|
||
_TASK_TYPES[task_cls.task_type] = task_cls
|
||
return task_cls
|
||
|
||
def get_task_class(task_type):
|
||
return _TASK_TYPES.get(task_type)
|
||
```
|
||
|
||
### 2.2 TaskJob — 任务计划
|
||
|
||
`TaskJob` 是"什么时候、在哪些设备上、用什么参数执行某个 TaskType"的持久化计划,存于 SQLite(`data/users.db` 的 `task_job` 表)。包含字段:
|
||
|
||
- `task_type` — 任务类型(对应 `_TASK_TYPES` 的 key)
|
||
- `target` — 目标设备:`{"mode": "all"|"group"|"serial", "group_name": "", "serial": ""}`
|
||
- `params` — 任务参数(与 `DEFAULT_PARAMS` 深合并)
|
||
- `schedule` — 调度策略:`{"mode": "once"|"cron", "cron": "0 9 * * *"}`
|
||
- `retry` — 重试策略:`{"max_attempts": 1, "delay": 60}`
|
||
- `enabled` — 是否启用
|
||
|
||
### 2.3 DeviceGroup — 设备分组
|
||
|
||
设备分组存于 SQLite(`device_group` 表),便于按批次/项目/客户分组下发任务。一个 Job 可指定 `target.mode="group"`,调度器展开为组内全部设备序列号。
|
||
|
||
### 2.4 Worker — 单设备执行线程
|
||
|
||
每个被调度的设备对应一个 `Worker` 实例,跑在独立线程中,继承 `BaseWorker`(`core/device_worker.py`)。Worker 负责一台设备的完整生命周期:申请设备 → 连接 u2 → setup → run_task → teardown → 释放设备。
|
||
|
||
### 2.5 Action — 操作
|
||
|
||
`Action` 是任务循环里执行的"原子操作"(点赞 / 评论 / 滑动 / 关注)。每个 app 有**独立的 Action 注册表**(通过 `create_action_registry()` 创建),互不污染。全局基类 `core/actions/base.py::BaseAction` 提供通用能力。
|
||
|
||
### 2.6 进度上报(通用,适配任意 app)
|
||
|
||
Worker 通过 `self.set_progress(**fields)` 上报进度,前端统一解析展示。**不再硬编码"已看视频数"等业务字段**。
|
||
|
||
**通用字段**:
|
||
|
||
| 字段 | 类型 | 说明 |
|
||
| --- | --- | --- |
|
||
| `done` | int | 已完成数量 |
|
||
| `total` | int | 总数量(0=不限数量,只显示已完成数) |
|
||
| `unit` | str | 计数单位("视频"/"轮次"/"条") |
|
||
| `action_counts` | dict | 操作计数 `{"like": 3, "comment": 1}` |
|
||
| `elapsed` | int | 已运行时长(秒,可选,前端显示为 "Xm Ys") |
|
||
|
||
**前端展示**:进度条(百分比,total>0 时)+ "done/total unit" + 运行时长 + 操作计数徽章
|
||
|
||
**示例**:
|
||
```python
|
||
# 抖音任务(有数量限制)
|
||
self.set_progress(done=5, total=80, unit="视频",
|
||
action_counts={"like": 3}, elapsed=120)
|
||
|
||
# 抖音任务(仅时长限制,无数量)
|
||
self.set_progress(done=5, total=0, unit="视频",
|
||
action_counts={"like": 3}, elapsed=120)
|
||
|
||
# 快手任务
|
||
self.set_progress(done=3, total=20, unit="轮次",
|
||
action_counts={"like": 2}, elapsed=60)
|
||
```
|
||
|
||
### 2.7 运行时长终止(通用,适配任意 app)
|
||
|
||
`BaseWorker` 提供运行时长终止能力,与"数量终止"配合使用。两者**哪个先到就停**。
|
||
|
||
| 成员 | 说明 |
|
||
| --- | --- |
|
||
| `self.max_duration` | 最大运行时长(秒),0=不限时 |
|
||
| `self._start_timer()` | 子类在 `run_task` 开头调用,启动计时 |
|
||
| `self.is_time_up()` | 是否已达 max_duration(max_duration=0 永远返回 False) |
|
||
| `self.elapsed()` | 已运行时长(秒) |
|
||
|
||
**循环条件模板**:
|
||
```python
|
||
while not self.stopped():
|
||
if watch_count > 0 and watched >= watch_count:
|
||
break # 数量终止
|
||
if self.is_time_up():
|
||
break # 时长终止
|
||
# ... 业务逻辑
|
||
```
|
||
|
||
**三种终止模式**:
|
||
- 仅数量:`watch_count=80, max_duration=0` → 看完 80 个视频停
|
||
- 仅时长:`watch_count=0, max_duration=1800` → 跑满 30 分钟停
|
||
- 双条件:`watch_count=80, max_duration=1800` → 哪个先到就停
|
||
- 都为 0:永不停止,需手动停止
|
||
|
||
### 2.8 任务调度模式(schedule)
|
||
|
||
任务计划 `TaskJob.schedule` 支持三种模式:
|
||
|
||
| mode | 字段 | 行为 |
|
||
| --- | --- | --- |
|
||
| `once` | 无 | 手动执行(前端点"立即执行"或调 `/api/jobs/<id>/run`) |
|
||
| `cron` | `cron` | 定时启动:到 cron 时间点自动启动 worker |
|
||
| `cron_stop` | `cron` + `stop_cron` | 定时启停:启动 cron 到点启动,停止 cron 到点停止本任务的 worker |
|
||
|
||
**cron_stop 模式**只停止**本 job 启动的 worker**,不影响其他正在运行的任务。典型用法:
|
||
```json
|
||
{
|
||
"schedule": {
|
||
"mode": "cron_stop",
|
||
"cron": "0 9 * * *",
|
||
"stop_cron": "0 18 * * *"
|
||
}
|
||
}
|
||
```
|
||
含义:每天 9 点自动启动任务,18 点自动停止。
|
||
|
||
### 2.9 心跳看门狗
|
||
|
||
每个 Worker 在 `set_action` / `set_progress` / `heartbeat` 时更新心跳时间戳。看门狗线程(`_Watchdog`)定期扫描,若超过 `_HEARTBEAT_TIMEOUT=120s` 未更新则判定卡死,标记 error。
|
||
|
||
**长耗时操作必须周期性调用 `self.heartbeat()`**,否则会被误杀。
|
||
|
||
---
|
||
|
||
## 3. 新增一个 app 任务(完整步骤)
|
||
|
||
以"快手养号"为例。完整步骤 6 步,全部代码可直接复制。
|
||
|
||
### 步骤 1:在 `tasks/` 下建 `kuaishou/` 子包
|
||
|
||
```
|
||
tasks/kuaishou/
|
||
├── __init__.py
|
||
├── task.py
|
||
└── actions/
|
||
├── __init__.py
|
||
├── base.py
|
||
└── like.py
|
||
```
|
||
|
||
### 步骤 2:写 `actions/base.py`(本任务的注册表)
|
||
|
||
```python
|
||
# tasks/kuaishou/actions/base.py
|
||
"""快手 Action 注册表。"""
|
||
from core.actions import (
|
||
BaseAction, register_action, create_action_registry,
|
||
list_actions, get_action, should_trigger,
|
||
)
|
||
|
||
# 快手专属操作注册表(独立 dict,不污染其他 app)
|
||
ACTIONS = create_action_registry()
|
||
|
||
|
||
def list_action_types():
|
||
"""返回所有已注册快手操作的元信息(供前端展示)。"""
|
||
return list_actions(ACTIONS)
|
||
|
||
|
||
def get_action_class(action_type):
|
||
"""按 action_type 取快手操作类。"""
|
||
return get_action(ACTIONS, action_type)
|
||
```
|
||
|
||
### 步骤 3:写 `actions/like.py`
|
||
|
||
```python
|
||
# tasks/kuaishou/actions/like.py
|
||
"""快手点赞 Action。"""
|
||
from core.actions import BaseAction, register_action, should_trigger
|
||
from core.logger import get_logger
|
||
from . import ACTIONS # 必须从 __init__ 导入注册表
|
||
|
||
_log = get_logger("action.kuaishou.like")
|
||
|
||
|
||
@register_action(ACTIONS)
|
||
class LikeAction(BaseAction):
|
||
action_type = "like"
|
||
name = "点赞"
|
||
description = "看完视频后随机点赞"
|
||
default_params = {
|
||
"rate": 0.8, # 触发概率 0~1
|
||
"method": "double_tap", # double_tap | heart_icon
|
||
}
|
||
|
||
def execute(self, d, params, worker):
|
||
rate = float(params.get("rate", 0.8))
|
||
if not should_trigger(rate):
|
||
return False
|
||
method = params.get("method", "double_tap")
|
||
try:
|
||
if method == "double_tap":
|
||
info = d.info
|
||
w, h = info["displayWidth"], info["displayHeight"]
|
||
d.double_click(int(w * 0.5), int(h * 0.5))
|
||
else:
|
||
el = d(description="点赞")
|
||
if not el.exists:
|
||
_log.info("未找到点赞按钮")
|
||
return False
|
||
el.click()
|
||
_log.info("点赞成功")
|
||
return True
|
||
except Exception as e:
|
||
_log.warning(f"点赞异常: {e}")
|
||
return False
|
||
```
|
||
|
||
### 步骤 4:写 `actions/__init__.py`(注意循环导入顺序)
|
||
|
||
```python
|
||
# tasks/kuaishou/actions/__init__.py
|
||
"""快手操作注册包。import 触发各操作注册。
|
||
|
||
⚠️ 循环导入坑:必须先从 base 导入 ACTIONS,再导入各操作模块!
|
||
"""
|
||
from .base import (
|
||
BaseAction, register_action, create_action_registry,
|
||
list_actions, get_action, should_trigger,
|
||
ACTIONS, list_action_types, get_action_class,
|
||
)
|
||
|
||
# 再导入各操作模块,触发 @register_action(ACTIONS) 注册
|
||
from . import like # noqa: F401
|
||
# from . import comment # 新增 action 在此 import
|
||
|
||
__all__ = [
|
||
"BaseAction", "register_action", "create_action_registry",
|
||
"list_actions", "get_action", "should_trigger",
|
||
"ACTIONS", "list_action_types", "get_action_class",
|
||
]
|
||
```
|
||
|
||
### 步骤 5:写 `task.py`
|
||
|
||
```python
|
||
# tasks/kuaishou/task.py
|
||
"""快手养号任务定义。
|
||
|
||
本文件自包含所有快手养号参数,不依赖 core 的业务配置。
|
||
快手专属操作(点赞/评论)在 actions/ 子包里,xpath 只适用于快手。
|
||
"""
|
||
import time
|
||
import random
|
||
|
||
from tasks.base import BaseTask, register_task
|
||
from .actions import list_action_types, get_action_class
|
||
from core.device_worker import BaseWorker, _update_status
|
||
from core.u2_helper import ensure_app_running, wait_for_app_home
|
||
from core.logger import get_logger
|
||
|
||
_log = get_logger("task.kuaishou")
|
||
|
||
KUAISHOU_PKG = "com.smile.gifmaker"
|
||
|
||
DEFAULT_PARAMS = {
|
||
"watch_count": 50, # 观看视频数量
|
||
"watch_min": 5.0, # 单个视频最短观看秒数
|
||
"watch_max": 30.0, # 单个视频最长观看秒数
|
||
"swipe_min": 0.25, # 上滑手势最短时长(秒)
|
||
"swipe_max": 0.50, # 上滑手势最长时长(秒)
|
||
"gap_min": 1.0, # 视频间隔最短秒数
|
||
"gap_max": 3.0, # 视频间隔最长秒数
|
||
"actions": {
|
||
"like": {
|
||
"enabled": True,
|
||
"params": {"rate": 0.3, "method": "double_tap"},
|
||
},
|
||
},
|
||
}
|
||
|
||
|
||
class KuaishouWorker(BaseWorker):
|
||
"""快手养号 worker。"""
|
||
|
||
def __init__(self, stf_client, serial, params=None):
|
||
super().__init__(stf_client, serial, params)
|
||
p = {**DEFAULT_PARAMS, **(self.params or {})}
|
||
self.watch_count = int(p["watch_count"])
|
||
self.watch_min = float(p["watch_min"])
|
||
self.watch_max = float(p["watch_max"])
|
||
self.swipe_min = float(p["swipe_min"])
|
||
self.swipe_max = float(p["swipe_max"])
|
||
self.gap_min = float(p["gap_min"])
|
||
self.gap_max = float(p["gap_max"])
|
||
self.actions_cfg = p.get("actions", {})
|
||
self._actions = []
|
||
for atype, cfg in self.actions_cfg.items():
|
||
if not cfg.get("enabled"):
|
||
continue
|
||
cls = get_action_class(atype)
|
||
if cls:
|
||
self._actions.append(cls())
|
||
if self._actions:
|
||
summary = ", ".join(
|
||
f"{a.action_type}(rate={self.actions_cfg.get(a.action_type, {}).get('params', {}).get('rate', '?')})"
|
||
for a in self._actions
|
||
)
|
||
_log.info(f"[{serial}] 启用操作: {summary}")
|
||
else:
|
||
_log.warning(f"[{serial}] 未启用任何操作")
|
||
|
||
def run_task(self, d):
|
||
"""快手养号主逻辑。d 是 u2.Device,已由基类连好。"""
|
||
def is_home(d):
|
||
return (d(descriptionContains="首页").exists
|
||
or d(descriptionContains="拍摄").exists)
|
||
|
||
d.app_start(KUAISHOU_PKG, wait=True)
|
||
if not wait_for_app_home(d, KUAISHOU_PKG, is_home, timeout=40):
|
||
self.set_action("首页加载超时,继续尝试")
|
||
|
||
watched = 0
|
||
action_counts = {a.action_type: 0 for a in self._actions}
|
||
# 初始化通用进度上报(前端会解析 done/total/unit + action_counts)
|
||
self.set_progress(done=0, total=self.watch_count, unit="视频",
|
||
action_counts=action_counts)
|
||
|
||
while not self.stopped() and watched < self.watch_count:
|
||
if not ensure_app_running(d, KUAISHOU_PKG):
|
||
_update_status(self.serial, status="error",
|
||
last_error="快手连续重启失败,放弃该设备")
|
||
return
|
||
|
||
watch = random.uniform(self.watch_min, self.watch_max)
|
||
self.set_action(f"观看视频 {watched+1}/{self.watch_count},{watch:.0f}s")
|
||
time.sleep(watch)
|
||
|
||
# 执行启用的操作
|
||
for action in self._actions:
|
||
if self.stopped():
|
||
break
|
||
cfg = self.actions_cfg.get(action.action_type, {})
|
||
params = {**action.default_params, **cfg.get("params", {})}
|
||
try:
|
||
ok = action.execute(d, params, self)
|
||
_log.info(f"[{self.serial}] 视频{watched+1}: {action.action_type} execute={ok}")
|
||
if ok:
|
||
action_counts[action.action_type] += 1
|
||
time.sleep(random.uniform(0.5, 1.5))
|
||
except Exception as e:
|
||
_log.error(f"[{self.serial}] 视频{watched+1}: 操作 {action.action_type} 异常: {e}")
|
||
|
||
if self.stopped():
|
||
break
|
||
d.swipe(500, 1000, 500, 300, random.uniform(self.swipe_min, self.swipe_max))
|
||
time.sleep(random.uniform(self.gap_min, self.gap_max))
|
||
watched += 1
|
||
# 上报通用进度
|
||
self.set_progress(done=watched, total=self.watch_count, unit="视频",
|
||
action_counts=dict(action_counts))
|
||
|
||
summary = f"完成 {watched} 个视频" + "".join(
|
||
f",{k} {v}次" for k, v in action_counts.items() if v
|
||
)
|
||
_update_status(self.serial, current_action=summary)
|
||
|
||
|
||
@register_task
|
||
class KuaishouTask(BaseTask):
|
||
"""快手养号任务。"""
|
||
task_type = "kuaishou_nurture"
|
||
name = "快手养号"
|
||
description = "自动观看快手视频,按配置执行点赞等操作"
|
||
default_params = dict(DEFAULT_PARAMS)
|
||
|
||
@classmethod
|
||
def list_action_types(cls):
|
||
return list_action_types()
|
||
|
||
@classmethod
|
||
def get_action_class(cls, action_type):
|
||
return get_action_class(action_type)
|
||
|
||
def create_worker(self, stf, serial, params):
|
||
merged = {**DEFAULT_PARAMS, **(params or {})}
|
||
# actions 字段参数级深合并(保留前端没传的操作默认值)
|
||
default_actions = DEFAULT_PARAMS["actions"]
|
||
merged_actions = params.get("actions", {}) if params else {}
|
||
for atype, dflt in default_actions.items():
|
||
if atype not in merged_actions:
|
||
merged_actions[atype] = dflt
|
||
else:
|
||
cfg = merged_actions[atype]
|
||
merged_cfg = {}
|
||
for k in ("enabled", "params"):
|
||
merged_cfg[k] = cfg.get(k, dflt.get(k))
|
||
merged_params = dict(dflt.get("params", {}))
|
||
merged_params.update(cfg.get("params", {}))
|
||
merged_cfg["params"] = merged_params
|
||
merged_actions[atype] = merged_cfg
|
||
merged["actions"] = merged_actions
|
||
return KuaishouWorker(stf, serial, params=merged)
|
||
```
|
||
|
||
### 步骤 6:注册任务包
|
||
|
||
`tasks/kuaishou/__init__.py`:
|
||
```python
|
||
# tasks/kuaishou/__init__.py
|
||
from . import task # noqa: F401 触发 @register_task 注册
|
||
```
|
||
|
||
`tasks/__init__.py` 加一行:
|
||
```python
|
||
# tasks/__init__.py
|
||
from .base import BaseTask, register_task, list_task_types, get_task_class
|
||
from . import douyin # noqa: F401
|
||
from . import kuaishou # noqa: F401 ← 新增这一行
|
||
```
|
||
|
||
完成。重启 web 后,前端任务类型下拉自动出现 `kuaishou_nurture`。
|
||
|
||
---
|
||
|
||
## 4. Worker 开发指南
|
||
|
||
`BaseWorker` 位于 `core/device_worker.py`,封装了设备生命周期、心跳、异常分类、与调度器的状态通信。
|
||
|
||
### 4.1 生命周期
|
||
|
||
```
|
||
acquire(serial) # 向 STF 申请设备占用
|
||
│
|
||
▼
|
||
adb connect + u2.connect # 连接 uiautomator2(带 30s 超时保护)
|
||
│
|
||
▼
|
||
setup(d) # 子类可选钩子(启动 app、授权、关闭弹窗)
|
||
│
|
||
▼
|
||
run_task(d) ◄── 必须实现 # 任务主循环
|
||
│
|
||
▼
|
||
teardown(d) # 子类可选钩子(退出 app、清理)
|
||
│
|
||
▼
|
||
release(serial) # 释放 STF 占用
|
||
```
|
||
|
||
任意阶段抛出 `DeviceOfflineError` → 立即终止,**不重试**。其他异常 → 按 Job 的 `retry` 策略重试。
|
||
|
||
### 4.2 必须实现 / 可选钩子
|
||
|
||
| 方法 | 是否必须 | 说明 |
|
||
| --- | --- | --- |
|
||
| `run_task(self, d)` | **必须** | 任务主循环,`d` 为 `uiautomator2.Device` |
|
||
| `setup(self, d)` | 可选 | 设备/应用初始化 |
|
||
| `teardown(self, d)` | 可选 | 收尾,即使出错也会执行 |
|
||
| `on_error(self, d, err)` | 可选 | 异常通知钩子 |
|
||
|
||
### 4.3 工具方法
|
||
|
||
| 方法 | 说明 |
|
||
| --- | --- |
|
||
| `self.stopped()` | **循环里必须检查**,返回 True 表示收到停止信号 |
|
||
| `self.set_action(s)` | 设置当前动作(前端大屏"当前动作"列可见),同时刷新心跳 |
|
||
| `self.set_progress(**fields)` | 上报进度(见 §2.6),同时刷新心跳 |
|
||
| `self.heartbeat()` | 手动刷新心跳(长操作中间调) |
|
||
| `self.params` | 已合并 `DEFAULT_PARAMS` 与 Job 参数后的最终参数 |
|
||
| `self.serial` | 当前设备序列号 |
|
||
| `self.d` | u2.Device(run_task 的 d 参数) |
|
||
|
||
### 4.4 进度上报规范(重要)
|
||
|
||
**所有 app 任务必须用 `set_progress` 上报进度**,前端会统一解析展示。
|
||
|
||
```python
|
||
# ✅ 正确:用通用字段
|
||
self.set_progress(done=5, total=80, unit="视频",
|
||
action_counts={"like": 3, "comment": 1})
|
||
|
||
# ❌ 错误:硬编码业务字段(前端无法识别)
|
||
self.set_progress(videos_watched=5) # 前端不认这个字段
|
||
```
|
||
|
||
前端展示效果:
|
||
- 进度条:`████████░░░░` (按 done/total 算百分比)
|
||
- 计数文本:`5/80 视频`
|
||
- 操作徽章:`点赞 3` `评论 1`
|
||
|
||
### 4.5 异常分类
|
||
|
||
| 异常 | 处理 |
|
||
| --- | --- |
|
||
| `DeviceOfflineError` | 设备掉线,**不重试**,立即释放 |
|
||
| 其他 `Exception` | 按 Job 的 `retry` 次数重试,退避后重新申请设备 |
|
||
|
||
### 4.6 `run_task` 模板(可直接复制)
|
||
|
||
```python
|
||
def run_task(self, d):
|
||
"""任务主循环模板。"""
|
||
# 1. 启动 app
|
||
d.app_start("com.xxx", wait=True)
|
||
if not wait_for_app_home(d, "com.xxx", lambda d: d(text="首页").exists, timeout=40):
|
||
self.set_action("首页加载超时,继续尝试")
|
||
|
||
# 2. 初始化进度上报
|
||
watched = 0
|
||
action_counts = {a.action_type: 0 for a in self._actions}
|
||
self.set_progress(done=0, total=self.watch_count, unit="视频",
|
||
action_counts=action_counts)
|
||
|
||
# 3. 主循环
|
||
while not self.stopped() and watched < self.watch_count:
|
||
# 3.1 确保 app 在前台
|
||
if not ensure_app_running(d, "com.xxx"):
|
||
_update_status(self.serial, status="error",
|
||
last_error="app 连续重启失败")
|
||
return
|
||
|
||
# 3.2 观看
|
||
watch = random.uniform(self.watch_min, self.watch_max)
|
||
self.set_action(f"观看 {watched+1}/{self.watch_count},{watch:.0f}s")
|
||
time.sleep(watch)
|
||
|
||
# 3.3 执行操作
|
||
for action in self._actions:
|
||
if self.stopped():
|
||
break
|
||
cfg = self.actions_cfg.get(action.action_type, {})
|
||
params = {**action.default_params, **cfg.get("params", {})}
|
||
try:
|
||
if action.execute(d, params, self):
|
||
action_counts[action.action_type] += 1
|
||
except Exception as e:
|
||
_log.error(f"[{self.serial}] 操作异常: {e}")
|
||
|
||
# 3.4 滑动
|
||
if self.stopped():
|
||
break
|
||
d.swipe(500, 1000, 500, 300, 0.3)
|
||
time.sleep(random.uniform(1.0, 3.0))
|
||
|
||
# 3.5 上报进度
|
||
watched += 1
|
||
self.set_progress(done=watched, total=self.watch_count, unit="视频",
|
||
action_counts=dict(action_counts))
|
||
|
||
# 4. 收尾
|
||
summary = f"完成 {watched} 个视频"
|
||
_update_status(self.serial, current_action=summary)
|
||
```
|
||
|
||
> **铁律**:循环里必须高频调用 `self.stopped()`,否则停止按钮无响应、看门狗误杀。
|
||
|
||
---
|
||
|
||
## 5. Action 开发指南
|
||
|
||
### 5.1 全局基类与独立注册表
|
||
|
||
- 全局基类:`core/actions/base.py::BaseAction`,提供 `should_trigger` 等通用能力
|
||
- 每个 app 通过 `create_action_registry()` 创建**独立注册表**,避免不同 app 的 `like` / `comment` 同名冲突
|
||
- 注册装饰器:`@register_action(ACTIONS)`,`ACTIONS` 为本 app 的注册表
|
||
|
||
### 5.2 `BaseAction` 关键 API
|
||
|
||
| 成员 | 说明 |
|
||
| --- | --- |
|
||
| `action_type` | 类属性,注册 key,必须与 `params["actions"]` 的 key 一致 |
|
||
| `name` | 类属性,中文名(前端展示) |
|
||
| `description` | 类属性,描述 |
|
||
| `default_params` | 类属性,自包含默认参数 |
|
||
| `execute(self, d, params, worker)` | **必须实现**,返回 `True`=成功 / `False`=跳过 |
|
||
| `should_trigger(rate)` | 按 `rate` 概率返回是否触发(`rate=0.8` → 80% 概率 True) |
|
||
|
||
### 5.3 完整 Action 模板
|
||
|
||
```python
|
||
# tasks/xxx/actions/comment.py
|
||
import time
|
||
import random
|
||
|
||
from core.actions import BaseAction, register_action, should_trigger
|
||
from core.logger import get_logger
|
||
from . import ACTIONS # 从 __init__ 导入本 app 注册表
|
||
|
||
_log = get_logger("action.xxx.comment")
|
||
|
||
|
||
@register_action(ACTIONS)
|
||
class CommentAction(BaseAction):
|
||
action_type = "comment"
|
||
name = "评论"
|
||
description = "看完视频后随机发一条评论"
|
||
default_params = {
|
||
"rate": 0.1,
|
||
"texts": ["不错", "666", "学到了"],
|
||
}
|
||
|
||
def execute(self, d, params, worker):
|
||
rate = float(params.get("rate", 0.1))
|
||
if not should_trigger(rate):
|
||
return False
|
||
|
||
texts = params.get("texts") or ["不错"]
|
||
text = random.choice(texts)
|
||
|
||
# 定位评论按钮(多策略组合,失败回退)
|
||
for desc in ("评论", "未评论", "comment"):
|
||
el = d(description=desc)
|
||
if el.exists:
|
||
el.click()
|
||
break
|
||
else:
|
||
_log.info("未找到评论按钮")
|
||
return False
|
||
|
||
time.sleep(1.5)
|
||
|
||
# 找输入框
|
||
for rid in ("com.xxx:id/comment_input",):
|
||
el = d(resourceId=rid)
|
||
if el.exists:
|
||
el.click()
|
||
break
|
||
else:
|
||
d.press("back")
|
||
return False
|
||
|
||
time.sleep(0.8)
|
||
|
||
# 输入文本(中文需切输入法)
|
||
try:
|
||
d.set_fastinput_ime(True)
|
||
except Exception:
|
||
pass
|
||
d.send_keys(text)
|
||
time.sleep(0.6)
|
||
|
||
# 发送
|
||
el = d(text="发送")
|
||
if el.exists:
|
||
el.click()
|
||
_log.info(f"评论成功: {text}")
|
||
time.sleep(1.0)
|
||
d.press("back")
|
||
return True
|
||
d.press("back")
|
||
return False
|
||
```
|
||
|
||
> **返回值约定**:`True`=成功执行;`False`=主动跳过(概率未中、元素不存在等);抛异常=执行失败,由 Worker 捕获并记录。
|
||
|
||
---
|
||
|
||
## 6. 参数设计规范
|
||
|
||
### 6.1 自包含
|
||
|
||
`DEFAULT_PARAMS` 放在 `task.py` 顶部,**所有**该任务需要的参数都要列出,包括每个 action 的子参数。不允许"隐式默认值"散落在 action 内部。
|
||
|
||
### 6.2 参数合并(参数级深合并)
|
||
|
||
`create_worker` 时执行三层合并:
|
||
|
||
```python
|
||
def create_worker(self, stf, serial, params):
|
||
merged = {**DEFAULT_PARAMS, **(params or {})}
|
||
# actions 字段参数级深合并
|
||
default_actions = DEFAULT_PARAMS["actions"]
|
||
merged_actions = params.get("actions", {}) if params else {}
|
||
for atype, dflt in default_actions.items():
|
||
if atype not in merged_actions:
|
||
merged_actions[atype] = dflt
|
||
else:
|
||
cfg = merged_actions[atype]
|
||
merged_cfg = {}
|
||
for k in ("enabled", "params"):
|
||
merged_cfg[k] = cfg.get(k, dflt.get(k))
|
||
# params 再深合并一层
|
||
merged_params = dict(dflt.get("params", {}))
|
||
merged_params.update(cfg.get("params", {}))
|
||
merged_cfg["params"] = merged_params
|
||
merged_actions[atype] = merged_cfg
|
||
merged["actions"] = merged_actions
|
||
return MyWorker(stf, serial, params=merged)
|
||
```
|
||
|
||
即:
|
||
- 顶层字段:Job 参数覆盖默认参数
|
||
- `actions` 字段:**参数级深合并**,前端可只覆盖某个 action 的某个子字段(如只改 `like.rate`)
|
||
|
||
### 6.3 前端任务参数 JSON 示例
|
||
|
||
Job 下发时只传**需要覆盖**的字段,调度器做深合并。例如只想把点赞概率从 0.3 调到 0.5,Job params 只需:
|
||
|
||
```json
|
||
{
|
||
"actions": {
|
||
"like": {"params": {"rate": 0.5}}
|
||
}
|
||
}
|
||
```
|
||
|
||
其余字段自动取 `DEFAULT_PARAMS`。**不要**在 Job 里传完整 params——升级默认值时会丢失新字段。
|
||
|
||
---
|
||
|
||
## 7. 日志规范
|
||
|
||
### 7.1 获取 logger
|
||
|
||
```python
|
||
from core.logger import get_logger
|
||
|
||
log = get_logger("task.kuaishou") # 任务日志 → logs/task.log
|
||
log = get_logger("action.kuaishou.like") # action 日志 → logs/action.log
|
||
log = get_logger("core.task_manager") # 核心日志 → logs/core.log
|
||
log = get_logger("web") # web 日志 → logs/web.log
|
||
```
|
||
|
||
logger 名前缀决定写入哪个文件:
|
||
|
||
| 前缀 | 文件 |
|
||
| --- | --- |
|
||
| `core.*` | `logs/core.log` |
|
||
| `task.*` | `logs/task.log` |
|
||
| `action.*` | `logs/action.log` |
|
||
| `web.*` | `logs/web.log` |
|
||
|
||
### 7.2 级别
|
||
|
||
- `DEBUG` — 详细元素查找、参数 dump(生产关闭)
|
||
- `INFO` — 正常流程节点(启动、轮次、action 结果)
|
||
- `WARNING` — 可恢复异常(元素找不到、action 失败)
|
||
- `ERROR` — 不可恢复错误(设备掉线、调度失败)
|
||
|
||
### 7.3 规则
|
||
|
||
- **禁止 `print`**,统一用 `get_logger`
|
||
- 日志里带 `[{self.serial}]` 设备前缀,多设备并发时才能区分
|
||
- 单文件 10MB 滚动,保留 5 份历史,无需手动清理
|
||
- 不要在循环里高频打 INFO(如每个 `exists()` 都打),用 DEBUG
|
||
|
||
---
|
||
|
||
## 8. STF 设备调试
|
||
|
||
### 8.1 常见错误
|
||
|
||
| 现象 | 原因 | 处理 |
|
||
| --- | --- | --- |
|
||
| HTTP 504 | 设备掉线 / STF 卡住 | 抛 `DeviceOfflineError`,不重试 |
|
||
| `DeviceOfflineError` | u2 连不上 / adb 远程不通 | 立即释放,跳过该设备 |
|
||
| `present=True` 但操作失败 | STF 状态有缓存,`present` 不代表真在线 | 用前台 App 扫描复测 |
|
||
| u2.connect 永久 hang | atx-agent 无响应 | 基类已加 30s 超时保护,超时抛异常 |
|
||
|
||
### 8.2 前台 App 扫描(不打扰设备)
|
||
|
||
Web 提供"扫描前台App"按钮(`/api/scan_foreground`),按设备状态分三类处理:
|
||
|
||
| 设备状态 | 处理方式 | 是否打扰 |
|
||
| --- | --- | --- |
|
||
| worker 运行中 | 复用已有 ADB 连接查询 | 否 |
|
||
| 完全空闲 | `adb connect` → `dumpsys` → `adb disconnect` | 否 |
|
||
| 他人占用 | 标记"(他人占用)" | 否 |
|
||
|
||
**绝不使用 STF occupy/release**——会唤醒 STF agent 导致设备退回桌面。
|
||
|
||
### 8.3 不重试原则
|
||
|
||
`DeviceOfflineError` 一律不重试——设备掉线后短时间内不会自愈,重试只会浪费配额并阻塞调度器。让设备进入冷却,由运维/STF 恢复后再启用。
|
||
|
||
---
|
||
|
||
## 9. 定位元素技巧
|
||
|
||
### 9.1 抓界面
|
||
|
||
用 [weditor](https://github.com/alibaba/web-editor)(`pip install weditor` → `python -m weditor`)实时查看 UI 树,复制定位表达式。
|
||
|
||
### 9.2 定位优先级
|
||
|
||
```
|
||
description > descriptionContains > resourceId > text/textContains > xpath
|
||
```
|
||
|
||
- **`description`** 最稳,开发者较少改动 contentDescription
|
||
- **`descriptionContains`** 模糊匹配,适配不同版本文案(如"点赞"/"未点赞")
|
||
- **`resourceId`** 注意带包名前缀(`com.xxx:id/...`),跨版本可能变,建议**多候选**
|
||
- **`xpath`** 用**相对定位**,禁止依赖 `FrameLayout[2]` / `LinearLayout[3]` 这类绝对序号
|
||
|
||
### 9.3 多策略组合 + 回退
|
||
|
||
```python
|
||
def find_like_button(d):
|
||
"""多策略定位点赞按钮,失败回退。"""
|
||
# 1. description 精确
|
||
for desc in ("点赞", "未点赞", "like"):
|
||
el = d(description=desc)
|
||
if el.exists:
|
||
return el
|
||
# 2. descriptionContains 模糊
|
||
for kw in ("赞", "like"):
|
||
el = d(descriptionContains=kw)
|
||
if el.exists:
|
||
return el
|
||
# 3. resourceId 列表(多候选)
|
||
for rid in ("com.xxx:id/aky", "com.xxx:id/d-like-view-icon"):
|
||
el = d(resourceId=rid)
|
||
if el.exists:
|
||
return el
|
||
return None
|
||
```
|
||
|
||
### 9.4 xpath 写法
|
||
|
||
```python
|
||
# ✅ 相对定位,稳
|
||
d.xpath('//android.widget.TextView[@text="关注"]').click()
|
||
|
||
# ❌ 绝对序号,UI 一变就崩
|
||
d.xpath('//FrameLayout[2]/LinearLayout[1]/TextView[3]').click()
|
||
```
|
||
|
||
---
|
||
|
||
## 10. 常见问题
|
||
|
||
### 10.1 循环导入
|
||
|
||
`actions/__init__.py` 必须**先导入 `base` 再导入各 action 模块**:
|
||
|
||
```python
|
||
# tasks/xxx/actions/__init__.py
|
||
from .base import ACTIONS, ... # 1. 先建注册表
|
||
from . import like, comment # 2. 再导入各 action,触发 @register_action
|
||
```
|
||
|
||
`tasks/__init__.py` 同理:先 `from .base import BaseTask`,再 `from . import douyin`。
|
||
|
||
### 10.2 中文输入
|
||
|
||
uiautomator2 默认 IME 不支持中文。需切到 fastinput:
|
||
|
||
```python
|
||
try:
|
||
d.set_fastinput_ime(True) # 切入
|
||
d.send_keys("中文内容")
|
||
finally:
|
||
try:
|
||
d.set_fastinput_ime(False) # 用完切回
|
||
except Exception:
|
||
pass
|
||
```
|
||
|
||
设备未装 FastInput 输入法时 `set_fastinput_ime` 会静默失败,建议加 try/except + 日志。
|
||
|
||
### 10.3 多设备并发
|
||
|
||
`adb_helper` 内置**全局锁**串行化所有 adb 调用(`adb connect` / `adb devices` 等)。原因:
|
||
|
||
- 多线程并发调 adb 会触发 adb server 竞争,导致连接抖动
|
||
- **禁止**在任务代码里调 `adb kill-server`——会踢掉所有设备的连接
|
||
- 设备申请/释放走 `stf_client`,与 adb 锁配合避免冲突
|
||
|
||
```python
|
||
# ✅ 正确:用 adb_helper 封装
|
||
from core.adb_helper import adb_connect
|
||
adb_connect(serial)
|
||
|
||
# ❌ 错误:自己起 subprocess 调 adb,绕过全局锁
|
||
import subprocess
|
||
subprocess.run(["adb", "connect", serial])
|
||
|
||
# ❌ 严禁
|
||
subprocess.run(["adb", "kill-server"])
|
||
```
|
||
|
||
### 10.4 看门狗误杀
|
||
|
||
若任务有长耗时操作(如长视频播放等待 5 分钟),看门狗可能误判卡死。解决:
|
||
|
||
- 在长操作内部**周期性调用 `self.heartbeat()`**(如每 30 秒一次),而不是只在整个操作前后调
|
||
- 不要调高看门狗阈值——真卡死的设备需要尽快释放
|
||
|
||
```python
|
||
# 长等待的正确写法
|
||
end = time.time() + 300
|
||
while time.time() < end:
|
||
if self.stopped():
|
||
return
|
||
self.heartbeat() # 长循环内部也要心跳
|
||
time.sleep(5)
|
||
```
|
||
|
||
### 10.5 u2.connect 卡死
|
||
|
||
`u2.connect()` 在 atx-agent 无响应时会永久 hang。基类已用 `ThreadPoolExecutor + future.result(timeout=30)` 包裹,超时返回 None 并抛异常。**子类无需处理**,但要避免在 run_task 里直接调 `u2.connect()`。
|
||
|
||
### 10.6 任务参数前端覆盖
|
||
|
||
Job 下发时只传**需要覆盖**的字段,调度器做深合并(见 §6.2)。例如只想把点赞概率从 0.8 调到 0.5,Job params 只需:
|
||
|
||
```json
|
||
{
|
||
"actions": {
|
||
"like": {"params": {"rate": 0.5}}
|
||
}
|
||
}
|
||
```
|
||
|
||
其余字段自动取 `DEFAULT_PARAMS`。**不要**在 Job 里传完整 params——升级默认值时会丢失新字段。
|
||
|
||
### 10.7 进度上报必须用通用字段
|
||
|
||
前端只认 `progress = {done, total, unit, action_counts}` 结构。**不要**用 `videos_watched`、`round_idx` 等业务字段名——前端不会识别。
|
||
|
||
```python
|
||
# ✅ 正确
|
||
self.set_progress(done=5, total=80, unit="视频",
|
||
action_counts={"like": 3})
|
||
|
||
# ❌ 错误(前端不认)
|
||
self.set_progress(videos_watched=5, round_idx=3)
|
||
```
|
||
|
||
---
|
||
|
||
## 附录:新增任务 Checklist
|
||
|
||
新建一个 app 任务时,按此清单逐项确认:
|
||
|
||
- [ ] `tasks/<app>/__init__.py` 有 `from . import task`
|
||
- [ ] `tasks/<app>/task.py` 有 `DEFAULT_PARAMS`(自包含)+ `Worker(BaseWorker)` + `Task(BaseTask)` + `@register_task`
|
||
- [ ] `Worker.run_task` 已实现,循环顶部和 action 之间都检查 `self.stopped()`
|
||
- [ ] `Worker.run_task` 用 `self.set_progress(done=, total=, unit=, action_counts=)` 上报进度
|
||
- [ ] 长循环内周期性调用 `self.heartbeat()`
|
||
- [ ] `tasks/<app>/actions/__init__.py` 先 `from .base import ACTIONS` 再导入各 action
|
||
- [ ] 每个 Action 有 `action_type` / `name` / `default_params` / `execute`,返回 `True/False`
|
||
- [ ] `Task.create_worker` 做 actions 参数级深合并(照抄抖音模板)
|
||
- [ ] `tasks/__init__.py` 已 `from . import <app>` 注册
|
||
- [ ] 日志用 `get_logger("task.<app>")` / `get_logger("action.<app>.<name>")`,无 `print`
|
||
- [ ] 定位元素优先 `description` / `descriptionContains`,resourceId 多候选,xpath 用相对定位
|
||
- [ ] 中文输入用 `set_fastinput_ime`,加 try/except
|
||
- [ ] adb 操作走 `adb_helper`,未自起 subprocess,未 `kill-server`
|
||
|
||
完成上述清单后,重启 web,前端单页应用即可看到新任务类型并可下发。
|