文本日志只能 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。
278 lines
11 KiB
Python
278 lines
11 KiB
Python
"""统一日志系统:按模块分文件,带时间戳和级别,自动滚动。
|
||
|
||
用法:
|
||
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}"
|