"""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/)。 回显、日志、发送记录一律用它——**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)。 界面上的提示文案也挂在这里(`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 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 {})} ADAPTERS = {WecomAdapter.name: WecomAdapter, JsonAdapter.name: JsonAdapter, BarkAdapter.name: BarkAdapter} # ================== 发送 ================== 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 里的错误码,但**成功码不一样** # (企业微信 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