fix(通知): 多设备批次不再逐台报成功;批次汇总瘦身并列出失败设备;设备只显示名字
用户反馈:一条任务覆盖 13 台设备,每台成功都推一条,群里被刷屏;批次汇总里
「任务:抖音养号」与标题重复、「样本:抖音养号」毫无信息量、「stopped」还是英文。
- task_manager:
· **多设备批次不发 `task.device.success`**(批次汇总里已有成功台数;单台任务照发);
失败仍逐台发——那是少数,且要知道是哪台;
· `_BatchTracker` 收集失败设备(`名字(型号):原因`,最多 5 条)并加进批次汇总;
· 设备级事件带 `model`:型号取自**设备池快照**(批次开始时一次查好带下去),
不是 worker 每次连接现采的(那次 `d.info()` 可能超时 → 空),也不查库;
- notifier:
· 聚合窗口口径修正:事件声明 `agg_window=0`(批次结束/服务启停/备份恢复等低频
高危事件)**不再被 hook 的窗口拖住**——原先事件级设置完全失效;
· 「样本」行只在真的合并了多条(>1)时出现,且按**设备**维度写(合并多台时写
任务名每条都一样);与标题重复的字段行(「任务:x」)不再重复渲染;
· 设备名与地址同时存在时**只显示名字**(IP 是给日志看的);
· 补 `stopped`/`failed_devices` 中文标签(原先直接显示英文 key);
- notify_events:批次事件补 `skipped`/`failed_devices` 字段,设备事件补 `model`。
实测(dev 真机 + 假接收端):
批次 2 台 → 只收到 batch.started/finished,无 device.success,汇总 总数2/成功1/跳过1;
单台任务 → 收到 device.success,显示「设备:cs1 / 型号:M2010J19SC」不含 IP。
文档:doc/NOTIFY.md §3(多设备批次怎么发 + 批次消息样例)、§6(窗口口径与注意事项)。
This commit is contained in:
+67
-16
@@ -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}")
|
||||
# 归还:本任务是抢占任务,结束后自动重新启动被抢占的原任务
|
||||
|
||||
Reference in New Issue
Block a user