""" 回复队列执行记录 ================ 待回复队列本身只存"现在还欠谁一条回复",不留痕迹:任务一完成就被删掉,出了 问题事后完全看不出它经历过什么。排查"为什么这个客户没收到回复"、"为什么同一 句话回了三遍"时,缺的就是这条时间线。 这里按时间顺序记下每个会话在队列里的每一步:入队、开始处理、生成回复、已发送、 完成、跳过、失败。机器人线程写、界面线程读,所以写入走临时文件加原子替换, 界面永远不会读到写了一半的 JSON。 """ import json import os import tempfile import threading import time from runtime_paths import application_data_dir # 队列事件 ENQUEUED = "入队" OPENED = "开始处理" GENERATED = "生成回复" SENT = "已发送" DONE = "完成" SKIPPED = "跳过" FAILED = "失败" class QueueLog: """回复队列的执行流水。""" # 只留最近这么多条。这是给人看的排查线索,不是审计账本; # 无上限的话跑上几天,界面每次刷新都要读一个几十兆的文件 MAX_EVENTS = 500 def __init__(self, path: str = ""): self.path = str(path or os.path.join(application_data_dir(), "queue_events.json")) self._lock = threading.Lock() # ── 写 ─────────────────────────────────────────────────────────────────── def append( self, session_id: str, name: str, event: str, detail: str = "", ) -> bool: """记一步。写失败绝不能影响回复,所以吞掉异常只返回 False。""" entry = { "ts": time.time(), "session_id": str(session_id or "")[:32], "name": str(name or ""), "event": str(event or ""), "detail": str(detail or "")[:200], } with self._lock: try: events = self._read_unlocked() events.append(entry) if len(events) > self.MAX_EVENTS: events = events[-self.MAX_EVENTS:] self._write_unlocked(events) return True except Exception: return False def clear(self) -> bool: with self._lock: try: self._write_unlocked([]) return True except Exception: return False # ── 读 ─────────────────────────────────────────────────────────────────── def recent(self, limit: int = 200) -> list[dict]: """最近的事件,最新的排在前面。""" with self._lock: try: events = self._read_unlocked() except Exception: return [] count = max(0, int(limit)) return list(reversed(events[-count:])) if count else [] # ── 落盘 ───────────────────────────────────────────────────────────────── def _read_unlocked(self) -> list[dict]: if not os.path.exists(self.path): return [] try: with open(self.path, encoding="utf-8") as handle: data = json.load(handle) except (OSError, ValueError): # 文件损坏时从头开始,总好过让界面和机器人一起崩 return [] if isinstance(data, dict): data = data.get("events", []) return [item for item in data if isinstance(item, dict)] def _write_unlocked(self, events: list) -> None: directory = os.path.dirname(os.path.abspath(self.path)) or "." os.makedirs(directory, exist_ok=True) handle = tempfile.NamedTemporaryFile( "w", encoding="utf-8", dir=directory, prefix=".queue_events_", suffix=".tmp", delete=False, ) try: with handle: json.dump({"events": events}, handle, ensure_ascii=False, indent=2) # 原子替换:界面随时可能在读,不能让它撞上写了一半的文件 os.replace(handle.name, self.path) except Exception: try: os.unlink(handle.name) except OSError: pass raise