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