feat(通知): 系统级 Webhook 通知子系统(事件目录 + 可插拔适配器 + 防刷屏)
平台此前出了问题只能靠人盯页面。现在各组件统一走 `notifier.notify(事件, **字段)`, 推到企业微信 / 自建服务;**所有可通知点都登记进事件目录,默认全关,用户按 webhook 勾选**。 ## 架构(`core/notifier.py` + `core/notify_events.py`) 业务线程调 `notify()` → 只做内存操作(读配置快照/匹配订阅/入队)→ 返回; 后台 1 个 dispatcher(聚合 + 每 hook 限流 + 折叠摘要)+ 3 个 sender(真实 HTTP、退避重试) 负责真正发出去。硬约束:**notify 零 DB、零 HTTP、零阻塞、异常不冒泡**——所以任务线程里 可以直接调(不用 app_context、不用 try/except),但**必须放在所有 `with self._lock` 之外**。 - **事件目录 36 条**(任务批次/单设备/设备/Worker/业务/安装/系统/AI),支持 `task.*` 通配订阅; 语义分工避免重复告警:`worker.*` 是单次尝试级,`task.device.success/failed` 是唯一权威结论。 - **适配器可插拔**:`wecom`(markdown,按 4096 **字节**截断、超限不截半个汉字)+ `json`(模板占位符,替换值按 JSON 转义,保存前干跑校验);钉钉/飞书留了插槽(前端置灰)。 - **防打爆四层**:聚合窗口(默认 30s,同批次合并成一条并带样本)→ 令牌桶限流(默认 18/分, 对齐企微硬限 20)→ 被限流的**折叠成摘要不丢弃** → 有界队列背压。取舍:失败通知最多延迟 一个窗口,换来群不被刷屏。 ## 安全与存储 - 配置只落 `app_meta.notify_webhooks` 一个键(**不建表** → 不涉及备份覆盖红线)。 - URL 本身就是凭据(企微 `?key=`)→ 接口回显/发送记录/日志一律 `mask_url()/scrub()`; 编辑时留空即不修改;secret 永不回显。DATA_MODEL 的明文凭据告警补上了这一条。 - 发送记录:内存环形缓冲 200 条(重启清空)+ 独立 `logs/notify.log`。 ## 接入点(每个都放在锁外、不改 return 顺序) task_manager(批次开始/结束用新增的 `_BatchTracker` 统一在 finally 计数、单设备成功/失败/ 离线/重试/停止/抢占/归还/cron 停止)、device_worker 心跳看门狗、generic 任务选择器连续失效、 apk 安装开始/完成、设备上下线(**状态沿检测**,只报新变化)、备份导出/恢复、经验巡检、 用户登录、服务启停。 ## 前端 系统 Tab 新增「通知 / Webhook」子分栏:多条 webhook 列表(URL 打码)+ 编辑弹窗(格式/URL/ 密钥/事件勾选树带 ★建议/聚合/限流/自定义模板/预览)+ 发送测试 + 发送记录。 ## 自测 - 进程内逻辑 10 组断言全绿:聚合合并、限流+折叠、无配置/全局关静默丢弃、未知事件、 内部异常不外泄、JSON 转义(标题含引号换行仍合法)、URL/异常消息脱敏、配置校验。 - 端到端(假 webhook 接收端)17 项断言全绿:真实事件投递(user.login / task.batch.no_device)、 企微请求体形状、**HTTP 200 + errcode 93000 判为失败**、500 重试 3 次、记录里 URL 打码。 - 韧性:webhook 指向黑洞地址时登录耗时 100~114ms(基线 107~133ms,**异步隔离生效**); 配置写成坏 JSON 服务照常启动、通知静默不发、日志有 error(服务端实测后已复原)。 - 页面:系统 → 通知 面板/弹窗/36 个事件复选框/预览全部正常,无 JS 报错。 - 自测数据已清理(webhook、自建任务、写坏又复原的配置键)。 文档:新增 doc/NOTIFY.md(事件表/配置/格式约束/防刷屏/加事件三步骤/排障)并登记进 doc/README; API.md §2.11;DATA_MODEL 的 app_meta 键表与明文凭据告警;ARCHITECTURE 线程表/分层/扩展点; 根 README 功能索引与日志表。
This commit is contained in:
@@ -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。
|
||||
|
||||
@@ -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),
|
||||
|
||||
+11
-1
@@ -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()
|
||||
|
||||
@@ -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"
|
||||
|
||||
@@ -0,0 +1,884 @@
|
||||
"""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,
|
||||
}
|
||||
|
||||
# 目前实现不了的格式:前端置灰,但用「通用 JSON」能手搓
|
||||
PLANNED_FORMATS = {"dingtalk": "钉钉", "feishu": "飞书"}
|
||||
|
||||
|
||||
# ================== URL / 文本脱敏 ==================
|
||||
_SECRET_KEYS = ("key", "access_token", "token", "secret", "sign", "apikey", "api_key")
|
||||
|
||||
|
||||
def mask_url(url):
|
||||
"""把 URL 里的凭据打码(企微 ?key=、钉钉 ?access_token=、飞书 /hook/<token>)。
|
||||
|
||||
回显、日志、发送记录一律用它——**URL 本身就是可直接发消息的凭据**。
|
||||
"""
|
||||
if not url:
|
||||
return ""
|
||||
try:
|
||||
from urllib.parse import urlsplit, urlunsplit, parse_qsl, urlencode
|
||||
sp = urlsplit(url)
|
||||
if not sp.query:
|
||||
return url
|
||||
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)。"""
|
||||
|
||||
name = "base"
|
||||
label = "基类"
|
||||
byte_limit = 4096 # 单个消息体的字节上限(0=不限)
|
||||
limit_default = 20 # 该格式每机器人每分钟的官方上限
|
||||
truncated_note = "\n…(已截断)"
|
||||
|
||||
@classmethod
|
||||
def render(cls, msg, hook):
|
||||
"""→ requests.post 的参数 dict(url/json/headers)。"""
|
||||
raise NotImplementedError
|
||||
|
||||
@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
|
||||
|
||||
@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
|
||||
|
||||
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}")
|
||||
return {"url": hook["url"], "json": parsed,
|
||||
"headers": dict(hook.get("headers") or {})}
|
||||
|
||||
@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, ""
|
||||
|
||||
|
||||
ADAPTERS = {WecomAdapter.name: WecomAdapter, JsonAdapter.name: JsonAdapter}
|
||||
|
||||
|
||||
# ================== 发送 ==================
|
||||
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 _post_once(adapter, msg, hook, timeout):
|
||||
kw = adapter.render(msg, hook)
|
||||
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 里的错误码
|
||||
ok = r.status_code == 200 and (code in (0, None))
|
||||
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
|
||||
@@ -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
|
||||
@@ -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
|
||||
|
||||
|
||||
|
||||
+142
-12
@@ -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):
|
||||
|
||||
Reference in New Issue
Block a user