Files
butubb 621d48c800 feat(日志): 任务步骤明细表 + 「日志 → 步骤明细」面板(按设备/任务/时间过滤、按运行归组、导出 CSV)
文本日志只能 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。
2026-09-16 08:58:43 +08:00

278 lines
11 KiB
Python
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""统一日志系统:按模块分文件,带时间戳和级别,自动滚动。
用法:
from core.logger import get_logger
log = get_logger("task") # 写 logs/task.log + 控制台
log.info("开始任务")
log.error("失败", exc_info=True)
模块划分:
core — 核心程序(STF/adb/worker/task_manager)
task — 任务执行(worker 业务逻辑)
web — web_server 请求/管理
action — 操作执行(点赞/评论等)
所有日志同时输出到控制台和对应文件,logs/ 目录自动创建。
文件按 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")
os.makedirs(_LOG_DIR, exist_ok=True)
# 已创建的 logger,避免重复添加 handler
_LOGGERS = {}
# 模块 → 文件名映射
_MODULE_FILES = {
"core": "core.log",
"task": "task.log",
"web": "web.log",
"action": "action.log",
# 通知单独一个文件:排障时"哪条通知发了/失败了/为什么"要能一眼翻到
"notify": "notify.log",
}
_FORMAT = "%(asctime)s [%(levelname)s] [%(name)s] %(message)s"
_DATE_FMT = "%Y-%m-%d %H:%M:%S"
def get_logger(name="core"):
"""获取指定模块的 logger。
name: 模块名(core/task/web/action),决定写哪个文件。
也可传子模块名如 "task.generic",会归到 task.log。
返回配置好的 logging.Logger。
"""
if name in _LOGGERS:
return _LOGGERS[name]
# 归类到对应文件:取顶层模块名
top = name.split(".")[0]
filename = _MODULE_FILES.get(top, "core.log")
log = logging.getLogger(name)
log.setLevel(logging.DEBUG)
# 避免向 root logger 传播导致重复输出
log.propagate = False
fmt = logging.Formatter(_FORMAT, _DATE_FMT)
# 文件 handler:10MB 滚动,保留 5 个
file_path = os.path.join(_LOG_DIR, filename)
fh = RotatingFileHandler(file_path, maxBytes=10 * 1024 * 1024,
backupCount=5, encoding="utf-8")
fh.setLevel(logging.DEBUG)
fh.setFormatter(fmt)
# 控制台 handler
ch = logging.StreamHandler()
ch.setLevel(logging.INFO)
ch.setFormatter(fmt)
log.addHandler(fh)
log.addHandler(ch)
_LOGGERS[name] = log
return log
# 兼容旧代码的 print 风格:紧急排查时可临时用
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}"