278 lines
9.8 KiB
Python
278 lines
9.8 KiB
Python
"""会话档案收口工具
|
||
==================
|
||
|
||
行锚点漂移曾让同一个联系人被铸造出多个会话键(头像哈希相差 9~31 位,容差只有
|
||
6 位)。锚点已按头像方块本身定位修好,但磁盘上残留的分裂档案和空壳待回复任务
|
||
不会自动消失,本工具做一次性收口:
|
||
|
||
1. 把只能唯一证明属于同一联系人的旧档案键合并进当前键,按时间戳归并历史,
|
||
保留当前键的画面快照(旧快照会污染增量比对基线);
|
||
2. 删除既没有聊天内容、也没有待发回复的空壳待回复任务——它们驱动不了任何
|
||
重试,只会让机器人以为有幽灵会话欠着回复。
|
||
|
||
判定“同一联系人”沿用 wechat_bot 的标准:两份档案各自的客户说话人集合都只有
|
||
一个人,且是同一个人。证据不唯一时一律保留,绝不猜测合并。
|
||
|
||
用法(先停掉监听,避免运行中的进程回写覆盖结果):
|
||
|
||
python reconcile_session_archive.py # 只报告,不改动
|
||
python reconcile_session_archive.py --apply # 执行收口,改动前自动备份
|
||
"""
|
||
|
||
import argparse
|
||
import json
|
||
import os
|
||
import re
|
||
import shutil
|
||
import time
|
||
|
||
_HERE = os.path.dirname(os.path.abspath(__file__))
|
||
_SPEAKER_HEADER = re.compile(
|
||
r"^(?P<speaker>.{1,40}?)\s+\d{1,2}/\d{1,2}\s+\d{1,2}:\d{2}(:\d{2})?$"
|
||
)
|
||
|
||
|
||
def _load(path: str) -> dict:
|
||
if not os.path.exists(path):
|
||
return {}
|
||
try:
|
||
with open(path, encoding="utf-8") as handle:
|
||
data = json.load(handle)
|
||
except (OSError, ValueError) as error:
|
||
print(f"[-] 读取 {os.path.basename(path)} 失败: {error}")
|
||
return {}
|
||
return data if isinstance(data, dict) else {}
|
||
|
||
|
||
def _agent_name() -> str:
|
||
try:
|
||
from ai_config import AI_AGENT_NAME
|
||
except ImportError:
|
||
return ""
|
||
return str(AI_AGENT_NAME or "").strip()
|
||
|
||
|
||
def _outgoing_speakers(archive: dict) -> set:
|
||
"""收集全部档案里已证实的我方发言人名。
|
||
|
||
企业微信里我方气泡用的是账号显示名(例如“高兴亮”),与配置的 AI 人设名
|
||
(AI_AGENT_NAME)通常不同;档案发送成功时会把它记进 outgoing_speakers。
|
||
这个名字对所有会话都一样,因此可以跨档案排除。
|
||
"""
|
||
names = set()
|
||
for entry in archive.values():
|
||
if not isinstance(entry, dict):
|
||
continue
|
||
names.update(
|
||
str(speaker).strip()
|
||
for speaker in entry.get("outgoing_speakers") or []
|
||
if str(speaker).strip()
|
||
)
|
||
return names
|
||
|
||
|
||
def _speakers(entry: dict, agent_name: str, outgoing: set = frozenset()) -> set:
|
||
"""收集该档案里出现过的客户说话人(排除我方坐席)。"""
|
||
history = [item for item in (entry.get("history") or []) if isinstance(item, dict)]
|
||
assistant_bodies = {
|
||
" ".join(str(item.get("content") or "").split())
|
||
for item in history
|
||
if item.get("role") == "assistant"
|
||
}
|
||
sources = ["\n".join(entry.get("last_lines") or [])]
|
||
sources.extend(str(item.get("content") or "") for item in history)
|
||
|
||
blocks = []
|
||
for source in sources:
|
||
speaker, body = "", []
|
||
for raw_line in str(source).splitlines():
|
||
line = raw_line.strip()
|
||
match = _SPEAKER_HEADER.match(line)
|
||
if match:
|
||
if speaker:
|
||
blocks.append((speaker, "\n".join(body).strip()))
|
||
speaker, body = match.group("speaker").strip(), []
|
||
elif speaker and line:
|
||
body.append(line)
|
||
if speaker:
|
||
blocks.append((speaker, "\n".join(body).strip()))
|
||
|
||
# 内容与我方回复对得上的说话人就是坐席自己,不能算客户。
|
||
agents = {
|
||
speaker
|
||
for speaker, body in blocks
|
||
if " ".join(body.split()) in assistant_bodies and body
|
||
}
|
||
if agent_name:
|
||
agents.update(speaker for speaker, _ in blocks if agent_name in speaker)
|
||
agents.update(outgoing)
|
||
return {
|
||
speaker
|
||
for speaker, _ in blocks
|
||
if speaker and speaker not in agents
|
||
and not (agent_name and agent_name in speaker)
|
||
}
|
||
|
||
|
||
def _describe(key: str, entry: dict) -> str:
|
||
updated = entry.get("updated") or 0
|
||
stamp = (
|
||
time.strftime("%m-%d %H:%M", time.localtime(updated)) if updated else "未记录"
|
||
)
|
||
return "%s…(%d 字节键) 最后更新 %s,历史 %d 条,快照 %d 行" % (
|
||
key[:16],
|
||
len(key) // 2,
|
||
stamp,
|
||
len(entry.get("history") or []),
|
||
len(entry.get("last_lines") or []),
|
||
)
|
||
|
||
|
||
def plan_archive_merges(archive: dict, agent_name: str) -> list:
|
||
"""找出可以唯一证明归属的分裂档案,返回 (旧键, 新键, 客户名) 列表。"""
|
||
outgoing = _outgoing_speakers(archive)
|
||
profiles = {}
|
||
for key, entry in archive.items():
|
||
if not isinstance(entry, dict):
|
||
continue
|
||
names = _speakers(entry, agent_name, outgoing)
|
||
if len(names) == 1:
|
||
profiles[key] = next(iter(names))
|
||
|
||
merges = []
|
||
for name in set(profiles.values()):
|
||
owners = sorted(
|
||
(key for key, value in profiles.items() if value == name),
|
||
key=lambda key: (
|
||
len(key),
|
||
float(archive[key].get("updated") or 0),
|
||
),
|
||
)
|
||
if len(owners) < 2:
|
||
continue
|
||
# 键最长、最近更新的那份是当前生效的身份,其余合并进它。
|
||
target = owners[-1]
|
||
merges.extend((old, target, name) for old in owners[:-1])
|
||
return merges
|
||
|
||
|
||
def merge_archive_entry(archive: dict, old_key: str, new_key: str) -> None:
|
||
"""把旧档案的历史按时间戳并入新档案,快照沿用新档案的。"""
|
||
old_entry = archive.get(old_key) or {}
|
||
new_entry = archive.get(new_key) or {}
|
||
history = [
|
||
item
|
||
for item in list(old_entry.get("history") or [])
|
||
+ list(new_entry.get("history") or [])
|
||
if isinstance(item, dict)
|
||
]
|
||
seen = set()
|
||
merged = []
|
||
for item in sorted(history, key=lambda item: float(item.get("ts") or 0)):
|
||
marker = (
|
||
str(item.get("role") or ""),
|
||
" ".join(str(item.get("content") or "").split()),
|
||
)
|
||
if marker in seen:
|
||
continue
|
||
seen.add(marker)
|
||
merged.append(item)
|
||
new_entry["history"] = merged[-200:]
|
||
archive[new_key] = new_entry
|
||
archive.pop(old_key, None)
|
||
|
||
|
||
def plan_pending_drops(pending: dict) -> list:
|
||
"""列出既无聊天内容、也无待发回复的空壳任务。"""
|
||
drops = []
|
||
for key, value in pending.items():
|
||
if not isinstance(value, dict):
|
||
continue
|
||
if str(value.get("send_state") or "").strip():
|
||
continue # 发送中的任务必须保留,可能已经发出去了
|
||
actionable = (
|
||
str(value.get("chat_text") or "").strip()
|
||
or (value.get("last_lines") or [])
|
||
or str(value.get("reply_text") or "").strip()
|
||
or str(value.get("staged_reply_text") or "").strip()
|
||
or str(value.get("staged_user_text") or "").strip()
|
||
or str(value.get("exchange_id") or "").strip()
|
||
)
|
||
if not actionable:
|
||
drops.append(key)
|
||
return drops
|
||
|
||
|
||
def main() -> None:
|
||
parser = argparse.ArgumentParser(description="会话档案与待回复任务收口")
|
||
parser.add_argument("--apply", action="store_true", help="执行改动(默认只报告)")
|
||
args = parser.parse_args()
|
||
|
||
archive_path = os.path.join(_HERE, "conversations.json")
|
||
pending_path = os.path.join(_HERE, "pending_replies.json")
|
||
archive = _load(archive_path)
|
||
pending = _load(pending_path)
|
||
agent_name = _agent_name()
|
||
|
||
merges = plan_archive_merges(archive, agent_name)
|
||
drops = plan_pending_drops(pending)
|
||
|
||
print("会话档案 %d 份,待回复任务 %d 个。" % (len(archive), len(pending)))
|
||
if merges:
|
||
print("\n[分裂档案] 可唯一证明属于同一联系人,将合并:")
|
||
for old, new, name in merges:
|
||
print(" 客户「%s」" % name)
|
||
print(" 旧 %s" % _describe(old, archive.get(old) or {}))
|
||
print(" 新 %s" % _describe(new, archive.get(new) or {}))
|
||
else:
|
||
print("\n[分裂档案] 没有能唯一证明归属的重复档案。")
|
||
if drops:
|
||
print("\n[空壳任务] 无内容也无待发回复,将删除:")
|
||
for key in drops:
|
||
created = (pending.get(key) or {}).get("created_at") or 0
|
||
print(
|
||
" %s… 建立于 %s"
|
||
% (
|
||
key[:16],
|
||
time.strftime("%m-%d %H:%M", time.localtime(created))
|
||
if created
|
||
else "未记录",
|
||
)
|
||
)
|
||
else:
|
||
print("\n[空壳任务] 没有需要清理的空壳任务。")
|
||
|
||
if not merges and not drops:
|
||
print("\n无需收口。")
|
||
return
|
||
if not args.apply:
|
||
print("\n以上为预演。加 --apply 执行(会先备份两个 json)。")
|
||
return
|
||
|
||
stamp = time.strftime("%Y%m%d_%H%M%S")
|
||
for path in (archive_path, pending_path):
|
||
if os.path.exists(path):
|
||
backup = f"{path}.{stamp}.bak"
|
||
shutil.copy2(path, backup)
|
||
print("已备份 %s" % os.path.basename(backup))
|
||
|
||
for old, new, _name in merges:
|
||
merge_archive_entry(archive, old, new)
|
||
for key in drops:
|
||
pending.pop(key, None)
|
||
|
||
for path, data in ((archive_path, archive), (pending_path, pending)):
|
||
tmp_path = f"{path}.tmp"
|
||
with open(tmp_path, "w", encoding="utf-8") as handle:
|
||
json.dump(data, handle, ensure_ascii=False, indent=2)
|
||
os.replace(tmp_path, path)
|
||
|
||
print(
|
||
"\n收口完成:合并 %d 份分裂档案,删除 %d 个空壳任务。"
|
||
% (len(merges), len(drops))
|
||
)
|
||
|
||
|
||
if __name__ == "__main__":
|
||
main()
|