408 lines
18 KiB
Python
408 lines
18 KiB
Python
"""
|
||
引擎 B:数据直读检测(第二套企业微信自动回复方案)
|
||
====================================================
|
||
与引擎 A(截图 RPA)并存的第二检测通道。引擎 B **全程不碰鼠标键盘**:
|
||
|
||
检测:**两路并行、互为兜底**——
|
||
① DB 直读(升级数据源,wxwork_db.WXWorkDB 解密 message.db 增量轮询)
|
||
② conversations.json(会话档案监听,引擎 A 视觉回写的独立通道)
|
||
每轮两路都独立检测,任一数据源有盲区另一路补漏;投递靠
|
||
dedup_key 去重(fp_hex:时间戳,300s 防抖 + 队列级去重)保证
|
||
同一条消息只投一次。
|
||
投递:经 WeChatBot.enqueue_detected() 写入共享队列(pending_replies.json),
|
||
带 dedup_key 去重、detected_by 标记来源;
|
||
发送:不自己发送。队列条目由引擎 A 的 `_resume_orphaned_pending_reply`
|
||
统一捡起(重新定位 → 视觉提取 → 生成 → send_reply),发送动作天然
|
||
收敛到唯一通道 + 发送互斥锁(send_lock),任意时刻只有一个发送者。
|
||
|
||
DB 直读的 fp_hex 由会话名派生(wxwork_db.session_fp_from_name,与引擎 A
|
||
identity_by_name 模式的 _fp_from_name 算法一致),引擎 A 可在会话列表中
|
||
按名称重新定位窗口。DB 数据源不可用(无密钥/解密失败)时该路自动停用,
|
||
不影响 JSON 路径继续独立工作。
|
||
|
||
这样两套检测(A 看图、B 读数据)互不干扰、互为备份,回复出口严格串行。
|
||
|
||
对 WeChatBot 实例的接口约定(由 wechat_bot.py 提供):
|
||
- enqueue_detected(fp_hex, dedup_key, detected_by, chat_text,
|
||
display_name, last_lines) -> (bool, reason)
|
||
- has_active_pending(fp_hex) -> bool
|
||
"""
|
||
|
||
import json
|
||
import os
|
||
import threading
|
||
import time
|
||
|
||
_SCRIPT_DIR = os.path.dirname(os.path.abspath(__file__))
|
||
CONVERSATIONS_PATH = os.path.join(_SCRIPT_DIR, "conversations.json")
|
||
|
||
# 启动时「最后一条客户消息在 N 秒内」的会话视为刚来的新消息,允许投递;
|
||
# 超过 N 秒的只建档不投递,避免机器人一启动就把全部历史当新消息回复一遍。
|
||
STARTUP_FRESH_WINDOW_SECONDS = 120.0
|
||
|
||
# 防抖窗口:同一 dedup_key(同一会话 + 同一消息时间戳)在 N 秒内只投一次。
|
||
DEDUP_WINDOW_SECONDS = 300.0
|
||
|
||
|
||
class DataEngine:
|
||
"""引擎 B:数据直读检测器。独立线程轮询会话档案,增量投递共享队列。
|
||
|
||
db_source: 可选 WXWorkDB 实例(DB 直读升级数据源)。传入时与
|
||
conversations.json 路径**并行双跑**,互为兜底;DB 异常只停用
|
||
DB 路(db_active=False),JSON 路径不受影响。
|
||
|
||
data_source_mode: 引擎 B 内部数据源策略,三选一:
|
||
- "parallel"(默认):DB 直读 与 conversations.json 每轮并行检测,
|
||
互为兜底;同一条消息由 dedup_key 去重。
|
||
- "db":仅 DB 直读(无 db_source 时自动退化为仅 JSON)。
|
||
- "json":仅 conversations.json 档案监听(DB 直读完全关闭)。
|
||
"""
|
||
|
||
def __init__(self, bot=None, conversations_path=None, poll_interval: float = 2.0,
|
||
db_source=None, data_source_mode: str = "parallel"):
|
||
self.bot = bot
|
||
self._conversations_path = conversations_path or CONVERSATIONS_PATH
|
||
self._poll_interval = max(0.5, float(poll_interval))
|
||
self._db_source = db_source
|
||
self.data_source_mode = str(data_source_mode or "parallel").lower()
|
||
if self.data_source_mode not in ("parallel", "db", "json"):
|
||
self.data_source_mode = "parallel"
|
||
# DB 直读游标:已见最大 send_time(秒),各账号共用墙钟。
|
||
# 初始化为「现在 - 新鲜窗口」:与 conversations.json 路径的启动语义一致
|
||
# (只关心启动前 120s 内刚来的消息),避免启动时全量扫描历史建档。
|
||
self._db_cursor_ts = time.time() - STARTUP_FRESH_WINDOW_SECONDS
|
||
# fp_hex -> 会话游标(JSON 路径专用)
|
||
self._seen: dict = {}
|
||
# fp_hex -> 会话游标(DB 路径专用)。两条路径游标完全独立,
|
||
# 各自建档/判定,互不污染;投递去重由 _recently(dedup_key) 统一兜底。
|
||
self._seen_db: dict = {}
|
||
# dedup_key -> wall ts
|
||
self._recently: dict = {}
|
||
self._stop = threading.Event()
|
||
self._thread: threading.Thread | None = None
|
||
self.last_error = ""
|
||
self.enqueued_count = 0
|
||
self.processed_count = 0
|
||
self.db_active = db_source is not None
|
||
|
||
# ── 生命周期 ────────────────────────────────────────────────────────────
|
||
def start(self) -> None:
|
||
if self._thread and self._thread.is_alive():
|
||
return
|
||
self._stop.clear()
|
||
self._thread = threading.Thread(
|
||
target=self._run, name="engine_b", daemon=True
|
||
)
|
||
self._thread.start()
|
||
print("[引擎B] 数据直读检测已启动")
|
||
|
||
def stop(self) -> None:
|
||
self._stop.set()
|
||
if self._thread:
|
||
self._thread.join(timeout=3)
|
||
self._thread = None
|
||
|
||
def is_alive(self) -> bool:
|
||
return bool(self._thread and self._thread.is_alive())
|
||
|
||
def _run(self) -> None:
|
||
while not self._stop.is_set():
|
||
try:
|
||
# 按数据源模式调度:
|
||
# parallel:DB 直读 与 conversations.json 每轮并行检测,互为兜底;
|
||
# db:仅 DB 直读(无数据源时退化 JSON);
|
||
# json:仅 conversations.json 档案监听。
|
||
if self.data_source_mode == "db":
|
||
self._poll_db_once()
|
||
if self._db_source is None:
|
||
self.poll_once() # 无 DB 时退化 JSON,避免静默无检测
|
||
elif self.data_source_mode == "json":
|
||
self.poll_once()
|
||
else:
|
||
self._poll_db_once()
|
||
self.poll_once()
|
||
except Exception as exc: # 检测线程绝不允许崩溃退出
|
||
self.last_error = str(exc)
|
||
print(f" [引擎B] [!] 检测异常: {exc}")
|
||
# 防抖表必须在这里清,不能只在 poll_once 里清:data_source_mode="db"
|
||
# 且 DB 数据源可用时 poll_once 整轮都不会被调用,防抖记录就只进不出,
|
||
# 一台跑上几天的机器最后会把内存吃光。
|
||
self._prune_dedup()
|
||
self._stop.wait(self._poll_interval)
|
||
|
||
def _prune_dedup(self) -> None:
|
||
"""丢弃已经过了防抖窗口的投递记录。"""
|
||
now = time.time()
|
||
for key in [
|
||
key
|
||
for key, ts in list(self._recently.items())
|
||
if now - float(ts or 0.0) > DEDUP_WINDOW_SECONDS
|
||
]:
|
||
self._recently.pop(key, None)
|
||
|
||
# ── DB 直读检测(升级数据源) ──────────────────────────────────────────
|
||
def _poll_db_once(self) -> None:
|
||
"""DB 直读路径:从解密后的 message.db 增量读取新客户消息并投递。
|
||
|
||
与 conversations.json 路径**并行独立**(由 _run 每轮分别调用):
|
||
本路失败只停用本路(db_active=False),不影响 JSON 路径。
|
||
注意:DB 增量返回的是窗口内**逐条**消息,同一会话可能有多条;
|
||
必须按 fp_hex 聚合、只以最后一条作为「最新客户消息」投递,
|
||
否则会把会话历史逐条误判为新消息(与 conversations.json 路径
|
||
只看 history 最后一条的语义保持一致)。
|
||
"""
|
||
db = self._db_source
|
||
if db is None:
|
||
return # 无 DB 数据源:仅 JSON 路径工作
|
||
try:
|
||
msgs = db.get_new_messages(self._db_cursor_ts)
|
||
except Exception as exc:
|
||
self.last_error = f"DB 直读失败: {exc}"
|
||
print(f" [引擎B] [!] DB 直读失败: {exc}(JSON 路径不受影响,继续独立检测)")
|
||
self.db_active = False
|
||
return
|
||
# 查询成功即视为 DB 数据源可用(无新消息也应恢复状态,
|
||
# 否则故障恢复后 db_active 会一直停留在 False)。
|
||
self.db_active = True
|
||
if not msgs:
|
||
return
|
||
now = time.time()
|
||
max_ts = self._db_cursor_ts
|
||
# 按会话聚合增量消息(max_ts 先对全部消息推进,含被过滤的,防重复扫描)
|
||
by_fp: dict[str, list] = {}
|
||
for m in msgs:
|
||
try:
|
||
ts = float(m.get("send_time") or 0.0)
|
||
except (TypeError, ValueError):
|
||
continue
|
||
if ts > max_ts:
|
||
max_ts = ts
|
||
if m.get("is_self"):
|
||
continue # 自己发的消息不投递
|
||
fp_hex = str(m.get("fp_hex") or "")
|
||
content = str(m.get("content") or "").strip()
|
||
if not fp_hex or not content:
|
||
continue
|
||
by_fp.setdefault(fp_hex, []).append(m)
|
||
for fp_hex, group in by_fp.items():
|
||
group.sort(key=lambda m: float(m.get("send_time") or 0.0))
|
||
last = group[-1]
|
||
entry = {
|
||
"display_name": str(last.get("display_name") or "").strip(),
|
||
"history": [
|
||
{"role": "user",
|
||
"content": str(m.get("content") or "").strip(),
|
||
"ts": float(m.get("send_time") or 0.0)}
|
||
for m in group
|
||
],
|
||
"last_lines": [str(m.get("content") or "").strip() for m in group],
|
||
}
|
||
self._process_entry(fp_hex, entry, now, seen=self._seen_db)
|
||
self._db_cursor_ts = max_ts
|
||
self.db_active = True
|
||
|
||
# ── 检测 ─────────────────────────────────────────────────────────────────
|
||
def poll_once(self) -> None:
|
||
"""读一次 conversations.json,对比增量,投递新消息到共享队列。"""
|
||
data = self._read_conversations()
|
||
if data is None:
|
||
return
|
||
now = time.time()
|
||
for fp_hex, entry in data.items():
|
||
if not isinstance(entry, dict):
|
||
continue
|
||
self._process_entry(fp_hex, entry, now)
|
||
self._prune_dedup()
|
||
|
||
def _process_entry(self, fp_hex: str, entry: dict, now: float,
|
||
seen: dict | None = None) -> None:
|
||
"""单会话建档/增量判定/投递(DB 直读与 conversations.json 共用)。
|
||
|
||
seen 参数指定该路径自己的游标表(DB 路径传 _seen_db,JSON 路径
|
||
用默认 _seen),保证两路独立判定、互不污染;投递去重统一走
|
||
_maybe_enqueue 的 dedup_key 防抖。
|
||
|
||
- 首见会话:只记录游标;最后一条客户消息在新鲜窗口内(120s)才投递,
|
||
避免机器人一启动就把全部历史当新消息回复一遍。
|
||
- 已知会话:最后一条客户消息更新/内容变化时投递。
|
||
"""
|
||
store = self._seen if seen is None else seen
|
||
last_user = self._last_user_message(entry.get("history") or [])
|
||
cur = store.get(fp_hex)
|
||
if cur is None:
|
||
store[fp_hex] = self._make_cursor(entry, last_user)
|
||
if (
|
||
last_user
|
||
and (now - self._msg_ts(last_user)) <= STARTUP_FRESH_WINDOW_SECONDS
|
||
):
|
||
self._maybe_enqueue(fp_hex, entry, last_user)
|
||
return
|
||
if last_user and self._is_new_message(cur, last_user):
|
||
self._maybe_enqueue(fp_hex, entry, last_user)
|
||
store[fp_hex] = self._make_cursor(entry, last_user)
|
||
|
||
def _read_conversations(self) -> dict | None:
|
||
try:
|
||
if not os.path.exists(self._conversations_path):
|
||
return {}
|
||
with open(self._conversations_path, encoding="utf-8") as handle:
|
||
raw = json.load(handle)
|
||
return raw if isinstance(raw, dict) else {}
|
||
except Exception as exc:
|
||
self.last_error = str(exc)
|
||
return None
|
||
|
||
@staticmethod
|
||
def _last_user_message(hist) -> dict | None:
|
||
"""返回历史中最后一条 role=user 的消息;没有返回 None。"""
|
||
for msg in reversed(list(hist or [])):
|
||
if not isinstance(msg, dict):
|
||
continue
|
||
if str(msg.get("role") or "").lower() == "user":
|
||
return msg
|
||
return None
|
||
|
||
@staticmethod
|
||
def _msg_ts(msg: dict | None) -> float:
|
||
try:
|
||
return float((msg or {}).get("ts") or 0.0)
|
||
except (TypeError, ValueError):
|
||
return 0.0
|
||
|
||
def _make_cursor(self, entry: dict, last_user: dict | None) -> dict:
|
||
return {
|
||
"last_user_ts": self._msg_ts(last_user),
|
||
"last_user_content": str((last_user or {}).get("content") or ""),
|
||
"last_msg_count": len(entry.get("history") or []),
|
||
}
|
||
|
||
def _is_new_message(self, seen: dict, last_user: dict | None) -> bool:
|
||
ts = self._msg_ts(last_user)
|
||
content = str((last_user or {}).get("content") or "")
|
||
if ts > float(seen.get("last_user_ts") or 0.0):
|
||
return True
|
||
# 时间戳相同但内容变了:档案被重建/回填,也视为新消息
|
||
if content and content != str(seen.get("last_user_content") or ""):
|
||
return True
|
||
return False
|
||
|
||
# ── 投递 ─────────────────────────────────────────────────────────────────
|
||
def _maybe_enqueue(self, fp_hex: str, entry: dict, last_user: dict) -> None:
|
||
content = str(last_user.get("content") or "").strip()
|
||
ts = self._msg_ts(last_user)
|
||
dedup_key = f"{fp_hex}:{ts}"
|
||
now = time.time()
|
||
if dedup_key in self._recently:
|
||
return
|
||
self._recently[dedup_key] = now
|
||
# 无文字的新消息(纯图片/语音/表情)v1 先不投递,等待引擎 A 视觉兜底
|
||
if not content:
|
||
return
|
||
bot = self.bot
|
||
if bot is None:
|
||
return
|
||
try:
|
||
has_active = getattr(bot, "has_active_pending", None)
|
||
if has_active is not None and has_active(fp_hex):
|
||
# 引擎 A 已在处理同一会话(或队列里已有活跃任务),只补来源标记
|
||
self._merge_detected_by(fp_hex)
|
||
return
|
||
enqueue = getattr(bot, "enqueue_detected", None)
|
||
if enqueue is None:
|
||
return
|
||
ok, reason = enqueue(
|
||
fp_hex=fp_hex,
|
||
dedup_key=dedup_key,
|
||
detected_by="engine_b",
|
||
chat_text=content,
|
||
display_name=str(entry.get("display_name") or "").strip(),
|
||
last_lines=list(entry.get("last_lines") or [])[-20:],
|
||
)
|
||
self.processed_count += 1
|
||
if ok:
|
||
self.enqueued_count += 1
|
||
name = str(entry.get("display_name") or fp_hex[:8])
|
||
print(f" [引擎B] 投递新消息 → {name}: {content[:40]}")
|
||
else:
|
||
print(f" [引擎B] [!] 投递未执行: {reason}")
|
||
except Exception as exc:
|
||
self.last_error = str(exc)
|
||
print(f" [引擎B] [!] 投递异常: {exc}")
|
||
|
||
def _merge_detected_by(self, fp_hex: str) -> None:
|
||
"""队列中已有该会话的活跃任务时,把 detected_by 补上 engine_b 标记。"""
|
||
try:
|
||
bot = self.bot
|
||
if bot is None:
|
||
return
|
||
merge = getattr(bot, "merge_detected_by", None)
|
||
if merge is not None:
|
||
merge(fp_hex, "engine_b")
|
||
except Exception:
|
||
pass
|
||
|
||
|
||
if __name__ == "__main__":
|
||
# 冒烟测试:构造最小桩 bot,验证检测逻辑能发现新消息并投递
|
||
class StubBot:
|
||
def __init__(self):
|
||
self.enqueued = []
|
||
|
||
def has_active_pending(self, fp_hex):
|
||
return False
|
||
|
||
def enqueue_detected(self, **kwargs):
|
||
self.enqueued.append(kwargs)
|
||
return True, "ok"
|
||
|
||
import tempfile
|
||
|
||
stub = StubBot()
|
||
path = os.path.join(tempfile.gettempdir(), "engine_b_smoke.json")
|
||
now = time.time()
|
||
with open(path, "w", encoding="utf-8") as handle:
|
||
json.dump(
|
||
{
|
||
"a" * 80: {
|
||
"display_name": "测试客户",
|
||
"history": [
|
||
{"role": "user", "content": "你好", "ts": now - 50},
|
||
{"role": "assistant", "content": "您好,很高兴为您服务", "ts": now - 49},
|
||
{"role": "user", "content": "我想挂号", "ts": now - 5},
|
||
],
|
||
"last_lines": ["我想挂号"],
|
||
}
|
||
},
|
||
handle,
|
||
ensure_ascii=False,
|
||
)
|
||
engine = DataEngine(bot=stub, conversations_path=path, poll_interval=0.5)
|
||
engine.poll_once()
|
||
assert engine.enqueued_count == 1, stub.enqueued
|
||
assert stub.enqueued[0]["dedup_key"] == "a" * 80 + f":{now - 5}"
|
||
assert stub.enqueued[0]["detected_by"] == "engine_b"
|
||
# 第二次轮询:无新消息,不应重复投递
|
||
engine.poll_once()
|
||
assert engine.enqueued_count == 1
|
||
# 模拟客户又发一条
|
||
with open(path, "w", encoding="utf-8") as handle:
|
||
json.dump(
|
||
{
|
||
"a" * 80: {
|
||
"display_name": "测试客户",
|
||
"history": [
|
||
{"role": "user", "content": "你好", "ts": now - 50},
|
||
{"role": "assistant", "content": "您好", "ts": now - 49},
|
||
{"role": "user", "content": "我想挂号", "ts": now - 5},
|
||
{"role": "user", "content": "请问几点上班", "ts": now - 1},
|
||
],
|
||
"last_lines": ["请问几点上班"],
|
||
}
|
||
},
|
||
handle,
|
||
ensure_ascii=False,
|
||
)
|
||
engine.poll_once()
|
||
assert engine.enqueued_count == 2, stub.enqueued
|
||
print("[OK] 引擎 B 冒烟测试通过:检测增量 → 去重 → 投递")
|