Files
kefu/wechat_rpa/engine_b.py
2026-08-27 14:04:28 +08:00

408 lines
18 KiB
Python
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""
引擎 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:
# 按数据源模式调度:
# parallelDB 直读 与 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_dbJSON 路径
用默认 _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 冒烟测试通过:检测增量 → 去重 → 投递")