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