diff --git a/core/notifier.py b/core/notifier.py index 1caa6d9..7f7a595 100644 --- a/core/notifier.py +++ b/core/notifier.py @@ -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: # 通知出问题绝不能影响业务 diff --git a/core/notify_events.py b/core/notify_events.py index 5c00803..b035bc8 100644 --- a/core/notify_events.py +++ b/core/notify_events.py @@ -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"], diff --git a/core/task_manager.py b/core/task_manager.py index 20f9af2..e71fdad 100644 --- a/core/task_manager.py +++ b/core/task_manager.py @@ -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}") # 归还:本任务是抢占任务,结束后自动重新启动被抢占的原任务 diff --git a/doc/NOTIFY.md b/doc/NOTIFY.md index 0c59fbc..85ba9f2 100644 --- a/doc/NOTIFY.md +++ b/doc/NOTIFY.md @@ -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 `/…/`、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. 开发:给新功能加通知