diff --git a/README.md b/README.md index 2c90cca..a0ad241 100644 --- a/README.md +++ b/README.md @@ -139,6 +139,7 @@ auto_control/ │ ├── apk_manager.py # APK 上传/解析/批量安装 │ ├── system_backup.py # 数据备份导出/导入(重启生效) │ ├── notifier.py # 通知分发(队列/聚合/限流/适配器)+ notify_events.py 事件目录 +│ ├── step_log.py # 任务步骤明细(异步写线程 + 保留期清理) │ ├── tailscale_client.py # Tailscale API v2 客户端 │ ├── ssh_client.py # SSH 封装(当前无人调用,预留) │ ├── logger.py # 分文件日志(core/task/web/action) @@ -184,7 +185,7 @@ auto_control/ ``` base.js → markdown.js → list.js → monitor.js → editor.js → tasks.js - → tools.js → apps.js → admin.js → agent.js → taskgen.js → system.js → notify.js + → tools.js → apps.js → admin.js → agent.js → taskgen.js → system.js → notify.js → steplog.js ``` --- @@ -246,7 +247,7 @@ self.set_progress(done=5, total=80, unit="视频", action_counts={"like": 3}, el |-----|------|--------| | **监控** | 统计卡片 + 设备表(在线/型号/任务状态/进度/前台 App/最近错误,支持排序、搜索、分页、多选批量操作)+ 异常汇总 + **任务运行概况**(每张任务卡:执行任务 / 停用任务 + 覆盖设备彩色 chip);导航栏「📺 大屏」打开 `/wall` | 所有登录用户(设备操作按钮需"设备控制") | | **任务** | 子分栏:**任务计划**(CRUD / 启停 / 执行 / 下次运行时间 / 离线跳过 / 步骤编辑器)、**自定义动作**(步骤打包复用) | 可看;写操作需"任务管理" | -| **日志** | 实时日志(按文件切换、可自动刷新) | 需"日志查看" | +| **日志** | 子分栏:**文件日志**(文件切换含滚动历史、关键字/级别/时间过滤、命中高亮、下载筛选结果、自动刷新)、**步骤明细**(按设备/任务/结果/时间过滤、按运行归组的概览、导出 CSV) | 需"日志查看" | | **用户** | 用户 CRUD、改密、分配权限 | 仅管理员 | | **工具** | 8 个子分栏:剪贴板注入 / adb 远程终端 / Tailscale 管理 / 应用管理 / 应用版本管理 / 设备已装应用 / **设备池管理**(含自动发现)/ 设备分组 | 仅管理员 | | **AI 控制台** | 会话列表 + 对话区(Markdown 渲染、推理链折叠、token 统计)+ 实时画面(MJPEG)+ 目标设备选择;右上角:经验库 / 动作库 / 模型配置 | 仅管理员 | @@ -357,7 +358,16 @@ from core.logger import get_logger log = get_logger("task.generic") # → logs/task.log ``` -「日志」Tab 可在线查看。 +「日志」Tab 可在线查看:文件下拉包含 `.1/.2…` 滚动历史(可回溯更早的日志), +支持关键字(命中高亮)/最低级别/时间范围过滤,点「下载」导出当前筛选结果 +(无筛选就是整个文件)。读取逻辑在 `core/logger.py` 的 `list_log_files()/query_log()/ +read_log_text()`,`file` 参数走白名单,不接受任意路径。 + +「日志 → 步骤明细」是**结构化**的另一半:每一次步骤执行落一行 `task_step_log` +(设备/任务/步骤路径/结果/耗时/选择器),因此能按设备、任务、结果、时间过滤, +也能按 `run_id` 归组看"这一次运行为什么失败"、导出 CSV。写入是异步的 +(`core/step_log.py`,任务线程只入队,不阻塞执行),保留 14 天、单次运行最多 +2000 条 —— 这两个上限决定这张表(以及备份包)能长多大。 --- diff --git a/core/device_worker.py b/core/device_worker.py index 8efab40..7342cd8 100644 --- a/core/device_worker.py +++ b/core/device_worker.py @@ -275,10 +275,13 @@ class BaseWorker(threading.Thread): 其他异常 — 可重试 """ - def __init__(self, serial, params=None, daemon=True): + def __init__(self, serial, params=None, daemon=True, ctx=None): super().__init__(daemon=daemon) self.serial = serial self.params = params or {} + # 本次运行的上下文(run_id / job_id / job_name / device_name),由 + # TaskManager 传入,供步骤明细等"结构化记录"标注来源。可空。 + self.ctx = dict(ctx or {}) self._stop_flag = threading.Event() self.d = None # u2.Device,run_task 里用 # 通用进度字段(子类通过 set_progress 上报) diff --git a/core/logger.py b/core/logger.py index cc30904..a40981e 100644 --- a/core/logger.py +++ b/core/logger.py @@ -16,7 +16,9 @@ 文件按 10MB 滚动,保留 5 个历史文件。 """ import os +import re import logging +from collections import deque from logging.handlers import RotatingFileHandler _LOG_DIR = os.path.join(os.path.dirname(os.path.dirname(os.path.abspath(__file__))), "logs") @@ -84,3 +86,192 @@ def log_print(msg, level="info", module="core"): """print 的替代品,转发到 logging。""" log = get_logger(module) getattr(log, level, log.info)(msg) + + +# ================== 日志读取(「日志」页的筛选 / 下载) ================== +# +# 读日志的代码放这里而不是 web 层:文件命名与滚动规则归本模块管 +# (_MODULE_FILES + RotatingFileHandler 的 backupCount),**白名单校验必须和 +# 写入侧用同一份事实**,否则改了文件名就会漏掉一处,变成目录穿越。 + +# 行格式与 _FORMAT 对应:2026-09-16 00:27:11 [INFO] [web.notify] 消息 +_LINE_RE = re.compile(r"^(\d{4}-\d{2}-\d{2} \d{2}:\d{2}:\d{2}) " + r"\[([A-Z]+)\] \[([^\]]*)\] ?(.*)$") +_LEVEL_ORDER = {"DEBUG": 10, "INFO": 20, "WARNING": 30, "ERROR": 40, "CRITICAL": 50} + +# 单次读取的行数上限:文件本身被限在 10MB(约 10 万行),这里再兜一层, +# 万一将来有人调大 maxBytes,也不会把内存吃爆。 +_MAX_SCAN_LINES = 500000 + +_MODULE_LABELS = { + "core.log": "核心(adb/调度)", + "task.log": "任务执行", + "web.log": "Web/AI", + "action.log": "动作执行", + "notify.log": "通知发送", +} + + +def _rotation_files(base): + """基础名 + 磁盘上实际存在的滚动副本(core.log → core.log.1/.2…)。 + + RotatingFileHandler 的备份编号**越大越旧**,所以 base 最新,接着 .1/.2… + """ + names = [] + if os.path.exists(os.path.join(_LOG_DIR, base)): + names.append(base) + i = 1 + while i <= 99 and os.path.exists(os.path.join(_LOG_DIR, f"{base}.{i}")): + names.append(f"{base}.{i}") + i += 1 + return names + + +def list_log_files(): + """可查看的日志文件(模块文件 + 滚动历史),最近写入的排前面。 + + `file` 参数只能取这里的 name(白名单),杜绝 `../../` 之类的路径穿越。 + """ + out = [] + for base in sorted(set(_MODULE_FILES.values())): + label = _MODULE_LABELS.get(base, base) + for name in _rotation_files(base): + try: + st = os.stat(os.path.join(_LOG_DIR, name)) + except OSError: + continue + rotated = name != base + out.append({ + "name": name, + "label": f"{label} · {name}" + ("(历史)" if rotated else ""), + "module": base, + "size": st.st_size, + "mtime": int(st.st_mtime), + "rotated": rotated, + }) + out.sort(key=lambda f: (-f["mtime"], f["name"])) + return out + + +def norm_ts(value, end=False): + """把前端的时间输入归一成 'YYYY-MM-DD HH:MM:SS'(可比字符串)。 + + 接受 datetime-local 的 'YYYY-MM-DDTHH:MM'、'YYYY-MM-DD HH:MM' 和纯日期; + 纯日期按整天算(end=True 取 23:59:59,否则 00:00:00)。 + """ + s = (value or "").strip().replace("T", " ") + if not s: + return "" + if len(s) == 10: # 只给了日期 → 按整天 + return s + (" 23:59:59" if end else " 00:00:00") + if len(s) == 16: # YYYY-MM-DD HH:MM + return s + (":59" if end else ":00") + return s[:19] + + +def _iter_filtered(fh, keyword, min_level, since, until, stats=None): + """逐行读文件,产出命中的 {ts, level, module, msg, cont}。 + + 过滤口径: + - 时间/级别用**继承值**——续行(traceback 的缩进行等)本身没有前缀, + 继承上一条带前缀的行,这样按 ERROR 筛选时整段堆栈不会被拆散; + - 关键字匹配本行原文(续行也参与),大小写不敏感。 + + stats(可选 dict)里回填 scanned:本次实际读了多少行(用于告诉用户 + "是读到文件开头了,还是被行数上限截断的")。 + """ + kw = (keyword or "").strip().lower() + floor = _LEVEL_ORDER.get((min_level or "").upper(), 0) + cur_ts = cur_level = cur_module = "" + for raw in fh: + if stats is not None: + stats["scanned"] += 1 + line = raw.rstrip("\r\n") + m = _LINE_RE.match(line) + if m: + cur_ts, cur_level, cur_module, msg = m.groups() + cont = False + else: + # 续行(traceback、被切成多行的长消息):沿用上一条的元信息 + msg, cont = line, True + if since and (not cur_ts or cur_ts < since): + continue + if until and (not cur_ts or cur_ts > until): + continue + if floor and _LEVEL_ORDER.get(cur_level, 0) < floor: + continue + if kw and kw not in line.lower(): + continue + yield {"ts": cur_ts, "level": cur_level, "module": cur_module, + "msg": msg, "cont": cont} + + +def query_log(name, keyword="", min_level="", since="", until="", + limit=300, max_scan=_MAX_SCAN_LINES): + """取日志文件里**最近** limit 条满足条件的行。 + + 返回 dict: + rows — [{ts, level, module, msg, cont}],保持文件原顺序(旧→新) + matched — 命中总数(可能大于 len(rows):更早的命中没返回) + scanned — 实际读了多少行(达到 max_scan 说明是被上限截断的) + truncated — matched > len(rows),即"还有更早的命中没显示" + error — 出错时的原因(如文件不存在) + + 正序读整文件 + deque 保留最后 N 条:文件被 RollingFileHandler 限在 10MB + (约 10 万行),整读一次百毫秒级;换来的是续行能正确继承时间戳——倒着 + 读的写法在遇到 traceback 时必须把续行攒着等前面那行,容易出错。 + """ + # 白名单:`name` 必须来自 list_log_files(),否则一律拒绝(防目录穿越) + if name not in {f["name"] for f in list_log_files()}: + return {"rows": [], "matched": 0, "scanned": 0, "truncated": False, + "error": "日志文件不存在或不可读取"} + since, until = norm_ts(since), norm_ts(until, end=True) + try: + keep = max(1, min(int(limit or 300), 5000)) + except (TypeError, ValueError): + keep = 300 + rows, matched, stats = deque(maxlen=keep), 0, {"scanned": 0} + try: + with open(os.path.join(_LOG_DIR, name), encoding="utf-8", + errors="replace") as fh: + for row in _iter_filtered(fh, keyword, min_level, since, until, stats): + matched += 1 + rows.append(row) + if stats["scanned"] >= max_scan: + break + except OSError as e: + return {"rows": [], "matched": 0, "scanned": 0, "truncated": False, + "error": f"读取失败: {e}"} + return {"rows": list(rows), "matched": matched, "scanned": stats["scanned"], + "truncated": matched > len(rows), "error": ""} + + +def read_log_text(name, keyword="", min_level="", since="", until="", + limit=200000): + """导出用的整段文本:命中行按原格式拼回去(下载按钮)。 + + 没给任何条件时直接读整个文件(含滚动历史由调用方选定的那一个)。 + """ + if name not in {f["name"] for f in list_log_files()}: + return None, "日志文件不存在或不可读取" + since, until = norm_ts(since), norm_ts(until, end=True) + any_filter = bool((keyword or "").strip() or min_level or since or until) + path = os.path.join(_LOG_DIR, name) + try: + with open(path, encoding="utf-8", errors="replace") as fh: + if not any_filter: + return fh.read(), "" + buf, n = [], 0 + for row in _iter_filtered(fh, keyword, min_level, since, until): + if row["cont"]: + buf.append(row["msg"]) # 续行原样输出 + else: + buf.append(f"{row['ts']} [{row['level']}] " + f"[{row['module']}] {row['msg']}") + n += 1 + if n >= limit: + buf.append(f"...(已达导出上限 {limit} 行,请缩小时间范围)") + break + return "\n".join(buf) + ("\n" if buf else ""), "" + except OSError as e: + return None, f"读取失败: {e}" diff --git a/core/models.py b/core/models.py index b053c2b..bddeb68 100644 --- a/core/models.py +++ b/core/models.py @@ -369,6 +369,40 @@ class DeviceInstallLog(db.Model): created_at = db.Column(db.String(20), default="") +class TaskStepLog(db.Model): + """任务步骤明细:每一次步骤执行的落库记录(「日志 → 步骤明细」页)。 + + 与 `logs/task.log` 的分工:文本日志是**排障时的原始现场**(什么都往里写、 + 10MB 滚动),本表是**结构化的一份**——设备/任务/步骤/结果/耗时都是列, + 所以能按设备、任务、时间、结果过滤和统计,文本日志只能 grep。 + + 写入方是 `core/step_log.py` 的专用写线程(异步批量落库),**任务线程不直接 + 写库**:一次运行可能上万步,每步一次 INSERT 会拖慢热路径。 + + 保留期由 `core/step_log.py` 的清理任务控制(默认 14 天,见该模块常量)。 + """ + __tablename__ = "task_step_log" + id = db.Column(db.Integer, primary_key=True, autoincrement=True) + run_id = db.Column(db.String(24), default="", index=True) # 一次运行=设备×任务×第几次尝试 + job_id = db.Column(db.String(32), default="") + job_name = db.Column(db.String(120), default="") + serial = db.Column(db.String(120), default="") + device_name = db.Column(db.String(80), default="") + step_path = db.Column(db.String(32), default="") # 嵌套位置,如 "2.1.3" + step_label = db.Column(db.String(120), default="") + step_type = db.Column(db.String(40), default="") # click_el / loop / input_text … + selector = db.Column(db.String(300), default="") # 元素选择器(长选择器截断) + result = db.Column(db.String(16), default="") # ok/skip/miss/error/unknown/cap + detail = db.Column(db.String(500), default="") # 异常消息、跳过原因等 + duration_ms = db.Column(db.Integer, default=0) + created_at = db.Column(db.String(20), default="", index=True) + + __table_args__ = ( + db.Index("ix_step_log_serial_ts", "serial", "created_at"), + db.Index("ix_step_log_job_ts", "job_id", "created_at"), + ) + + class AgentConversation(db.Model): """AI 控制台会话:整个消息序列以 JSON 存在一行里(单会话几十 KB,够用)。""" __tablename__ = "agent_conversation" 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/step_log.py b/core/step_log.py new file mode 100644 index 0000000..1f9899f --- /dev/null +++ b/core/step_log.py @@ -0,0 +1,247 @@ +"""任务步骤明细:异步落库(专用写线程)+ 保留期清理。 + +为什么不让任务线程直接写库:步骤执行是**热路径**(一次运行可能上万步),每步 +一次 INSERT 要建会话、等 MySQL 往返,还会和业务事务抢连接。这里改成「专用写 +线程 + 有界队列」:任务线程只做一次 put_nowait(微秒级),落库、重排、批量都 +在写线程里做。 + +铁律(与 `core/notifier.py` 一致): + - `record()` 零阻塞、零 DB、**绝不抛异常** —— 记录绝不影响任务执行; + - 落库失败只写日志、丢弃该批,**不重排**(不为补数据把内存撑爆)。 + +线程用裸 `threading.Thread(daemon=True)`(与项目其它后台线程一致;线程池的 +非 daemon 线程会在退出时被 atexit join 卡住)。 +""" +import queue +import threading +import time + +from core.logger import get_logger + +_log = get_logger("task.step") + +# 一次运行最多记录多少条:防「2 小时无限循环」这类任务把表写爆。 +# 达到上限后本次运行不再逐条记,只补一条 cap 说明行。 +MAX_ROWS_PER_RUN = 2000 + +# 保留天数:每日清理任务删除更早的记录。**这个值直接决定备份包体积** +# (task_step_log 会进整库备份,见 doc/DATA_MODEL.md §6),改大前先想清楚。 +KEEP_DAYS = 14 + +# 队列上限与批量:队列满 = 丢弃并计数(宁可丢明细,也不能拖住任务线程) +_QUEUE_MAX = 5000 +_BATCH = 200 +_FLUSH_INTERVAL = 1.0 + +_app = None +_q = queue.Queue(maxsize=_QUEUE_MAX) +_writer = None +_stop = threading.Event() +_stats = {"queued": 0, "written": 0, "dropped": 0, "failed": 0} + + +# ================== 生命周期 ================== + +def init_app(app): + """启动写线程(web_server 启动时调用;重复调用安全)。""" + global _app, _writer + _app = app + if _writer is not None and _writer.is_alive(): + return + _stop.clear() + _writer = threading.Thread(target=_writer_loop, name="step-log-writer", + daemon=True) + _writer.start() + _log.info(f"步骤明细写线程已启动(保留 {KEEP_DAYS} 天,单次运行上限 " + f"{MAX_ROWS_PER_RUN} 条)") + + +def shutdown(timeout=3.0): + """退出前停写线程并尽量排空队列(超时就算了,明细不是关键数据)。""" + _stop.set() + if _writer is not None: + _writer.join(timeout) + + +def stats(): + """本进程内的计数(排障用:看有没有在丢记录)。""" + return dict(_stats, queue_size=_q.qsize()) + + +# ================== 写入侧 ================== + +def record(run_id="", job_id="", job_name="", serial="", device_name="", + step_path="", step_label="", step_type="", selector="", + result="", detail="", duration_ms=0): + """记录一次步骤执行。**任务线程直接调用**——零阻塞、零 DB、不抛异常。 + + 队列满时丢弃并计数(dropped),不会阻塞调用方。 + """ + try: + row = { + "run_id": (run_id or "")[:24], + "job_id": (job_id or "")[:32], + "job_name": (job_name or "")[:120], + "serial": (serial or "")[:120], + "device_name": (device_name or "")[:80], + "step_path": (step_path or "")[:32], + "step_label": (step_label or "")[:120], + "step_type": (step_type or "")[:40], + "selector": (selector or "")[:300], + "result": (result or "")[:16], + "detail": (detail or "")[:500], + "duration_ms": int(duration_ms or 0), + "created_at": time.strftime("%Y-%m-%d %H:%M:%S"), + } + _q.put_nowait(row) + _stats["queued"] += 1 + except queue.Full: + _stats["dropped"] += 1 + except Exception: # 截断/类型异常都不该冒泡到任务线程 + _stats["dropped"] += 1 + + +def _collect(): + """取一批待写记录:最多等 _FLUSH_INTERVAL,凑够 _BATCH 就走。""" + rows = [] + try: + rows.append(_q.get(timeout=_FLUSH_INTERVAL)) + except queue.Empty: + return rows + while len(rows) < _BATCH: + try: + rows.append(_q.get_nowait()) + except queue.Empty: + break + return rows + + +def _insert(rows): + """批量落库。失败只记日志——明细丢了可以接受,任务不能被拖住。""" + if not rows or _app is None: + return + from core.models import TaskStepLog, db + try: + with _app.app_context(): + db.session.execute(db.insert(TaskStepLog), rows) + db.session.commit() + _stats["written"] += len(rows) + except Exception as e: + _stats["failed"] += len(rows) + try: + db.session.rollback() + except Exception: + pass + _log.error(f"步骤明细落库失败(丢弃 {len(rows)} 条): {e}") + + +def _writer_loop(): + while not _stop.is_set(): + batch = _collect() + if batch: + _insert(batch) + # 退出前再排空一次,尽量不丢最后几条 + while True: + batch = _collect() + if not batch: + break + _insert(batch) + + +# ================== 查询(给 API 用,在请求线程里跑) ================== + +def query(serial="", job_id="", result="", run_id="", keyword="", + since="", until="", limit=200, offset=0): + """按条件查明细,返回 (rows, total)。**倒序**(最近的在前)。""" + from sqlalchemy import or_ + from core.models import TaskStepLog + q = TaskStepLog.query + if serial: + q = q.filter(TaskStepLog.serial == serial) + if job_id: + q = q.filter(TaskStepLog.job_id == job_id) + if result: + q = q.filter(TaskStepLog.result == result) + if run_id: + q = q.filter(TaskStepLog.run_id == run_id) + if since: + q = q.filter(TaskStepLog.created_at >= since) + if until: + q = q.filter(TaskStepLog.created_at <= until) + if keyword: + kw = f"%{keyword}%" + q = q.filter(or_(TaskStepLog.detail.like(kw), + TaskStepLog.step_label.like(kw), + TaskStepLog.step_type.like(kw), + TaskStepLog.selector.like(kw), + TaskStepLog.job_name.like(kw))) + total = q.count() + rows = (q.order_by(TaskStepLog.id.desc()) + .offset(max(0, int(offset or 0))) + .limit(max(1, min(int(limit or 200), 2000))) + .all()) + return rows, total + + +def runs(serial="", job_id="", since="", until="", limit=100): + """按 run_id 归组的一次运行概览(哪次运行失败、失败几步)。 + + 用一次 GROUP BY 查出来,避免前端按明细自己聚合。 + """ + from sqlalchemy import case, func + from core.models import TaskStepLog, db + q = db.session.query( + TaskStepLog.run_id, + func.min(TaskStepLog.created_at).label("started_at"), + func.max(TaskStepLog.created_at).label("ended_at"), + func.max(TaskStepLog.job_name).label("job_name"), + func.max(TaskStepLog.job_id).label("job_id"), + func.max(TaskStepLog.serial).label("serial"), + func.max(TaskStepLog.device_name).label("device_name"), + func.count(TaskStepLog.id).label("steps"), + func.sum(case((TaskStepLog.result.in_(("error", "miss", "unknown")), 1), + else_=0)).label("failures"), + ) + if serial: + q = q.filter(TaskStepLog.serial == serial) + if job_id: + q = q.filter(TaskStepLog.job_id == job_id) + if since: + q = q.filter(TaskStepLog.created_at >= since) + if until: + q = q.filter(TaskStepLog.created_at <= until) + q = (q.filter(TaskStepLog.run_id != "") + .group_by(TaskStepLog.run_id) + .order_by(func.max(TaskStepLog.id).desc()) + .limit(max(1, min(int(limit or 100), 500)))) + return q.all() + + +def purge_old(keep_days=None, batch=5000): + """删除超过保留期的明细。返回删除行数。 + + 分批删(先取 id 再按 id 删,**不用 DELETE ... LIMIT**——SQLite 默认编译 + 选项不支持它),避免一条长事务把表锁住。 + """ + from core.models import TaskStepLog, db + days = KEEP_DAYS if keep_days is None else int(keep_days) + cutoff = time.strftime("%Y-%m-%d %H:%M:%S", + time.localtime(time.time() - days * 86400)) + deleted = 0 + try: + while True: + ids = [r.id for r in db.session.query(TaskStepLog.id) + .filter(TaskStepLog.created_at < cutoff).limit(batch).all()] + if not ids: + break + deleted += db.session.query(TaskStepLog) \ + .filter(TaskStepLog.id.in_(ids)) \ + .delete(synchronize_session=False) + db.session.commit() + except Exception as e: + db.session.rollback() + _log.error(f"步骤明细清理失败(已删 {deleted} 条): {e}") + return deleted + if deleted: + _log.info(f"步骤明细清理: 删除 {deleted} 条(保留 {days} 天)") + return deleted diff --git a/core/system_backup.py b/core/system_backup.py index d07cd52..3ebbac8 100644 --- a/core/system_backup.py +++ b/core/system_backup.py @@ -56,6 +56,7 @@ TABLE_LABELS = { "pending_device": "待连接设备", "agent_conversation": "AI 会话", "agent_experience": "经验库", "experience_audit": "经验巡检", "agent_action": "动作库", "device_install_log": "设备端安装记录", + "task_step_log": "任务步骤明细", } _STAGE_TTL = 1800 # 导入暂存有效期(秒) diff --git a/core/task_manager.py b/core/task_manager.py index 5a68894..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): # 每台设备一个重试循环线程,互不影响;错峰延迟在各自线程内等待 @@ -825,8 +865,13 @@ class TaskManager: worker = None t_attempt = time.time() # 单次 attempt 耗时(通知里带) + # 本次运行的上下文:步骤明细按 run_id 归组,把同一设备同一次尝试 + # 执行的步骤串起来(重试会生成新的 run_id,两次尝试不混在一起) + run_ctx = {"run_id": uuid.uuid4().hex[:12], "job_id": job.id, + "job_name": job.name, "device_name": dname, + "attempt": attempt} try: - worker = task.create_worker(serial, job.params) + worker = task.create_worker(serial, job.params, ctx=run_ctx) with self._lock: self._running[serial]["worker"] = worker _log.info(f"{serial} 开始任务 {job.name} (第{attempt}/{max_attempts}次)") @@ -843,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: @@ -869,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) @@ -877,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]) @@ -920,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: # 清除停止标志:整个重试循环结束(成功/失败/停止)后允许下次任务 @@ -929,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/API.md b/doc/API.md index 3694fda..c210210 100644 --- a/doc/API.md +++ b/doc/API.md @@ -139,7 +139,12 @@ | POST | `/api/users` | Admin | 新建用户 | | PUT | `/api/users/` | Admin | 更新用户(改密/权限/管理员) | | DELETE | `/api/users/` | Admin | 删除用户 | -| GET | `/api/logs` | G | 读日志文件尾部 N 行 | +| GET | `/api/logs` | G | 读日志(关键字/级别/时间过滤,返回结构化行) | +| GET | `/api/logs/download` | G | 下载日志(同样支持过滤;无过滤即整个文件) | +| GET | `/api/step_logs` | G | 任务步骤明细(按设备/任务/结果/时间/关键字过滤,分页) | +| GET | `/api/step_logs/runs` | G | 按 `run_id` 归组的一次运行概览(几步、失败几步) | +| GET | `/api/step_logs/filters` | G | 步骤明细的筛选项(明细里出现过的设备/任务 + 结果枚举) | +| GET | `/api/step_logs/download` | G | 导出步骤明细 CSV(带 UTF-8 BOM,Excel 直接打开) | ### 2.5 devices(`web/devices_api.py`) @@ -615,7 +620,59 @@ | `POST /api/users` | Admin | `{"username","password","is_admin"?,"perms"?}` | 用户名/密码空、重名 → 400 | | `PUT /api/users/` | Admin | `{"password"?,"is_admin"?,"perms"?}` | 取消最后一个管理员 → 400;不存在 404 | | `DELETE /api/users/` | Admin | — | 删 `admin`/删自己/删最后一个管理员 → 400 | -| `GET /api/logs?file=core.log&lines=300` | G | — | 返回日志尾部 + 可选文件清单 | +| `GET /api/logs?file=core.log&lines=300` | G | `q` 关键字、`level` 最低级别、`since`/`until` 时间 | 见下 §10.1 | +| `GET /api/logs/download?file=core.log` | G | 同上(参数一致) | 以附件返回 `.txt`,文件名带时间戳 | + +#### 10.1 日志查询(`GET /api/logs`) + +参数: + +| 参数 | 说明 | +|------|------| +| `file` | 日志文件名,**白名单校验**(只能取接口返回的 `files[].name`,含 `.1/.2` 滚动历史);非法值返回 `ok:false` | +| `lines` | 返回**最近**多少条命中(默认 300,上限 5000) | +| `q` | 关键字,大小写不敏感,匹配整行原文 | +| `level` | 最低级别:`ERROR` / `WARNING` / `INFO` / `DEBUG`(含以上);留空为全部 | +| `since` / `until` | 时间范围,接受 `YYYY-MM-DDTHH:MM`(datetime-local)或 `YYYY-MM-DD`(按整天) | + +返回: + +```json +{"ok": true, "file": "core.log", + "files": [{"name": "core.log", "label": "核心(adb/调度) · core.log", + "size": 1175910, "mtime": 1789519242, "rotated": false}], + "rows": [{"ts": "2026-09-16 08:26:50", "level": "INFO", + "module": "core.disc", "msg": "发现: …", "cont": false}], + "matched": 11423, "scanned": 11423, "truncated": true} +``` + +- `rows` 保持文件原顺序(旧→新),即"最近 N 条命中"按时间正序排列; +- `cont=true` 是**续行**(traceback 的缩进行等):它本身没有时间/级别前缀, + 过滤时**继承上一条带前缀的行**,所以按 ERROR 筛选不会把堆栈拆散; +- `matched` 是命中总数,`truncated=true` 表示更早的命中没返回(应缩小时间范围或加关键字); +- `scanned` 是实际扫描行数(上限 50 万行)。 + +#### 10.2 任务步骤明细(`/api/step_logs*`) + +数据源是 `task_step_log` 表(**结构化**,与上面的文本日志不是一回事):每一次步骤 +执行一条,能按设备/任务/结果/时间过滤、能按运行归组、能导出 CSV。写入侧见 +[DATA_MODEL.md](DATA_MODEL.md) §2.8 与 `core/step_log.py`。 + +| 接口 | 参数 | 返回 | +|------|------|------| +| `GET /api/step_logs` | `serial` `job_id` `result` `run_id` `q` `since` `until` `limit`(≤2000,默认 200) `offset` | `rows`(**最近的在前**)、`total`、`stats`(本次进程的 queued/written/dropped/failed)、`keep_days`、`max_rows_per_run` | +| `GET /api/step_logs/runs` | `serial` `job_id` `since` `until` `limit`(≤500) | `runs[]`:`run_id` `started_at` `ended_at` `steps` `failures` | +| `GET /api/step_logs/filters` | — | `devices[]` / `jobs[]`(**只列明细里真的出现过的**,不依赖设备池与任务表)、`results[]` | +| `GET /api/step_logs/download` | 同列表接口 | 附件 `.csv`(UTF-8 **带 BOM**),最多 2 万行,按时间正序 | + +`result` 取值:`ok` / `miss`(handler 返回 False,如元素没找到)/ `error`(抛异常)/ +`unknown`(未知步骤类型)/ `skip`(概率未触发)/ `cap`(本次运行已达记录上限)。 + +`run_id` 是"设备 × 任务 × 第几次尝试",同一行里能拿到 `step_path`(如 `2.1.3`) +还原嵌套结构——**步骤明细面板的「最近运行概览 → 查看」就是按它过滤**。 + +> 稳定性:记录走异步队列(`record()` 零阻塞),**队列满会丢弃**(计数在 `stats.dropped`, +> 页面提示栏会显示),所以明细允许缺条 —— 权威结论仍看任务状态与 [NOTIFY.md](NOTIFY.md) 的事件。 --- diff --git a/doc/ARCHITECTURE.md b/doc/ARCHITECTURE.md index 2a6a8ff..645de38 100644 --- a/doc/ARCHITECTURE.md +++ b/doc/ARCHITECTURE.md @@ -31,6 +31,7 @@ ┌───────────────────────────▼──────────────────────────────────────────┐ │ 基础层 adb_helper · u2_helper · uiauto_helper · ocr · clipboard │ │ notifier(通知分发:队列/聚合/限流/适配器) │ +│ step_log(步骤明细:队列 + 批量落库 + 保留期清理) │ │ models(SQLite)· logger · config │ └──────────────────────────────────────────────────────────────────────┘ ▲ @@ -78,15 +79,15 @@ |------|------|---------| | **B** | `:28-58` | `Flask(__name__)`;会话密钥(`.env` 的 `WEB_SECRET_KEY`,缺失则随机生成并 warning);`TEMPLATES_AUTO_RELOAD=True`;**数据库目标由 `core/db_config` 装配**(`.env` 的 `DEPLOY_ENV`/`DB_*` → URI + 引擎参数),配置错直接 `SystemExit(2)`;`LoginManager` + `login_view="auth.login"` | | **C** | `:60-67` | **恢复任务消费** `consume_pending_restore()`。SQLite 时代它必须在 engine 首次打开 `users.db` **之前**(Windows 无法替换被持有的文件);改用 MySQL 后这一步的语义会变成"启动期事务替换",见 §7 | -| **D** | `:69-80` | `init_db(app)`(建表 → 补列 → 版本账本 → 唯一索引 → 默认管理员 → 旧 JSON 迁移)→ **库环境标签校验 + 启动横幅**(`db_config.verify_deployment_label/print_banner`,不符拒绝启动);`device_pool.init_app`(**起一次性线程**,3s 后采集型号);`device_discovery.init_app`(**起常驻扫描线程**);`TaskManager(app=app)`(APScheduler + 看门狗 + 从库加载分组/任务 + 重注册 cron);`ApkManager(app=app)` | +| **D** | `:69-88` | `init_db(app)`(建表 → 补列 → 版本账本 → 唯一索引 → 默认管理员 → 旧 JSON 迁移)→ `notifier.init_app`(通知 dispatcher/sender 线程)与 `step_log.init_app`(步骤明细写线程)→ **恢复任务消费** `consume_pending_restore()` → **库环境标签校验 + 启动横幅**(`db_config.verify_deployment_label/print_banner`,不符拒绝启动);`device_pool.init_app`(**起一次性线程**,3s 后采集型号);`device_discovery.init_app`(**起常驻扫描线程**);`TaskManager(app=app)`(APScheduler + 看门狗 + 从库加载分组/任务 + 重注册 cron);`ApkManager(app=app)` | ### 2.3 阶段 E~G:蓝图、巡检调度器、真正启动 | 阶段 | 位置 | 做了什么 | |------|------|---------| | **E** | `:64-67` | `context.init(...)`;`register_blueprints(app)`(10 个蓝图);`agent_api.set_app(app)`(供后台线程推 app context) | -| **F** | `:71-80` | **第二个独立 APScheduler**:`CronTrigger(hour=3, minute=47)` 挂经验库巡检;失败仅 warning | -| **G** | `:253-265`(`__main__`) | `_ensure_uiauto_running()`(拉起 uiautodev:20242,写 `data/uiauto.pid`,`atexit` 清理)→ `_preconnect_pool_devices()`(后台并发 connect 池内网络设备)→ `_run_server()`(候选端口依次 bind:`0.0.0.0:18050` → `127.0.0.1:18050` → `127.0.0.1:18051..18055`);退出时 `mgr.shutdown()` + `device_discovery.shutdown()` + 停 uiautodev | +| **F** | 同上附近 | **第二个独立 APScheduler**:`CronTrigger(hour=3, minute=47)` 挂经验库巡检、`hour=4, minute=13` 挂步骤明细清理(`_purge_step_log`,自建 app context);失败仅 warning | +| **G** | `__main__` | `_ensure_uiauto_running()`(拉起 uiautodev:20242,写 `data/uiauto.pid`,`atexit` 清理)→ `_preconnect_pool_devices()`(后台并发 connect 池内网络设备)→ `_purge_step_log_async()`(后台清理超期步骤明细)→ `_run_server()`(候选端口依次 bind:`0.0.0.0:18050` → `127.0.0.1:18050` → `127.0.0.1:18051..18055`);退出时 `notifier.shutdown()` + `step_log.shutdown()` + `mgr.shutdown()` + `device_discovery.shutdown()` + 停 uiautodev | > ⚠️ **阶段 A~F 在 import 期就会起线程/调度器**,只有 uiautodev 拉起与预连接在 `__main__` 分支。以 WSGI 方式 import 本模块会得到"半个启动"的进程——本地调试请直接 `python web_server.py`。 @@ -112,6 +113,7 @@ | 巡检手动线程 | `web/agent_api.py` | 手动触发巡检 | 按需 | | **通知 dispatcher** | `core/notifier.init_app` | 通知聚合 + 每 hook 限流 + 折叠摘要 | 常驻 1 个 | | **通知 sender ×3** | 同上 | 真实发 webhook(退避重试、环形记录) | 常驻 3 个 | +| **步骤明细写线程** | `core/step_log.init_app` | 批量落库 `task_step_log`(队列满丢弃并计数) | 常驻 1 个 | 两个 APScheduler 相互独立,时区均固定 `Asia/Shanghai`。 diff --git a/doc/DATA_MODEL.md b/doc/DATA_MODEL.md index e0dc455..5c4d486 100644 --- a/doc/DATA_MODEL.md +++ b/doc/DATA_MODEL.md @@ -27,7 +27,7 @@ 启动时校验「`.env` 声明」与「库名」「库中登记的 `app_meta.deployment_env`」三方一致,不符**拒绝启动**。 -**表清单(12 张,全部是 `core/models.py` 里的 ORM 模型)** +**表清单(14 张,全部是 `core/models.py` 里的 ORM 模型)** | # | 表 | 用途 | |---|---|------| @@ -44,6 +44,7 @@ | 11 | `experience_audit` | 经验巡检结论 | | 12 | `agent_action` | 动作库(命名动作) | | 13 | `device_install_log` | 设备端应用商店的下载/安装记录(设备上报,见 [DEVICE_AGENT.md](DEVICE_AGENT.md)) | +| 14 | `task_step_log` | 任务步骤明细(每次步骤执行一条,见 §2.8;**唯一有无界增长风险的表**,靠保留期清理) | > 2026-09-13 之前,`app_meta` 与 4 张 `agent_*` 表是各模块里的裸 `CREATE TABLE` > (不进模型层)。迁 MySQL 时那批 SQL 的 `AUTOINCREMENT`/`TEXT DEFAULT ''`/`TEXT PRIMARY KEY` @@ -136,6 +137,41 @@ > ⚠️ SQLAlchemy 模型的 `default=` 是 **Python 侧默认值**,SQLite 建表语句里没有 `DEFAULT` 子句;只有原生建表的表才有真正的 SQL DEFAULT。 +### 2.8 `task_step_log` — 任务步骤明细 + +每一次步骤执行一条(`tasks/generic/task.py:_exec_one` 里记录),是「日志 → 步骤明细」 +页的数据源。与 `logs/task.log` 的分工:那边是**排障原文**(什么都写、10MB 滚动), +这边是**结构化的一份**——设备/任务/步骤/结果/耗时都是列,能过滤、能统计、能导出 CSV。 + +| 列 | 类型 | 默认 | 说明 | +|----|------|------|------| +| `id` | Integer | — | 主键 | +| `run_id` | String(24) | `""` | 一次运行 = 设备 × 任务 × 第几次尝试;`TaskManager._run_with_retry` 每次尝试生成一个(12 位 hex),把这次尝试的所有步骤串起来 | +| `job_id` / `job_name` | String(32/120) | `""` | 任务快照(任务删了明细还在,名字仍可读) | +| `serial` / `device_name` | String(120/80) | `""` | 设备地址与**当时**的名字(快照,改名不影响历史) | +| `step_path` | String(32) | `""` | 嵌套位置,如 `2.1.3`;容器步骤(loop/group/if_el)会记自己那条,children 追加一级 | +| `step_label` / `step_type` | String(120/40) | `""` | 步骤标签与类型(`click_el`/`loop`/`wait`…) | +| `selector` | String(300) | `""` | 元素选择器(长选择器截断) | +| `result` | String(16) | `""` | `ok` / `miss`(handler 返回 False)/ `error`(抛异常)/ `unknown`(未知步骤类型)/ `skip`(概率未触发)/ `cap`(本次运行已达上限) | +| `detail` | String(500) | `""` | 异常消息、跳过原因等 | +| `duration_ms` | Integer | `0` | 本步耗时(慢步骤一眼可辨) | +| `created_at` | String(20) | `""` | 执行时刻(**保留期按它算**) | + +索引:`run_id`、`created_at`、`(serial, created_at)`、`(job_id, created_at)`。 + +**两条硬边界**(都在 `core/step_log.py`): + +| 常量 | 默认 | 作用 | +|---|---|---| +| `MAX_ROWS_PER_RUN` | 2000 | 单次运行最多记 2000 条,超出只补一条 `cap` 说明行。**没有它,`forever` 循环任务会瞬间写爆这张表** | +| `KEEP_DAYS` | 14 | 保留期:每天 04:13(+ 每次启动)清理更早的记录 | + +> ⚠️ **这张表是唯一有无界增长风险的表**,而它会**自动进整库备份**(§6 的派生规则), +> 所以 `KEEP_DAYS` 直接决定备份包体积。调大之前先想清楚导出的 zip 会有多大。 + +写入走 `core/step_log.py` 的**专用写线程 + 有界队列**(任务线程只 `put_nowait`, +微秒级;队列满丢弃并计数)——步骤执行是热路径,绝不能在任务线程里同步写库。 + --- ## 3. 非模型表 diff --git a/doc/DEPLOY.md b/doc/DEPLOY.md index 27baa35..9372b92 100644 --- a/doc/DEPLOY.md +++ b/doc/DEPLOY.md @@ -212,9 +212,14 @@ tail -20 logs/web.log # 无 ERROR/Traceback SUMMARY_TABLES = tuple(sorted(t.name for t in db.metadata.tables.values())) ``` -当前 12 张表:`app_meta` / `user` / `device_group` / `task_job` / `custom_action` / +当前 14 张表:`app_meta` / `user` / `device_group` / `task_job` / `custom_action` / `apk_file` / `device` / `pending_device` / `agent_conversation` / `agent_experience` / -`experience_audit` / `agent_action`。完整说明见 [DATA_MODEL.md](DATA_MODEL.md) §6。 +`experience_audit` / `agent_action` / `device_install_log` / `task_step_log`。 +完整说明见 [DATA_MODEL.md](DATA_MODEL.md) §6。 + +> ⚠️ `task_step_log`(任务步骤明细)是**唯一会持续增长**的表:它按 `KEEP_DAYS` +> (默认 14 天,见 `core/step_log.py`)自动清理,但备份包里会带上保留期内的全部行。 +> 设备多、任务密时导出 zip 会明显变大——需要更小的包就调小那个常量。 > **新增一张 ORM 表,就自动进了覆盖清单**,不可能再漏(2026-09-10 动作库 `agent_action` > 曾因手工维护漏登记,数据其实在快照里,只是清单没列 → 被误判为"没有备份")。 diff --git a/doc/DEVELOPMENT.md b/doc/DEVELOPMENT.md index c956f43..984f6a2 100644 --- a/doc/DEVELOPMENT.md +++ b/doc/DEVELOPMENT.md @@ -56,7 +56,7 @@ | 3 | **空闲设备扫描不主动 connect/disconnect** | 避免扰动共享连接 | 前台扫描对空闲设备直接返回"空闲";设备发现用 socket 探测 | | 4 | **adb key 保持历史 key 不变** | 设备信任该 key,换 key 全部 `unauthorized` | 部署沿用 `~/.android/adbkey` | | 5 | **生产(220)默认只读** | 生产事故成本高 | 任何写操作(pull/重启/改文件)都需负责人确认 | -| 6 | **新增持久化表必须登记备份覆盖清单** | 漏登记 = 等于没备份 | `core/system_backup.py` 的 `SUMMARY_TABLES` + `TABLE_LABELS`,详见 [DEPLOY.md](DEPLOY.md) §5.2 | +| 6 | **新增持久化表必须进备份覆盖清单** | 漏登记 = 等于没备份 | `SUMMARY_TABLES` 已由模型元数据自动派生(加了 ORM 表就进清单);**手工要做的只有补 `TABLE_LABELS` 中文标签**,详见 [DEPLOY.md](DEPLOY.md) §5.2 | | 7 | **功能/配置/接口改动必须同步文档** | 文档落后会误导开发与运维 | 见 §6;索引 [doc/README.md](README.md) | ### 其它开发约束 @@ -184,9 +184,10 @@ MCP_ALLOW_WRITE=1 MCP_PLATFORM_USER=admin MCP_PLATFORM_PASS=<密码> \ ### 5.2 新增数据库字段/表 -- 模型改 `core/models.py`;新表 `create_all()` 会建 +- 模型改 `core/models.py`;新表 `create_all()` 会建(**幂等**:已是模型就会自动建/补列) - **老库**要在 `SCHEMA_MIGRATIONS` 里加迁移(版本号递增 + SQL) -- **新增表**:登记进 `core/system_backup.py` 的 `SUMMARY_TABLES` + `TABLE_LABELS`(**红线**) +- **新增表**:`SUMMARY_TABLES` 由模型元数据**自动派生**(=自动进备份覆盖清单), + 仍需手工做的是补 `TABLE_LABELS` 的中文标签(**红线**) - 更新 [DATA_MODEL.md](DATA_MODEL.md) 与 [DEPLOY.md](DEPLOY.md) §5.2 ### 5.3 修改前端 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. 开发:给新功能加通知 diff --git a/doc/TASK_DEV.md b/doc/TASK_DEV.md index 106d8ea..14414f8 100644 --- a/doc/TASK_DEV.md +++ b/doc/TASK_DEV.md @@ -348,9 +348,9 @@ class MyTask(BaseTask): def get_action_class(cls, action_type): return None - def create_worker(self, serial, params): + def create_worker(self, serial, params, ctx=None): merged = {**DEFAULT_PARAMS, **(params or {})} - return MyWorker(serial, params=merged) + return MyWorker(serial, params=merged, ctx=ctx) ``` ```python @@ -369,8 +369,11 @@ from .myapp import task # ← 新增 ### 8.3 签名约定 -- `BaseWorker.__init__(self, serial, params=None, daemon=True)` -- `Task.create_worker(self, serial, params)` —— **不要带 `stf_client` / `stf` 形参**(STF 已摘除,历史签名已清理) +- `BaseWorker.__init__(self, serial, params=None, daemon=True, ctx=None)` +- `Task.create_worker(self, serial, params, ctx=None)` —— **不要带 `stf_client` / `stf` 形参**(STF 已摘除,历史签名已清理) +- `ctx` 是本次运行的上下文(`run_id` / `job_id` / `job_name` / `device_name`),由 + `TaskManager._run_with_retry` 传入;**只用于结构化记录**(如步骤明细 + `core/step_log.py`),不参与业务逻辑。不关心就原样透传给 `BaseWorker` 即可。 --- diff --git a/static/admin/admin.js b/static/admin/admin.js index 63a6fea..0af633e 100644 --- a/static/admin/admin.js +++ b/static/admin/admin.js @@ -91,22 +91,85 @@ async function deleteGroup(name){ } // ================== Tab 4: 日志 ================== +// 条件(文件/关键字/级别/时间)→ /api/logs(结构化行);下载走 /api/logs/download +function _logParams(){ + const v=id=>document.getElementById(id); + const p=new URLSearchParams(); + p.set('file', v('log-file').value||''); + p.set('lines', v('log-lines').value); + const q=v('log-q').value.trim(); if(q)p.set('q',q); + const lv=v('log-level').value; if(lv)p.set('level',lv); + const since=v('log-since').value; if(since)p.set('since',since); + const until=v('log-until').value; if(until)p.set('until',until); + return p; +} + +// 文件下拉由接口返回的清单渲染(写死的选项会漏掉新增模块,如 notify.log) +function _renderLogFiles(files,current){ + const sel=document.getElementById('log-file'); + const sig=(files||[]).map(f=>f.name).join(','); + if(sel.dataset.sig!==sig){ // 只在清单变化时重建,避免打断当前选择 + sel.dataset.sig=sig; + sel.innerHTML=(files||[]).map(f=> + '').join(''); + } + if(current)sel.value=current; +} + async function loadLogs(){ - const file=document.getElementById('log-file').value; - const lines=document.getElementById('log-lines').value; - const r=await apiGet('/api/logs?file='+encodeURIComponent(file)+'&lines='+lines); - if(!r||!r.ok)return; - const content=r.content||''; - // 简单着色 - const colored=esc(content) - .replace(/\[ERROR\]/g,'[ERROR]') - .replace(/\[WARNING\]/g,'[WARNING]') - .replace(/\[INFO\]/g,'[INFO]'); - document.getElementById('log-content').innerHTML=colored||'(空)'; - document.getElementById('log-hint').textContent=file+' · '+content.split('\n').length+' 行'; - // 自动滚动到底部 const el=document.getElementById('log-content'); - el.scrollTop=el.scrollHeight; + if(!el)return; + // 用户正往上翻看历史时不要被自动刷新拽回底部 + const atBottom=el.scrollHeight-el.scrollTop-el.clientHeight<40; + const r=await apiGet('/api/logs?'+_logParams().toString()); + if(!r)return; + _renderLogFiles(r.files,r.file); + const hint=document.getElementById('log-hint'); + if(!r.ok){ + el.innerHTML=''+esc(r.error||'读取失败')+''; + hint.textContent=''; + return; + } + const rows=r.rows||[], kw=document.getElementById('log-q').value.trim(); + el.innerHTML=rows.length + ? rows.map(row=>_logRowHtml(row,kw)).join('') + : '(无匹配行)'; + let text=r.file+' · 命中 '+r.matched+' 行' + +(rows.length['+esc(lvl)+'] [' + +esc(r.module)+'] '; + return '
'+head+_logHl(r.msg,kw)+'
'; +} + +// 关键字高亮:**先转义再匹配**——反过来的话消息里的 < > 会先被吃成实体 +function _logHl(text,kw){ + const safe=esc(text||''); + if(!kw)return safe; + const k=esc(kw); + if(!k)return safe; + const re=new RegExp(k.replace(/[.*+?^${}()|[\]\\]/g,'\\$&'),'gi'); + return safe.replace(re,m=>''+m+''); +} + +function downloadLogs(){ + window.location='/api/logs/download?'+_logParams().toString(); +} + +function resetLogFilters(){ + ['log-q','log-since','log-until'].forEach(id=>{document.getElementById(id).value='';}); + document.getElementById('log-level').value=''; + loadLogs(); } function toggleLogAuto(){ diff --git a/static/admin/base.js b/static/admin/base.js index 590d518..06379f8 100644 --- a/static/admin/base.js +++ b/static/admin/base.js @@ -106,7 +106,11 @@ function showTab(name){ if(name==='monitor'){loadMonitor();_monitorTimer=setInterval(loadMonitor,5000);} if(name==='tasks'){showSubTab('tasks',_activeSubs.tasks);loadTasks();loadCustomActions();} if(name==='tools'){showSubTab('tools',_activeSubs.tools);loadToolsDevices();loadAdbDevices();loadTailscaleDevices();loadApks();} - if(name==='logs'){loadLogs();if(document.getElementById('log-auto').checked)_logTimer=setInterval(loadLogs,3000);} + if(name==='logs'){ + showSubTab('logs', _activeSubs.logs||'files'); // 文件日志 / 步骤明细 + if((_activeSubs.logs||'files')==='files' + && document.getElementById('log-auto').checked)_logTimer=setInterval(loadLogs,3000); + } if(name==='agent'){ showSubTab('agent', _activeSubs.agent||'chat'); // 聊天 / AI 建任务 if(typeof initAgentChat==='function') initAgentChat(); @@ -117,7 +121,8 @@ function showTab(name){ // ================== 页内子分栏(任务/工具 通用) ================== // 每个带子分栏的 Tab 记住上次选中的子分栏,切走再切回来保持原位 -let _activeSubs = {tasks: 'plan', tools: 'clipboard', system: 'backup', agent: 'chat'}; +let _activeSubs = {tasks: 'plan', tools: 'clipboard', system: 'backup', agent: 'chat', + logs: 'files'}; let _discoveryTimer = null; // 设备自动发现 10s 轮询(仅 devpool 子分栏激活时) function showSubTab(tabId, name){ @@ -127,6 +132,8 @@ function showSubTab(tabId, name){ tab.querySelectorAll('.sub-panel').forEach(p=>p.classList.toggle('active', p.id===tabId+'-sub-'+name)); if(tabId==='agent' && name==='taskgen' && typeof initTaskGen==='function') initTaskGen(); if(tabId==='system' && name==='notify' && typeof loadNotifyPanel==='function') loadNotifyPanel(); + if(tabId==='logs' && name==='files' && typeof loadLogs==='function') loadLogs(); + if(tabId==='logs' && name==='steps' && typeof loadStepLogs==='function') loadStepLogs(true); if(name==='groups' && typeof loadGroups==='function') loadGroups(); if(name==='apks' && typeof loadAgentStore==='function') loadAgentStore(); if(name==='devpool' && typeof loadDevPool==='function'){ diff --git a/static/admin/steplog.js b/static/admin/steplog.js new file mode 100644 index 0000000..17aed00 --- /dev/null +++ b/static/admin/steplog.js @@ -0,0 +1,143 @@ +// ================== 「日志 → 步骤明细」子分栏 ================== +// 数据源:task_step_log 表(结构化),接口 /api/step_logs*(见 web/admin_api.py)。 +// 与「文件日志」的分工:那边是原始文本,这边能按设备/任务/时间/结果过滤、能导出。 + +let _slOffset = 0; // 分页偏移 +let _slTotal = 0; // 命中总数 +let _slRunId = ''; // "只看某次运行"时的 run_id(空 = 不限) +let _slInited = false; // 下拉选项只拉一次 + +function _slVal(id){const el=document.getElementById(id);return el?el.value:'';} + +function _slParams(){ + const p=new URLSearchParams(); + const put=(k,v)=>{if(v)p.set(k,v);}; + put('serial', _slVal('sl-device')); + put('job_id', _slVal('sl-job')); + put('result', _slVal('sl-result')); + put('q', _slVal('sl-q').trim()); + put('since', _slVal('sl-since')); + put('until', _slVal('sl-until')); + put('run_id', _slRunId); + p.set('limit', _slVal('sl-limit')||200); + p.set('offset', _slOffset); + return p; +} + +async function loadStepLogs(reset){ + const body=document.getElementById('sl-body'); + if(!body)return; + if(reset)_slOffset=0; + const p=_slParams(); + const r=await apiGet('/api/step_logs?'+p.toString()); + if(!r||!r.ok){body.innerHTML='读取失败';return;} + _slTotal=r.total||0; + const rows=r.rows||[]; + body.innerHTML = rows.length ? rows.map(_slRowHtml).join('') + : '(没有匹配的步骤记录——先跑一个任务,' + +'或在上面放宽筛选条件)'; + // 提示栏:分页位置 + 保留期 + 运行期间丢弃计数(排障用) + const from=_slTotal?(_slOffset+1):0, to=_slOffset+rows.length; + let hint=`命中 ${_slTotal} 条 · 显示 ${from}-${to}`; + if(_slRunId)hint=`只看运行 ${_slRunId} · `+hint; + const st=r.stats||{}; + if(st.dropped)hint+=` · 队列溢出丢弃 ${st.dropped} 条`; + if(st.failed)hint+=` · 落库失败 ${st.failed} 条`; + document.getElementById('sl-hint').textContent=hint; + document.getElementById('sl-page').textContent= + `第 ${Math.floor(_slOffset/(parseInt(_slVal('sl-limit'),10)||200))+1} 页`; + document.getElementById('sl-prev').disabled = _slOffset<=0; + document.getElementById('sl-next').disabled = _slOffset+rows.length>=_slTotal; + if(!_slInited){_slInited=true;loadStepFilters();} + if(reset)loadStepRuns(); +} + +// 结果 → 颜色。注意不能复用 .log-err/.log-warn/.log-info:那套的作用域是 +// .log-content(文件日志面板),表格里得用下面这套 res-* +const _SL_RESULT={error:'res-error',unknown:'res-error',miss:'res-miss', + cap:'res-miss',skip:'res-skip',ok:'res-ok'}; + +function _slRowHtml(r){ + const cls=_SL_RESULT[r.result]||''; + const dev=esc(r.device_name||r.serial||''); + const dur=r.duration_ms>=1000?(r.duration_ms/1000).toFixed(1)+'s':(r.duration_ms||0)+'ms'; + return '' + +''+esc(r.created_at||'')+'' + +''+dev+'' + +''+esc(r.job_name||'')+'' + +''+esc(r.step_path||'')+'' + +''+esc(r.step_label||'')+'' + +''+esc(r.step_type||'')+'' + +''+esc(r.result||'')+'' + +''+dur+'' + +''+esc(_slTrim(r.detail,120))+'' + +''; +} + +function _slTrim(s,n){s=s||'';return s.length>n?s.slice(0,n)+'…':s;} + +function stepLogPage(delta){ + const size=parseInt(_slVal('sl-limit'),10)||200; + _slOffset=Math.max(0,_slOffset+delta*size); + loadStepLogs(false); +} + +function resetStepFilters(){ + ['sl-q','sl-since','sl-until'].forEach(id=>{const el=document.getElementById(id);if(el)el.value='';}); + ['sl-device','sl-job','sl-result'].forEach(id=>{const el=document.getElementById(id);if(el)el.value='';}); + _slRunId=''; + loadStepLogs(true); +} + +function downloadStepLogs(){ + const p=_slParams(); + p.delete('limit');p.delete('offset'); + window.location='/api/step_logs/download?'+p.toString(); +} + +// ---------- 筛选下拉(只列"明细里真的出现过"的设备/任务) ---------- +async function loadStepFilters(){ + const r=await apiGet('/api/step_logs/filters'); + if(!r||!r.ok)return; + const dev=document.getElementById('sl-device'); + const job=document.getElementById('sl-job'); + const keep=(sel,val)=>{const cur=sel.value;sel.innerHTML=val;sel.value=cur;}; + keep(dev,''+(r.devices||[]).map(d=> + '').join('')); + keep(job,''+(r.jobs||[]).map(j=> + '').join('')); + document.getElementById('sl-note').textContent= + `明细保留 ${r.keep_days} 天(超期自动清理);单次运行最多记录 ` + +`${r.max_rows_per_run} 条(无限循环任务不会写爆这张表)。`; +} + +// ---------- 最近运行概览(点一行 = 只看那次运行) ---------- +async function loadStepRuns(){ + const p=new URLSearchParams(); + if(_slVal('sl-device'))p.set('serial',_slVal('sl-device')); + if(_slVal('sl-job'))p.set('job_id',_slVal('sl-job')); + if(_slVal('sl-since'))p.set('since',_slVal('sl-since')); + if(_slVal('sl-until'))p.set('until',_slVal('sl-until')); + p.set('limit','20'); + const r=await apiGet('/api/step_logs/runs?'+p.toString()); + const body=document.getElementById('sl-runs-body'); + if(!r||!r.ok){body.innerHTML='';return;} + const runs=r.runs||[]; + document.getElementById('sl-runs-summary').textContent= + `最近运行概览(${runs.length} 次,点「查看」只看那次)`; + body.innerHTML = runs.length ? runs.map(x=>'' + +''+esc(x.started_at||'')+'' + +''+esc(x.device_name||x.serial||'')+'' + +''+esc(x.job_name||'')+'' + +''+x.steps+'' + +''+(x.failures||0)+'' + +'' + +'').join('') + : '(暂无运行记录)'; +} + +function filterByRun(runId){ + _slRunId=runId||''; + document.getElementById('sl-runs-box').open=false; + loadStepLogs(true); +} diff --git a/tasks/base.py b/tasks/base.py index 1f46061..56a432b 100644 --- a/tasks/base.py +++ b/tasks/base.py @@ -49,8 +49,8 @@ from .actions import get_action_class return get_action_class(action_type) - def create_worker(self, serial, params): - return MyWorker(serial, params=params) + def create_worker(self, serial, params, ctx=None): + return MyWorker(serial, params=params, ctx=ctx) 5. 在 my_task/__init__.py 加:from . import task (触发注册) 6. 在 tasks/__init__.py 加:from .my_task import task (触发注册) @@ -85,8 +85,12 @@ class BaseTask: description = "" default_params = {} - def create_worker(self, serial, params): - """返回一个 threading.Thread(已启动或待启动),执行实际任务。""" + def create_worker(self, serial, params, ctx=None): + """返回一个 threading.Thread(已启动或待启动),执行实际任务。 + + ctx:本次运行的上下文(`run_id`/`job_id`/`job_name`/`device_name`), + 由 TaskManager 传入,供步骤明细等结构化记录标注来源;可为 None。 + """ raise NotImplementedError @classmethod diff --git a/tasks/generic/task.py b/tasks/generic/task.py index dfefb72..b9bbc7d 100644 --- a/tasks/generic/task.py +++ b/tasks/generic/task.py @@ -37,7 +37,7 @@ from tasks.base import BaseTask, register_task from core.device_worker import BaseWorker, _update_status from core.u2_helper import ensure_app_running, wait_for_app_home, random_sleep from core.logger import get_logger -from core import notifier +from core import notifier, step_log _log = get_logger("task.generic") @@ -124,13 +124,18 @@ class GenericStepsWorker(BaseWorker): 只实现 run_task,按步骤类型分发到 _exec_ 方法。 """ - def __init__(self, serial, params=None): - super().__init__(serial, params) + def __init__(self, serial, params=None, ctx=None): + super().__init__(serial, params, ctx=ctx) p = {**DEFAULT_PARAMS, **(self.params or {})} self.max_duration = int(p.get("max_duration", 0)) self.steps = p.get("steps", []) self._action_counts = {} # 步骤执行计数 {step_label: count} self._miss_counts = {} # 选择器健康:selector -> 连续未命中次数 + # 步骤路径:容器步骤(loop/group/if_el)进出时压栈,得到 "2.1.3" 这种位置, + # 写进「步骤明细」用来还原嵌套结构。worker 每台设备一个线程,实例级状态安全。 + self._path = [] + self._step_rows = 0 # 本次运行已记录的明细条数(封顶见 _record_step) + self._step_capped = False # 某 click 选择器连续未找到元素的次数达到该值,判定可能失效(App 改版) _MAX_CONSECUTIVE_MISS = 10 @@ -177,10 +182,16 @@ class GenericStepsWorker(BaseWorker): if depth > 5: _log.warning(f"[{self.serial}] 步骤嵌套深度超限(>5),跳过") return - for step in steps: + for i, step in enumerate(steps): if self.stopped() or self.is_time_up(): return - self._exec_one(d, step, depth) + # 压入本步在兄弟里的序号(1 起):容器步骤执行 children 时会在其后 + # 继续追加,得到 "2.1" 这样的嵌套路径 + self._path.append(str(i + 1)) + try: + self._exec_one(d, step, depth) + finally: + self._path.pop() def _exec_one(self, d, step, depth=0): """执行单个步骤。""" @@ -197,19 +208,29 @@ class GenericStepsWorker(BaseWorker): prob = float(params.get("probability", 100)) if prob < 100 and random.random() * 100 > prob: _log.info(f"[{self.serial}] 步骤 '{label}' 概率 {prob}% 未触发,跳过") + self._record_step(step, "skip", f"概率 {prob}% 未触发", 0) return self.set_action(f"执行: {label}") _log.info(f"[{self.serial}] 步骤: {label}({stype}) params={params}") ret = None + # 结果归类:异常 > 未命中(handler 返回 False)> 正常 + # (True/None 都算正常——None 是"没有判断语义"的步骤,如点击坐标) + result, detail = "ok", "" + t0 = time.time() handler = getattr(self, f"_exec_{stype}", None) if handler: try: ret = handler(d, params, depth) + if ret is False: + result, detail = "miss", "未命中/超时" except Exception as e: _log.error(f"[{self.serial}] 步骤 {label}({stype}) 异常: {e}") + result, detail = "error", str(e) else: _log.warning(f"[{self.serial}] 未知步骤类型: {stype}") + result, detail = "unknown", f"未知步骤类型: {stype}" + self._record_step(step, result, detail, int((time.time() - t0) * 1000)) # 容器步骤(loop/group/if_el)不计入操作计数,避免监控显示噪音 if stype not in ("loop", "group", "if_el"): self._action_counts[label] = self._action_counts.get(label, 0) + 1 @@ -220,6 +241,44 @@ class GenericStepsWorker(BaseWorker): action_counts=dict(self._action_counts), elapsed=elapsed) return ret # 命中结果(True/False/None),"测试此步骤"用 + def _record_step(self, step, result, detail, duration_ms): + """把这一步写进「步骤明细」(异步落库,不阻塞任务线程)。 + + 单次运行条数封顶:`forever` 循环任务会执行成千上万步,不封顶这张表 + 会被瞬间写爆。到顶后只补一条 cap 说明行,此后本次运行不再记录。 + """ + try: + # 只记"任务运行"里的步骤:编辑器里的「测试此步骤」也会走 _exec_one, + # 但没有 run_id(不归任何一次运行),记进来只是噪音 + if not self.ctx.get("run_id") or self._step_capped: + return + params = step.get("params") or {} + label = step.get("label") or step.get("type", "") + if self._step_rows >= step_log.MAX_ROWS_PER_RUN: + self._step_capped = True + step_log.record( + run_id=self.ctx.get("run_id", ""), serial=self.serial, + device_name=self.ctx.get("device_name", ""), + job_id=self.ctx.get("job_id", ""), + job_name=self.ctx.get("job_name", ""), + step_path=".".join(self._path), step_label=label, + step_type=step.get("type", ""), result="cap", + detail=f"本次运行明细已达上限 {step_log.MAX_ROWS_PER_RUN} 条," + f"后续步骤不再记录(保留期与上限见 core/step_log.py)") + return + self._step_rows += 1 + step_log.record( + run_id=self.ctx.get("run_id", ""), serial=self.serial, + device_name=self.ctx.get("device_name", ""), + job_id=self.ctx.get("job_id", ""), + job_name=self.ctx.get("job_name", ""), + step_path=".".join(self._path), step_label=label, + step_type=step.get("type", ""), + selector=str(params.get("selector_value", "") or ""), + result=result, detail=detail, duration_ms=duration_ms) + except Exception: + pass # 记录失败绝不能影响任务执行 + # ================== 步骤执行器 ================== def _exec_screen_on(self, d, params, depth=0): """亮屏:息屏时唤醒并滑动解锁。任务执行前用(息屏时 u2 无法操作)。""" @@ -671,9 +730,9 @@ class GenericStepsTask(BaseTask): def get_action_class(cls, action_type): return None - def create_worker(self, serial, params): + def create_worker(self, serial, params, ctx=None): merged = {**DEFAULT_PARAMS, **(params or {})} - return GenericStepsWorker(serial, params=merged) + return GenericStepsWorker(serial, params=merged, ctx=ctx) def test_step(serial, step): diff --git a/templates/admin/monitor.html b/templates/admin/monitor.html index 6cfb923..29f10f9 100644 --- a/templates/admin/monitor.html +++ b/templates/admin/monitor.html @@ -112,7 +112,19 @@ body{background:var(--bg);font-family:var(--body);color:var(--text);font-size:14 .log-content .log-err{color:#f87171} .log-content .log-warn{color:#fbbf24} .log-content .log-info{color:#67e8f9} -.log-toolbar{display:flex;gap:10px;margin-bottom:12px;align-items:center} +.log-row{white-space:pre-wrap} +.log-cont{color:#8a94a6;padding-left:2.2em} /* 续行(traceback 缩进行):淡色缩进 */ +.log-hit{background:#fbbf2433} /* 关键字命中高亮 */ + +/* 「步骤明细」表格的结果列配色(.log-* 那套作用域在 .log-content 里,表格用这套) */ +.res-error{color:#dc2626;font-weight:600} /* 异常 / 未知步骤类型 */ +.res-miss{color:#d97706} /* 未命中 / 超时 / 已达上限 */ +.res-skip{color:var(--text-light)} /* 概率未触发 */ +.res-ok{color:#059669} +/* 工具栏:控件较多,允许换行;控件本身不参与收缩(否则会被挤成竖排文字) */ +.log-toolbar{display:flex;flex-wrap:wrap;gap:8px 10px;margin-bottom:12px;align-items:center} +.log-toolbar>*{flex:0 0 auto} +.log-toolbar .hint{flex:1 1 200px;min-width:160px} /* 区块标题 */ .section-title{font-family:var(--display);font-size:15px;font-weight:700;color:var(--text);margin:0 0 12px;padding-bottom:8px;border-bottom:2px solid var(--primary);display:inline-block;letter-spacing:.2px} @@ -778,25 +790,103 @@ body{background:var(--bg);font-family:var(--body);color:var(--text);font-size:14
日志查看
-
查看系统各模块运行日志
+
原始文本日志 + 任务步骤明细(结构化,可按设备/任务/时间过滤)
+ +
+ + +
+ + +
- + + + + ~ + + +
+
+ + +
+
+ + + + + + ~ + + + + + + +
+
+ +
+ 最近运行概览 + + + + + + +
开始设备任务步骤数失败
+
+ + + + + + + + +
时间设备任务步骤标签类型结果耗时详情
+
+ + + +
+
@@ -1266,5 +1356,6 @@ body{background:var(--bg);font-family:var(--body);color:var(--text);font-size:14 + diff --git a/web/admin_api.py b/web/admin_api.py index be9151c..1f5dc8a 100644 --- a/web/admin_api.py +++ b/web/admin_api.py @@ -1,17 +1,19 @@ """管理域 API:用户管理/日志查看。""" -import os import time +from io import BytesIO -from flask import Blueprint, jsonify, request +from flask import Blueprint, jsonify, request, send_file from flask_login import current_user from flask import session -from core.logger import _LOG_DIR, _MODULE_FILES +from core import logger +from core.logger import get_logger from core.models import db, User from web import context from web.auth import (admin_required, perm_required, _validate_perms, PERM_LOGS) +_log = get_logger("web") bp = Blueprint("admin", __name__) @bp.route("/api/users") @@ -79,22 +81,173 @@ def api_users_delete(uid): # ================== API:日志查看 ================== +# +# 过滤/白名单都在 core.logger 里做(读取侧与写入侧共用同一份文件清单, +# 见 core/logger.py「日志读取」一节);这里只负责取参数 + 组装响应。 + + +def _log_filters(): + """从 query string 取过滤条件(/api/logs 与 /api/logs/download 共用)。""" + return dict(keyword=request.args.get("q", ""), + min_level=request.args.get("level", ""), + since=request.args.get("since", ""), + until=request.args.get("until", "")) + @bp.route("/api/logs") @perm_required(PERM_LOGS) def api_logs(): - files = list(_MODULE_FILES.values()) - current = request.args.get("file", "core.log") - lines = int(request.args.get("lines", 300)) - content = "" - path = os.path.join(_LOG_DIR, current) - if os.path.exists(path): - try: - with open(path, encoding="utf-8") as f: - content = "".join(f.readlines()[-lines:]) - except Exception as e: - content = f"读取失败: {e}" - return jsonify({"ok": True, "content": content, "file": current, "files": files}) + """读日志:可按关键字 / 最低级别 / 时间范围过滤,返回结构化行。 + + 参数:file 文件名(白名单)、q 关键字、level 最低级别(ERROR 含以上…)、 + since/until 时间、lines 返回条数(取**最近** N 条命中)。 + """ + files = logger.list_log_files() + names = [f["name"] for f in files] + # 默认仍落在 core.log(接口按 mtime 倒序返回,names[0] 会随当前哪个模块在 + # 写而漂移,作为"打开日志页时看哪个"不稳定) + default = "core.log" if "core.log" in names else (names[0] if names else "") + current = request.args.get("file") or default + try: + lines = int(request.args.get("lines", 300)) + except (TypeError, ValueError): + lines = 300 + res = logger.query_log(current, limit=lines, **_log_filters()) + return jsonify({"ok": not res["error"], "error": res["error"], "file": current, + "files": files, "rows": res["rows"], "matched": res["matched"], + "scanned": res["scanned"], "truncated": res["truncated"]}) + + +@bp.route("/api/logs/download") +@perm_required(PERM_LOGS) +def api_logs_download(): + """下载日志文件(带同样的过滤条件;无条件时就是整个文件)。 + + 用 BytesIO 发送而不是 send_file(路径):Windows 上流式发送时文件句柄可能 + 到 close 仍未释放,而日志文件正被日志线程持续写入,按路径发容易踩锁。 + """ + name = request.args.get("file", "") + text, err = logger.read_log_text(name, **_log_filters()) + if err: + return jsonify({"ok": False, "error": err}), 404 + data = text.encode("utf-8") + fname = f"{name}_{time.strftime('%Y%m%d_%H%M%S')}.txt" + _log.info("下载日志: %s(%d 字节)by %s", name, len(data), + getattr(current_user, "username", "")) + return send_file(BytesIO(data), as_attachment=True, download_name=fname, + mimetype="text/plain; charset=utf-8") + + +# ================== API:任务步骤明细 ================== +# +# 数据源是 task_step_log 表(结构化),不是 logs/task.log 文本——所以能按设备/ +# 任务/时间/结果过滤。写入侧见 core/step_log.py(异步写线程 + 保留期清理)。 + +def _step_filters(): + """从 query string 取步骤明细的过滤条件(列表与下载共用)。""" + return dict(serial=request.args.get("serial", ""), + job_id=request.args.get("job_id", ""), + result=request.args.get("result", ""), + run_id=request.args.get("run_id", ""), + keyword=request.args.get("q", ""), + since=logger.norm_ts(request.args.get("since", "")), + until=logger.norm_ts(request.args.get("until", ""), end=True)) + + +def _step_to_dict(r): + return {"id": r.id, "run_id": r.run_id, "job_id": r.job_id, + "job_name": r.job_name, "serial": r.serial, + "device_name": r.device_name, "step_path": r.step_path, + "step_label": r.step_label, "step_type": r.step_type, + "selector": r.selector, "result": r.result, "detail": r.detail, + "duration_ms": r.duration_ms, "created_at": r.created_at} + + +@bp.route("/api/step_logs") +@perm_required(PERM_LOGS) +def api_step_logs(): + """任务步骤明细(分页,最近的在前)。""" + from core import step_log + try: + limit = int(request.args.get("limit", 200)) + offset = int(request.args.get("offset", 0)) + except (TypeError, ValueError): + limit, offset = 200, 0 + rows, total = step_log.query(limit=limit, offset=offset, **_step_filters()) + return jsonify({"ok": True, "rows": [_step_to_dict(r) for r in rows], + "total": total, "limit": limit, "offset": offset, + "stats": step_log.stats(), + "keep_days": step_log.KEEP_DAYS, + "max_rows_per_run": step_log.MAX_ROWS_PER_RUN}) + + +@bp.route("/api/step_logs/runs") +@perm_required(PERM_LOGS) +def api_step_logs_runs(): + """按 run_id 归组的一次运行概览(哪次运行、几步、失败几步)。""" + from core import step_log + f = _step_filters() + rows = step_log.runs(serial=f["serial"], job_id=f["job_id"], + since=f["since"], until=f["until"], + limit=request.args.get("limit", 100)) + return jsonify({"ok": True, "runs": [ + {"run_id": r.run_id, "job_id": r.job_id, "job_name": r.job_name, + "serial": r.serial, "device_name": r.device_name, + "started_at": r.started_at, "ended_at": r.ended_at, + "steps": int(r.steps or 0), "failures": int(r.failures or 0)} + for r in rows]}) + + +@bp.route("/api/step_logs/filters") +@perm_required(PERM_LOGS) +def api_step_logs_filters(): + """筛选下拉的可选项:**只在明细里出现过的**设备与任务(按数据自洽,不依赖 + 设备池/任务表,也不要求调用方另有 task/device 权限)。""" + from core import step_log + from core.models import TaskStepLog, db + devices = (db.session.query(TaskStepLog.serial, TaskStepLog.device_name) + .filter(TaskStepLog.serial != "").distinct().limit(500).all()) + jobs = (db.session.query(TaskStepLog.job_id, TaskStepLog.job_name) + .filter(TaskStepLog.job_id != "").distinct().limit(500).all()) + return jsonify({"ok": True, + "devices": sorted(({"serial": s, "device_name": n} + for s, n in devices), + key=lambda x: x["device_name"] or x["serial"]), + "jobs": sorted(({"job_id": i, "job_name": n} for i, n in jobs), + key=lambda x: x["job_name"] or x["job_id"]), + "results": ["ok", "miss", "error", "unknown", "skip", "cap"], + "keep_days": step_log.KEEP_DAYS, + "max_rows_per_run": step_log.MAX_ROWS_PER_RUN}) + + +@bp.route("/api/step_logs/download") +@perm_required(PERM_LOGS) +def api_step_logs_download(): + """导出步骤明细为 CSV(带同样的过滤条件)。 + + 加 UTF-8 BOM:不加的话 Excel 打开中文是乱码(这是给运维看的表, + 大概率会被 Excel 打开)。 + """ + import csv + from io import StringIO + from core import step_log + f = _step_filters() + # 导出上限:和列表页共用同一次查询,最多 2 万行(再多请缩小时间范围) + rows, total = step_log.query(limit=20000, offset=0, **f) + buf = StringIO() + w = csv.writer(buf) + w.writerow(["时间", "设备", "设备名", "任务", "运行ID", "步骤路径", "步骤", + "类型", "选择器", "结果", "耗时(ms)", "详情"]) + for r in reversed(rows): # 导出的时间序与页面相反:文件里按正序更好读 + w.writerow([r.created_at, r.serial, r.device_name, r.job_name, + r.run_id, r.step_path, r.step_label, r.step_type, + r.selector, r.result, r.duration_ms, r.detail]) + data = buf.getvalue().encode("utf-8-sig") # utf-8-sig = UTF-8 带 BOM(Excel 中文不乱码) + fname = f"step_log_{time.strftime('%Y%m%d_%H%M%S')}.csv" + _log.info("导出步骤明细: %d 行(命中 %d)by %s", len(rows), total, + getattr(current_user, "username", "")) + return send_file(BytesIO(data), as_attachment=True, download_name=fname, + mimetype="text/csv; charset=utf-8") # ================== API:运行控制 ================== diff --git a/web_server.py b/web_server.py index 31be35c..6bd926c 100644 --- a/web_server.py +++ b/web_server.py @@ -73,6 +73,12 @@ try: _notifier.init_app(app) except Exception as _e: # 通知出问题不影响启动 _log.warning(f"通知模块初始化失败(不影响启动): {_e}") +# 任务步骤明细:起一个专用写线程,任务线程只入队(见 core/step_log.py) +try: + from core import step_log as _step_log + _step_log.init_app(app) +except Exception as _e: + _log.warning(f"步骤明细模块初始化失败(不影响启动): {_e}") # 消费「待生效的备份恢复」——必须在 init_db 之后(表已建好,且现在是事务替换而非换文件) _consume_pending_restore() # 库环境标签校验 + 启动横幅:每次启动都明确写出「现在连的是哪个库」, @@ -104,6 +110,24 @@ from web.auth import _csrf_protect, _csrf_token context.init(mgr, apk_mgr, device_pool) register_blueprints(app) +def _purge_step_log(): + """清理超期的任务步骤明细(启动时 + 每日 04:13)。 + + app context 由本函数自建:APScheduler 的作业跑在调度线程里,没有请求上下文。 + """ + try: + from core.step_log import purge_old + with app.app_context(): + purge_old() + except Exception as e: + _log.warning(f"步骤明细清理失败(不影响主服务): {e}") + + +def _purge_step_log_async(): + """启动时后台清理一次(不阻塞启动;服务若频繁重启,定时任务可能一直轮不到)。""" + threading.Thread(target=_purge_step_log, daemon=True).start() + + # 经验库每日 AI 巡检(凌晨 03:47):评审标记疑似问题经验,删除只走人工确认。 # 巡检线程在 run_experience_audit 内自建 app context,不依赖这里。 try: @@ -112,6 +136,9 @@ try: from web.agent_api import run_experience_audit _audit_sched = BackgroundScheduler(timezone="Asia/Shanghai") _audit_sched.add_job(run_experience_audit, CronTrigger(hour=3, minute=47)) + # 步骤明细清理:超过保留期(core/step_log.KEEP_DAYS)的记录每天删一次。 + # 不清理的话这张表是唯一无界增长的表,会一路顶大备份包。 + _audit_sched.add_job(_purge_step_log, CronTrigger(hour=4, minute=13)) _audit_sched.start() _log.info("经验库每日巡检已注册(03:47 Asia/Shanghai)") except Exception as e: @@ -294,6 +321,8 @@ if __name__ == "__main__": _ensure_uiauto_running() # 预连接设备池:adb server 重启后设备全掉线,后台并发重连加速恢复 _preconnect_pool_devices() + # 步骤明细清理:服务频繁重启时定时任务可能一直轮不到,启动也清一次 + _purge_step_log_async() _log.info(f"管理后台: http://localhost:{WEB_PORT}/ (admin/admin123)") # 启动通知:放在 __main__ 里是刻意的——以 WSGI 方式 import 本模块只有"半个启动" # (阶段 A~F 在 import 期就跑完了),不该发"服务已启动"。 @@ -316,6 +345,11 @@ if __name__ == "__main__": _notifier.shutdown(timeout=3.0) except Exception: pass + try: + from core import step_log as _step_log + _step_log.shutdown(timeout=3.0) # 排空队列后再退 + except Exception: + pass mgr.shutdown() device_discovery.shutdown() _stop_uiauto()