WechatOnCloud/bridge/woc_bridge/messaging/friend_watcher.py
Kris c45282f094 feat: 新增好友自动通过、UI自动化能力与多分辨率适配
- 新增登录状态守卫后台任务
- 新增好友申请自动通过规则引擎
- 新增多分辨率UI配置与模板资源
- 新增消息拉取复合游标支持
- 优化发送队列与UI自动化逻辑
- 新增批量发送日志与错误处理
- 优化Docker镜像构建与ptrace初始化
- 新增联系人名称缓存预热
2026-07-17 18:10:31 +08:00

558 lines
23 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.

"""好友申请监听器:轮询 fmessage 系统消息,解析后触发自动通过。
独立后台协程,不侵入 MessageStreamer 架构。复用其 mtime 感知 +
复合游标 (create_time, local_id) 增量检测模式,但独立游标与频率控制。
auto_accept 关闭时挂起在 asyncio.Event 上,不轮询 DB。
DB_ENCRYPTED 时 _init_cursor 返回 False每轮重试初始化。
"""
from __future__ import annotations
import asyncio
import logging
import time
from dataclasses import dataclass
from typing import Optional
from woc_bridge.models import AcceptDecision, BridgeError
from woc_bridge.models.contact import FriendRequestInfo
from woc_bridge.messaging.friend_parser import parse_friend_request
logger = logging.getLogger("woc-bridge")
class FriendRequestWatcher:
"""好友申请监听器:轮询 fmessage 系统消息,解析后触发自动通过。"""
def __init__(
self,
db_reader, # DbReader避免类型注解循环 import
rule_engine, # AcceptRuleEngine
send_queue, # SendQueue
idem_cache, # IdemCache
xdotool_driver, # XdotoolDriver
breaker, # CircuitBreaker
poll_interval: float = 3.0,
cursor_create_time: int = 0,
cursor_local_id: int = 0,
login_guard=None, # LoginGuard登录守卫None 时跳过登录检查)
) -> None:
self._db_reader = db_reader
self._rule_engine = rule_engine
self._send_queue = send_queue
self._idem_cache = idem_cache
self._xdotool = xdotool_driver
self._breaker = breaker
self._poll_interval = poll_interval
self._login_guard = login_guard
# 复合游标 (create_time, local_id),与 get_friend_requests_since 返回值对齐
self._cursor_create_time: int = cursor_create_time
self._cursor_local_id: int = cursor_local_id
self._last_db_mtime: float = 0.0
self._last_wal_mtime: float = 0.0
# contact.db 独立 mtime好友申请可能只写入 contact.db 的 ticket_info
# 不触发 message_0.db 变化,必须单独感知
self._last_contact_db_mtime: float = 0.0
self._last_contact_wal_mtime: float = 0.0
# ticket_info 增量游标:避免每次全量读取历史已通过/残留 ticket
self._ticket_cursor_local_id: int = 0
self._ticket_cursor_inited: bool = False
self._watcher: Optional[asyncio.Task] = None
self._enabled_event: asyncio.Event = asyncio.Event() # auto_accept 开关事件
self._cursor_inited: bool = False
# 上轮满载标志:为 True 时跳过 mtime 一级过滤,避免批量积压(>50 条)时
# 剩余批次因 DB mtime 未变被卡住
self._has_more_pending: bool = False
self._has_more_ticket_pending: bool = False
# 状态统计(供 /api/friends/auto_accept/status 查询)
self._processed_count: int = 0
self._accepted_count: int = 0
self._rejected_count: int = 0
self._last_processed_time: Optional[int] = None
# 短期去重:同一个 wxid 在 N 秒内不重复入队执行 UI 通过,
# 避免已通过但 verify 校验超时/失败时被反复处理导致乱点
self._recent_wxids: dict[str, float] = {}
self._recent_wxid_ttl: float = 60.0
@property
def is_running(self) -> bool:
return self._watcher is not None and not self._watcher.done()
@property
def is_enabled(self) -> bool:
return self._enabled_event.is_set()
@property
def processed_count(self) -> int:
return self._processed_count
@property
def accepted_count(self) -> int:
return self._accepted_count
@property
def rejected_count(self) -> int:
return self._rejected_count
@property
def last_processed_time(self) -> Optional[int]:
return self._last_processed_time
@property
def cursor_create_time(self) -> int:
return self._cursor_create_time
@property
def cursor_local_id(self) -> int:
return self._cursor_local_id
async def start(self) -> None:
if self._watcher is None or self._watcher.done():
self._watcher = asyncio.create_task(self._watch_loop())
logger.info("FriendRequestWatcher 已启动")
async def stop(self) -> None:
if self._watcher is not None and not self._watcher.done():
self._watcher.cancel()
try:
await self._watcher
except asyncio.CancelledError:
pass
self._watcher = None
def set_enabled(self, enabled: bool) -> None:
if enabled:
# 开启时重置游标,避免批量处理积压申请
self._cursor_inited = False
self._ticket_cursor_inited = False
self._has_more_pending = False
self._has_more_ticket_pending = False
self._enabled_event.set()
else:
self._enabled_event.clear()
logger.info("FriendRequestWatcher enabled=%s", enabled)
async def _init_cursor(self) -> bool:
"""对齐游标到当前 fmessage 表最大 (create_time, local_id)。
Returns:
True 表示已对齐(或 DB 不可读但标记为已初始化,避免无限重试)
False 表示异常,下一轮重试
"""
try:
max_cursor = await asyncio.to_thread(
self._db_reader.get_max_create_time_for_talker, "fmessage"
)
if max_cursor is None:
# DB 不可读_ensure_decrypted 抛 BridgeError 已被吞掉返回 None
# 不标记 _cursor_inited下一轮重试
logger.warning("FriendRequestWatcher 游标初始化: DB 不可读,将在下轮重试")
return False
if isinstance(max_cursor, tuple) and max_cursor[0] > 0:
self._cursor_create_time = max_cursor[0]
self._cursor_local_id = max_cursor[1]
logger.info(
"FriendRequestWatcher 游标对齐到 (create_time=%d, local_id=%d)(跳过历史申请)",
self._cursor_create_time, self._cursor_local_id,
)
# 无论 max_cursor 是否为 (0, 0)(空表),都标记为已初始化,避免空表时无限重试
self._cursor_inited = True
return True
except Exception as e:
logger.warning("FriendRequestWatcher 游标初始化失败: %s", e)
return False
async def _init_ticket_cursor(self) -> bool:
"""初始化 ticket_info 增量游标到当前最大 id跳过历史数据。
Returns:
True 表示已初始化None 时按 0 处理并标记成功。
"""
try:
max_id = await asyncio.to_thread(
self._db_reader.get_max_ticket_info_id
)
if max_id is None:
logger.warning(
"FriendRequestWatcher ticket 游标初始化: DB 不可读,将在下轮重试"
)
return False
self._ticket_cursor_local_id = max_id
self._ticket_cursor_inited = True
logger.info(
"FriendRequestWatcher ticket 游标对齐到 max_id=%d(跳过历史 ticket",
self._ticket_cursor_local_id,
)
return True
except Exception as e:
logger.warning("FriendRequestWatcher ticket 游标初始化失败: %s", e)
return False
async def _poll_once(self) -> None:
# 1. mtime 感知(一级过滤,避免空 SQL
# 必须同时感知 message_0.db 与 contact.db微信 4.x Linux 好友申请可能
# 只写入 contact.db 的 ticket_info不触发 message_0.db 变化。
# 上轮满载时跳过过滤(仍有未读批次),否则会因 DB mtime 未变卡住剩余申请。
msg_changed = False
contact_changed = False
mtimes = await asyncio.to_thread(
self._db_reader.get_db_mtime, "message/message_0.db"
)
if mtimes is None:
logger.debug("FriendRequestWatcher _poll_once: message_0.db mtime 不可读")
else:
db_mtime, wal_mtime = mtimes
if (
self._has_more_pending
or db_mtime != self._last_db_mtime
or wal_mtime != self._last_wal_mtime
):
msg_changed = True
self._last_db_mtime = db_mtime
self._last_wal_mtime = wal_mtime
mtimes_contact = await asyncio.to_thread(
self._db_reader.get_db_mtime, "contact/contact.db"
)
if mtimes_contact is None:
logger.debug("FriendRequestWatcher _poll_once: contact.db mtime 不可读")
else:
cdb_mtime, cwal_mtime = mtimes_contact
if (
self._has_more_ticket_pending
or cdb_mtime != self._last_contact_db_mtime
or cwal_mtime != self._last_contact_wal_mtime
):
contact_changed = True
self._last_contact_db_mtime = cdb_mtime
self._last_contact_wal_mtime = cwal_mtime
logger.debug(
"FriendRequestWatcher _poll_once: 开始轮询 "
"cursor=(%d,%d) msg_changed=%s contact_changed=%s "
"last_msg_mtime=(%.3f,%.3f) last_contact_mtime=(%.3f,%.3f)",
self._cursor_create_time, self._cursor_local_id,
msg_changed, contact_changed,
self._last_db_mtime, self._last_wal_mtime,
self._last_contact_db_mtime, self._last_contact_wal_mtime,
)
if not msg_changed and not contact_changed:
logger.debug(
"FriendRequestWatcher _poll_once: message_0.db/contact.db mtime 均未变化,跳过本次轮询"
)
return
# 2. 查询 fmessage 分片表增量消息(复合游标)
result = await asyncio.to_thread(
self._db_reader.get_friend_requests_since,
self._cursor_create_time,
self._cursor_local_id,
50, # limit
)
if result is None:
# DB 不可读_ensure_decrypted 抛 BridgeError
logger.warning(
"FriendRequestWatcher _poll_once: get_friend_requests_since 返回 NoneDB 不可读)"
)
return
raw_items: list[dict] = list(result["requests"])
# 2.1 fallback微信 4.x Linux 好友申请可能仅存于 contact.db 的 ticket_info
# fmessage 表可能为空。额外读取 ticket_info通过 since_id 游标跳过历史数据。
ticket_items = await asyncio.to_thread(
self._db_reader.get_pending_requests_from_ticket_info,
50,
self._ticket_cursor_local_id,
)
if ticket_items:
logger.info(
"FriendRequestWatcher _poll_once: 从 ticket_info 读到 %d 条新申请 "
"(since_id=%d)",
len(ticket_items), self._ticket_cursor_local_id,
)
raw_items.extend(ticket_items)
# 推进 ticket_info 游标到本批最大 id
max_ticket_id = max(item["local_id"] for item in ticket_items)
if max_ticket_id > self._ticket_cursor_local_id:
self._ticket_cursor_local_id = max_ticket_id
self._has_more_ticket_pending = len(ticket_items) >= 50
logger.info(
"FriendRequestWatcher _poll_once: 本轮检测到 %d 条原始申请 "
"(cursor=%d,%d → next=%d,%d)",
len(raw_items),
self._cursor_create_time, self._cursor_local_id,
result["next_create_time"], result["next_local_id"],
)
# 3. 解析 XML + 逐条处理
parsed_count = 0
for item in raw_items:
info = parse_friend_request(
item["content"], item["create_time"], item["local_id"]
)
if info is None:
# 非 verifyUser 类型或解析失败,记录日志便于排查
logger.debug(
"FriendRequestWatcher: 解析失败或非 verifyUsercontent=%s",
repr(item["content"][:200]) if item.get("content") else "",
)
continue
parsed_count += 1
logger.info(
"FriendRequestWatcher: 解析到申请 wxid=%s nickname=%s scene=%s",
info.stranger_wxid, info.nickname, info.scene,
)
await self._handle_request(info)
# 4. 推进复合游标(仅由 fmessage 结果决定ticket_info 无时间戳,不参与游标)
self._cursor_create_time = result["next_create_time"]
self._cursor_local_id = result["next_local_id"]
# 5. 满载标志:仅由 fmessage 本批是否达 limit 决定
self._has_more_pending = len(result["requests"]) >= 50
logger.debug(
"FriendRequestWatcher _poll_once: 解析成功 %d/%d 条,满载=%s",
parsed_count, len(raw_items), self._has_more_pending,
)
async def _watch_loop(self) -> None:
while True:
try:
# auto_accept 关闭时挂起,避免空轮询
await self._enabled_event.wait()
# 登录守卫:微信未登录时跳过轮询,不推进游标
# 登录后自然从上次游标位置重新读取DB 中积累的申请不丢失
if self._login_guard is not None and not self._login_guard.is_logged_in:
logger.info(
"FriendRequestWatcher: 微信未登录 (state=%s),暂停轮询等待登录恢复",
self._login_guard.current_state,
)
# 等待登录恢复,带超时避免紧密循环
# 超时后 continue 回到 while 顶部重新检查状态
await self._login_guard.wait_for_login(
timeout=self._poll_interval
)
continue
# 游标未初始化时先对齐DB 恢复可读后自动补齐)
# DB_ENCRYPTED 时 _init_cursor 返回 False下一轮仍会重试
if not self._cursor_inited:
await self._init_cursor()
if not self._cursor_inited:
# DB 仍不可读,本轮跳过 _poll_once
await asyncio.sleep(self._poll_interval)
continue
# fmessage 游标初始化成功后,再初始化 ticket_info 游标
if not self._ticket_cursor_inited:
await self._init_ticket_cursor()
if not self._ticket_cursor_inited:
await asyncio.sleep(self._poll_interval)
continue
await asyncio.sleep(self._poll_interval)
await self._poll_once()
except asyncio.CancelledError:
raise
except Exception as e:
logger.exception("FriendRequestWatcher 异常: %s", e)
await asyncio.sleep(5.0)
async def _verify_with_retry(self, stranger_wxid: str) -> bool:
"""UI 操作完成后轮询 DB 校验好友是否已通过。
先等待 1.5s 让微信刷盘,之后每隔 1s 查询一次,最多 10s。
期间任一次校验成功即返回 TrueDB 异常视为未通过并继续重试。
"""
await asyncio.sleep(1.5)
deadline = time.monotonic() + 10.0
attempt = 0
while time.monotonic() < deadline:
attempt += 1
try:
verified = await asyncio.to_thread(
self._db_reader.verify_friend_accepted, stranger_wxid
)
logger.info(
"FriendRequestWatcher: verify 尝试 #%d (wxid=%s) → verified=%s",
attempt, stranger_wxid, verified,
)
if verified:
return True
except Exception as exc:
logger.warning(
"FriendRequestWatcher: verify_friend_accepted 异常 #%d (wxid=%s): %s",
attempt, stranger_wxid, exc,
)
await asyncio.sleep(1.0)
logger.warning(
"FriendRequestWatcher: verify 超时 (wxid=%s, 尝试 %d 次)",
stranger_wxid, attempt,
)
return False
async def _handle_request(self, req: FriendRequestInfo) -> None:
self._processed_count += 1
self._last_processed_time = req.create_time
# 步骤 1收到申请打印原始信息
logger.info(
"[自动通过][步骤1/8] 收到好友申请 wxid=%s nickname=%s scene=%s verify=%s",
req.stranger_wxid,
req.nickname,
req.scene,
repr(req.verify_message[:80]) if req.verify_message else "",
)
# 步骤 2规则引擎决策
logger.info("[自动通过][步骤2/8] 进入规则引擎评估 wxid=%s", req.stranger_wxid)
decision = await self._rule_engine.evaluate(req)
if decision == AcceptDecision.REJECT:
self._rejected_count += 1
logger.info(
"[自动通过][步骤2/8] 决策结果=REJECT拒绝处理结束 wxid=%s",
req.stranger_wxid,
)
return
if decision != AcceptDecision.ACCEPT:
logger.info(
"[自动通过][步骤2/8] 决策结果=SKIP跳过处理结束 wxid=%s"
"原因未命中任何通过规则accept_all=false、无白名单/关键词匹配)",
req.stranger_wxid,
)
return
logger.info(
"[自动通过][步骤2/8] 决策结果=ACCEPT允许通过 wxid=%s",
req.stranger_wxid,
)
# 步骤 2.5:前置 DB 检查,避免已经通过的好友被重复执行 UI
logger.info("[自动通过][步骤2.5/8] 前置校验是否已是好友 wxid=%s", req.stranger_wxid)
try:
already_accepted = await asyncio.to_thread(
self._db_reader.verify_friend_accepted, req.stranger_wxid
)
except Exception as exc:
logger.warning(
"[自动通过][步骤2.5/8] verify_friend_accepted 异常,按未通过继续 wxid=%s: %s",
req.stranger_wxid, exc,
)
already_accepted = False
if already_accepted:
self._idem_cache.set(
"friend_accept", req.stranger_wxid, "",
{"success": True, "verified": True, "pre_check": True}, "",
)
self._accepted_count += 1
logger.info(
"[自动通过][步骤2.5/8] 已是好友,跳过 UI 操作并写入幂等缓存 wxid=%s",
req.stranger_wxid,
)
return
logger.info("[自动通过][步骤2.5/8] 尚未通过,继续执行 UI wxid=%s", req.stranger_wxid)
# 步骤 3幂等去重TTL=300s5 分钟内不重复处理同一申请人)
logger.info("[自动通过][步骤3/8] 检查幂等缓存 wxid=%s", req.stranger_wxid)
cached = self._idem_cache.get(
"friend_accept", req.stranger_wxid, "", ""
)
if cached is not None:
logger.info(
"[自动通过][步骤3/8] 幂等命中,跳过 wxid=%s (cached=%s)",
req.stranger_wxid, cached,
)
return
# 步骤 3.5:短期去重,避免 verify 超时/失败后在短时间内重复入队乱点
now = time.monotonic()
last_at = self._recent_wxids.get(req.stranger_wxid)
if last_at is not None and now - last_at < self._recent_wxid_ttl:
logger.info(
"[自动通过][步骤3.5/8] %.0fs 内已处理过,跳过 wxid=%s",
self._recent_wxid_ttl, req.stranger_wxid,
)
return
logger.info("[自动通过][步骤3/8] 幂等未命中,继续处理 wxid=%s", req.stranger_wxid)
# 步骤 4熔断检查
logger.info("[自动通过][步骤4/8] 检查熔断器 wxid=%s", req.stranger_wxid)
breaker_allowed = self._breaker.allow()
breaker_state = self._breaker.state.value
logger.info(
"[自动通过][步骤4/8] 熔断器状态 allow=%s state=%s wxid=%s",
breaker_allowed, breaker_state, req.stranger_wxid,
)
if not breaker_allowed:
logger.warning(
"[自动通过][步骤4/8] 熔断器 OPEN跳过 wxid=%s",
req.stranger_wxid,
)
return
# 步骤 5/6入队执行 UI 通过
logger.info(
"[自动通过][步骤5/8] 准备入队执行 UI 通过 wxid=%s",
req.stranger_wxid,
)
# 记录本次处理时间,用于短期去重(即使后续 verify 失败也不立即重试)
self._recent_wxids[req.stranger_wxid] = time.monotonic()
try:
result = await self._send_queue.enqueue(
lambda req=req: self._xdotool.accept_friend_request(
stranger_wxid=req.stranger_wxid,
nickname=req.nickname,
),
delay_ms=None, # 使用 send_queue 默认间隔
)
logger.info(
"[自动通过][步骤6/8] UI 操作完成 result=%s,开始 DB 校验 wxid=%s",
result, req.stranger_wxid,
)
# 步骤 7等待微信 DB WAL 刷盘后校验
verified = await self._verify_with_retry(req.stranger_wxid)
# 步骤 8结果处理
if verified:
# 只有真正通过并验证成功才写幂等缓存,避免坐标不准等临时失败
# 导致该申请人 5 分钟内无法再次处理
self._idem_cache.set(
"friend_accept", req.stranger_wxid, "",
{"success": True, "verified": True}, "",
)
self._breaker.record_success()
self._accepted_count += 1
logger.info(
"[自动通过][步骤8/8] 通过并验证成功 wxid=%s nickname=%s",
req.stranger_wxid, req.nickname,
)
else:
# UI 操作完成但 DB 校验未通过(可能有延迟),不写入幂等缓存,
# 下轮仍可尝试;仅记录熔断器失败。
self._breaker.record_failure()
logger.warning(
"[自动通过][步骤8/8] UI 操作完成但 DB 校验未通过,"
"不写入幂等缓存,允许下轮重试 wxid=%s",
req.stranger_wxid,
)
except BridgeError as e:
# 透传 BridgeError 错误码RATE_LIMITED / SEND_FAILED / WINDOW_NOT_FOUND 等)
self._breaker.record_failure()
logger.error(
"[自动通过][步骤6/8] BridgeError code=%s: %s wxid=%s",
getattr(e, "code", "UNKNOWN"), e, req.stranger_wxid,
)
except Exception as e:
self._breaker.record_failure()
logger.error(
"[自动通过][步骤6/8] 异常: %s wxid=%s",
e, req.stranger_wxid,
)