"""线程安全的系统诊断日志缓冲区。 任何上下文(异步协程、同步函数、WebSocket 回调所在的执行器线程)都可以安全调用 ``record(...)`` 写入一条诊断日志,用于排查私信收发失败的原因。 设计要点: - 内存环形缓冲区是“实时”数据源(API 直接读取,零延迟),任何线程都能写入; - ``_pending`` 队列由 main.py 的后台任务定期落库到 ``system_logs`` 表做持久化; - 启动时可用 ``seed(...)`` 把历史记录读回缓冲区。 """ import logging import threading from collections import deque from datetime import datetime from typing import Optional logger = logging.getLogger("douyin_im.system") _VALID_LEVELS = ("info", "success", "warning", "error") _lock = threading.Lock() _buffer: deque = deque(maxlen=3000) _pending: list = [] _seq = 0 def record( event: str, detail: str = "", level: str = "info", category: str = "system", account_id: Optional[int] = None, ) -> dict: """写入一条诊断日志,返回该条记录。""" global _seq if level not in _VALID_LEVELS: level = "info" with _lock: _seq += 1 entry = { "id": _seq, "account_id": account_id, "level": level, "category": category, "event": str(event or ""), "detail": str(detail or ""), "created_at": datetime.utcnow().isoformat(), } _buffer.appendleft(entry) _pending.append(entry) msg = f"[{category}] {event}" + (f" | {detail}" if detail else "") if level == "error": logger.error(msg) elif level == "warning": logger.warning(msg) else: logger.info(msg) return entry def get_logs( account_id: Optional[int] = None, level: Optional[str] = None, category: Optional[str] = None, limit: int = 200, ) -> list: """按条件读取最近的诊断日志(最新优先)。""" with _lock: items = list(_buffer) out = [] for e in items: if account_id is not None and e["account_id"] != account_id: continue if level and e["level"] != level: continue if category and e["category"] != category: continue out.append(e) if len(out) >= limit: break return out def drain_pending() -> list: """取出尚未落库的记录(供后台 flush 任务持久化)。""" with _lock: items = _pending[:] _pending.clear() return items def seed(entries: list) -> None: """启动时把持久化的历史记录读回缓冲区(不会重复落库)。""" global _seq with _lock: for e in sorted(entries, key=lambda x: x.get("id") or 0): _buffer.appendleft(e) if (e.get("id") or 0) > _seq: _seq = e["id"] def clear() -> None: with _lock: _buffer.clear() _pending.clear()