Merge branch 'fix/notify-noise'——通知降噪:多设备批次不逐台报成功、批次汇总瘦身、设备只显示名字

This commit is contained in:
2026-09-16 10:32:51 +08:00
4 changed files with 150 additions and 31 deletions
+31 -5
View File
@@ -331,6 +331,7 @@ _FIELD_LABELS = {
"timeout_s": "超时(秒)", "task_job": "原任务", "model": "型号",
"apk_name": "应用", "apk_id": "应用ID", "package_name": "包名",
"success": "成功", "failed": "失败", "skipped": "跳过", "total": "总数",
"stopped": "停止", "failed_devices": "失败设备",
"failed_items": "失败明细", "filename": "文件", "size": "大小(字节)",
"tables": "表数", "include_apk": "含APK", "user": "操作人", "operator": "操作人",
"env": "环境", "db_target": "数据库", "device_count": "设备数",
@@ -379,13 +380,25 @@ def build_message(ev, fields, hook=None):
v = fields.get(k)
if v in (None, "", [], {}):
continue
# 标题里已经写出来的主体(任务名/设备名)不再重复占一行:
# "任务批次结束:抖音养号" 下面再来一行 "任务:抖音养号" 是纯噪音
if k in ("job_name", "device_name", "apk_name", "hook_name") \
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]))
if merged and merged.get("samples"):
rows.append(("样本", ";".join(merged["samples"][:3])))
# 样本只在**真的合并了多条**时才有意义:合并 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']} 条被限流折叠"))
@@ -913,7 +926,16 @@ def _dispatcher_loop():
def _sample(ev, fields):
subj = _subject(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 ''}"
@@ -976,9 +998,13 @@ def notify(event, **fields):
return
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,
int(h.get("agg_window") or 0)))
_event_q.put_nowait((h["id"], event, dict(fields), keys, window))
except queue.Full:
_dropped["event_q"] += 1
except Exception as e: # 通知出问题绝不能影响业务
+14 -9
View File
@@ -38,8 +38,10 @@ EVENTS = [
["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),
["job_id", "job_name", "total", "success", "failed", "stopped", "skipped",
"failed_devices", "duration_s"],
"整批跑完(成功/失败/停止台数 + 耗时;有失败时列出失败设备与型号)",
agg_window=0, recommend=True),
_e("task.batch.no_device", "任务无可用设备", "任务批次",
["job_id", "job_name", "target"],
"触发时一台可用设备都没有(任务空跑)", agg_window=0, recommend=True),
@@ -52,22 +54,25 @@ EVENTS = [
# ---------------- 任务 · 单设备(权威结论点) ----------------
_e("task.device.success", "设备任务成功", "任务·单设备",
["serial", "device_name", "job_id", "job_name", "attempt", "duration_s"],
"某台设备上的任务最终成功", agg_key=_BY_JOB, recommend=True),
["serial", "device_name", "model", "job_id", "job_name", "attempt", "duration_s"],
"某台设备上的任务最终成功。**只在目标设备只有 1 台时发送**——多设备批次看"
"「任务批次结束」就够了,13 台设备就是 13 条刷屏", agg_key=_BY_JOB),
_e("task.device.failed", "设备任务失败", "任务·单设备",
["serial", "device_name", "job_id", "job_name", "attempts", "cause", "msg"],
"某台设备上的任务最终失败(带真实原因)", agg_key=_BY_JOB, recommend=True),
["serial", "device_name", "model", "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"],
["serial", "device_name", "model", "job_id", "job_name", "error"],
"设备离线,不重试直接失败", agg_key=_BY_JOB, recommend=True),
_e("task.device.error", "单次尝试异常", "任务·单设备",
["serial", "device_name", "job_name", "attempt", "error"],
["serial", "device_name", "model", "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"],
["serial", "device_name", "model", "job_name", "attempt", "phase"],
"用户手动停止 / cron 停止", agg_key=_BY_JOB),
_e("task.device.preempted", "设备被抢占", "任务·单设备",
["serial", "device_name", "job_name", "preempted_job_id", "preempted_job_name"],
+67 -16
View File
@@ -49,6 +49,34 @@ _log = get_logger("core.tm")
_START_STAGGER_SEC = 0.2
def _dev_model(serial, tracker=None):
"""设备型号:优先用批次开始时查好的设备池快照,退回 worker 上报的状态表。
为什么优先快照:worker 每次启动都会 `_update_status(model=…)`,但那次采集
(`d.info()`)可能超时/失败 → 型号是空;而设备池的型号是后台统一采集的、
稳定得多。快照还有个好处:不查库(通知路径上不碰 DB)。
"""
if tracker is not None:
m = (tracker.models or {}).get(serial)
if m:
return m
with _WORKERS_LOCK:
return _WORKERS.get(serial, {}).get("model") or ""
def _dev_fail_note(serial, dname, tracker=None):
"""失败设备的一行说明:`名字(型号):原因`——批次汇总里直接列出来。
型号同上用快照;原因取 worker 上报的 last_error(刚写进去,还在)。
"""
with _WORKERS_LOCK:
err = _WORKERS.get(serial, {}).get("last_error") or ""
cause = err.strip().splitlines()[0] if err.strip() else ""
model = _dev_model(serial, tracker)
head = f"{dname}({model})" if model else str(dname)
return f"{head}:{cause[:60]}" if cause else head
class _BatchTracker:
"""一次任务批次的收尾统计:**所有设备都出结果**后发一条 `task.batch.finished`。
@@ -59,29 +87,40 @@ class _BatchTracker:
放在 finally 里是为了保证"任何一个 return 分支都会被计数一次",不会漏也不会重。
"""
def __init__(self, job, total, names=None):
def __init__(self, job, total, names=None, models=None):
self.job = job
self.total = total
# serial -> 设备名(通知里显示名字而不是裸地址);取不到就是空
self.names = dict(names or {})
# serial -> 型号(批次开始时从设备池查一次带下来,通知里直接用)
self.models = dict(models or {})
self._lock = threading.Lock()
self._left = total
self._stats = Counter()
self._failed = [] # 失败设备明细(型号 + 原因),汇总里有失败时列出
self._t0 = time.time()
def done(self, serial, outcome="failed"):
"""登记一台设备的结果(success/failed/stopped/skipped)。"""
def done(self, serial, outcome="failed", detail=""):
"""登记一台设备的结果(success/failed/stopped/skipped)。
detail 只对 failed 有意义("设备名(型号):原因"),最多留 5 条 ——
批次汇总里列出来,用户不用再去翻日志。
"""
with self._lock:
self._left -= 1
self._stats[outcome] += 1
if outcome == "failed" and detail and len(self._failed) < 5:
self._failed.append(detail)
left = self._left
s = dict(self._stats)
fails = list(self._failed)
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),
failed_devices=fails,
duration_s=int(time.time() - self._t0))
@@ -722,14 +761,15 @@ class TaskManager:
serials_preview=serials[:5])
# 批次收尾统计:所有设备都出结果后发一条汇总(见 _BatchTracker)
# 设备名一次查好带下去(worker 线程里没有 app context,查库要显式包 context)
names = {}
names, models = {}, {}
try:
with self._db():
names = {d["serial"]: (d.get("name") or "")
for d in device_pool.list_devices()}
rows = device_pool.list_devices()
names = {d["serial"]: (d.get("name") or "") for d in rows}
models = {d["serial"]: (d.get("model") or "") for d in rows}
except Exception as e:
_log.warning(f"读取设备名失败(通知里将显示地址): {e}")
tracker = _BatchTracker(job, len(serials), names)
_log.warning(f"读取设备名/型号失败(通知里将显示地址): {e}")
tracker = _BatchTracker(job, len(serials), names, models)
for idx, serial in enumerate(serials):
# 每台设备一个重试循环线程,互不影响;错峰延迟在各自线程内等待
@@ -848,11 +888,16 @@ class TaskManager:
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()))
# 单设备成功:**只在"这一批本来就只有一台设备"时发**。
# 多设备批次里逐台报成功没有信息量(批次汇总里已有成功台数),
# 13 台就是 13 条刷屏;失败仍然逐台发——那是少数且要知道是哪台。
if tracker is None or tracker.total <= 1:
notifier.notify("task.device.success", serial=serial,
device_name=dname,
model=_dev_model(serial, tracker),
job_id=job.id, job_name=job.name,
attempt=attempt,
duration_s=int(time.time() - t_attempt))
return
# 用户请求停止(无论 attempt 第几次、status 是什么)→ 不重试
with self._lock:
@@ -874,7 +919,9 @@ class TaskManager:
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])
model=_dev_model(serial, tracker),
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)
@@ -882,6 +929,7 @@ class TaskManager:
self._running.pop(serial, None)
_update_status(serial, task_job="")
notifier.notify("task.device.error", serial=serial, device_name=dname,
model=_dev_model(serial, tracker),
job_name=job.name, attempt=attempt,
error=str(e)[:200])
@@ -925,7 +973,8 @@ class TaskManager:
_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,
model=_dev_model(serial, tracker), job_id=job.id,
job_name=job.name, attempts=max_attempts,
cause=cause, msg=msg)
finally:
# 清除停止标志:整个重试循环结束(成功/失败/停止)后允许下次任务
@@ -934,7 +983,9 @@ class TaskManager:
# 批次统计:**每台设备只在这里上报一次**(所有 return 分支都会走到 finally)
if tracker is not None:
try:
tracker.done(serial, outcome)
tracker.done(serial, outcome,
detail=_dev_fail_note(serial, dname, tracker)
if outcome == "failed" else "")
except Exception as e:
_log.warning(f"批次统计上报失败(不影响任务): {e}")
# 归还:本任务是抢占任务,结束后自动重新启动被抢占的原任务
+38 -1
View File
@@ -65,6 +65,32 @@ daemon 线程。所以:
> 是"这台设备最终成功/失败"的**唯一权威点**。所以一个设备重试 3 次后失败,只会收到 **1 条**
> `task.device.failed`,不会收到 3 条噪音。
**多设备批次怎么发**(设备一多,逐台发就等于刷屏):
| 事件 | 多设备批次(>1 台) | 单台任务 |
|---|---|---|
| `task.batch.started` | 发 1 条(N 台) | 发 |
| `task.device.success` | **不发**(批次汇总里已有成功台数) | 发 |
| `task.device.failed` / `.offline` | 逐台发(失败是少数,且要知道是哪台) | 发 |
| `task.batch.finished` | **发 1 条汇总** | 发 |
批次汇总长这样(有失败时才会多出一行「失败设备」,列出名字(型号):原因,最多 5 条):
```
### ✅ 任务批次结束:抖音养号
> **任务ID**:b370bbbc
> **总数**:13
> **成功**:13
> **失败**:0
> **停止**:0
> **跳过**:0
> **耗时(秒)**:1347
> **时间**:2026-09-16 09:22:57
```
单设备事件里的设备一律**显示名字**(`cs1`)而不是地址;只有设备没命名时才退回 IP。
`型号` 取自设备池的快照(后台统一采集的那份),比 worker 每次连接时现采的稳。
---
## 4. 配置
@@ -147,13 +173,24 @@ Bark `/…/<key>`、Slack `/services/T…/B…/X…`)→ 接口回显、发送
| 层 | 机制 | 默认 |
|---|---|---|
| L1 聚合 | 同 webhook、同事件、同聚合键(如 `job_id`)在一个窗口内合并成一条,保留前 3 个样本 | 30s(低频高危事件设 0,立即发) |
| L1 聚合 | 同 webhook、同事件、同聚合键(如 `job_id`)在一个窗口内合并成一条,保留前 3 个样本 | 30s(**事件声明 `agg_window=0` 的一律立即发**,如批次结束 / 服务启停 / 备份恢复) |
| L2 限流 | 每 webhook 一个令牌桶 | 18 条/分(企业微信硬限 20,留余量) |
| L3 折叠 | 被限流的事件**不丢弃**,压成一条「被限流折叠 N 条」摘要 | 最多 60s 一条 |
| L4 背压 | 有界队列(event 2000 / send 1000),满了丢弃并计数 | 溢出会告警一次 |
取舍写明白:**失败通知最多延迟一个聚合窗口(默认 30s)**,换来群不被刷屏。
窗口口径(代码在 `core/notifier.py` 的 `notify()`):
- 事件自己写了 `agg_window=0` → **立即发**,不受 webhook 的窗口影响;
- 其余事件 → 用该 webhook 配置的 `agg_window`(它表示"这个群最多等多久合并")。
「样本」行只在**真的合并了多条**(>1)时出现,且按**设备**维度写(`cs1 · 超时`)——
合并多台设备时写任务名每条都一样,等于没写。
> ⚠️ 保存通知配置(`save_config`)会清空**待发聚合**与限流令牌桶。清理 hook 时
> 顺手丢掉的正是还没到窗口的聚合事件——排障时别把它当成"没发"。
---
## 7. 开发:给新功能加通知