Files
auto_control/core/notifier.py
T
butubb 8fdde978ea fix(设备名): 选择设备与所有通知都显示名称而不是 IP
用户反馈:① 元素抓取的选择设备列表显示的是 IP;② 所有 webhook 通知里都是 IP;
要都显示设备名。

- 通知链路统一补名字(调用点不用动):`core/notifier.py`
  · 新增 `set_device_name_resolver()` + `fill_device_names()`,在 `notify()` 与
    `build_message()`(预览/测试发送也走它)里把 `serial` 补成 `device_name`、
    把 `serials`/`devices` 列表逐项换成名字。
  · 补一处就全带名字了——标题主体、字段表、聚合样本认的都是 `device_name`。
  · **做成"纯内存回调"是刻意的**:notify() 的硬红线是零 DB,不能为了取个名字去查库
    (那等于在业务线程里加一次阻塞查询)。
  · 降级规则:调用点自己传了 device_name 就用它的;查不到名字(没命名/不在池里)
    保留原地址;解析器缺失或抛异常都只是降级,绝不影响发送。
- `core/device_pool.py`:维护 `serial → 名称` 内存快照——`init_app` 同步刷一次、
  增删改名/迁址后各刷一次(改名立刻生效)、`device-names` 线程每 60s 兜底刷一次
  (覆盖整库恢复这类进程外改动)。`name_of()` 只读内存,可在通知路径上安全调用。
- `web_server.py`:装配层接上 `notifier.set_device_name_resolver(device_pool.name_of)`。
- 元素抓取/测试此步骤的设备列表:`GET /api/uiauto/devices` 的 `name` 改用**平台名**,
  前端 `editor.js` 新增共用的 `_devCard()`——名字做主标题(粗体),
  `型号 · 地址` 作副标题。
  · 池外设备**退回地址而不是 uiautodev 的 name**:实测那份 name 是设备 codename
    (一柜子机器全叫 "earth"),拿它认设备等于没名字,地址至少唯一。
  · 「测试此步骤」的设备列表原来连状态角标都没有,一并统一成同一个卡片。
- 顺带修一个**既有 bug(不是本次需求)**:`/locate` 设备端定位页从 ce47a5b 那次
  web 蓝图拆分起就一直 500——拆分时漏掉了 `render_template_string` 与
  `markupsafe.escape as _esc` 两个 import。后者不只是缺个名字:定位页把 query 里的
  serial 拼进 HTML,而 `render_template_string` 的模板名是 "<template>"、
  **不会自动转义**,所以那还是个反射 XSS。已显式转义(已用 `<img onerror>` 验证)。

自测:device_pool 快照/改名即时生效、fill_device_names 全分支(含列表、
未命名保留地址、不在池保留地址、解析器缺失/抛异常降级)、真消息渲染断言
**通篇不含 IP**;GET 路由冒烟 54 个 0 个 500;前端语法 + 卡片渲染截图核对。
2026-09-24 08:55:10 +08:00

