124 lines
4.6 KiB
Python
124 lines
4.6 KiB
Python
"""
|
|
回复队列执行记录
|
|
================
|
|
待回复队列本身只存"现在还欠谁一条回复",不留痕迹:任务一完成就被删掉,出了
|
|
问题事后完全看不出它经历过什么。排查"为什么这个客户没收到回复"、"为什么同一
|
|
句话回了三遍"时,缺的就是这条时间线。
|
|
|
|
这里按时间顺序记下每个会话在队列里的每一步:入队、开始处理、生成回复、已发送、
|
|
完成、跳过、失败。机器人线程写、界面线程读,所以写入走临时文件加原子替换,
|
|
界面永远不会读到写了一半的 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
|