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。
This commit is contained in:
2026-09-16 08:58:43 +08:00
parent 2c32c397c8
commit 621d48c800
20 changed files with 838 additions and 38 deletions
+247
View File
@@ -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