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

248 lines
8.9 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.
"""任务步骤明细:异步落库(专用写线程)+ 保留期清理。
为什么不让任务线程直接写库:步骤执行是**热路径**(一次运行可能上万步),每步
一次 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