1168 lines
49 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""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>`
(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, fields=None):
"""消息级别:**事件自己带的 `level` 优先**(步骤「发通知」/巡检都能自己选),
否则按事件 key 猜(failed→error、success→success…),猜不到算 info。"""
lv = (fields or {}).get("level")
if lv in ("info", "success", "warning", "error"):
return lv
for hints, lv2 in _LEVEL_BY_HINT:
if any(h in event_key for h in hints):
return lv2
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": "总数",
"title": "标题", "message": "内容", "patrol_name": "巡检",
"check_label": "检查", "action_label": "动作", "detail": "结果",
"stopped": "停止", "failed_devices": "失败设备",
"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)
# ================== 设备名解析(让通知显示名字而不是 IP) ==================
# 由装配层注册(web_server.py):`notifier.set_device_name_resolver(device_pool.name_of)`。
# **刻意做成"纯内存回调"**:notify() 的红线是零 DB——不能为了取个名字去查库。
# 快照由 device_pool 维护(启动/池子变动/每 60s 兜底刷新)。
_device_name_fn = None
def set_device_name_resolver(fn):
"""注册 serial → 设备名 的解析器。必须是**纯内存**实现(不得查库/发请求)。"""
global _device_name_fn
_device_name_fn = fn
def _dev_name(serial):
"""尽力把地址解析成设备名;解析不到返回 ""(调用方保留原值)。"""
if not serial or _device_name_fn is None:
return ""
try:
return str(_device_name_fn(str(serial)) or "")
except Exception:
return ""
def fill_device_names(fields):
"""把 fields 里的设备地址补成设备名(**原地改**,返回同一个 dict)。
两类字段:
· `serial` → 补 `device_name`(原来没填才补)——标题主体、字段表、
聚合样本认的都是 `device_name`,补上这一处就全带名字了;
· `serials` / `devices`(多设备事件,如设备断联)→ 逐项换成名字,
通知里不该出现一串 IP。
查不到名字(设备没命名 / 已不在池里)就**保留原值**,绝不把内容弄丢。
"""
if not isinstance(fields, dict):
return fields
s = fields.get("serial")
if s and not fields.get("device_name"):
n = _dev_name(s)
if n:
fields["device_name"] = n
for k in ("serials", "devices"):
v = fields.get(k)
if isinstance(v, (list, tuple)) and v:
fields[k] = [_dev_name(x) or x for x in v]
return fields
def _subject(ev, fields):
"""标题主体:**用户自己写的 title 优先**,其次任务名,再次设备名/serial。"""
if fields.get("title"):
return str(fields["title"])
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 {})
fill_device_names(fields) # 同 notify():预览/测试发送也要显示设备名
merged = fields.pop("_merged", None)
level = _level_of(ev.key, fields)
subj = _subject(ev, fields)
# 「发通知」步骤自带标题时**直接用它**:用户写的是"抖音号1 掉线了",
# 前面再挂一个"任务自定义通知:"纯属噪音(其它事件照旧带事件名)
title = str(fields["title"]) if fields.get("title") else 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 k in ("job_name", "device_name", "apk_name", "hook_name", "title") \
and str(v) == str(subj):
continue
# 设备名优先:两个都带时只显示名字——IP 是给日志看的,通知里没人想读地址
if k == "serial" and fields.get("device_name"):
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]))
# 样本只在**真的合并了多条**时才有意义:合并 1 条时样本就是标题主体本身
if merged and merged.get("merged_count", 0) > 1:
samples = [s for s in (merged.get("samples") or [])[:3]
if s and str(s) != str(subj)]
if samples:
rows.append(("样本", ";".join(samples)))
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):
"""聚合样本:一句话说明"合并进来的是哪几条"。
**优先设备维度**(与标题相反):聚合键是 job_id 时,合并的是"同一任务的多台
设备",样本再写任务名就每条都一样(标题里已经有了),写设备名才看得出是哪几台。
"""
subj = ""
for k in ("device_name", "serial", "job_name", "apk_name", "hook_name"):
if fields.get(k):
subj = str(fields[k])
break
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
fill_device_names(fields) # 地址 → 设备名(纯内存快照,不查库)
for h in hooks:
keys = tuple(str(fields.get(k, "")) for k in (ev.agg_key or ()))
# 聚合窗口:事件自己声明「立即发」(agg_window=0,低频高危事件,如批次结束 /
# 服务启停 / 备份恢复)时**不受 hook 窗口影响**——否则白等一个窗口,还会
# 挂上一行毫无信息量的「样本」(单条事件的样本=标题主体,纯重复)。
# 其余事件听 hook 的:那个值表达的是"这个群最多等多久合并"。
window = 0 if ev.agg_window == 0 else int(h.get("agg_window") or 0)
try:
_event_q.put_nowait((h["id"], event, dict(fields), keys, window))
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