From 621d48c8006bf2021188b19741c8fd85eeb7a2fa Mon Sep 17 00:00:00 2001 From: butubb <1422726308@qq.com> Date: Wed, 16 Sep 2026 08:58:43 +0800 Subject: [PATCH] =?UTF-8?q?feat(=E6=97=A5=E5=BF=97):=20=E4=BB=BB=E5=8A=A1?= =?UTF-8?q?=E6=AD=A5=E9=AA=A4=E6=98=8E=E7=BB=86=E8=A1=A8=20+=20=E3=80=8C?= =?UTF-8?q?=E6=97=A5=E5=BF=97=20=E2=86=92=20=E6=AD=A5=E9=AA=A4=E6=98=8E?= =?UTF-8?q?=E7=BB=86=E3=80=8D=E9=9D=A2=E6=9D=BF=EF=BC=88=E6=8C=89=E8=AE=BE?= =?UTF-8?q?=E5=A4=87/=E4=BB=BB=E5=8A=A1/=E6=97=B6=E9=97=B4=E8=BF=87?= =?UTF-8?q?=E6=BB=A4=E3=80=81=E6=8C=89=E8=BF=90=E8=A1=8C=E5=BD=92=E7=BB=84?= =?UTF-8?q?=E3=80=81=E5=AF=BC=E5=87=BA=20CSV=EF=BC=89?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 文本日志只能 grep,"这台设备这次运行为什么失败"翻起来很费劲。新增一张 **结构化**的步骤明细表,把"哪一步、什么类型、哪个选择器、结果、耗时"落库。 - core/models.py:新增 task_step_log(run_id/job/设备/step_path/selector/ result/detail/duration_ms),索引 run_id、created_at、(serial,created_at)、 (job_id,created_at); - core/step_log.py(新):**专用写线程 + 有界队列**批量落库——步骤执行是热路径, 任务线程只 put_nowait(实测 9000 次入队 31ms),队列满丢弃并计数,绝不阻塞; 另有保留期清理(默认 14 天,每日 04:13 + 每次启动); - tasks/generic/task.py:_exec_one 记一条(异常=error / handler 返回 False=miss / 未知类型=unknown / 概率未触发=skip);_exec_steps 维护路径栈得到 "2.1.3" 这样的嵌套位置;**单次运行封顶 2000 条**——forever 循环任务否则会写爆表; - 运行上下文 ctx(run_id/job_id/job_name/device_name)由 TaskManager 生成, 经 create_worker(serial, params, ctx=None) 传入 worker(扩展点向后兼容); - 接口:/api/step_logs(过滤+分页)、/runs(按运行归组)、/filters(下拉选项)、 /download(CSV,带 BOM); - 前端:日志页拆成「文件日志 / 步骤明细」子分栏 + static/admin/steplog.js。 红线:新表自动进备份覆盖清单(SUMMARY_TABLES 由元数据派生),已补 TABLE_LABELS 中文标签,导出实测 14 张表、coverage_missing 为空。 文档:DATA_MODEL §2.8/§1、API §10.2、ARCHITECTURE §1.1/§2.2/§3.1、 DEVELOPMENT §5.2 与红线表、DEPLOY §5.2、TASK_DEV §8.3、README。 --- README.md | 11 +- core/device_worker.py | 5 +- core/logger.py | 7 +- core/models.py | 34 +++++ core/step_log.py | 247 +++++++++++++++++++++++++++++++++++ core/system_backup.py | 1 + core/task_manager.py | 7 +- doc/API.md | 26 ++++ doc/ARCHITECTURE.md | 8 +- doc/DATA_MODEL.md | 38 +++++- doc/DEPLOY.md | 9 +- doc/DEVELOPMENT.md | 7 +- doc/TASK_DEV.md | 11 +- static/admin/base.js | 11 +- static/admin/steplog.js | 143 ++++++++++++++++++++ tasks/base.py | 12 +- tasks/generic/task.py | 73 ++++++++++- templates/admin/monitor.html | 75 ++++++++++- web/admin_api.py | 117 ++++++++++++++++- web_server.py | 34 +++++ 20 files changed, 838 insertions(+), 38 deletions(-) create mode 100644 core/step_log.py create mode 100644 static/admin/steplog.js diff --git a/README.md b/README.md index 73ec219..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)+ 目标设备选择;右上角:经验库 / 动作库 / 模型配置 | 仅管理员 | @@ -362,6 +363,12 @@ log = get_logger("task.generic") # → logs/task.log (无筛选就是整个文件)。读取逻辑在 `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 b5e926f..a40981e 100644 --- a/core/logger.py +++ b/core/logger.py @@ -18,7 +18,6 @@ import os import re import logging -import datetime from collections import deque from logging.handlers import RotatingFileHandler @@ -154,7 +153,7 @@ def list_log_files(): return out -def _norm_ts(value, end=False): +def norm_ts(value, end=False): """把前端的时间输入归一成 'YYYY-MM-DD HH:MM:SS'(可比字符串)。 接受 datetime-local 的 'YYYY-MM-DDTHH:MM'、'YYYY-MM-DD HH:MM' 和纯日期; @@ -226,7 +225,7 @@ def query_log(name, keyword="", min_level="", since="", until="", 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) + since, until = norm_ts(since), norm_ts(until, end=True) try: keep = max(1, min(int(limit or 300), 5000)) except (TypeError, ValueError): @@ -255,7 +254,7 @@ def read_log_text(name, keyword="", min_level="", since="", until="", """ if name not in {f["name"] for f in list_log_files()}: return None, "日志文件不存在或不可读取" - since, until = _norm_ts(since), _norm_ts(until, end=True) + 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: 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/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..20f9af2 100644 --- a/core/task_manager.py +++ b/core/task_manager.py @@ -825,8 +825,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}次)") diff --git a/doc/API.md b/doc/API.md index f1ca9d0..c210210 100644 --- a/doc/API.md +++ b/doc/API.md @@ -141,6 +141,10 @@ | DELETE | `/api/users/` | Admin | 删除用户 | | 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`) @@ -648,6 +652,28 @@ - `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) 的事件。 + --- ## 11. 运维工具(adb / 剪贴板 / 应用版本 / Tailscale) 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/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/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 938098e..29f10f9 100644 --- a/templates/admin/monitor.html +++ b/templates/admin/monitor.html @@ -115,6 +115,12 @@ body{background:var(--bg);font-family:var(--body);color:var(--text);font-size:14 .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} @@ -784,7 +790,15 @@ body{background:var(--bg);font-family:var(--body);color:var(--text);font-size:14
日志查看
-
查看系统各模块运行日志
+
原始文本日志 + 任务步骤明细(结构化,可按设备/任务/时间过滤)
+ +
+ + +
+ + +
@@ -815,6 +829,64 @@ body{background:var(--bg);font-family:var(--body);color:var(--text);font-size:14
+
+ + +
+
+ + + + + + ~ + + + + + + +
+
+ +
+ 最近运行概览 + + + + + + +
开始设备任务步骤数失败
+
+ + + + + + + + +
时间设备任务步骤标签类型结果耗时详情
+
+ + + +
+
@@ -1284,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 4d55d43..1f5dc8a 100644 --- a/web/admin_api.py +++ b/web/admin_api.py @@ -1,7 +1,8 @@ """管理域 API:用户管理/日志查看。""" 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 @@ -125,8 +126,6 @@ def api_logs_download(): 用 BytesIO 发送而不是 send_file(路径):Windows 上流式发送时文件句柄可能 到 close 仍未释放,而日志文件正被日志线程持续写入,按路径发容易踩锁。 """ - from io import BytesIO - from flask import send_file name = request.args.get("file", "") text, err = logger.read_log_text(name, **_log_filters()) if err: @@ -139,5 +138,117 @@ def api_logs_download(): 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()