diff --git a/README.md b/README.md index fa4db85..2c90cca 100644 --- a/README.md +++ b/README.md @@ -53,6 +53,7 @@ | **AI 控制台** | 用自然语言驱动 AI 操作指定设备(MCP 工具 + 截图),流式输出、Markdown 渲染、推理链折叠、token 统计;成功操作自动沉淀「经验库 / 动作库」并在相似任务中召回 | | **MCP 接入** | 20 个 `de_*` 工具,把手机控制开放给外部 AI;写操作有开关、设备忙时拒绝、全量审计 | | **备份导出/导入** | 一键导出 zip(库快照 + manifest + 可选 APK),导入前校验预览、自动预备份、重启生效 | +| **通知 / Webhook** | 任务成功失败、设备上下线、安装完成、备份恢复等事件推送到企业微信 / 钉钉 / 飞书 / Bark / 自建服务(Slack 用通用 JSON);多条 webhook 各自订阅;聚合+限流防刷屏(见 [doc/NOTIFY.md](doc/NOTIFY.md)) | --- @@ -120,6 +121,7 @@ auto_control/ │ ├── tailscale_api.py # Tailscale 管理(改名/授权/密钥/IP/auth key) │ ├── agent_api.py # AI 控制台(会话/SSE/经验库/动作库/巡检) │ ├── system_api.py # 系统数据备份导出/导入 +│ ├── notify_api.py # 通知 / Webhook 配置(系统 Tab) │ ├── common.py # 跨模块共享工具(合并设备列表、屏幕状态) │ └── context.py # 共享对象注入(mgr / apk_mgr / device_pool) │ @@ -136,6 +138,7 @@ auto_control/ │ ├── clipboard_helper.py # 剪贴板注入(ClipInject 通道) │ ├── apk_manager.py # APK 上传/解析/批量安装 │ ├── system_backup.py # 数据备份导出/导入(重启生效) +│ ├── notifier.py # 通知分发(队列/聚合/限流/适配器)+ notify_events.py 事件目录 │ ├── tailscale_client.py # Tailscale API v2 客户端 │ ├── ssh_client.py # SSH 封装(当前无人调用,预留) │ ├── logger.py # 分文件日志(core/task/web/action) @@ -162,7 +165,7 @@ auto_control/ │ ├── monitor.html # 主单页应用(7 个顶级 Tab + 模态框 + 内联样式) │ ├── login.html # 登录页 │ └── wall.html # 监控大屏(独立页面,自包含) -├── static/admin/ # 前端 JS(11 个文件,按顺序同步加载,见下) +├── static/admin/ # 前端 JS(13 个文件,按顺序同步加载,见下) ├── static/fonts/ # 自托管字体(Bricolage Grotesque + IBM Plex Mono) │ ├── data/ # 运行时数据(不入 git) @@ -181,7 +184,7 @@ auto_control/ ``` base.js → markdown.js → list.js → monitor.js → editor.js → tasks.js - → tools.js → apps.js → admin.js → agent.js → system.js + → tools.js → apps.js → admin.js → agent.js → taskgen.js → system.js → notify.js ``` --- @@ -347,6 +350,7 @@ MCP_ALLOW_WRITE=1 MCP_PLATFORM_USER=admin MCP_PLATFORM_PASS=<密码> \ | `logs/task.log` | `task.*` | 任务执行(worker 业务逻辑) | | `logs/web.log` | `web.*` | Web 请求与管理操作 | | `logs/action.log` | `action.*` | 操作执行 | +| `logs/notify.log` | `notify` | 通知发送(成功/失败/平台错误码) | ```python from core.logger import get_logger diff --git a/core/apk_manager.py b/core/apk_manager.py index ec282ee..ace0169 100644 --- a/core/apk_manager.py +++ b/core/apk_manager.py @@ -30,6 +30,7 @@ from concurrent.futures import (ThreadPoolExecutor, as_completed, from config import APK_DIR, ADB_PATH from core.logger import get_logger +from core import notifier from core.models import db, ApkFile as ApkRow from .device_worker import get_all_worker_status @@ -320,6 +321,8 @@ class ApkManager: name="apk-install", daemon=True) t.start() _log.info(f"开始安装 {apk_name} 到 {len(serials)} 台设备") + notifier.notify("apk.install.started", apk_id=apk_id, apk_name=apk_name, + package_name=package_name, total=len(serials)) return True, f"开始安装 {apk_name} 到 {len(serials)} 台设备" def _install_worker(self, apk_id, apk_path, apk_name, package_name, @@ -378,6 +381,16 @@ class ApkManager: skipped = sum(1 for v in self._install_task["items"].values() if v["status"] == "skipped") _log.info(f"安装完成 {apk_name}: 成功={success}, 失败={failed}, 跳过={skipped}") + # 通知:批量安装结束(带失败设备与原因,运维要知道哪几台没装上) + try: + items = ((self._install_task or {}).get("items") or {}) + failed_items = [f"{v.get('name') or s}·{v.get('msg') or ''}" + for s, v in items.items() if v.get("status") == "failed"][:10] + notifier.notify("apk.install.finished", apk_name=apk_name, + success=success, failed=failed, skipped=skipped, + total=len(items), failed_items=failed_items) + except Exception as e: + _log.warning(f"发送安装完成通知失败(不影响安装): {e}") def _install_one(self, serial, apk_path, package_name=""): """直连设备安装 APK。 diff --git a/core/device_discovery.py b/core/device_discovery.py index f3799f1..3d0e1dc 100644 --- a/core/device_discovery.py +++ b/core/device_discovery.py @@ -22,6 +22,7 @@ from config import (USB_ADB_HOST, DISCOVERY_PORT, DISCOVERY_SUBNETS, DISCOVERY_INTERVAL) from core.adb_helper import _adb, adb_connect_light from core.logger import get_logger +from core import notifier from core.models import db, Device, PendingDevice _log = get_logger("core.disc") @@ -42,6 +43,8 @@ _PROBE_TIMEOUT = 0.4 _app = None _scan_lock = threading.Lock() # 定时/手动扫描互斥 +# 上一轮"池内离线"集合:通知只报**状态沿**(这一轮新变成离线的),不每轮刷屏 +_last_offline = set() _stop_event = threading.Event() # shutdown 用 _scanning = False # 状态快照(API/前端) _last_scan = None # (时间串, 开放数, 可连数, 新增数) @@ -266,6 +269,7 @@ def scan_once(manual=False): now = _fmt() existing = {p.serial for p in PendingDevice.query.all()} added = 0 + added_serials = [] from core import device_pool fps = {} for serial in verified: @@ -283,6 +287,7 @@ def scan_once(manual=False): first_seen=now, last_seen=now, fingerprint=fp)) added += 1 + added_serials.append(serial) db.session.commit() # 5.5 自动认领(可选,默认关):指纹命中池中已有设备 → 直接把记录迁到新地址。 # 默认关是因为认领会改写分组/任务引用(数据结构变动),交人工点一下更稳妥; @@ -293,11 +298,34 @@ def scan_once(manual=False): # 全局锁串行),连上即恢复在线,无需人工干预。pending 池是给 # 「未授权新设备」的,正式池设备断联不进 pending,而是自动重连。 back = _reconnect_offline(configured) + # 7. 离线集合(供"状态沿"通知用:只报**这一轮新变成离线**的,不每轮刷屏) + offline_now = {s for s in configured if not device_pool.is_online(s)} _log.info(f"发现: 探测开放 {len(open_ips)} 台,可连 {len(verified)} 台," f"新增待连接 {added} 台" + (f",自动重连恢复 {len(back)} 台 {back}" if back else "") + (f",自动认领 {len(claimed)} 台 {[c[0] + '→' + c[1] for c in claimed]}" if claimed else "")) + # 通知统一放在 `with _ctx():` 之后(别在 DB 事务期间做额外的事) + global _last_offline + try: + newly_offline = sorted(offline_now - _last_offline) + _last_offline = set(offline_now) + if newly_offline: + notifier.notify("device.offline", serials=newly_offline, + count=len(newly_offline), + devices=newly_offline[:10]) + if back: + notifier.notify("device.online", serials=list(back), count=len(back)) + if added_serials: + notifier.notify("device.discovered", serials=added_serials[:10], + count=len(added_serials)) + if claimed: + notifier.notify("device.claimed", + pairs=[f"{c[0]}→{c[1]}" for c in claimed][:10], + count=len(claimed)) + except Exception as e: + _log.warning(f"发送设备上下线通知失败(不影响扫描): {e}") + _last_scan = (now, len(open_ips), len(verified), added) _last_error = "" return True, {"found": len(open_ips), "verified": len(verified), diff --git a/core/device_worker.py b/core/device_worker.py index a613e2a..8efab40 100644 --- a/core/device_worker.py +++ b/core/device_worker.py @@ -19,6 +19,7 @@ import uiautomator2 as u2 from config import USB_ADB_HOST, USB_ADB_PORT from core.logger import get_logger from core import device_pool +from core import notifier from .adb_helper import _adb, adb_connect, _adb_remote _log = get_logger("core.worker") @@ -188,9 +189,18 @@ class _Watchdog(threading.Thread): if hb and now - hb > _HEARTBEAT_TIMEOUT: stale.append(serial) for serial in stale: - _log.error(f"[{serial}] 心跳超时 {now - get_worker_heartbeat(serial):.0f}s,标记卡死") + hb_age = now - get_worker_heartbeat(serial) + _log.error(f"[{serial}] 心跳超时 {hb_age:.0f}s,标记卡死") + with _WORKERS_LOCK: + cur = dict(_WORKERS.get(serial, {})) _update_status(serial, status="error", last_error=f"心跳超时 {_HEARTBEAT_TIMEOUT}s,worker 可能卡死") + # 设备卡死是最需要人工介入的静默故障——通知一次 + # (天然只报一次:下面的检查只针对 running/connecting,标成 error 后不再命中) + notifier.notify("device.heartbeat_timeout", serial=serial, + device_name=cur.get("device_name") or "", + timeout_s=int(hb_age), task_job=cur.get("task_job") or "", + model=cur.get("model") or "") def stop(self): self._stop.set() diff --git a/core/logger.py b/core/logger.py index bc37006..cc30904 100644 --- a/core/logger.py +++ b/core/logger.py @@ -31,6 +31,8 @@ _MODULE_FILES = { "task": "task.log", "web": "web.log", "action": "action.log", + # 通知单独一个文件:排障时"哪条通知发了/失败了/为什么"要能一眼翻到 + "notify": "notify.log", } _FORMAT = "%(asctime)s [%(levelname)s] [%(name)s] %(message)s" diff --git a/core/notifier.py b/core/notifier.py new file mode 100644 index 0000000..1caa6d9 --- /dev/null +++ b/core/notifier.py @@ -0,0 +1,1081 @@ +"""Webhook 通知运行时:平台各组件通过 `notify(event, **fields)` 发通知。 + +## 设计底线(改动前务必先读) + +**通知绝不能影响主流程。** `notify()` 的保证是:**零 DB 访问、零 HTTP、零阻塞、 +异常绝不冒泡**——它只做「读内存配置快照 → 匹配订阅 → 丢进队列」这几步内存操作, +真正发 HTTP 的是后台 daemon 线程。所以: + +- 任务线程里可以直接调,**不用包 app_context**、不用 try/except(notify 自己兜) +- 但**必须放在所有 `with self._lock` 之外**(锁内只做状态变更),别让通知拖住调度锁 +- 任何一处 notify 出问题,只会在日志里看到警告,任务照跑 + +## 数据流 + + notify() ──入队──> _event_q ──> notify-dispatcher(1 线程) + ① 聚合:同 (hook, 事件, 聚合键) 窗口内合并 + ② 限流:每 hook 令牌桶(默认 18/分) + ③ 折叠:被限流的不丢,压成摘要(最多 60s 一条) + └──> _send_q ──> notify-sender × 3(真实 HTTP) + └──> 环形缓冲 + logs/notify.log + +防打爆四层(聚合/限流/折叠/背压)的取舍:失败通知最多延迟一个聚合窗口(默认 30s), +换来群不被刷屏——**刷屏会让通知彻底失效**。 + +## 配置 + +只落 `app_meta.notify_webhooks` 一个键(JSON),**不新建表**(因此不涉及备份覆盖红线)。 +读写只走 `core.db_config.meta_get/meta_set`(方言中立;`app_meta.key` 是 MySQL 保留字)。 +URL 里带凭据(企微 `?key=`),对外一律 `mask_url()`,日志/异常消息过 `scrub()`。 +""" +import hmac # noqa: F401 (签名插槽用,钉钉/飞书实现时启用) +import json +import queue +import random +import threading +import time +from collections import Counter, deque +from datetime import datetime + +import requests + +from core import notify_events +from core.logger import get_logger + +_log = get_logger("notify") # 单独写 logs/notify.log(见 core/logger.py 的 _MODULE_FILES) + +# ---------------- 常量 ---------------- +META_KEY = "notify_webhooks" +MAX_HOOKS = 20 +MAX_CONFIG_CHARS = 60000 # app_meta.value 是 TEXT,对齐 agent_task_draft 的口径 +MAX_URL_LEN = 2048 +MAX_EVENTS_PER_HOOK = 200 +EVENT_Q_SIZE = 2000 +SEND_Q_SIZE = 1000 +SENDER_THREADS = 3 +RING_SIZE = 200 +FOLD_INTERVAL = 60 # 被限流折叠后,多久汇总报一次 +TEST_PER_MIN = 10 # 「发送测试」按钮自身的限流(不占业务令牌桶) +_HTTP_TIMEOUT = (3, 5) # (连接, 读取) 秒 +_RETRY_DELAYS = (1, 4) # 失败退避:最多重试 2 次(合计 3 次尝试) + +_DEFAULT_SETTINGS = { + "global_enabled": True, + "default_agg_window": 30, # 秒;0=不聚合 + "default_rate_limit": 18, # 条/分钟(企业微信硬限 20,留余量) + "log_keep": RING_SIZE, + "http_timeout": 5, +} + +# 计划中但还没实现的格式(前端置灰)——钉钉/飞书已于 2026-09-15 实现,这里是留给后续的 +PLANNED_FORMATS = {} + + +# ================== URL / 文本脱敏 ================== +_SECRET_KEYS = ("key", "access_token", "token", "secret", "sign", "apikey", "api_key") + + +def mask_url(url): + """把 URL 里的凭据打码。 + + 两类都要处理(**URL 本身就是可直接发消息的凭据**): + - query 参数:企微 `?key=`、钉钉 `?access_token=` → 值换成 `***` + - **路径段的 token**:飞书 `https://open.feishu.cn/open-apis/bot/v2/hook/` + (token 在路径最后一段,只处理 query 会漏掉它) + 回显、发送记录、日志一律用它。 + """ + if not url: + return "" + try: + from urllib.parse import urlsplit, urlunsplit, parse_qsl, urlencode + sp = urlsplit(url) + segs = sp.path.split("/") if sp.path else [] + # 最后一段是长随机串(飞书 hook token 那种)→ 打码;普通路径(/robot/send)不动 + if segs and len(segs[-1]) >= 20 and segs[-1].replace("-", "").replace("_", "").isalnum(): + segs[-1] = "***" + sp = sp._replace(path="/".join(segs)) + # Slack 的凭据在路径三连(/services/T…/B…/X…)——整段打码 + if sp.netloc in ("hooks.slack.com", "hooks.slack-gov.com") and "services" in segs: + i = segs.index("services") + 1 + sp = sp._replace(path="/".join(segs[:i] + ["***"] * (len(segs) - i))) + if not sp.query: + return urlunsplit((sp.scheme, sp.netloc, sp.path, "", sp.fragment)) + q = [(k, ("***" if k.lower() in _SECRET_KEYS else v)) + for k, v in parse_qsl(sp.query, keep_blank_values=True)] + # safe='*' 让打码后的 *** 保持原样(否则会被编码成 %2A%2A%2A,看着像乱码) + return urlunsplit((sp.scheme, sp.netloc, sp.path, + urlencode(q, safe="*"), sp.fragment)) + except Exception: + return url[:60] + "…" + + +_SCRUB_RE = None + + +def scrub(text): + """清洗异常消息里的凭据(requests 的 str(e) 会带完整 URL,最容易漏的泄漏口)。""" + global _SCRUB_RE + if not text: + return "" + if _SCRUB_RE is None: + import re + _SCRUB_RE = re.compile( + r"(?i)\b(key|access_token|token|secret|sign)=(?![*\s&])[^&\s\"']+") + return _SCRUB_RE.sub(lambda m: m.group(1) + "=***", str(text)) + + +# ================== 模块状态 ================== +_app = None # Flask app(只在 _load_config/save_config 里用) +_ready = False # init_app 是否完成 +_config = None # 内存配置快照;None=未就绪(notify 直接丢) +_cfg_lock = threading.RLock() +_event_q = queue.Queue(maxsize=EVENT_Q_SIZE) +_send_q = queue.Queue(maxsize=SEND_Q_SIZE) +_stop = threading.Event() +_threads = [] +_ring = deque(maxlen=RING_SIZE) +_ring_lock = threading.Lock() +_dropped = {"event_q": 0, "send_q": 0} +_buckets = {} # hook_id -> {"tokens":, "ts":} +_folds = {} # hook_id -> {"count":, "by_event": Counter, "last": ts} +_agg = {} # (hook_id, event, aggkey) -> {"deadline","count","samples",...} +_test_hits = deque(maxlen=TEST_PER_MIN + 1) + + +# ================== 配置读写 ================== +def _normalize_hook(h, settings): + """补齐/纠正单条 webhook 配置(保存与读回都走它,保证内存与库一致)。""" + h = dict(h or {}) + events = [str(e).strip() for e in (h.get("events") or []) if str(e).strip()] + return { + "id": str(h.get("id") or ("wh_" + random.getrandbits(32).to_bytes(4, "big").hex())), + "name": str(h.get("name") or "未命名").strip()[:40], + "enabled": bool(h.get("enabled", True)), + "format": str(h.get("format") or "wecom"), + "url": str(h.get("url") or "").strip(), + "secret": str(h.get("secret") or ""), + "events": events[:MAX_EVENTS_PER_HOOK], + "agg_window": int(h.get("agg_window", settings["default_agg_window"]) or 0), + "rate_limit_per_min": int(h.get("rate_limit_per_min", + settings["default_rate_limit"]) or 0), + "title_template": str(h.get("title_template") or ""), + "body_template": str(h.get("body_template") or ""), + "headers": dict(h.get("headers") or {}), + "created_at": h.get("created_at") or datetime.now().strftime("%Y-%m-%d %H:%M:%S"), + "updated_at": h.get("updated_at") or datetime.now().strftime("%Y-%m-%d %H:%M:%S"), + } + + +def _normalize_config(raw): + if not isinstance(raw, dict): + raw = {} + settings = dict(_DEFAULT_SETTINGS) + if isinstance(raw.get("settings"), dict): + for k, v in raw["settings"].items(): + if k in settings and isinstance(v, (int, bool)): + settings[k] = v + hooks = [_normalize_hook(h, settings) for h in (raw.get("webhooks") or [])] + return {"version": 1, "settings": settings, "webhooks": hooks[:MAX_HOOKS]} + + +def _load_config(): + """从 app_meta 读配置(**唯一允许碰 DB 的地方之一**,必须在 app context 内调)。""" + from core.db_config import meta_get + raw = meta_get(META_KEY) or "" + if not raw: + return _normalize_config({}) + try: + obj = json.loads(raw) + except (ValueError, TypeError) as e: + _log.error("通知配置 JSON 解析失败(按空配置处理,不回写覆盖): %s", scrub(e)) + return _normalize_config({}) + return _normalize_config(obj) + + +def reload_config(): + """重新读配置(保存后、或外部改动 app_meta 后调用)。""" + global _config + if _app is None: + return + try: + with _app.app_context(): + cfg = _load_config() + except Exception as e: + _log.warning("重新加载通知配置失败(保持旧配置): %s", scrub(e)) + return + with _cfg_lock: + _config = cfg + _buckets.clear() + _agg.clear() + _log.info("通知配置已加载: %d 条 webhook,全局开关=%s", + len(cfg["webhooks"]), cfg["settings"]["global_enabled"]) + + +def get_config(): + """内存配置快照(含明文 url,**只在内部用**)。""" + with _cfg_lock: + return json.loads(json.dumps(_config or _normalize_config({}))) + + +def get_public_config(): + """给接口的配置:url 打码、secret 只回是否配置过。""" + cfg = get_config() + cfg["webhooks"] = [dict(h, url=mask_url(h["url"]), + secret="", secret_set=bool(h.get("secret"))) + for h in cfg["webhooks"]] + return cfg + + +def save_config(new_cfg): + """校验并保存配置。返回 (ok, msg|errors)。**保存后才 reload。**""" + global _config + if not isinstance(new_cfg, dict): + return False, "配置必须是对象" + hooks = new_cfg.get("webhooks") + if not isinstance(hooks, list): + return False, "webhooks 必须是数组" + if len(hooks) > MAX_HOOKS: + return False, f"webhook 最多 {MAX_HOOKS} 条(当前 {len(hooks)})" + + settings = dict(_DEFAULT_SETTINGS) + if isinstance(new_cfg.get("settings"), dict): + settings.update({k: v for k, v in new_cfg["settings"].items() if k in settings}) + + old = {h["id"]: h for h in (get_config()["webhooks"])} + norm = [] + errs = [] + for i, h in enumerate(hooks, 1): + nh = _normalize_hook(h, settings) + # url 省略/为空/是打码值 → 保持原值(前端"留空不修改"的兜底) + if not nh["url"] or nh["url"] == mask_url(old.get(nh["id"], {}).get("url", "")): + nh["url"] = old.get(nh["id"], {}).get("url", nh["url"]) + if not nh["url"].startswith(("http://", "https://")): + errs.append(f"第 {i} 条「{nh['name']}」: URL 必须是 http/https") + elif len(nh["url"]) > MAX_URL_LEN: + errs.append(f"第 {i} 条「{nh['name']}」: URL 太长(>{MAX_URL_LEN})") + if nh["format"] not in ADAPTERS and nh["format"] not in PLANNED_FORMATS: + errs.append(f"第 {i} 条「{nh['name']}」: 未知格式 {nh['format']}") + if not nh["events"]: + errs.append(f"第 {i} 条「{nh['name']}」: 至少要勾一个事件") + for p in nh["events"]: + if p != "*" and not any(notify_events.match(p, e.key) + for e in notify_events.EVENTS): + errs.append(f"第 {i} 条「{nh['name']}」: 事件 {p} 不存在") + if nh["format"] == "json": + ok, msg = JsonAdapter.validate_template(nh["body_template"]) + if not ok: + errs.append(f"第 {i} 条「{nh['name']}」: 自定义模板有问题——{msg}") + if not nh.get("secret") and old.get(nh["id"], {}).get("secret"): + nh["secret"] = old[nh["id"]]["secret"] # 密钥留空 = 不修改 + norm.append(nh) + if errs: + return False, errs + + cfg = {"version": 1, "settings": settings, "webhooks": norm} + raw = json.dumps(cfg, ensure_ascii=False) + if len(raw) > MAX_CONFIG_CHARS: + return False, f"配置太大({len(raw)} 字符 > {MAX_CONFIG_CHARS}),请精简模板/事件" + + try: + from core.db_config import meta_set + with _app.app_context(): + meta_set(META_KEY, raw) + except Exception as e: + _log.exception("保存通知配置失败") + return False, f"保存失败: {scrub(e)}" + with _cfg_lock: + _config = cfg + _buckets.clear() + _agg.clear() + return True, "已保存" + + +def set_global_enabled(on): + """只翻全局开关(列表页那个开关用)。""" + with _cfg_lock: + if not _config: + return False, "通知配置未就绪" + _config["settings"]["global_enabled"] = bool(on) + cfg = json.loads(json.dumps(_config)) + return save_config(cfg) + + +# ================== 消息渲染 ================== +_LEVEL_BY_HINT = ( + (("failed", "failure", "error", "timeout", "restore_failed", "invalid", + "no_device", "offline", "unknown_type"), "error"), + (("retry", "preempt", "stopped", "stopping", "skipped"), "warning"), + (("success", "finished", "started", "online", "done", "restored", + "exported", "claimed", "imported"), "success"), +) + + +def _level_of(event_key): + for hints, lv in _LEVEL_BY_HINT: + if any(h in event_key for h in hints): + return lv + return "info" + + +def _fmt_ts(ts=None): + return datetime.fromtimestamp(ts or time.time()).strftime("%Y-%m-%d %H:%M:%S") + + +# 字段名 → 中文标签(通知是发给人看的群消息,别甩英文字段名) +_FIELD_LABELS = { + "serial": "设备地址", "device_name": "设备", "job_name": "任务", "job_id": "任务ID", + "task_type": "任务类型", "attempt": "第几次", "attempts": "尝试次数", + "duration_s": "耗时(秒)", "cause": "失败原因", "msg": "说明", "error": "错误", + "error_msg": "错误", "count": "数量", "serials": "设备", "devices": "设备", + "pairs": "认领", "target": "目标", "selector": "选择器", "miss_count": "连续未命中", + "timeout_s": "超时(秒)", "task_job": "原任务", "model": "型号", + "apk_name": "应用", "apk_id": "应用ID", "package_name": "包名", + "success": "成功", "failed": "失败", "skipped": "跳过", "total": "总数", + "failed_items": "失败明细", "filename": "文件", "size": "大小(字节)", + "tables": "表数", "include_apk": "含APK", "user": "操作人", "operator": "操作人", + "env": "环境", "db_target": "数据库", "device_count": "设备数", + "job_count": "任务数", "version": "版本", "uptime_s": "运行时长(秒)", + "username": "用户", "reviewed": "评审条数", "suggested": "建议删除", + "summary": "结论", "stopped_count": "停止台数", "preempted_job_name": "被抢占任务", + "returned": "已归还", "reason": "原因", "transient": "可重试", + "max_duration_hit": "到时停止", "applied_rows": "恢复表数", + "schema_version": "schema 版本", "hook_name": "Webhook", "phase": "阶段", + "next_attempt": "下一次", "delay_s": "间隔(秒)", "miss": "未命中", +} + + +def _label(k): + return _FIELD_LABELS.get(k, k) + + +def _subject(ev, fields): + """标题主体:优先任务名,其次设备名/serial。""" + for k in ("job_name", "apk_name", "hook_name"): + if fields.get(k): + return str(fields[k]) + if fields.get("device_name") or fields.get("serial"): + return str(fields.get("device_name") or fields.get("serial")) + if fields.get("count") is not None: + return f"{fields.get('count')} 台设备" + return "" + + +def build_message(ev, fields, hook=None): + """把事件渲染成统一消息体(title/summary/fields/level/markdown)。""" + fields = dict(fields or {}) + merged = fields.pop("_merged", None) + level = _level_of(ev.key) + subj = _subject(ev, fields) + title = ev.label + (f":{subj}" if subj else "") + if merged and merged.get("merged_count", 0) > 1: + title = f"{ev.label}:{subj}" if subj else ev.label + title += f"({merged['merged_count']} 次)" + + rows = [] + if merged and merged.get("by_event"): + rows.append(("汇总", "、".join(f"{k}×{v}" for k, v in + merged["by_event"].most_common(6)))) + for k in ev.fields: + v = fields.get(k) + if v in (None, "", [], {}): + continue + if isinstance(v, (list, tuple)): + v = "、".join(str(x) for x in list(v)[:6]) + ("…" if len(v) > 6 else "") + elif isinstance(v, dict): + v = json.dumps(v, ensure_ascii=False)[:200] + rows.append((_label(k), str(v)[:400])) + if merged and merged.get("samples"): + rows.append(("样本", ";".join(merged["samples"][:3]))) + if merged and merged.get("folded"): + rows.append(("限流折叠", f"另有 {merged['folded']} 条被限流折叠")) + + icon = {"success": "✅", "error": "❌", "warning": "⚠️"}.get(level, "ℹ️") + lines = [f"### {icon} {title}"] + for k, v in rows: + lines.append(f"> **{k}**:{v}") + lines.append(f"> **时间**:{_fmt_ts()}") + md = "\n".join(lines) + summary = f"{title}" + (f"({rows[0][1][:60]})" if rows else "") + return {"event": ev.key, "title": title, "summary": summary, "level": level, + "fields": rows, "fields_raw": fields, "markdown": md, "ts": time.time()} + + +def _apply_template(tpl, msg, hook_name=""): + """通用 JSON 模板的占位符替换。 + + 值一律用 `json.dumps(v)[1:-1]`(**已转义的 JSON 字符串片段**)——这样标题里的 + 引号/换行不会把外层 JSON 打坏。 + """ + raw = dict(msg.get("fields_raw") or {}) + mapping = { + "event": msg["event"], "title": msg["title"], "summary": msg["summary"], + "ts": _fmt_ts(msg["ts"]), "level": msg["level"], "markdown": msg["markdown"], + "hook_name": hook_name, "fields_json": json.dumps(raw, ensure_ascii=False), + "fields": " / ".join(f"{k}={v}" for k, v in msg["fields"]), + } + out = tpl + for k, v in mapping.items(): + out = out.replace("{{%s}}" % k, json.dumps(str(v), ensure_ascii=False)[1:-1]) + for k, v in raw.items(): # {{field.xxx}} + out = out.replace("{{field.%s}}" % k, + json.dumps(str(v), ensure_ascii=False)[1:-1]) + return out + + +# ================== 适配器 ================== +class BaseAdapter: + """格式适配器基类(可插拔:新增格式 = 加一个类 + 注册进 ADAPTERS)。 + + 界面上的提示文案也挂在这里(`url_hint`/`url_help`/`secret_*`/`limit_help`)—— + **格式一换,界面上的说明要跟着换**,别把企业微信的说明写死在页面上。 + """ + + name = "base" + label = "基类" + byte_limit = 4096 # 单个消息体的字节上限(0=不限) + limit_default = 20 # 该格式每机器人每分钟的官方上限 + ok_codes = (0,) # 响应体里表示成功的错误码(None=只看 HTTP 状态) + truncated_note = "\n…(已截断)" + needs_template = False # 是否需要用户自定义请求体模板 + secret_label = "签名密钥(可选)" + secret_help = "" + url_hint = "https://…" + url_help = "" + limit_help = "" + + @classmethod + def render(cls, msg, hook): + """→ requests.post 的参数 dict(url/json/headers)。""" + raise NotImplementedError + + @classmethod + def sign_request(cls, url, secret, body): + """加签钩子:返回 (url, body)。 + + 默认不改动;钉钉(拼 URL,毫秒)和飞书(放 body,秒)各写各的 —— **两家的 + 签名算法不一样**,照抄另一家的会静默失败(钉钉 errcode 310000、飞书 code 19021)。 + """ + return url, body + + @classmethod + def meta(cls): + """给前端用的格式元信息(下拉项 + 随格式变化的提示文案)。""" + return {"name": cls.name, "label": cls.label, "implemented": True, + "byte_limit": cls.byte_limit, "limit_default": cls.limit_default, + "needs_template": cls.needs_template, + "secret_label": cls.secret_label, "secret_help": cls.secret_help, + "url_hint": cls.url_hint, "url_help": cls.url_help, + "limit_help": cls.limit_help} + + @classmethod + def _cut(cls, text): + """按 **UTF-8 字节**截断(企业微信限 4096 字节,不是字符)。""" + if not cls.byte_limit: + return text + b = text.encode("utf-8") + if len(b) <= cls.byte_limit: + return text + keep = cls.byte_limit - len(cls.truncated_note.encode("utf-8")) + # 走 encode/decode,天然不会截出半个汉字 + return b[:keep].decode("utf-8", "ignore") + cls.truncated_note + + +class WecomAdapter(BaseAdapter): + """企业微信群机器人:markdown 消息。 + + 真实约束(官方):content ≤ **4096 UTF-8 字节**;**每机器人每分钟 20 条** + (超限返回 errcode 45009);成功必须 `errcode == 0`。 + 另外企微 markdown 子集**不支持表格** —— 所以字段一律渲染成引用行。 + """ + + name = "wecom" + label = "企业微信" + byte_limit = 4096 + limit_default = 20 + ok_codes = (0,) + url_hint = "https://qyapi.weixin.qq.com/cgi-bin/webhook/send?key=…" + url_help = "企业微信:群设置 → 群机器人 → 添加 → 复制 Webhook 地址" + secret_help = "企业微信不需要密钥,留空即可" + limit_help = "企业微信硬限 20 条/分" + + @classmethod + def render(cls, msg, hook): + return { + "url": hook["url"], + "json": {"msgtype": "markdown", + "markdown": {"content": cls._cut(msg["markdown"])}}, + "headers": dict(hook.get("headers") or {}), + } + + +class JsonAdapter(BaseAdapter): + """通用 JSON:模板即请求体(Slack 也用它)。 + + 占位符:{{event}} {{title}} {{summary}} {{ts}} {{level}} {{markdown}} + {{fields}} {{fields_json}} {{hook_name}} {{field.<字段名>}} + """ + + name = "json" + label = "通用 JSON / Slack" + byte_limit = 0 + limit_default = 60 + needs_template = True + url_hint = "https://your-service/hook" + url_help = "你自己的接收端地址;Slack 填它的 Incoming Webhook URL" + secret_label = "自定义密钥(可选)" + secret_help = "填了会作为请求头 X-Webhook-Secret 一并发出,供你的服务校验" + limit_help = "由你自己的服务决定;默认 60 条/分" + + DEFAULT_TEMPLATE = json.dumps({ + "event": "{{event}}", + "title": "{{title}}", + "summary": "{{summary}}", + "level": "{{level}}", + "text": "{{markdown}}", + "ts": "{{ts}}", + }, ensure_ascii=False, indent=2) + + @classmethod + def render(cls, msg, hook): + tpl = (hook.get("body_template") or "").strip() or cls.DEFAULT_TEMPLATE + body = _apply_template(tpl, msg, hook.get("name", "")) + try: + parsed = json.loads(body) + except ValueError as e: + raise ValueError(f"模板渲染结果不是合法 JSON: {e}") + headers = dict(hook.get("headers") or {}) + if hook.get("secret"): + headers.setdefault("X-Webhook-Secret", hook["secret"]) + return {"url": hook["url"], "json": parsed, "headers": headers} + + @classmethod + def validate_template(cls, tpl): + """保存前的干跑校验:模板必须能渲染成合法 JSON。""" + tpl = (tpl or "").strip() or cls.DEFAULT_TEMPLATE + fake = {"event": "notify.test", "title": "标题\"含引号\n换行", "summary": "摘要", + "level": "info", "markdown": "# 正文", "ts": time.time(), + "fields": [("k", "v")], "fields_raw": {"k": "v"}} + try: + json.loads(_apply_template(tpl, fake, "测试")) + except ValueError as e: + return False, f"模板渲染后不是合法 JSON({e})" + except Exception as e: + return False, str(e)[:120] + return True, "" + + +class BarkAdapter(BaseAdapter): + """Bark(iOS 推送 App):`POST https://api.day.app/push`。 + + 成功判定**和企业微信不一样**:Bark 返回 `{"code":200,"message":"success"}` —— + 200 才是成功(企微是 errcode 0)。所以成功码挂在适配器上(ok_codes),不能写死。 + + 设备 key 两种填法都支持: + ① URL 直接粘贴 Bark 里复制的那串(`https://api.day.app/<你的key>`)——key 在路径里; + ② URL 填 `https://api.day.app/push`,把 key 填到「设备 Key」里(作为 device_key 字段发)。 + """ + + name = "bark" + label = "Bark(iOS 推送)" + byte_limit = 2048 # 走 APNs,单条体量有限,超了截断 + limit_default = 60 + ok_codes = (200,) + url_hint = "https://api.day.app/push" + url_help = ("Bark App 里复制的那串地址:可直接粘 https://api.day.app/<你的key>," + "或填 https://api.day.app/push 并把 key 填到下面的「设备 Key」") + secret_label = "设备 Key(可选)" + secret_help = "地址里已含 key 就留空;否则填 Bark App 里的那串 key(会以 device_key 发送)" + limit_help = "Bark 官方没有明确的每分钟上限;默认 60 条/分" + + @classmethod + def render(cls, msg, hook): + body = { + "title": msg["title"], + # body 给纯文本(老版本 App 不认 markdown 字段时也能看清),markdown 给富文本 + "body": cls._cut(msg["summary"]), + "markdown": cls._cut(msg["markdown"]), + "group": "auto_control", + "level": "active", + } + if hook.get("secret"): + body["device_key"] = hook["secret"] + return {"url": hook["url"], "json": body, + "headers": dict(hook.get("headers") or {})} + + +class DingtalkAdapter(BaseAdapter): + """钉钉群机器人:markdown 消息 + 加签。 + + 加签口径(官方,**和飞书不一样,别照抄**): + timestamp 是**毫秒**;`sign = urlencode(base64(HMAC-SHA256(key=secret, + message=f"{timestamp}\\n{secret}")))`;**timestamp/sign 拼到 URL 上**(不是 body)。 + 错了会返回 `errcode 310000 invalid signature`。 + 机器人若用"自定义关键词/IP 白名单"做安全设置,则**不需要** secret(留空即可)。 + """ + + name = "dingtalk" + label = "钉钉" + byte_limit = 20000 # markdown 正文上限 20000 字节 + limit_default = 15 # 官方 20/分,但超限会被限流 10 分钟 —— 留足余量 + ok_codes = (0,) + url_hint = "https://oapi.dingtalk.com/robot/send?access_token=…" + url_help = "钉钉:群设置 → 智能群助手 → 添加机器人 → 自定义 → 复制 Webhook 地址" + secret_label = "加签密钥(可选)" + secret_help = ("机器人安全设置选「加签」时,填 SEC 开头的那串;" + "选「自定义关键词/IP 白名单」则留空") + limit_help = "官方 20 条/分,超限会被限流 10 分钟(这里默认 15 留余量)" + + @classmethod + def render(cls, msg, hook): + return {"url": hook["url"], + "json": {"msgtype": "markdown", + "markdown": {"title": msg["title"], + "text": cls._cut(msg["markdown"])}}, + "headers": dict(hook.get("headers") or {})} + + @classmethod + def sign_request(cls, url, secret, body): + from urllib.parse import quote_plus + ts = str(int(time.time() * 1000)) # 毫秒 + sign = quote_plus(_hmac_b64(secret.encode("utf-8"), + f"{ts}\n{secret}".encode("utf-8"))) + sep = "&" if "?" in url else "?" + return f"{url}{sep}timestamp={ts}&sign={sign}", body + + +class FeishuAdapter(BaseAdapter): + """飞书群机器人:交互式卡片(markdown 元素)+ 加签。 + + 加签口径(官方,**和钉钉不一样**): + timestamp 是**秒**;`sign = base64(HMAC-SHA256(key=f"{timestamp}\\n{secret}", + message=空))`;**timestamp/sign 放在 JSON body 里**(拼 URL 通不过校验)。 + 错了会返回 `code 19021`。 + """ + + name = "feishu" + label = "飞书" + byte_limit = 20000 # 请求体上限 20KB + limit_default = 60 + ok_codes = (0,) + url_hint = "https://open.feishu.cn/open-apis/bot/v2/hook/…" + url_help = "飞书:群设置 → 群机器人 → 添加机器人 → 自定义机器人 → 复制 Webhook 地址" + secret_label = "签名校验密钥(可选)" + secret_help = "机器人安全设置勾了「签名校验」才需要填;否则留空" + limit_help = "飞书官方 100 条/分;这里默认 60 留余量" + + @classmethod + def render(cls, msg, hook): + body = { + "msg_type": "interactive", + "card": { + "config": {"wide_screen_mode": True}, + "header": { + "title": {"tag": "plain_text", "content": msg["title"][:120]}, + "template": {"error": "red", "warning": "orange", + "success": "green"}.get(msg["level"], "blue"), + }, + "elements": [{"tag": "markdown", + "content": cls._cut(msg["markdown"])}], + }, + } + return {"url": hook["url"], "json": body, + "headers": dict(hook.get("headers") or {})} + + @classmethod + def sign_request(cls, url, secret, body): + ts = str(int(time.time())) # 秒 + body = dict(body) + body["timestamp"] = ts + body["sign"] = _hmac_b64(f"{ts}\n{secret}".encode("utf-8"), b"") + return url, body + + +ADAPTERS = {WecomAdapter.name: WecomAdapter, JsonAdapter.name: JsonAdapter, + BarkAdapter.name: BarkAdapter, DingtalkAdapter.name: DingtalkAdapter, + FeishuAdapter.name: FeishuAdapter} + + +# ================== 发送 ================== +def _record(rec): + with _ring_lock: + _ring.append(rec) + lv = "info" if rec.get("ok") else "warning" + getattr(_log, lv)( + "通知发送 %s | %s | %s | %s | %sms%s", + "成功" if rec.get("ok") else "失败", rec.get("hook_name"), + rec.get("event"), mask_url(rec.get("url", "")), rec.get("elapsed_ms"), + ("" if rec.get("ok") else " | " + scrub(rec.get("error", "")))) + + +def _hmac_b64(key: bytes, message: bytes) -> str: + """HMAC-SHA256 → Base64(钉钉/飞书加签共用这一步)。""" + import base64 + import hashlib + import hmac as _hmac + return base64.b64encode( + _hmac.new(key, message, digestmod=hashlib.sha256).digest()).decode() + + +def _post_once(adapter, msg, hook, timeout): + kw = adapter.render(msg, hook) + if hook.get("secret"): + kw["url"], kw["json"] = adapter.sign_request(kw["url"], hook["secret"], + kw.get("json") or {}) + kw["timeout"] = timeout + kw["allow_redirects"] = False # 防重定向把消息带到别处 + r = requests.post(**kw) + body = "" + code = None + try: + data = r.json() + if isinstance(data, dict): + if "errcode" in data: + code = data.get("errcode") + elif "code" in data: + code = data.get("code") + body = json.dumps(data, ensure_ascii=False)[:200] + except ValueError: + body = (r.text or "")[:200] + # HTTP 200 不等于成功:各平台都用 body 里的错误码,但**成功码不一样** + # (企业微信 errcode=0、Bark code=200),所以按适配器声明的 ok_codes 判 + ok = r.status_code == 200 and (code is None or code in adapter.ok_codes) + return ok, r.status_code, code, body + + +def send_now(hook, msg, timeout=_HTTP_TIMEOUT): + """同步发送一条(供 dispatcher / 测试按钮用)。返回结果 dict(不抛异常)。""" + adapter = ADAPTERS.get(hook.get("format")) + name = hook.get("name", "") + t0 = time.time() + if adapter is None: + rec = {"ts": t0, "hook_id": hook.get("id"), "hook_name": name, + "event": msg.get("event"), "title": msg.get("title"), "ok": False, + "url": hook.get("url", ""), "http_status": None, "errcode": None, + "elapsed_ms": 0, "attempts": 0, + "error": f"格式 {hook.get('format')} 尚未实现"} + _record(rec) + return rec + last = "" + last_code = None + last_status = None + for attempt in range(1, len(_RETRY_DELAYS) + 2): + try: + ok, status, code, body = _post_once(adapter, msg, hook, timeout) + if ok: + rec = {"ts": t0, "hook_id": hook.get("id"), "hook_name": name, + "event": msg.get("event"), "title": msg.get("title"), "ok": True, + "url": hook.get("url", ""), "http_status": status, + "errcode": code, "elapsed_ms": int((time.time() - t0) * 1000), + "attempts": attempt, "error": ""} + _record(rec) + return rec + last = f"HTTP {status} errcode={code} {scrub(body)}" + last_code, last_status = code, status + # 4xx(除 45009 限流)不重试:我们的请求本身有问题 + if status and 400 <= status < 500 and status != 429 and code != 45009: + break + except Exception as e: + last = scrub(e) + if attempt <= len(_RETRY_DELAYS): + time.sleep(_RETRY_DELAYS[attempt - 1]) + rec = {"ts": t0, "hook_id": hook.get("id"), "hook_name": name, + "event": msg.get("event"), "title": msg.get("title"), "ok": False, + "url": hook.get("url", ""), "http_status": last_status, + "errcode": last_code, # 带上平台错误码(45009=限流,UI 要能看出来) + "elapsed_ms": int((time.time() - t0) * 1000), + "attempts": len(_RETRY_DELAYS) + 1, "error": last[:300]} + _record(rec) + return rec + + +def get_logs(limit=50): + """最近的发送记录(内存环形缓冲,重启清空)。""" + with _ring_lock: + rows = list(_ring)[-max(1, min(int(limit or 50), RING_SIZE)):] + out = [] + for r in reversed(rows): + out.append({k: v for k, v in r.items() if k != "url"} | {"url": mask_url(r.get("url", ""))}) + return out + + +# ================== 限流 / 聚合 ================== +def _rate_allow(hook): + """每 hook 一个令牌桶;返回 False 表示本次该被折叠。""" + rate = int(hook.get("rate_limit_per_min") or 0) + if rate <= 0: + return True + cap = min(rate, 5) + now = time.time() + b = _buckets.get(hook["id"]) + if b is None: + b = {"tokens": float(cap), "ts": now} + _buckets[hook["id"]] = b + b["tokens"] = min(cap, b["tokens"] + (now - b["ts"]) * rate / 60.0) + b["ts"] = now + if b["tokens"] >= 1: + b["tokens"] -= 1 + return True + return False + + +def _fold(hook, ev, msg): + f = _folds.setdefault(hook["id"], {"count": 0, "by_event": Counter(), + "last": time.time(), "ev": ev, "hook": hook}) + f["count"] += 1 + f["by_event"][ev.key] += 1 + if time.time() - f["last"] >= FOLD_INTERVAL: + _flush_fold(hook["id"]) + + +def _flush_fold(hook_id): + f = _folds.get(hook_id) + if not f or not f["count"]: + return + ev, hook = f["ev"], f["hook"] + merged = {"merged_count": f["count"], "by_event": dict(f["by_event"]), "folded": 0} + fields = {"_merged": {"merged_count": f["count"], "by_event": f["by_event"], + "folded": f["count"]}} + msg = build_message(ev, fields, hook) + msg["title"] = f"通知被限流折叠:{f['count']} 条" + msg["level"] = "warning" + msg["markdown"] = ("### ⚠️ 通知被限流折叠\n" + + "".join(f"> **{k}**:{v}\n" for k, v in f["by_event"].most_common(8)) + + f"> **时间**:{_fmt_ts()}\n" + f"> 完整记录见 系统 → 通知 → 发送记录 / logs/notify.log") + _enqueue_send(hook, msg, bypass_limit=True) + _folds[hook_id] = {"count": 0, "by_event": Counter(), "last": time.time(), + "ev": ev, "hook": hook} + + +def _enqueue_send(hook, msg, bypass_limit=False): + try: + _send_q.put_nowait((hook, msg, bypass_limit)) + except queue.Full: + _dropped["send_q"] += 1 + + +def _dispatch_one(hook, ev, fields): + """聚合窗口到期后真正投递一条(先过限流)。""" + if _rate_allow(hook): + _enqueue_send(hook, build_message(ev, fields, hook)) + else: + _fold(hook, ev, build_message(ev, fields, hook)) + + +def _dispatcher_loop(): + while not _stop.is_set(): + try: + item = _event_q.get(timeout=0.5) + except queue.Empty: + item = None + if item is not None and item is not _SENTINEL: + hook_id, event_key, fields, agg_key_vals, window = item + ev = notify_events.get(event_key) + hook = _hook_by_id(hook_id) + if ev and hook: + if window <= 0: + _dispatch_one(hook, ev, fields) + else: + k = (hook_id, event_key, agg_key_vals) + a = _agg.get(k) + if a is None: + _agg[k] = {"deadline": time.time() + window, "count": 1, + "fields": fields, "samples": [_sample(ev, fields)], + "ev": ev, "hook": hook} + else: + a["count"] += 1 + if len(a["samples"]) < 3: + a["samples"].append(_sample(ev, fields)) + # 到期的聚合窗口 + now = time.time() + for k in [k for k, a in list(_agg.items()) if a["deadline"] <= now]: + a = _agg.pop(k, None) + if not a: + continue + fields = dict(a["fields"]) + fields["_merged"] = {"merged_count": a["count"], + "samples": a["samples"]} + with _cfg_lock: + configured = bool(_config and _config["settings"]["global_enabled"]) + if configured: + _dispatch_one(a["hook"], a["ev"], fields) + # 折叠摘要 + for hid in list(_folds.keys()): + f = _folds.get(hid) + if f and f["count"] and time.time() - f["last"] >= FOLD_INTERVAL: + _flush_fold(hid) + # 队列丢弃提示(每 60s 最多一次) + if (_dropped["event_q"] or _dropped["send_q"]) and int(now) % 60 == 0: + n1, n2 = _dropped["event_q"], _dropped["send_q"] + _dropped["event_q"] = _dropped["send_q"] = 0 + _log.warning("通知队列溢出丢弃: event_q=%s send_q=%s(60s 内)", n1, n2) + _log.info("通知 dispatcher 退出") + + +def _sample(ev, fields): + subj = _subject(ev, fields) + reason = fields.get("cause") or fields.get("error") or fields.get("error_msg") or "" + return f"{subj}{' · ' + str(reason)[:40] if reason else ''}" + + +def _sender_loop(): + while not _stop.is_set(): + try: + hook, msg, bypass = _send_q.get(timeout=0.5) + except queue.Empty: + continue + if hook is _SENTINEL: + break + try: + timeout = _HTTP_TIMEOUT + with _cfg_lock: + if _config: + t = int(_config["settings"].get("http_timeout") or 0) + if t > 0: + timeout = (min(3, t), t) + send_now(hook, msg, timeout=timeout) + except Exception as e: # 兜底:sender 绝不能死 + _log.warning("通知发送线程异常: %s", scrub(e)) + _log.info("通知 sender 退出") + + +def _hook_by_id(hook_id): + with _cfg_lock: + if not _config: + return None + if _config["settings"]["global_enabled"] is False: + return None + for h in _config["webhooks"]: + if h["id"] == hook_id: + return h if h.get("enabled") else None + return None + + +# ================== 对外入口 ================== +def notify(event, **fields): + """发一条通知。**永不阻塞、永不抛异常、不碰 DB。** + + 调用方(任务/设备/安装/备份…)直接调即可;内部只做内存匹配 + 入队。 + ⚠ 调用点务必在 `with self._lock` **之外**。 + """ + try: + if not _ready or _config is None: + return + ev = notify_events.get(event) + if ev is None: + _log.warning("未知通知事件: %s(请在 core/notify_events.py 登记)", event) + return + with _cfg_lock: + cfg = _config + if not cfg["settings"]["global_enabled"]: + return + hooks = [h for h in cfg["webhooks"] + if h.get("enabled") and any(notify_events.match(p, event) + for p in h["events"])] + if not hooks: + return + for h in hooks: + keys = tuple(str(fields.get(k, "")) for k in (ev.agg_key or ())) + try: + _event_q.put_nowait((h["id"], event, dict(fields), keys, + int(h.get("agg_window") or 0))) + except queue.Full: + _dropped["event_q"] += 1 + except Exception as e: # 通知出问题绝不能影响业务 + try: + _log.warning("notify(%s) 内部异常: %s", event, e) + except Exception: + pass + + +def send_test(hook, event="notify.test", operator=""): + """「发送测试」按钮:同步发一条(不占业务令牌桶,单独限流)。""" + now = time.time() + while _test_hits and now - _test_hits[0] > 60: + _test_hits.popleft() + if len(_test_hits) >= TEST_PER_MIN: + return {"ok": False, "error": f"测试发送太频繁(每分钟最多 {TEST_PER_MIN} 次)"} + _test_hits.append(now) + ev = notify_events.get(event) or notify_events.get("notify.test") + fields = {"hook_name": hook.get("name", ""), "operator": operator} + msg = build_message(ev, fields, hook) + return send_now(hook, msg) + + +def preview(hook_cfg, event="notify.test"): + """保存前预览:真实请求体 + UTF-8 字节数 + 是否截断。""" + hook = _normalize_hook(hook_cfg, dict(_DEFAULT_SETTINGS)) + ev = notify_events.get(event) or notify_events.get("notify.test") + msg = build_message(ev, {"hook_name": hook["name"], "operator": "预览"}, hook) + out = {"event": ev.key, "title": msg["title"], "markdown": msg["markdown"]} + adapter = ADAPTERS.get(hook["format"]) + if adapter is None: + out["error"] = f"格式 {hook['format']} 尚未实现" + return out + try: + kw = adapter.render(msg, hook) + body = json.dumps(kw.get("json"), ensure_ascii=False) + limit = adapter.byte_limit or 0 + out.update(request_body=body, bytes=len(body.encode("utf-8")), + byte_limit=limit, + truncated=bool(limit and len(body.encode("utf-8")) > limit)) + except Exception as e: + out["error"] = scrub(e) + return out + + +# ================== 生命周期 ================== +_SENTINEL = object() + + +def init_app(app): + """启动时注入 app 并拉起后台线程。**必须早于 `consume_pending_restore()`**, + 否则"恢复完成"这类启动期事件会丢(notify 在未就绪时静默丢弃)。""" + global _app, _ready + _app = app + try: + reload_config() + except Exception as e: + _log.warning("初始化通知配置失败(先按空配置跑): %s", scrub(e)) + if _ready: + return + _ready = True + for i in range(SENDER_THREADS): + t = threading.Thread(target=_sender_loop, name=f"notify-sender-{i}", daemon=True) + t.start() + _threads.append(t) + t = threading.Thread(target=_dispatcher_loop, name="notify-dispatcher", daemon=True) + t.start() + _threads.append(t) + _log.info("通知模块已启动(%d 条 webhook,%d 个发送线程)", + len((get_config().get("webhooks") or [])), SENDER_THREADS) + + +def shutdown(timeout=3.0): + """退出时收尾:尽量把在途通知发完(最多等 timeout 秒)。""" + try: + for hid in list(_folds.keys()): + _flush_fold(hid) + for k in list(_agg.keys()): + a = _agg.pop(k, None) + if a: + fields = dict(a["fields"]) + fields["_merged"] = {"merged_count": a["count"], "samples": a["samples"]} + _dispatch_one(a["hook"], a["ev"], fields) + deadline = time.time() + timeout + while not _send_q.empty() and time.time() < deadline: + time.sleep(0.1) + except Exception: + pass + _stop.set() + try: + _event_q.put_nowait(_SENTINEL) + for _ in range(SENDER_THREADS): + _send_q.put_nowait((_SENTINEL, None, True)) + except Exception: + pass + for t in _threads: + try: + t.join(timeout=1.0) + except Exception: + pass diff --git a/core/notify_events.py b/core/notify_events.py new file mode 100644 index 0000000..5c00803 --- /dev/null +++ b/core/notify_events.py @@ -0,0 +1,203 @@ +"""通知事件目录(**唯一真相**)。 + +约定: +- **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 diff --git a/core/system_backup.py b/core/system_backup.py index 66c01dc..d07cd52 100644 --- a/core/system_backup.py +++ b/core/system_backup.py @@ -36,6 +36,7 @@ from config import (DATA_DIR, APK_DIR, BACKUP_DIR, RESTORE_STAGING_DIR, RESTORE_PENDING_DIR) from core.logger import get_logger from core.models import CURRENT_SCHEMA_VERSION, db +from core import notifier _log = get_logger("core.backup") @@ -547,6 +548,8 @@ def consume_pending_restore(): shutil.move(pending_db, os.path.join(fail_dir, "users.db")) remove_quiet(RESTORE_PENDING_DIR) _log.error(f"恢复任务校验失败,已搁置到 {fail_dir},继续使用当前库: {msg}") + notifier.notify("system.backup.restore_failed", error=str(msg)[:200], + fail_dir=fail_dir) return False try: @@ -569,6 +572,9 @@ def consume_pending_restore(): remove_quiet(RESTORE_PENDING_DIR) schema_version = msg.get("schema_version", "?") if isinstance(msg, dict) else "?" _log.info(f"已应用备份恢复(单事务替换): {applied},schema v{schema_version},apks 已合并") + # 通知在"下一次启动"才发得出去(恢复本身是重启生效的),文案里说明白 + notifier.notify("system.backup.restored", applied_rows=applied, + schema_version=schema_version) return True diff --git a/core/task_manager.py b/core/task_manager.py index 7113984..5a68894 100644 --- a/core/task_manager.py +++ b/core/task_manager.py @@ -21,6 +21,7 @@ import json import time import uuid import threading +from collections import Counter from datetime import datetime from concurrent.futures import ThreadPoolExecutor, as_completed @@ -31,6 +32,7 @@ from config import DATA_DIR from core.logger import get_logger from core.models import db, DeviceGroup as GroupRow, TaskJob as JobRow from core import device_pool +from core import notifier from .adb_helper import get_foreground_app, get_foreground_app_remote from .device_worker import ( DeviceOfflineError, @@ -47,6 +49,42 @@ _log = get_logger("core.tm") _START_STAGGER_SEC = 0.2 +class _BatchTracker: + """一次任务批次的收尾统计:**所有设备都出结果**后发一条 `task.batch.finished`。 + + 为什么不每台设备各发一条"批次结束":100 台设备就是 100 条刷屏。这里只计数, + 归零那一刻发一条汇总(锁外发,见 done())。 + + 每台设备的结果由 `_run_with_retry` 的 finally 调一次 `done(serial, outcome)`—— + 放在 finally 里是为了保证"任何一个 return 分支都会被计数一次",不会漏也不会重。 + """ + + def __init__(self, job, total, names=None): + self.job = job + self.total = total + # serial -> 设备名(通知里显示名字而不是裸地址);取不到就是空 + self.names = dict(names or {}) + self._lock = threading.Lock() + self._left = total + self._stats = Counter() + self._t0 = time.time() + + def done(self, serial, outcome="failed"): + """登记一台设备的结果(success/failed/stopped/skipped)。""" + with self._lock: + self._left -= 1 + self._stats[outcome] += 1 + left = self._left + s = dict(self._stats) + if left > 0: + return + notifier.notify("task.batch.finished", job_id=self.job.id, + job_name=self.job.name, total=self.total, + success=s.get("success", 0), failed=s.get("failed", 0), + stopped=s.get("stopped", 0), skipped=s.get("skipped", 0), + duration_s=int(time.time() - self._t0)) + + def _in_run_window(schedule, now=None): """是否在当前运行窗口内。 @@ -633,6 +671,8 @@ class TaskManager: w.stop() stopped.append(serial) _log.info(f"定时停止 {job.name}({job.id}): 停止 {len(stopped)} 台设备 {stopped}") + notifier.notify("task.cron.stopped", job_id=job.id, job_name=job.name, + stopped_count=len(stopped), serials=stopped[:10]) def next_run_of(self, job): """任务下次真正执行的时间(考虑运行窗口)。停用/无 cron 返回 None。""" @@ -664,23 +704,42 @@ class TaskManager: serials = job.resolve_serials(self) if not serials: _log.warning(f"任务 {job.name} 无可用设备") + notifier.notify("task.batch.no_device", job_id=job.id, job_name=job.name, + target=job.target) return task_cls = get_task_class(job.task_type) if not task_cls: _log.error(f"未知任务类型: {job.task_type}") + notifier.notify("task.batch.unknown_type", job_id=job.id, job_name=job.name, + task_type=job.task_type) return task = task_cls() max_attempts = max(1, job.retry.get("max_attempts", 1)) delay = job.retry.get("delay", 60) + notifier.notify("task.batch.started", job_id=job.id, job_name=job.name, + task_type=job.task_type, device_count=len(serials), + serials_preview=serials[:5]) + # 批次收尾统计:所有设备都出结果后发一条汇总(见 _BatchTracker) + # 设备名一次查好带下去(worker 线程里没有 app context,查库要显式包 context) + names = {} + try: + with self._db(): + names = {d["serial"]: (d.get("name") or "") + for d in device_pool.list_devices()} + except Exception as e: + _log.warning(f"读取设备名失败(通知里将显示地址): {e}") + tracker = _BatchTracker(job, len(serials), names) + for idx, serial in enumerate(serials): # 每台设备一个重试循环线程,互不影响;错峰延迟在各自线程内等待 t = threading.Thread(target=self._run_with_retry, args=(task, serial, job, max_attempts, delay, - idx * _START_STAGGER_SEC), daemon=True) + idx * _START_STAGGER_SEC, tracker), daemon=True) t.start() - def _run_with_retry(self, task, serial, job, max_attempts, delay, start_delay=0): + def _run_with_retry(self, task, serial, job, max_attempts, delay, start_delay=0, + tracker=None): """单设备任务执行 + 重试。 异常分类: @@ -691,36 +750,58 @@ class TaskManager: # 错峰启动:在各自线程内等待,分摊批量启动的连接/占用压力 if start_delay > 0: time.sleep(start_delay) + # 设备名(通知里显示"设备:A08"而不是裸地址) + dname = (tracker.names.get(serial) if tracker is not None else "") or serial # 被抢占的原任务 id(抢占结束后归还)——必须在循环外: # 重试时若重置为 None,finally 归还逻辑会丢失信息,被抢占任务永不恢复 preempted_job = None + # 本台设备的结果(供批次统计与通知用)。默认 "failed":走到重试耗尽就是失败; + # 各提前 return 的分支会覆盖它。**在 finally 里统一上报**,保证不漏不重。 + outcome = "failed" try: for attempt in range(1, max_attempts + 1): # 用户已请求停止 → 不再启动新 attempt + stopped_by_user = False with self._lock: if serial in self._stop_requested: _log.info(f"{serial} 用户已请求停止,取消重试 (job={job.name})") - return + stopped_by_user = True + if stopped_by_user: + outcome = "stopped" + notifier.notify("task.device.stopped", serial=serial, device_name=dname, + job_id=job.id, job_name=job.name, attempt=attempt, + phase="start") + return # 同一 serial 同时只能一个 worker;不同任务可配置抢占 preempt = False + skip_reason = "" with self._lock: if serial in self._running: cur = self._running[serial] if cur.get("job_id") == job.id: # 同一任务重复触发:跳过(原行为) _log.warning(f"{serial} 已有任务在跑,跳过 (job={job.name})") - return - if not job.params.get("preempt"): + skip_reason = "同一任务已在运行" + elif not job.params.get("preempt"): # 不同任务且本任务未开启抢占:跳过 _log.warning(f"{serial} 已有任务在跑,跳过 (job={job.name})") - return - preempt = True - if not preempted_job: - preempted_job = cur.get("job_id") + skip_reason = "设备被其他任务占用" + else: + preempt = True + if not preempted_job: + preempted_job = cur.get("job_id") + if skip_reason: + outcome = "skipped" + return if preempt: # 抢占:锁外停止该设备上的其他任务(stop_device 内部拿同一把锁, # 在锁内调用会死锁!),等其释放后接管 _log.warning(f"{serial} 抢占:停止任务 {preempted_job} 后执行 {job.name}") + notifier.notify("task.device.preempted", serial=serial, device_name=dname, + job_id=job.id, job_name=job.name, + preempted_job_id=preempted_job, + preempted_job_name=(self.jobs.get(preempted_job).name + if self.jobs.get(preempted_job) else "")) self.stop_device(serial) deadline = time.time() + 30 while time.time() < deadline: @@ -730,6 +811,10 @@ class TaskManager: time.sleep(0.5) else: _log.warning(f"{serial} 抢占超时(旧任务 30s 未退出),跳过 (job={job.name})") + outcome = "skipped" + notifier.notify("task.device.preempt_timeout", serial=serial, + device_name=dname, job_id=job.id, job_name=job.name, + preempted_job_id=preempted_job) return with self._lock: self._running[serial] = {"job_id": job.id, "started_at": time.time(), @@ -739,6 +824,7 @@ class TaskManager: max_attempts=max_attempts) worker = None + t_attempt = time.time() # 单次 attempt 耗时(通知里带) try: worker = task.create_worker(serial, job.params) with self._lock: @@ -756,11 +842,22 @@ class TaskManager: _update_status(serial, task_job="") if st == "done": _log.info(f"{serial} 任务 {job.name} 成功完成") + outcome = "success" + notifier.notify("task.device.success", serial=serial, + device_name=dname, job_id=job.id, + job_name=job.name, attempt=attempt, + duration_s=int(time.time() - t_attempt), + is_time_up=bool(worker.is_time_up())) return # 用户请求停止(无论 attempt 第几次、status 是什么)→ 不重试 with self._lock: if serial in self._stop_requested: _log.info(f"{serial} 用户已请求停止,不再重试 (job={job.name})") + outcome = "stopped" + notifier.notify("task.device.stopped", serial=serial, + device_name=dname, job_id=job.id, + job_name=job.name, attempt=attempt, + phase="after") return _log.warning(f"{serial} 任务未成功(status={st})") except DeviceOfflineError as e: @@ -770,30 +867,46 @@ class TaskManager: self._running.pop(serial, None) _update_status(serial, status="failed", task_job="", last_error=f"设备离线: {e}") + outcome = "failed" + notifier.notify("task.device.offline", serial=serial, device_name=dname, + job_id=job.id, job_name=job.name, error=str(e)[:200]) return except Exception as e: _log.error(f"{serial} 执行异常: {e}", exc_info=True) with self._lock: self._running.pop(serial, None) _update_status(serial, task_job="") + notifier.notify("task.device.error", serial=serial, device_name=dname, + job_name=job.name, attempt=attempt, + error=str(e)[:200]) if attempt < max_attempts: # 用户在 sleep 期间点停止也能中断 with self._lock: if serial in self._stop_requested: _log.info(f"{serial} 用户已请求停止,取消重试 (job={job.name})") + outcome = "stopped" + notifier.notify("task.device.stopped", serial=serial, + device_name=dname, job_name=job.name, + attempt=attempt, phase="retry-wait") return # 检查是否是端口耗尽类临时错误,需要更长退避等端口释放 with _WORKERS_LOCK: err = _WORKERS.get(serial, {}).get("last_error", "") - if err.startswith("[transient]"): + transient = err.startswith("[transient]") + if transient: # Windows TCP 端口耗尽,TIME_WAIT 默认 2-4 分钟,等 120 秒 extra_delay = max(delay, 120) _log.info(f"{serial} ADB 连接临时错误(端口耗尽),{extra_delay}s 后重试 ({attempt+1}/{max_attempts})") - time.sleep(extra_delay) + wait = extra_delay else: _log.info(f"{serial} {delay}s 后重试 ({attempt+1}/{max_attempts})") - time.sleep(delay) + wait = delay + notifier.notify("task.device.retry", serial=serial, device_name=dname, + job_name=job.name, attempt=attempt, + next_attempt=attempt + 1, delay_s=int(wait), + reason="transient" if transient else "normal") + time.sleep(wait) _log.error(f"{serial} 任务 {job.name} 重试耗尽,放弃") # 带上最后一次的真实失败原因——只写"重试N次失败"会让用户看不到为什么失败 @@ -805,20 +918,37 @@ class TaskManager: if cause and cause != msg: msg += f":{cause}" _update_status(serial, status="failed", last_error=msg[:200]) + outcome = "failed" + notifier.notify("task.device.failed", serial=serial, device_name=dname, + job_id=job.id, job_name=job.name, attempts=max_attempts, + cause=cause, msg=msg) finally: # 清除停止标志:整个重试循环结束(成功/失败/停止)后允许下次任务 with self._lock: self._stop_requested.discard(serial) + # 批次统计:**每台设备只在这里上报一次**(所有 return 分支都会走到 finally) + if tracker is not None: + try: + tracker.done(serial, outcome) + except Exception as e: + _log.warning(f"批次统计上报失败(不影响任务): {e}") # 归还:本任务是抢占任务,结束后自动重新启动被抢占的原任务 if preempted_job: j = self.jobs.get(preempted_job) if j and j.enabled: _log.info(f"{serial} 抢占任务 {job.name} 结束,归还设备给任务 {j.name}") + notifier.notify("task.device.released", serial=serial, device_name=dname, + job_name=job.name, preempted_job_id=preempted_job, + returned=True, reason="归还并重启被抢占任务") threading.Thread(target=self._run_job, args=(j,), daemon=True).start() else: # 打日志便于排查:被抢占任务已删除/停用时不会归还,但要知道原因 _log.info(f"{serial} 抢占任务 {job.name} 结束," f"被抢占任务 {preempted_job} {'已停用,不归还' if j else '已不存在,不归还'}") + notifier.notify("task.device.released", serial=serial, device_name=dname, + job_name=job.name, preempted_job_id=preempted_job, + returned=False, + reason="被抢占任务已停用" if j else "被抢占任务已不存在") # ---- 运行控制 ---- def stop_device(self, serial): diff --git a/doc/API.md b/doc/API.md index 193bfd4..3694fda 100644 --- a/doc/API.md +++ b/doc/API.md @@ -238,7 +238,23 @@ | POST | `/api/system/backup/preview` | 上传备份校验预览(multipart,字段 `file`) | | POST | `/api/system/backup/apply` | 应用恢复(重启生效) | -### 2.11 device_agent(`web/device_agent_api.py`) +### 2.11 notify(`web/notify_api.py`,全部 Admin) + +通知 / Webhook 配置与发送记录(事件目录、推送格式、限流语义见 [NOTIFY.md](NOTIFY.md)): + +| 方法 | 路径 | 功能 | +|------|------|------| +| GET | `/api/notify/webhooks` | 全部 webhook(`url` 打码、`secret` 只回 `secret_set`)+ 格式清单 | +| POST | `/api/notify/webhooks` | 新建(body 见 NOTIFY.md §4) | +| PUT | `/api/notify/webhooks/` | 更新(`url`/`secret` 省略或为打码值 → **保持原值**) | +| DELETE | `/api/notify/webhooks/` | 删除 | +| POST | `/api/notify/webhooks//test` | **同步**发一条测试消息(不占业务令牌桶,单独限 10 次/分) | +| POST | `/api/notify/preview` | 预览真实请求体 + UTF-8 字节数 + 是否截断 | +| GET | `/api/notify/events` | 事件目录(含字段清单) | +| GET | `/api/notify/logs?limit=` | 最近发送记录(内存 200 条,重启清空) | +| POST | `/api/notify/settings` | 全局开关与默认聚合/限流 | + +### 2.12 device_agent(`web/device_agent_api.py`) **设备端专用**(无登录会话,靠 `X-Device-Token` 鉴权;平台未启用时统一 404): @@ -321,7 +337,7 @@ | `POST /api/stop_all` | — | 停止全部 | | `POST /api/device/clear_error` | `{"serial":"…"}` | 运行中/重试等待中拒绝清理 | | `POST /api/device/clear_all_errors` | — | 跳过 running/connecting | -| `POST /api/device/screen_all` | `{"mode":"on"\|"off", "serials":[…]}` | 不传 serials 则对全部在线设备 | +| `POST /api/device/screen_all` | `{"mode":"on"\|"off", "serials":[…]}` | 不传 serials 则对全部在线设备。**熄屏是单向的**(`KEYCODE_SLEEP`),不会把已息屏的设备唤醒;亮屏用 `KEYCODE_WAKEUP` + `dismiss-keyguard` | | `POST /api/scan_foreground` | — | 已扫描中时返回 `{"ok":false}`(**HTTP 200**) | | `POST /api/device/locate` | `{"serial":"…","show":true}` | `show=true` 时在设备上打开 `/locate` | | `POST /api/device/locate/stop` | `{"serial":"…"}` | 结束定位 | diff --git a/doc/ARCHITECTURE.md b/doc/ARCHITECTURE.md index ac0b12a..2a6a8ff 100644 --- a/doc/ARCHITECTURE.md +++ b/doc/ARCHITECTURE.md @@ -30,6 +30,7 @@ │ ┌───────────────────────────▼──────────────────────────────────────────┐ │ 基础层 adb_helper · u2_helper · uiauto_helper · ocr · clipboard │ +│ notifier(通知分发:队列/聚合/限流/适配器) │ │ models(SQLite)· logger · config │ └──────────────────────────────────────────────────────────────────────┘ ▲ @@ -109,6 +110,8 @@ | `apk-install` | `ApkManager.install` | 并发 5 台安装 APK | 安装期,全局单任务 | | Agent 执行线程 | `web/agent_api.py` | AI 控制台一轮会话 | 按需 | | 巡检手动线程 | `web/agent_api.py` | 手动触发巡检 | 按需 | +| **通知 dispatcher** | `core/notifier.init_app` | 通知聚合 + 每 hook 限流 + 折叠摘要 | 常驻 1 个 | +| **通知 sender ×3** | 同上 | 真实发 webhook(退避重试、环形记录) | 常驻 3 个 | 两个 APScheduler 相互独立,时区均固定 `Asia/Shanghai`。 @@ -391,5 +394,7 @@ connecting ──获取设备──▶ u2 连接 ──▶ running ──▶ set | 加一张表 | `core/models.py` 模型 + `SCHEMA_MIGRATIONS`(老库) | **登记进备份覆盖清单** + [DATA_MODEL.md](DATA_MODEL.md)(红线) | | 加一个常驻线程 | 参考 `device_discovery` 的 `init_app` / `shutdown` 模式 | 更新本文 §3 与 [DEVELOPMENT.md](DEVELOPMENT.md) | | 加一个前端模块 | `static/admin/.js` + 在 `monitor.html` 按序引入 | 更新本文 §6 与 [DEVELOPMENT.md](DEVELOPMENT.md) | +| **加一个通知事件** | `core/notify_events.py` 的 `EVENTS` 加一条 + 触发点调 `notifier.notify(key, **fields)` | 事件目录要做全、默认不启用;**notify 必须放在锁外**(见 [NOTIFY.md](NOTIFY.md) §7) | +| 加一种通知格式 | `core/notifier.py` 的 `ADAPTERS` 注册一个 `BaseAdapter` 子类 | 补字节上限/默认限流;更新 [NOTIFY.md](NOTIFY.md) §5 | | 加一个 MCP 工具 | `mcp_server/mcp_server.py`(需要时加 `direct_ops.py`) | 写操作要挂门控三连;更新 [MCP.md](MCP.md) | | 加一个配置键 | `config.py`(程序级)或 `.env`(密钥类) | 更新 `.env.example` + [DEVELOPMENT.md](DEVELOPMENT.md) 配置速查 | diff --git a/doc/DATA_MODEL.md b/doc/DATA_MODEL.md index 4e99e47..e0dc455 100644 --- a/doc/DATA_MODEL.md +++ b/doc/DATA_MODEL.md @@ -240,8 +240,12 @@ UTF-8 等价于字节序)。 | `deployment_env` | **库环境标签**(`dev`/`prod`),启动时与 `.env` 比对 | `core/db_config.py`(首次连接)/ 迁移脚本 | | `deployment_id` | 库唯一标识(uuid),用于识别"这份备份来自哪个库" | 同上 | | `deployment_claimed_at` | 标签写入时间 | 同上 | +| `notify_webhooks` | **通知 / Webhook 全部配置**(JSON:`{version, settings, webhooks[]}`,见 [NOTIFY.md](NOTIFY.md) §4)。**不建表**——新增/删除 webhook 都只改这一个键 | 系统 → 通知 页 | > ⚠️ `agent_api_key` 是**明文存储**,导出备份的 zip 里也含它——备份预览会固定给出"含敏感信息"告警。 +> **同理 `notify_webhooks` 里的 webhook URL 本身就是凭据**(企业微信 `?key=`、钉钉 `?access_token=`、 +> 飞书 `/hook/`):拿到它就能往群里发消息。接口回显/发送记录/日志一律走 +> `notifier.mask_url()/scrub()` 打码,导出备份时也按敏感信息对待。 > `app_meta` 的列名 `key` 在 MySQL 里是保留字,**不要直接拼裸 SQL**,统一走 > `core/db_config.meta_get / meta_set`(方言中立、自动加引号)。 diff --git a/doc/NOTIFY.md b/doc/NOTIFY.md new file mode 100644 index 0000000..0c59fbc --- /dev/null +++ b/doc/NOTIFY.md @@ -0,0 +1,193 @@ +# 通知 / Webhook(NOTIFY) + +> 面向:要给平台接告警的运维、以及往各组件加通知点的开发。 +> 相关:[API.md](API.md)(接口)、[DATA_MODEL.md](DATA_MODEL.md)(配置存哪)、 +> [ARCHITECTURE.md](ARCHITECTURE.md)(线程模型)。 + +--- + +## 1. 它是什么 + +平台各组件(任务、设备、安装、备份、AI 巡检…)在关键时刻调一个统一入口: + +```python +from core import notifier +notifier.notify("task.device.failed", serial="192.168.20.71:5555", + device_name="A02", job_name="抖音养号", cause="选择器连续 10 次未命中") +``` + +`notify()` **只做内存操作**(读配置快照 → 匹配订阅 → 丢进队列),真正发 HTTP 的是后台 +daemon 线程。所以: + +- 任务线程里可以直接调,**不用包 app_context、不用 try/except**(内部全兜住了) +- **但必须放在所有 `with self._lock` 之外**——别让通知拖住调度锁 +- 通知模块自己出问题(地址写错、对方挂了、配置坏了)**绝不影响任务/设备/备份**, + 最多在日志里留一条警告 + +``` +业务线程 notify() ──入队──▶ dispatcher(1 线程)──▶ sender ×3 ──▶ 企业微信/自建服务 + ① 聚合 ② 限流 ③ 折叠 + └─▶ 发送记录(内存 200 条 + logs/notify.log) +``` + +--- + +## 2. 快速上手(企业微信) + +1. 企业微信群里 → 群机器人 → 添加 → 复制 Webhook 地址(形如 + `https://qyapi.weixin.qq.com/cgi-bin/webhook/send?key=xxxx`) +2. 后台 → **系统 → 通知 / Webhook → + 新建 Webhook**:名称随便起、格式选「企业微信」、 + 粘贴地址、在事件树里勾上关心的事件(★ 是建议开的) +3. 点该行的 **发送测试** —— 群里立刻能收到一条 markdown;收不到就看返回的 HTTP/errcode + +--- + +## 3. 事件目录 + +**权威定义在 `core/notify_events.py`**(后台「通知」页也是从它渲染的,不会漂移)。 +所有事件**默认不启用**——登记不等于推送,勾了才发。 + +| 类别 | 事件 | +|---|---| +| 任务批次 | `task.batch.started` · `.finished` · `.no_device` · `.unknown_type` · `task.cron.stopped` | +| 任务·单设备 | `task.device.success` · `.failed` · `.offline` · `.error` · `.retry` · `.stopped` · `.preempted` · `.preempt_timeout` · `.released` | +| 业务 | `task.selector.invalid`(选择器连续 10 次未命中——"任务成功但什么都没做"的隐蔽故障) | +| Worker | `worker.connected` · `.attempt.done` · `.attempt.error`(单次尝试级,噪音大,默认没人勾) | +| 设备 | `device.online` · `.offline` · `.discovered` · `.claimed` · `device.heartbeat_timeout` | +| 安装 | `apk.install.started` · `.finished` | +| 系统 | `system.backup.exported` · `.imported` · `.restored` · `.restore_failed` · `service.started` · `.stopping` · `user.login` | +| AI | `ai.audit.finished` · `.failed` · `.skipped` | +| 其它 | `notify.test`(测试按钮专用) | + +订阅支持通配:`*`(全部)、`task.*`、`device.*`。 + +> 语义分工(**避免重复告警**):`worker.*` 是**单次尝试**层面;`task.device.success/failed` +> 是"这台设备最终成功/失败"的**唯一权威点**。所以一个设备重试 3 次后失败,只会收到 **1 条** +> `task.device.failed`,不会收到 3 条噪音。 + +--- + +## 4. 配置 + +存在 `app_meta.notify_webhooks` 一个键里(JSON,**不建表**,因此不涉及备份覆盖清单): + +```json +{ + "version": 1, + "settings": {"global_enabled": true, "default_agg_window": 30, + "default_rate_limit": 18, "log_keep": 200, "http_timeout": 5}, + "webhooks": [ + {"id": "wh_ab12cd34", "name": "运维群", "enabled": true, "format": "wecom", + "url": "https://qyapi.weixin.qq.com/cgi-bin/webhook/send?key=…", + "secret": "", "events": ["task.device.failed", "device.*"], + "agg_window": 30, "rate_limit_per_min": 18, + "title_template": "", "body_template": "", "headers": {}} + ] +} +``` + +上限(超了**拒绝保存**,不静默截断):webhook ≤ 20 条、整体 JSON ≤ 60000 字符、URL ≤ 2048。 + +**安全**:URL 里带凭据(企微 `?key=`、钉钉 `?access_token=`、飞书 `/hook/`、 +Bark `/…/`、Slack `/services/T…/B…/X…`)→ 接口回显、发送记录、日志一律打码 +(`mask_url` 同时处理 query 与**路径里的 token**;异常消息过 `scrub`); +编辑时**留空即不修改**;`secret` 永不回显(只回"已配置")。 + +--- + +## 5. 推送格式 + +| 格式 | 请求体 | URL 怎么填 | 成功判定 / 关键约束 | +|---|---|---|---| +| `wecom` 企业微信 | `{"msgtype":"markdown","markdown":{"content":"…"}}` | 群机器人 → 复制的 Webhook 地址(含 `?key=`) | **`errcode==0`**;content **≤4096 字节**(按字节截断、不会截出半个汉字);**每机器人每分钟 20 条**,超限 `45009` | +| `dingtalk` 钉钉 | `{"msgtype":"markdown","markdown":{"title","text"}}` | 群设置 → 智能群助手 → 添加机器人 → 自定义 → 复制的 Webhook 地址(含 `?access_token=`) | **`errcode==0`**;正文 ≤20000 字节;官方 20 条/分,**超限会被限流 10 分钟**(这里默认 15 留余量);机器人安全设置选「加签」时密钥必填 | +| `feishu` 飞书 | `{"msg_type":"interactive","card":{header + markdown 元素[,"timestamp","sign"]}}` | 群设置 → 群机器人 → 添加 → 自定义机器人 → 复制的 Webhook 地址(`…/bot/v2/hook/`) | **`code==0`**;请求体 ≤20KB;官方 100 条/分(这里默认 60);开了「签名校验」时密钥必填 | +| `bark` iOS 推送 | `{"title","body","markdown","group","level"[,"device_key"]}` | Bark App 里复制的那串(`https://api.day.app/`)**或** `https://api.day.app/push` + 设备 Key 填到「设备 Key」 | **`code==200`**(注意和企业微信不一样!);走 APNs,正文按 2048 字节截断;`markdown` 传富文本、`body` 传纯文本兜底 | +| `json` 通用 / Slack | 由 `body_template` 决定 | 你自己的接收端;Slack 填它的 Incoming Webhook URL | HTTP 200(响应体里有 `code`/`errcode` 时必须为 0);模板保存前**干跑校验**;`secret` 会作为 `X-Webhook-Secret` 头发出 | + +> 成功码**按格式**判定(企微/钉钉/飞书 = 0,Bark = 200),写在适配器的 `ok_codes` 上—— +> 加新格式时别忘了一起定。 + +### 加签:钉钉和飞书**算法不一样**,别互相照抄 + +| | 钉钉 | 飞书 | +|---|---|---| +| key(HMAC 的密钥) | `secret` | `"{timestamp}\n{secret}"` | +| message(被签内容) | `"{timestamp}\n{secret}"` | **空** | +| timestamp 单位 | **毫秒** | **秒** | +| 拼在哪 | **URL query**(`×tamp=…&sign=…`) | **JSON body** 的 `timestamp`/`sign` | +| 结果编码 | Base64 后再 **urlencode** | Base64 | +| 出错码 | `errcode 310000 invalid signature` | `code 19021 sign match fail` | + +实现见 `core/notifier.py` 的 `DingtalkAdapter.sign_request` / `FeishuAdapter.sign_request`; +两家的算法都用官方示例代码对拍过(钉钉逐字节一致,含 urlencode)。 + +**字段一律渲染成 `> **字段**:值` 引用行,不用 Markdown 表格**——企业微信/钉钉的 markdown +子集不支持表格,表格会原样吐出来。 + +通用 JSON 的占位符:`{{event}} {{title}} {{summary}} {{ts}} {{level}} {{markdown}} +{{fields}} {{fields_json}} {{hook_name}} {{field.<字段名>}}`。 +替换值按 JSON 字符串转义,所以标题里带引号/换行也不会打坏请求体。Slack 直接写 +`{"text":"{{markdown}}"}` 就行。 + +**格式相关的提示文案(URL 示例/说明、密钥叫什么、限流上限)都挂在适配器上** +(`BaseAdapter.url_hint/url_help/secret_label/secret_help/limit_help`),界面按当前格式渲染、 +切格式即时更新——新增格式时把这些一起填上,别把某个平台的说明写死在页面上。 + +**再加新格式**:写一个 `BaseAdapter` 子类(`render` + 可选 `sign_request`)并把元信息 +(`byte_limit`/`ok_codes`/`limit_default`/`url_hint`/`url_help`/`secret_*`/`limit_help`)填全, +然后注册进 `ADAPTERS` —— 界面上的下拉项与提示文案会自动跟着出来,**不用改前端**。 +`PLANNED_FORMATS` 是"还没实现、下拉里置灰"的占位表(目前为空)。 + +--- + +## 6. 防打爆(四层) + +100 台设备同时失败 = 100 条消息,群里会被刷到静音——**刷屏会让通知彻底失效**。所以: + +| 层 | 机制 | 默认 | +|---|---|---| +| L1 聚合 | 同 webhook、同事件、同聚合键(如 `job_id`)在一个窗口内合并成一条,保留前 3 个样本 | 30s(低频高危事件设 0,立即发) | +| L2 限流 | 每 webhook 一个令牌桶 | 18 条/分(企业微信硬限 20,留余量) | +| L3 折叠 | 被限流的事件**不丢弃**,压成一条「被限流折叠 N 条」摘要 | 最多 60s 一条 | +| L4 背压 | 有界队列(event 2000 / send 1000),满了丢弃并计数 | 溢出会告警一次 | + +取舍写明白:**失败通知最多延迟一个聚合窗口(默认 30s)**,换来群不被刷屏。 + +--- + +## 7. 开发:给新功能加通知 + +1. 在 `core/notify_events.py` 的 `EVENTS` 里加一条 `_e("模块.对象.动作", "中文标签", + "分类", ["字段1", "字段2"], "什么时候发", agg_window=…, recommend=…)` +2. 在触发点调 `notifier.notify("模块.对象.动作", 字段1=…, 字段2=…)` + —— **放在所有 `with self._lock` 之外**,且不要改变原有 `return` 的顺序 +3. 把事件补进本文档 §3 的表格 + +约定: + +- `notify()` **不阻塞、不抛异常、不碰 DB**——这是硬约束,别在它里面加 HTTP 或查库 +- 事件的 `fields` 是前端字段表与模板占位符的白名单,只写真正有用的 +- 高频事件(每台设备/每次尝试都会发生的)把 `agg_window` 设大一点或 `recommend=False` + +--- + +## 8. 排障 + +| 现象 | 看哪里 | +|---|---| +| 完全没收到 | ① 全局开关是不是关了(列表页「启用通知」)② 该 webhook 是否启用 ③ 事件勾了没(`task.device.failed` 是**单设备最终失败**,不是每次尝试) | +| 收到但内容不全 | 消息被 4096 字节截断了(末尾有「…(已截断)」);把 `cause` 之类长字段在事件侧截短 | +| 只在群里看到「被限流折叠」 | 短时间内同类事件太多,触发了 L2/L3;调大该 webhook 的限流值(企微上限 20)或调大聚合窗口 | +| 失败原因 | 后台「通知」页的**发送记录**(内存,重启清空)或 `logs/notify.log`(完整历史)。`errcode 45009`=企微限流、`93000`=Webhook 地址无效、`HTTP 200 + errcode≠0` **也算失败** | +| 配置坏了 | 服务照常启动(启动日志有 error),通知静默不发;在后台删掉坏配置或直接改 `app_meta.notify_webhooks` | + +--- + +## 9. 已知限制 + +- 发送记录在**内存**(最近 200 条,重启清空);持久历史只有 `logs/notify.log` 文本 +- **不支持自定义 webhook 请求头**(`headers` 字段留着但界面没暴露)——需要的话说一声 +- 只做 http/https,**不做内网 IP 黑名单**(内网自建 webhook 是合法用法),但禁止重定向 + (`allow_redirects=False`) +- `system.backup.restored` 是**重启后**才发(恢复本身就是重启生效的) diff --git a/doc/README.md b/doc/README.md index 8051a94..ac2299f 100644 --- a/doc/README.md +++ b/doc/README.md @@ -18,6 +18,7 @@ | [MCP_DESIGN.md](MCP_DESIGN.md) | **MCP 设计文档**:边界划分、错误码、白名单/审计设计、演进方向 | 平台开发者 | | [AI_CONSOLE.md](AI_CONSOLE.md) | **AI 控制台**:会话/SSE、经验库、动作库、巡检、Markdown 渲染、推理链、token 统计 | 使用者、平台开发者 | | [AI_TASK_GEN.md](AI_TASK_GEN.md) | **AI 建任务**:AI 自己在真机探索 → 写出可调度任务 → 人工确认入库(P0 已实现;契约与红线) | 平台开发者、使用者 | +| [NOTIFY.md](NOTIFY.md) | **通知 / Webhook**:事件目录、企业微信/通用 JSON 适配、防刷屏(聚合/限流/折叠)、排障、怎么加事件 | 运维、平台开发者 | | [DEVICE_AGENT.md](DEVICE_AGENT.md) | **设备端 Agent 接口契约**:应用商店的设备专用接口(清单/下载/上报)、adb 指令协议、版本约定 —— 与设备端 APK 仓库共享的契约 | 设备端开发者、平台开发者 | | [DEPLOY.md](DEPLOY.md) | **部署与运维**:环境准备、生产容器、数据备份导出/导入、故障排查 | 运维、部署者 | | [DEVELOPMENT.md](DEVELOPMENT.md) | **开发手册**:git 流程、技术红线、本地开发与调试、常见开发任务、文档同步约定 | 所有开发者 | diff --git a/doc/TASK_DEV.md b/doc/TASK_DEV.md index b7fc697..106d8ea 100644 --- a/doc/TASK_DEV.md +++ b/doc/TASK_DEV.md @@ -118,7 +118,7 @@ from .generic import task # 触发 @register_task(当前唯一任务类型) | 2 | `stop_app` | 结束 App | `package` | — | — | `app_stop`(`am force-stop`,下次打开是冷启动) | | 3 | `screen_on` | 亮屏 | — | — | — | `d.unlock()`(息屏时唤醒并滑动解锁) | | 4 | `screen_off` | 息屏 | — | — | — | `d.screen_off()` | -| 5 | `keep_screen` | 保持亮屏 | `mode`("on") | — | — | `mode="off"` → 恢复自动息屏;否则 `svc power stayon true`(充电时常亮,适合长任务) | +| 5 | `keep_screen` | 保持亮屏 | `mode`("on") | — | — | `mode="off"` → 恢复自动息屏;否则**把系统息屏超时顶到最大**(`screen_off_timeout=2147483647`)+ `svc power stayon true` + 唤醒一次。⚠️ 只发 `svc power stayon true` 对**没插充电器**的设备(走 WiFi 的那些)**完全无效**——它管的是"充电时屏幕常亮",2026-09-14 就是因此被反馈"保持亮屏没用" | | 6 | `key_event` | 按键 | `key`("back") | — | — | `d.press(key)`;**值直传,无白名单校验** | | 7 | `swipe` | 滑动 | `direction`("up") | `duration_min`(0.25)、`duration_max`(0.50) | — | 时长取随机值;up = 中轴从 0.8h → 0.2h(down 反向,left/right 同理);**非法 direction 静默不滑** | | 8 | `swipe_until` | 滑动直到元素 | `selector_type`、`selector_value` | `direction`("up")、`max_swipes`(8)、`click_when_found`(True) | — | 每轮滑动后 sleep 0.5s;命中即(可选)点击并返回 True;**只有 `direction="down"` 才是向下滑**,其余按上滑 | diff --git a/static/admin/base.js b/static/admin/base.js index 120dd5b..590d518 100644 --- a/static/admin/base.js +++ b/static/admin/base.js @@ -126,6 +126,7 @@ function showSubTab(tabId, name){ tab.querySelectorAll('.sub-tab').forEach(b=>b.classList.toggle('active', b.dataset.sub===name)); tab.querySelectorAll('.sub-panel').forEach(p=>p.classList.toggle('active', p.id===tabId+'-sub-'+name)); if(tabId==='agent' && name==='taskgen' && typeof initTaskGen==='function') initTaskGen(); + if(tabId==='system' && name==='notify' && typeof loadNotifyPanel==='function') loadNotifyPanel(); if(name==='groups' && typeof loadGroups==='function') loadGroups(); if(name==='apks' && typeof loadAgentStore==='function') loadAgentStore(); if(name==='devpool' && typeof loadDevPool==='function'){ diff --git a/static/admin/editor.js b/static/admin/editor.js index 2b6b057..257df75 100644 --- a/static/admin/editor.js +++ b/static/admin/editor.js @@ -291,7 +291,7 @@ var _stepEditor={ h+='
息屏。任务结束后用,避免长时间亮屏导致烧屏/发热
'; }else if(step.type==='keep_screen'){ h+='
'; h+='
长任务(如养号看视频几小时)放"保持亮屏"防屏幕超时熄灭;任务结束前放"恢复自动息屏"
'; }else if(step.type==='swipe'){ diff --git a/static/admin/notify.js b/static/admin/notify.js new file mode 100644 index 0000000..d542e33 --- /dev/null +++ b/static/admin/notify.js @@ -0,0 +1,384 @@ +// 系统 → 通知:多 webhook 配置(每条自己订阅事件)+ 测试发送 + 发送记录 +// 后端:web/notify_api.py(全部 admin_required);配置落 app_meta.notify_webhooks 单键。 +// 安全口径:URL 里有凭据(企微 ?key=)→ 列表/记录里一律打码;编辑时**留空即不修改**。 +let _notifyHooks = []; +let _notifyEvents = []; +let _notifyFormats = []; +let _notifySettings = {}; +let _notifyEditingId = null; // null=新建 +let _notifyPatterns = []; // 编辑中的通配订阅(如 task.*) +let _notifyPrevFormat = ''; // 上一个选中格式(切换时判断限流值要不要跟着走) +let _notifyRateTouched = false; // 用户是否手动改过限流值(没改过就跟随格式的推荐值) + +// ================== 加载与渲染 ================== +function loadNotifyPanel(){ + Promise.all([apiGet('/api/notify/webhooks'), + apiGet('/api/notify/events'), + apiGet('/api/notify/logs?limit=50')]).then(function(res){ + const [hooks, events, logs] = res; + if(hooks && hooks.ok){ + _notifyHooks = hooks.webhooks || []; + _notifyFormats = hooks.formats || []; + _notifySettings = hooks.settings || {}; + } + if(events && events.ok) _notifyEvents = events.events || []; + renderNotifySettings(); + renderNotifyHooks(); + renderNotifyLogs(logs && logs.ok ? (logs.logs || []) : []); + }); +} + +function renderNotifySettings(){ + const on = document.getElementById('notify-global'); + if(on) on.checked = _notifySettings.global_enabled !== false; + const agg = document.getElementById('notify-agg'); + if(agg) agg.value = _notifySettings.default_agg_window != null ? _notifySettings.default_agg_window : 30; + const rate = document.getElementById('notify-rate'); + if(rate) rate.value = _notifySettings.default_rate_limit != null ? _notifySettings.default_rate_limit : 18; +} + +function renderNotifyHooks(){ + const tb = document.getElementById('tb-notify-hooks'); + if(!tb) return; + if(!_notifyHooks.length){ + tb.innerHTML = '还没有配置 webhook。点「+ 新建 Webhook」加一条——' + + '比如把「任务失败」推到企业微信群。'; + return; + } + tb.innerHTML = _notifyHooks.map(function(h){ + const fmt = (_notifyFormats.find(f => f.name === h.format) || {}).label || h.format; + const evs = (h.events || []).length; + return '' + + '' + esc(h.name) + (h.secret_set ? ' 🔑' : '') + '' + + '' + esc(fmt) + '' + + '' + + esc(h.url || '') + '' + + '' + evs + ' 个事件' + + '' + (h.enabled ? '启用' + : '停用') + '' + + '' + + ' ' + + ' ' + + ' ' + + '' + + ''; + }).join(''); +} + +function renderNotifyLogs(logs){ + const tb = document.getElementById('tb-notify-logs'); + if(!tb) return; + if(!logs.length){ + tb.innerHTML = '本次运行还没有发送记录(重启后清空,历史见 logs/notify.log)'; + return; + } + tb.innerHTML = logs.map(function(r){ + const t = r.ts ? new Date(r.ts * 1000).toLocaleString('zh-CN', {hour12:false}) : ''; + const st = r.ok ? '成功' + : '失败'; + const detail = r.ok ? ('HTTP ' + (r.http_status == null ? '' : r.http_status) + + (r.errcode ? ' errcode ' + r.errcode : '')) + : esc(r.error || ''); + return '' + esc(t) + '' + + '' + esc(r.hook_name || '') + '' + + '' + esc(r.event || '') + '' + + '' + st + '' + + '' + detail + (r.elapsed_ms != null ? ' · ' + r.elapsed_ms + 'ms' : '') + ''; + }).join(''); +} + +// ================== 编辑弹窗 ================== +function notifyEventMatches(pattern, key){ + if(!pattern) return false; + if(pattern === '*' || pattern === key) return true; + if(pattern.slice(-2) === '.*') return key.indexOf(pattern.slice(0, -1)) === 0; + return false; +} + +function openNotifyModal(id){ + _notifyEditingId = id || null; + const h = id ? (_notifyHooks.find(x => x.id === id) || {}) : {}; + _notifyPatterns = (h.events || []).filter(p => p.indexOf('*') >= 0); + const checked = (h.events || []).filter(p => p.indexOf('*') < 0); + const fmtOpts = _notifyFormats.map(function(f){ + const sel = (h.format || 'wecom') === f.name ? ' selected' : ''; + return ''; + }).join(''); + const box = document.getElementById('notify-modal-box'); + if(!box) return; + box.innerHTML = + '
' + + '
' + + (id ? '编辑 Webhook' : '新建 Webhook') + '
' + + '' + + '
' + + '
' + + '
' + + '
' + + '
' + + '
' + + '
' + + '
' + // 占位与说明都由**当前格式**决定(见 _notifyFormatChanged),切换格式会跟着变 + + '' + + '
' + + '
' + + '
' + + '' + + '
' + + '
' + + '' + + '
0=不聚合
' + + '
' + + '' + + '
' + + '
' + + '
' + + '
按类别分的完整事件目录;' + + '带 ★ 的是建议开启的。也可以直接写通配(如 task.*)——在下面的芯片框里加。
' + + '
' + + '
' + + '
' + + '
' + + '
' + + '
' + + '' + + '
' + + '
点「预览请求体」看真实会发出去的内容与字节数。
' + + '
' + // 底部操作栏(缺了它就只能靠刷新页面逃出去) + + '
' + + '' + + '' + + '' + + 'Esc 或点空白处也能关' + + '
'; + box.querySelector('#nf-body').value = h.body_template || ''; + renderEventTree(checked); + renderNotifyPatterns(); + _notifyPrevFormat = (h.format || 'wecom'); + _notifyRateTouched = (h.rate_limit_per_min != null); // 编辑已有配置时视为"已定" + _notifyFormatChanged(); + document.getElementById('notify-modal-overlay').style.display = 'flex'; +} + +function _notifyFormatChanged(){ + const name = (document.getElementById('nf-format') || {}).value; + const meta = _notifyFormats.find(x => x.name === name) || {}; + const h = _notifyEditingId + ? (_notifyHooks.find(x => x.id === _notifyEditingId) || {}) : {}; + + // URL:编辑时占位显示现有(打码)地址,新建时显示该格式的示例 + const url = document.getElementById('nf-url'); + if(url){ + url.placeholder = (h.url && h.url !== undefined) ? h.url : (meta.url_hint || 'https://…'); + } + const urlHelp = document.getElementById('nf-url-help'); + if(urlHelp){ + urlHelp.innerHTML = (h.url ? ('已配置:' + esc(h.url) + '(留空=不修改)
') : '') + + esc(meta.url_help || ''); + } + + // 密钥:不同格式叫法/用途不同(企业微信不需要、Bark 是设备 Key、通用 JSON 是自定义头) + const sl = document.getElementById('nf-secret-label'); + if(sl) sl.textContent = (h.secret_set ? '🔑 ' : '') + (meta.secret_label || '密钥(可选)'); + const sh = document.getElementById('nf-secret-help'); + if(sh) sh.innerHTML = esc(meta.secret_help || '') + + (h.secret_set ? '
已配置过:留空=不修改,想清除请点上面的 🔑 后填写新值' : ''); + + // 限流:上限因格式而异(企微 20/分、Bark 无硬限) + const lh = document.getElementById('nf-limit-help'); + if(lh) lh.textContent = meta.limit_help || ''; + // 用户没手动改过限流值时,跟着格式的推荐值走(企微 20 / 通用 JSON 60 / Bark 60) + const rate = document.getElementById('nf-rate'); + if(rate && !_notifyRateTouched && meta.limit_default){ + rate.value = meta.limit_default; + } + // 请求体模板:只有需要模板的格式(通用 JSON)才显示 + const bodyWrap = document.getElementById('nf-body-wrap'); + if(bodyWrap) bodyWrap.style.display = meta.needs_template ? 'block' : 'none'; + + _notifyPrevFormat = name; +} + +// Esc 关闭弹窗(点空白处关闭在 overlay 的 onclick 上) +document.addEventListener('keydown', function(ev){ + if(ev.key !== 'Escape') return; + const ov = document.getElementById('notify-modal-overlay'); + if(ov && ov.style.display !== 'none') closeNotifyModal(); +}); + +function renderEventTree(checked){ + const box = document.getElementById('nf-events'); + if(!box) return; + const cats = []; + _notifyEvents.forEach(function(e){ if(cats.indexOf(e.category) < 0) cats.push(e.category); }); + box.innerHTML = cats.map(function(cat){ + const rows = _notifyEvents.filter(e => e.category === cat); + const body = rows.map(function(e){ + const id = 'nfe-' + e.key.replace(/\./g, '-'); + return ''; + }).join(''); + return '
' + + esc(cat) + '
' + + body + '
'; + }).join(''); +} + +function _notifyToggleCat(cat){ + const keys = _notifyEvents.filter(e => e.category === cat).map(e => e.key); + const boxes = [...document.querySelectorAll('#nf-events .nf-ev')].filter(b => keys.indexOf(b.value) >= 0); + const allOn = boxes.every(b => b.checked); + boxes.forEach(b => { b.checked = !allOn; }); +} + +function addNotifyPattern(){ + const inp = document.getElementById('nf-pattern-input'); + const v = (inp.value || '').trim(); + if(!v) return; + if(v !== '*' && v.slice(-2) !== '.*'){ showToast('通配请写成 task.* 或 device.* 或 *', 'error'); return; } + if(_notifyPatterns.indexOf(v) < 0) _notifyPatterns.push(v); + inp.value = ''; + renderNotifyPatterns(); +} + +function renderNotifyPatterns(){ + const box = document.getElementById('nf-patterns'); + if(!box) return; + box.innerHTML = _notifyPatterns.map(function(p, i){ + return '' + + esc(p) + ' ×'; + }).join('') || '(无)'; +} + +function removeNotifyPattern(i){ + _notifyPatterns.splice(i, 1); + renderNotifyPatterns(); +} + +function closeNotifyModal(){ + const ov = document.getElementById('notify-modal-overlay'); + if(ov) ov.style.display = 'none'; +} + +// ================== 保存 / 删除 / 启停 / 测试 ================== +function collectNotifyForm(){ + const events = [...document.querySelectorAll('#nf-events .nf-ev')] + .filter(b => b.checked).map(b => b.value); + _notifyPatterns.forEach(function(p){ if(events.indexOf(p) < 0) events.push(p); }); + const body = {name: (document.getElementById('nf-name').value || '').trim(), + format: document.getElementById('nf-format').value, + events: events, + agg_window: parseInt(document.getElementById('nf-agg').value, 10) || 0, + rate_limit_per_min: parseInt(document.getElementById('nf-rate').value, 10) || 0, + body_template: (document.getElementById('nf-body').value || '')}; + const url = (document.getElementById('nf-url').value || '').trim(); + if(url) body.url = url; // 留空 = 不修改(后端保留原值) + const sec = (document.getElementById('nf-secret').value || '').trim(); + if(sec) body.secret = sec; + return body; +} + +function saveNotifyHook(){ + const body = collectNotifyForm(); + if(!body.name){ showToast('请填名称', 'error'); return; } + if(!body.events.length){ showToast('至少要勾一个事件', 'error'); return; } + const id = _notifyEditingId; + const req = id ? apiPut('/api/notify/webhooks/' + id, body) + : apiPost('/api/notify/webhooks', body); + req.then(function(r){ + if(!r || !r.ok){ + let msg = (r && r.error) || '保存失败'; + if(r && r.errors && r.errors.length) msg = r.errors.join(';'); + showToast(msg, 'error'); + return; + } + showToast('已保存', 'success'); + closeNotifyModal(); + loadNotifyPanel(); + }); +} + +function deleteNotifyHook(id){ + const h = _notifyHooks.find(x => x.id === id) || {}; + if(!confirm('删除 webhook「' + (h.name || '') + '」?\n(只是不再推送,不影响任何任务/设备数据)')) return; + apiDelete('/api/notify/webhooks/' + id).then(function(r){ + if(!r || !r.ok){ showToast((r && r.error) || '删除失败', 'error'); return; } + showToast('已删除', 'success'); + loadNotifyPanel(); + }); +} + +function toggleNotifyHook(id, on){ + apiPut('/api/notify/webhooks/' + id, {enabled: on}).then(function(r){ + if(!r || !r.ok){ showToast((r && r.error) || '操作失败', 'error'); return; } + loadNotifyPanel(); + }); +} + +function testNotifyHook(id){ + showToast('正在发送测试消息…', 'success'); + apiPost('/api/notify/webhooks/' + id + '/test', {}).then(function(r){ + if(!r){ showToast('请求失败', 'error'); return; } + const detail = 'HTTP ' + (r.http_status == null ? '-' : r.http_status) + + (r.errcode ? ' · errcode ' + r.errcode : '') + + (r.elapsed_ms != null ? ' · ' + r.elapsed_ms + 'ms' : '') + + (r.attempts > 1 ? ' · 重试 ' + r.attempts + ' 次' : '') + + (r.error ? '\n' + r.error : ''); + showToast(r.ok ? ('测试消息已发出(' + detail + ')') : ('发送失败(' + detail + ')'), + r.ok ? 'success' : 'error'); + loadNotifyPanel(); + }); +} + +function previewNotifyHook(){ + const body = collectNotifyForm(); + apiPost('/api/notify/preview', {hook: body, event: 'notify.test'}).then(function(r){ + const box = document.getElementById('nf-preview'); + if(!box) return; + if(!r || !r.ok){ box.innerHTML = '预览失败:' + esc((r && r.error) || '') + ''; return; } + let html = '
UTF-8 字节数:' + (r.bytes == null ? '-' : r.bytes) + + (r.byte_limit ? ' / 上限 ' + r.byte_limit : '') + + (r.truncated ? ' (会截断)' : '') + '
'; + html += '
'
+      + esc(r.request_body || '') + '
'; + box.innerHTML = html; + }); +} + +function saveNotifyGlobal(){ + const on = document.getElementById('notify-global').checked; + const body = {global_enabled: on, + default_agg_window: parseInt(document.getElementById('notify-agg').value, 10) || 0, + default_rate_limit: parseInt(document.getElementById('notify-rate').value, 10) || 18}; + apiPost('/api/notify/settings', body).then(function(r){ + if(!r || !r.ok){ showToast((r && r.error) || '保存失败', 'error'); return; } + showToast('已保存(新配置对之后的事件生效)', 'success'); + loadNotifyPanel(); + }); +} + +function toggleNotifyGlobal(el){ + apiPost('/api/notify/settings', {global_enabled: el.checked}).then(function(r){ + if(!r || !r.ok){ el.checked = !el.checked; showToast((r && r.error) || '操作失败', 'error'); return; } + showToast(el.checked ? '通知总开关:已开启' : '通知总开关:已关闭(所有 webhook 暂停推送)', 'success'); + }); +} diff --git a/tasks/generic/task.py b/tasks/generic/task.py index c859417..dfefb72 100644 --- a/tasks/generic/task.py +++ b/tasks/generic/task.py @@ -37,6 +37,7 @@ 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") @@ -45,6 +46,13 @@ _log = get_logger("task.generic") _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]`(只改整体前缀,保留后续子路径)。 @@ -141,6 +149,10 @@ class GenericStepsWorker(BaseWorker): _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。""" @@ -228,17 +240,35 @@ class GenericStepsWorker(BaseWorker): def _exec_keep_screen(self, d, params, depth=0): """保持亮屏 / 恢复自动息屏。 - 长任务(养号看视频数小时)时屏幕会按系统超时自动熄灭导致任务无法执行, - 用 svc power stayon 让充电时屏幕常亮;任务结束前用 mode=off 恢复自动息屏。 + ⚠️ 只用 `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") - _log.info(f"[{self.serial}] 恢复自动息屏") + 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: - d.shell("svc power stayon true") - _log.info(f"[{self.serial}] 保持亮屏(充电时屏幕常亮)") + 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}") diff --git a/templates/admin/monitor.html b/templates/admin/monitor.html index 4392129..6cfb923 100644 --- a/templates/admin/monitor.html +++ b/templates/admin/monitor.html @@ -1112,6 +1112,7 @@ body{background:var(--bg);font-family:var(--body);color:var(--text);font-size:14
+
@@ -1151,6 +1152,51 @@ body{background:var(--bg);font-family:var(--body);color:var(--text);font-size:14
+ + +
+
+ 通知:把平台事件(任务失败/成功、设备上下线、安装完成、备份恢复…)推到 webhook, + 例如企业微信群机器人。每条 webhook 自己勾选关心哪些事件——不勾就不推。 +
为防刷屏:同一条 webhook 默认每 30 秒合并一次同类事件、每分钟最多 18 条 + (企业微信硬限 20),超出的会压成一条「被限流折叠」摘要而不是丢掉。详见 doc/NOTIFY.md。 +
+
+ + + + 默认聚合 + 秒 · + 限流 条/分 + + +
+ + + +
名称格式URL订阅状态操作
+ +
+ 发送记录 + 仅本次运行期间(重启清空);完整历史见 logs/notify.log + +
+ + + +
时间Webhook事件结果详情
+
+ + + + @@ -1219,5 +1265,6 @@ body{background:var(--bg);font-family:var(--body);color:var(--text);font-size:14 + diff --git a/web/__init__.py b/web/__init__.py index 58ca3ea..9f15e35 100644 --- a/web/__init__.py +++ b/web/__init__.py @@ -11,6 +11,7 @@ tailscale — Tailscale 管理 system — 系统数据备份导出/导入 device_agent— 设备端 Agent 专用接口(应用商店:清单/下载/上报)+ 其管理端配置 + notify — 通知 / Webhook 管理(配置、测试发送、发送记录、事件目录) """ from flask import Blueprint @@ -30,7 +31,8 @@ def register_blueprints(app): _agent_mod.set_app(app) from .system_api import bp as system_bp from .device_agent_api import bp as devagent_bp + from .notify_api import bp as notify_bp for bp in (auth_bp, monitor_bp, tasks_bp, admin_bp, tools_bp, devices_bp, apks_bp, tailscale_bp, agent_bp, system_bp, - devagent_bp): + devagent_bp, notify_bp): app.register_blueprint(bp) diff --git a/web/agent_api.py b/web/agent_api.py index 62ff460..fa4afe8 100644 --- a/web/agent_api.py +++ b/web/agent_api.py @@ -836,6 +836,8 @@ def run_experience_audit(): cfg = _read_cfg() if not cfg.get("api_key"): _log.info("经验巡检跳过:未配置 API Key") + from core import notifier + notifier.notify("ai.audit.skipped", reason="未配置模型 API Key") return rows = db.session.execute(db.text( "SELECT id, task_prompt, recipe, tool_seq, hits FROM agent_experience " @@ -868,8 +870,16 @@ def run_experience_audit(): summary = f"评审 {reviewed} 条,建议删除 {suggested} 条(待人工确认)" _audit_state.update(last=now, last_summary=summary) _log.info("经验巡检完成: %s", summary) + from core import notifier + notifier.notify("ai.audit.finished", reviewed=reviewed, + suggested=suggested, summary=summary) except Exception as e: _log.warning(f"经验巡检异常: {e}") + try: + from core import notifier + notifier.notify("ai.audit.failed", error=str(e)[:200]) + except Exception: + pass finally: _audit_state["running"] = False diff --git a/web/auth.py b/web/auth.py index 62e3663..52a1b55 100644 --- a/web/auth.py +++ b/web/auth.py @@ -144,6 +144,11 @@ def login(): login_user(user) _csrf_token() # 建立 CSRF token,前端通过 /api/csrf 获取 _log.info(f"用户 {username} 登录") + try: + from core import notifier + notifier.notify("user.login", username=username) + except Exception: + pass return redirect(request.args.get("next") or url_for("auth.index")) return render_template("admin/login.html", error="用户名或密码错误") return render_template("admin/login.html", error=None) diff --git a/web/monitor.py b/web/monitor.py index 144d548..7bf07d6 100644 --- a/web/monitor.py +++ b/web/monitor.py @@ -239,8 +239,11 @@ def api_device_screen_all(): if ":" in serial: adb_connect(serial) if mode == "off": - subprocess.run([ADB_PATH, "-s", serial, "shell", "input", "keyevent", "26"], - capture_output=True, timeout=15) # KEYCODE_POWER 息屏 + # KEYCODE_SLEEP(223) 是**单向**熄屏;不能用 KEYCODE_POWER(26)—— + # 那是电源键**开关**,对已经息屏的设备会把屏幕点亮(2026-09-14 用户反馈: + # "一键息屏"把睡着的设备唤醒了,看着像按钮反了) + subprocess.run([ADB_PATH, "-s", serial, "shell", "input", "keyevent", "223"], + capture_output=True, timeout=15) else: subprocess.run([ADB_PATH, "-s", serial, "shell", "input", "keyevent", "224"], capture_output=True, timeout=15) # WAKEUP diff --git a/web/notify_api.py b/web/notify_api.py new file mode 100644 index 0000000..567f0a7 --- /dev/null +++ b/web/notify_api.py @@ -0,0 +1,190 @@ +"""通知 / Webhook 管理接口(系统级设置,全部 admin_required)。 + +配置读写统一走 `core/notifier`(它落 `app_meta.notify_webhooks` 一个键)——本模块 +只做"取配置 → 改一处 → 存回去"的编排,**校验与持久化都在 notifier 里**,不在 web 层 +重复实现(避免两份规则漂移)。 + +安全口径: +- URL 里有凭据(企微 `?key=`)→ 回显一律 `notifier.mask_url()`; +- 前端"留空不修改":提交上来的 url 若是打码值或空,`save_config` 会保留原值; +- `secret` 永不回显,只回 `secret_set`。 +""" +import uuid + +from flask import Blueprint, jsonify, request +from flask_login import current_user + +from core import notify_events, notifier +from core.logger import get_logger +from web.auth import admin_required + +_log = get_logger("web.notify") +bp = Blueprint("notify", __name__) + + +def _operator(): + try: + return current_user.username or "admin" + except Exception: + return "admin" + + +def _formats(): + """可用格式(含未实现的,前端据此置灰)。 + + 每项带**该格式自己的提示文案**(URL 占位/说明、密钥标签、限流说明)—— + 这些必须跟着格式走,不能把企业微信的说明写死在页面上。 + """ + out = [cls.meta() for _, cls in notifier.ADAPTERS.items()] + for name, label in notifier.PLANNED_FORMATS.items(): + out.append({"name": name, "label": label + "(未实现)", "implemented": False, + "byte_limit": 0, "limit_default": 20, "needs_template": False, + "secret_label": "签名密钥(可选)", "secret_help": "", + "url_hint": "", "url_help": "该格式尚未实现,可先用「通用 JSON」手搓", + "limit_help": ""}) + return out + + +@bp.route("/api/notify/webhooks") +@admin_required +def api_notify_webhooks(): + """全部 webhook 配置(url 打码、secret 只回是否配置过)+ 格式清单 + 事件目录。""" + cfg = notifier.get_public_config() + return jsonify({"ok": True, "webhooks": cfg["webhooks"], "settings": cfg["settings"], + "formats": _formats(), "meta_key": notifier.META_KEY, + "max_hooks": notifier.MAX_HOOKS, + "planned": list(notifier.PLANNED_FORMATS)}) + + +@bp.route("/api/notify/webhooks", methods=["POST"]) +@admin_required +def api_notify_webhook_create(): + """新建一条 webhook。body 即该条配置(见 doc/NOTIFY.md 的字段表)。""" + data = request.json or {} + cfg = notifier.get_config() + if len(cfg["webhooks"]) >= notifier.MAX_HOOKS: + return jsonify({"ok": False, "error": f"最多 {notifier.MAX_HOOKS} 条"}), 400 + new_id = "wh_" + uuid.uuid4().hex[:8] + hook = dict(data) + hook["id"] = new_id + cfg["webhooks"].append(hook) + ok, msg = notifier.save_config(cfg) + if not ok: + return jsonify({"ok": False, "error": msg, "errors": msg + if isinstance(msg, list) else None}), 400 + _log.info("新增通知 webhook: %s(%s)by %s", hook.get("name"), new_id, _operator()) + return jsonify({"ok": True, "msg": "已创建", "id": new_id}) + + +@bp.route("/api/notify/webhooks/", methods=["PUT"]) +@admin_required +def api_notify_webhook_update(hook_id): + """更新一条(部分字段;url/secret 省略或为打码值时保持原值)。""" + data = request.json or {} + cfg = notifier.get_config() + target = next((h for h in cfg["webhooks"] if h["id"] == hook_id), None) + if not target: + return jsonify({"ok": False, "error": "webhook 不存在"}), 404 + for k, v in data.items(): + if k in ("id", "created_at"): + continue + target[k] = v + ok, msg = notifier.save_config(cfg) + if not ok: + return jsonify({"ok": False, "error": msg, "errors": msg + if isinstance(msg, list) else None}), 400 + _log.info("更新通知 webhook: %s by %s", hook_id, _operator()) + return jsonify({"ok": True, "msg": "已保存"}) + + +@bp.route("/api/notify/webhooks/", methods=["DELETE"]) +@admin_required +def api_notify_webhook_delete(hook_id): + cfg = notifier.get_config() + before = len(cfg["webhooks"]) + cfg["webhooks"] = [h for h in cfg["webhooks"] if h["id"] != hook_id] + if len(cfg["webhooks"]) == before: + return jsonify({"ok": False, "error": "webhook 不存在"}), 404 + ok, msg = notifier.save_config(cfg) + if not ok: + return jsonify({"ok": False, "error": msg}), 400 + _log.info("删除通知 webhook: %s by %s", hook_id, _operator()) + return jsonify({"ok": True, "msg": "已删除"}) + + +@bp.route("/api/notify/webhooks//test", methods=["POST"]) +@admin_required +def api_notify_webhook_test(hook_id): + """同步发一条测试消息(用户等结果),返回真实 HTTP 状态与平台错误码。 + + ⚠ 不占业务令牌桶:否则管理员点两下测试就把业务通知的配额吃掉了。 + """ + cfg = notifier.get_config() + hook = next((h for h in cfg["webhooks"] if h["id"] == hook_id), None) + if not hook: + return jsonify({"ok": False, "error": "webhook 不存在"}), 404 + event = ((request.json or {}).get("event") or "notify.test").strip() + rec = notifier.send_test(hook, event=event, operator=_operator()) + return jsonify({"ok": bool(rec.get("ok")), + "msg": "测试消息已发出" if rec.get("ok") else "发送失败", + "http_status": rec.get("http_status"), + "errcode": rec.get("errcode"), + "elapsed_ms": rec.get("elapsed_ms"), + "attempts": rec.get("attempts"), + "error": rec.get("error") or "", + "url": notifier.mask_url(hook.get("url", ""))}) + + +@bp.route("/api/notify/preview", methods=["POST"]) +@admin_required +def api_notify_preview(): + """保存前预览:真实请求体 + UTF-8 字节数 + 是否会被截断。""" + data = request.json or {} + hook = data.get("hook") or {} + # 预览时 URL 可能还没填:给个占位,只关心渲染结果 + hook = dict(hook) + hook.setdefault("url", "https://example.invalid/hook") + out = notifier.preview(hook, event=(data.get("event") or "notify.test")) + return jsonify({"ok": not out.get("error"), **out}) + + +@bp.route("/api/notify/events") +@admin_required +def api_notify_events(): + """事件目录(前端画勾选树 / 模板占位符速查)。""" + rows = notify_events.list_events() + cats = [] + for e in rows: + if e["category"] not in cats: + cats.append(e["category"]) + return jsonify({"ok": True, "events": rows, "categories": cats}) + + +@bp.route("/api/notify/logs") +@admin_required +def api_notify_logs(): + """最近发送记录(内存环形缓冲,重启清空;历史见 logs/notify.log)。""" + limit = request.args.get("limit", 50) + try: + limit = max(1, min(int(limit), 200)) + except (TypeError, ValueError): + limit = 50 + return jsonify({"ok": True, "logs": notifier.get_logs(limit), + "note": "仅显示本次运行期间记录;历史见 logs/notify.log"}) + + +@bp.route("/api/notify/settings", methods=["POST"]) +@admin_required +def api_notify_settings(): + """全局设置:global_enabled / default_agg_window / default_rate_limit / log_keep。""" + data = request.json or {} + cfg = notifier.get_config() + for k in ("global_enabled", "default_agg_window", "default_rate_limit", + "log_keep", "http_timeout"): + if k in data and data[k] is not None: + cfg["settings"][k] = data[k] + ok, msg = notifier.save_config(cfg) + if not ok: + return jsonify({"ok": False, "error": msg}), 400 + _log.info("更新通知全局设置 by %s: %s", _operator(), data) + return jsonify({"ok": True, "msg": "已保存"}) diff --git a/web/system_api.py b/web/system_api.py index bbed2ac..aca3031 100644 --- a/web/system_api.py +++ b/web/system_api.py @@ -9,6 +9,7 @@ from io import BytesIO from flask import Blueprint, jsonify, request, send_file +from flask_login import current_user from web.auth import admin_required from core import system_backup as sb @@ -39,6 +40,14 @@ def api_system_backup_export(): with open(zip_path, "rb") as f: data = f.read() sb.remove_quiet(zip_path) + try: + from core import notifier + notifier.notify("system.backup.exported", filename=fname, size=len(data), + tables=len(manifest.get("tables") or []), + include_apk=include_apk, + user=getattr(current_user, "username", "")) + except Exception: + pass # 通知失败不影响导出 resp = send_file(BytesIO(data), as_attachment=True, download_name=fname, mimetype="application/zip") return resp diff --git a/web_server.py b/web_server.py index 51ffd8c..31be35c 100644 --- a/web_server.py +++ b/web_server.py @@ -66,6 +66,13 @@ def _consume_pending_restore(): # 先初始化数据库(含旧 JSON 迁移),再创建 TaskManager(需要 app context 读写 DB) init_db(app) +# 通知模块:**必须早于 _consume_pending_restore()**——启动期就要能发「备份已恢复」 +# 这类通知(notify 在未就绪时静默丢弃,晚初始化就丢事件了) +try: + from core import notifier as _notifier + _notifier.init_app(app) +except Exception as _e: # 通知出问题不影响启动 + _log.warning(f"通知模块初始化失败(不影响启动): {_e}") # 消费「待生效的备份恢复」——必须在 init_db 之后(表已建好,且现在是事务替换而非换文件) _consume_pending_restore() # 库环境标签校验 + 启动横幅:每次启动都明确写出「现在连的是哪个库」, @@ -288,9 +295,27 @@ if __name__ == "__main__": # 预连接设备池:adb server 重启后设备全掉线,后台并发重连加速恢复 _preconnect_pool_devices() _log.info(f"管理后台: http://localhost:{WEB_PORT}/ (admin/admin123)") + # 启动通知:放在 __main__ 里是刻意的——以 WSGI 方式 import 本模块只有"半个启动" + # (阶段 A~F 在 import 期就跑完了),不该发"服务已启动"。 + try: + from core import notifier as _notifier + with app.app_context(): + _devs = len(device_pool.list_configured()) + _notifier.notify("service.started", env=db_config.env_badge().get("env", ""), + db_target=db_config.describe_target(), + device_count=_devs, job_count=len(mgr.jobs)) + except Exception as _e: + _log.warning(f"发送启动通知失败(不影响启动): {_e}") + _t0 = time.time() try: _run_server(WEB_HOST, WEB_PORT) finally: + try: + from core import notifier as _notifier + _notifier.notify("service.stopping", uptime_s=int(time.time() - _t0)) + _notifier.shutdown(timeout=3.0) + except Exception: + pass mgr.shutdown() device_discovery.shutdown() _stop_uiauto()