Files
auto_control/core/notify_events.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

204 lines
10 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.
"""通知事件目录(**唯一真相**)。
约定:
- **key 点分层**:`模块.对象.动作`(如 `task.device.failed`),支持通配订阅
(`task.*`、`device.*`、`*`),匹配规则见 `match()`。
- **登记 ≠ 会发**:所有事件 `default` 一律 False —— 用户在「系统 → 通知」里
按 webhook 勾选才会推送。`recommend` 只用于界面高亮"建议开启",不改默认值。
- **聚合**:`agg_window` 秒内的同 `(hook, event, agg_key 取值)` 事件合并成一条
(保留前几个样本);`agg_window=0` = 不聚合,立即发(低频高危事件用它)。
- **fields**:该事件保证携带的上下文字段,供前端画字段表、也是模板占位符白名单。
新增事件:在 `EVENTS` 里加一条 + 在触发点调 `notifier.notify(key, **fields)`,
并同步 `doc/NOTIFY.md` 的事件表(见 doc/README.md 的文档同步红线)。
"""
from collections import namedtuple
# key, label, category, default, recommend, agg_window, agg_key, fields, desc
EventDef = namedtuple(
"EventDef",
"key label category default recommend agg_window agg_key fields desc")
def _e(key, label, category, fields, desc="", agg_window=30, agg_key=None,
recommend=False):
"""构造一条事件定义(default 恒为 False:登记不等于推送)。"""
return EventDef(key, label, category, False, recommend, agg_window,
agg_key, tuple(fields), desc)
# 聚合键:按"任务批次"合并(同一次任务的多台设备结果合成一条)
_BY_JOB = ("job_id",)
# 按"设备"合并(一台设备的重复抖动合成一条)
_BY_SERIAL = ("serial",)
EVENTS = [
# ---------------- 任务批次 ----------------
_e("task.batch.started", "任务批次开始", "任务批次",
["job_id", "job_name", "task_type", "device_count", "serials_preview"],
"一次任务开始铺开到 N 台设备", agg_window=0, recommend=True),
_e("task.batch.finished", "任务批次结束", "任务批次",
["job_id", "job_name", "total", "success", "failed", "stopped", "duration_s"],
"整批跑完(含成功/失败/被停台数与耗时)", agg_window=0, recommend=True),
_e("task.batch.no_device", "任务无可用设备", "任务批次",
["job_id", "job_name", "target"],
"触发时一台可用设备都没有(任务空跑)", agg_window=0, recommend=True),
_e("task.batch.unknown_type", "任务类型不存在", "任务批次",
["job_id", "job_name", "task_type"],
"任务类型已从代码里删掉,任务永远不会执行", agg_window=0, recommend=True),
_e("task.cron.stopped", "定时停止任务", "任务批次",
["job_id", "job_name", "stopped_count", "serials"],
"cron_stop 到点,停掉了正在跑的设备", agg_window=0),
# ---------------- 任务 · 单设备(权威结论点) ----------------
_e("task.device.success", "设备任务成功", "任务·单设备",
["serial", "device_name", "job_id", "job_name", "attempt", "duration_s"],
"某台设备上的任务最终成功", agg_key=_BY_JOB, recommend=True),
_e("task.device.failed", "设备任务失败", "任务·单设备",
["serial", "device_name", "job_id", "job_name", "attempts", "cause", "msg"],
"某台设备上的任务最终失败(带真实原因)", agg_key=_BY_JOB, recommend=True),
_e("task.device.offline", "设备离线放弃", "任务·单设备",
["serial", "device_name", "job_id", "job_name", "error"],
"设备离线,不重试直接失败", agg_key=_BY_JOB, recommend=True),
_e("task.device.error", "单次尝试异常", "任务·单设备",
["serial", "device_name", "job_name", "attempt", "error"],
"某次尝试抛异常(后面还会重试,噪音较大)", agg_key=_BY_JOB),
_e("task.device.retry", "设备任务重试", "任务·单设备",
["serial", "device_name", "job_name", "attempt", "next_attempt", "delay_s", "reason"],
"即将重试", agg_key=_BY_JOB),
_e("task.device.stopped", "设备任务被停止", "任务·单设备",
["serial", "device_name", "job_name", "attempt", "phase"],
"用户手动停止 / cron 停止", agg_key=_BY_JOB),
_e("task.device.preempted", "设备被抢占", "任务·单设备",
["serial", "device_name", "job_name", "preempted_job_id", "preempted_job_name"],
"本任务抢占了该设备上正在跑的其他任务", agg_key=_BY_JOB, recommend=True),
_e("task.device.preempt_timeout", "抢占超时跳过", "任务·单设备",
["serial", "device_name", "job_name", "preempted_job_id"],
"等被抢占任务退出超时,本设备放弃执行", agg_window=0, recommend=True),
_e("task.device.released", "抢占结束归还", "任务·单设备",
["serial", "device_name", "job_name", "preempted_job_id", "returned", "reason"],
"抢占结束时重新拉起被抢占的任务", agg_key=_BY_JOB),
# ---------------- Worker 状态机(单次 attempt 级,默认关) ----------------
_e("worker.connected", "设备已连接", "Worker",
["serial", "model", "remote_adb_url"],
"worker 连上设备(单次尝试级,噪音大)"),
_e("worker.attempt.done", "单次执行完成", "Worker",
["serial", "task_type", "duration_s", "max_duration_hit"],
"单次 attempt 正常跑完(不代表任务最终成功)"),
_e("worker.attempt.error", "单次执行出错", "Worker",
["serial", "task_type", "error", "transient"],
"单次 attempt 抛异常"),
# ---------------- 业务(任务内容) ----------------
_e("task.selector.invalid", "选择器连续失效", "业务",
["serial", "device_name", "selector", "miss_count"],
"某选择器连续 10 次未命中——任务可能显示成功但什么都没做", agg_window=0,
recommend=True),
# ---------------- 设备 ----------------
_e("device.online", "设备恢复在线", "设备",
["serials", "count"],
"原本断联的设备又能连上了", agg_window=0, recommend=True),
_e("device.offline", "设备断联", "设备",
["serials", "count", "devices"],
"设备池里的设备连不上了", agg_window=0, recommend=True),
_e("device.discovered", "发现新设备", "设备",
["serials", "count"],
"扫描到尚未在池内的设备(待认领)", agg_window=0),
_e("device.claimed", "设备自动认领", "设备",
["pairs", "count"],
"指纹匹配成功,自动把旧记录迁到新地址", agg_window=0),
_e("device.heartbeat_timeout", "心跳超时", "设备",
["serial", "device_name", "timeout_s", "task_job", "model"],
"设备卡死(长时间没心跳),任务可能已中断", agg_window=0, recommend=True),
# ---------------- 应用安装 ----------------
_e("apk.install.started", "应用安装开始", "安装",
["apk_id", "apk_name", "package_name", "total"],
"开始往 N 台设备推装", agg_window=0),
_e("apk.install.finished", "应用安装完成", "安装",
["apk_name", "success", "failed", "skipped", "total", "failed_items"],
"批量安装结束(含失败台数与原因)", agg_window=0, recommend=True),
# ---------------- 系统 ----------------
_e("system.backup.exported", "备份已导出", "系统",
["filename", "size", "tables", "include_apk", "user"],
"有人导出了整库备份", agg_window=0, recommend=True),
_e("system.backup.imported", "备份导入已挂起", "系统",
["token_prefix", "env_label", "force"],
"上传了备份并确认导入(重启后生效)", agg_window=0),
_e("system.backup.restored", "备份已恢复", "系统",
["applied_rows", "schema_version"],
"启动时应用了待恢复的备份(所以这条是重启后才发)", agg_window=0,
recommend=True),
_e("system.backup.restore_failed", "备份恢复失败", "系统",
["error", "fail_dir"],
"恢复校验没过,已搁置(数据未被改动)", agg_window=0, recommend=True),
_e("service.started", "服务已启动", "系统",
["version", "env", "db_target", "device_count", "job_count"],
"平台进程起来了", agg_window=0, recommend=True),
_e("service.stopping", "服务正在停止", "系统",
["uptime_s"],
"平台进程收到退出信号", agg_window=0),
_e("user.login", "用户登录", "系统",
["username"],
"有人登录了后台", agg_window=0),
# ---------------- AI ----------------
_e("ai.audit.finished", "经验巡检完成", "AI",
["reviewed", "suggested", "summary"],
"每日经验库巡检跑完(或手动触发)", agg_window=0, recommend=True),
_e("ai.audit.failed", "经验巡检异常", "AI",
["error"],
"巡检过程中报错", agg_window=0),
_e("ai.audit.skipped", "经验巡检跳过", "AI",
["reason"],
"没配模型 Key 等原因跳过巡检", agg_window=0),
# ---------------- 其它 ----------------
_e("notify.test", "测试通知", "其它",
["hook_name", "operator"],
"「发送测试」按钮专用", agg_window=0),
]
BY_KEY = {e.key: e for e in EVENTS}
# 界面上的分组顺序
CATEGORIES = ("任务批次", "任务·单设备", "业务", "设备", "安装", "系统", "AI",
"Worker", "其它")
def get(key):
"""按 key 取事件定义;未知 key 返回 None(未知事件会被 notify 丢弃)。"""
return BY_KEY.get(key)
def match(pattern, key):
"""事件订阅匹配。
- `*` → 全部
- `task.*` → task 下所有(含 task.device.failed)
- `task.device.*` → 任意层级前缀匹配
- `task.device.failed` → 精确
"""
if not pattern:
return False
if pattern == "*" or pattern == key:
return True
if pattern.endswith(".*"):
return key.startswith(pattern[:-1]) # "task.*" → "task."
return False
def list_events():
"""给接口/前端用的事件目录(按 CATEGORIES 排序)。"""
order = {c: i for i, c in enumerate(CATEGORIES)}
rows = [{"key": e.key, "label": e.label, "category": e.category,
"default": e.default, "recommend": e.recommend,
"agg_window": e.agg_window, "agg_key": list(e.agg_key or ()),
"fields": list(e.fields), "desc": e.desc}
for e in EVENTS]
rows.sort(key=lambda r: (order.get(r["category"], 99), r["key"]))
return rows