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