76 KiB
UIActionScheduler 统一调度方案
文档版本:v1.0 创建日期:2026-07-17 范围:
bridge/woc_bridge全链路 目标:将散落在 5 个子系统的 UI 操作收口到统一调度器,消除绕过队列的竞态风险,引入优先级与可观测性
一、背景与现状
1.1 当前架构
当前 UI 操作分散在 5 个子系统,串行化覆盖不完整:
| 子系统 | 入口 | 是否经 SendQueue | 互斥情况 |
|---|---|---|---|
FlowOrchestrator.send_text/send_file |
HTTP + BatchWorker | ✅ 入队 | SendQueue 串行 + _ui_lock 单命令互斥 |
routes/contacts.py (set_remark/add_friend/accept) |
HTTP | ✅ 入队 | 同上 |
routes/moments.py (publish/share/like/comment/delete) |
HTTP | ✅ 入队 | 同上 |
routes/send.py legacy (revoke/forward) |
HTTP | ✅ 入队 | 同上 |
FriendRequestWatcher._handle_request |
后台轮询 3s | ✅ 入队 | 同上 |
routes/login.py (logout/restart) |
HTTP | ❌ 直接调 xdotool | 仅 _ui_lock 单命令互斥 |
routes/login.py (qr_start/qr_wait) |
HTTP | ❌ 直接调 xdotool | qr_capture 无锁 |
routes/screenshot.py |
HTTP | ❌ 直接调 qr_capture | 无锁 |
routes/diagnostic.py (autofix) |
HTTP | ❌ 内联 pkill | 无锁 |
LoginGuard._check_loop |
后台轮询 5s | ❌ 直接调 xdotool | 仅 _ui_lock |
WeChatWatchdog._run |
后台轮询 10s | ❌ 直接调 pgrep/kill | 无锁 |
1.2 已确认的问题
P0 竞态风险(已发生过实际故障)
| 问题 | 根因 | 后果 | 引用 |
|---|---|---|---|
| friend_watcher 与 send_text 竞态 | 早期未入队,现已修复入队 | erratic clicking、重复拨号 | 项目记忆记录 |
| logout 多步操作被 send_queue 任务穿插 | logout 不入队,仅靠 _ui_lock 单命令互斥 |
菜单状态错乱、登出失败 | routes/login.py:201 |
| restart kill 微信与并发 UI 操作冲突 | restart 不入队 | xdotool 对空窗口操作、TimeoutError 雪崩 | routes/login.py:281 |
| diagnostic autofix pkill 无锁 | 直接 pkill -x wechat |
与并发 UI 操作冲突高风险 | routes/diagnostic.py:232-239 |
| QrCapture 无锁无超时 | proc.communicate() 无 timeout |
与 send_text 截图竞争 X11、可能无限阻塞 | ui/qr_capture.py:48-57 |
| WeChatWatchdog autofix 无锁 | pgrep+kill 不持 _ui_lock |
kill 微信时 UI 操作中途失败 | ui/watchdog.py:99-191 |
P1 设计债务
| 问题 | 根因 | 后果 |
|---|---|---|
| 无优先级机制 | SendQueue 是纯 FIFO | 紧急操作(logout/restart)无法插队,100 个 send 排队时 logout 延迟分钟级 |
| 命名与实现不符 | 类名 SendQueue,实际是通用 UI 队列 |
误导维护者认为只管 send |
路由层前置 detect_login_state 不入队 |
多处路由在 enqueue 前直接调 xdotool | 与 send_queue 内任务竞争 _ui_lock,增加排队延迟 |
| 启动清场不入队 | _full_cleanup_on_startup 直接调 |
若 friend_watcher 已启动可能并发(实际启动顺序规避了此风险) |
| 无统一可观测性 | 各子系统独立打日志,无统一 metrics | UI 操作队列深度、等待时长、执行时长无法监控 |
P1.1 detect_login_state 散落调用分析
detect_login_state 是只读探测(截图 + 模板匹配),但散落在 14 处直接调用,每处都会与 send_queue 内任务竞争 L1 _ui_lock:
| 文件 | 调用位置 | 调用频率 | 与队列竞争影响 |
|---|---|---|---|
routes/send.py |
L307 / L413 / L587 / L690 / L774 / L852(6 处) | 每次 HTTP 请求 | 每次持锁 ~200-500ms(screenshot+模板匹配),100 并发 send 时累计 12-30s 额外延迟 |
routes/moments.py |
L151 / L255 / L346 / L447 / L545 / L677(6 处) | 每次 HTTP 请求 | 同上 |
routes/contacts.py |
L236 / L313 / L433(3 处) | 每次 HTTP 请求 | 同上 |
routes/login.py |
L117 / L190(2 处) | 每次 HTTP 请求 | 同上 |
routes/diagnostic.py |
L114(1 处) | 诊断触发 | 频率低,影响小 |
routes/status.py |
L51-56(内联判定,未直接调) | 每次 status 请求 | 已优化为内联 pgrep/find_window,不持 _ui_lock |
LoginGuard._check_loop |
后台 5s 轮询 | 持续运行 | 每 5s 一次,频率低 |
friend_watcher._verify_with_retry |
后台触发 | 偶发 | 频率低 |
关键发现:routes/status.py 已经通过内联 pgrep+find_window 避开了 detect_login_state,证明这种"只读探测不持 _ui_lock"的优化路径是可行且被项目采纳的。
处理方案(本方案不强制收口,但记录优化路径):
- 短期(本方案不实施):保持现状。
detect_login_state是单次截图+模板匹配,持锁时间可控(< 1s)。在 send_queue 任务执行期间(通常 1-5s),最多 1-2 次 detect_login_state 抢占 _ui_lock,影响有限。 - 中期(独立优化项):参照
routes/status.py的优化模式,将路由前置detect_login_state改为内联 pgrep + find_window 判定(不持 _ui_lock),仅在判定为 "logged_in" 后才入队执行真正的 UI 操作。 - 长期(架构演进):LoginGuard 维护登录态缓存,路由层读缓存而非每次调
detect_login_state。
为什么 UIActionScheduler 不收口 detect_login_state:
- 入队会被 100 个 send 排队阻塞,导致登录态检测延迟分钟级,影响 friend_watcher 等依赖登录态的后台任务
- detect_login_state 是只读探测,不修改 UI 状态,与 send_queue 任务的"互斥"是性能问题而非正确性问题
- 收口 detect_login_state 会引入"队列内任务触发队列外只读探测"的循环依赖
风险评估:保持现状的代价是每个 send 请求多 ~300ms 延迟(detect_login_state 持锁 1 次),在 100 并发场景下累计 ~30s。可通过中期方案消除。
1.3 SendQueue 现状评估
SendQueue 实际上已经是事实上的通用 UI 调度器:
- 16 处
enqueue调用点,覆盖 send(4)/contacts(3)/moments(6)/friend_watcher(1)/orchestrator(2) - 接收任意
CoroFactory,不限定 send 语义 - 单 worker 串行 + 1 秒滑窗限流 + 队列满拒绝 + 等待超时
可直接复用为基础,但需扩展:
- 增加
priority参数(当前队列元素是tuple[CoroFactory, Future, Optional[int]],第三段是custom_delay_ms,无 priority 槽位) - 重命名以反映通用语义(向后兼容保留别名)
- 收口绕过队列的调用点
1.4 双层保护模型(L1 _ui_lock + L2 SendQueue)
当前架构存在两层互斥机制,理解其分工是设计 UIActionScheduler 的前提:
| 层级 | 互斥粒度 | 实现位置 | 持锁时长 | 保护语义 |
|---|---|---|---|---|
L1 _ui_lock |
单命令(一次 click/type/screenshot) | ui/backends/xdotool.py:50 asyncio.Lock |
单次 xdotool 子进程(~50ms-5s) | 防止两个 asyncio 协程同时调 xdotool 导致 X11 焦点错乱 |
L2 SendQueue |
多命令 Flow(整个 send_text/logout 流程) | messaging/send_queue.py worker 串行 |
整个 Flow(~1-30s) | 防止多步操作之间被其他操作步骤穿插 |
关键区别:
send_text Flow(多步操作,L2 保护):
activate → click_search_box → type_query → open_session →
focus_input → type_text → click_send → verify
↑ 每一步内部都 acquire/release L1 _ui_lock(单命令互斥)
↑ 整个 Flow 由 L2 SendQueue 串行执行(多步不被穿插)
logout Flow(多步操作,当前仅 L1 保护):
activate → click_main_menu → key_down × N → enter → click_confirm
↑ 每一步内部 acquire/release L1 _ui_lock
↑ 整个 Flow 不在 SendQueue 中 → send_text 的步骤可以穿插进来!
核心问题:L1 只保护单次命令,无法防止多步 Flow 被穿插。例如 logout 点击主菜单后释放 _ui_lock,send_text 立即获得 _ui_lock 执行 click_search_box,导致 logout 的下一步 key_down 落在错误的焦点上。
UIActionScheduler 的职责:提供 L2 层统一保护,所有多步 UI 操作必须入队,避免步骤穿插。L1 保持不变作为底层单命令互斥的第二道防线(防止绕过队列的极端情况,如 LoginGuard/detect_login_state)。
1.5 UI 操作调用点全景图
修改前(当前状态,散落 22 个调用点,仅 16 个入队):
┌─ HTTP 路由层 ────────────────────────────────────────────────────────────┐
│ │
│ routes/send.py │
│ ├─ send_text (L337) ── enqueue ──→ SendQueue [✅ 入队] │
│ ├─ send_file (L703) ── enqueue ──→ SendQueue [✅ 入队] │
│ ├─ revoke_message (L787) ── enqueue ──→ SendQueue [✅ 入队] │
│ ├─ forward_message (L881) ── enqueue ──→ SendQueue [✅ 入队] │
│ ├─ detect_login_state × 6 ── direct ──→ xdotool [⚠️ 仅 L1] │
│ ├─ post_verify_revoke ── direct ──→ DB only [✅ 无 UI] │
│ ├─ capture_forward_baseline ── direct ──→ DB only [✅ 无 UI] │
│ └─ post_verify_forward ── direct ──→ DB only [✅ 无 UI] │
│ │
│ routes/contacts.py │
│ ├─ set_remark (L249) ── enqueue ──→ SendQueue [✅ 入队] │
│ ├─ add_friend (L321) ── enqueue ──→ SendQueue [✅ 入队] │
│ ├─ accept_friend (L452) ── enqueue ──→ SendQueue [✅ 入队] │
│ ├─ detect_login_state × 3 ── direct ──→ xdotool [⚠️ 仅 L1] │
│ └─ post_verify_set_remark ── direct ──→ DB only [✅ 无 UI] │
│ post_verify_add_friend │
│ │
│ routes/moments.py │
│ ├─ publish (L167) ── enqueue ──→ SendQueue [✅ 入队] │
│ ├─ share (L271) ── enqueue ──→ SendQueue [✅ 入队] │
│ ├─ like (L379) ── enqueue ──→ SendQueue [✅ 入队] │
│ │ └─ _like_and_verify (含 post_verify_moment_like 含截图) │
│ ├─ comment (L480) ── enqueue ──→ SendQueue [✅ 入队] │
│ │ └─ _comment_and_verify (含 post_verify_moment_comment 含截图) │
│ ├─ delete (L583) ── enqueue ──→ SendQueue [✅ 入队] │
│ ├─ forward_moment (L692) ── enqueue ──→ SendQueue [✅ 入队] │
│ └─ detect_login_state × 6 ── direct ──→ xdotool [⚠️ 仅 L1] │
│ │
│ routes/login.py │
│ ├─ qr_start (L58-59) ── direct ──→ xdotool+qr [❌ P0 无锁] │
│ ├─ qr_wait (L117) ── direct ──→ xdotool [⚠️ 仅 L1] │
│ ├─ logout (L201) ── direct ──→ xdotool [❌ P0 仅 L1] │
│ └─ wechat_restart (L281) ── direct ──→ xdotool [❌ P0 仅 L1] │
│ │
│ routes/screenshot.py │
│ └─ screenshot (L42) ── direct ──→ qr_capture [❌ P0 无锁] │
│ │
│ routes/diagnostic.py │
│ ├─ autofix_wechat_running ── direct ──→ pkill [❌ P0 无锁] │
│ └─ detect_login_state (L114)── direct ──→ xdotool [⚠️ 仅 L1] │
│ │
│ routes/status.py │
│ └─ 内联 pgrep+find_window (不持 _ui_lock,已优化) [✅ 无 UI] │
│ │
└──────────────────────────────────────────────────────────────────────────┘
┌─ 后台任务层 ────────────────────────────────────────────────────────────┐
│ │
│ FlowOrchestrator (send_text/send_file) │
│ ├─ enqueue (L231, L364) ── enqueue ──→ SendQueue [✅ 入队] │
│ └─ SessionCache LRU │
│ │
│ FriendRequestWatcher._handle_request │
│ └─ enqueue (L506) ── enqueue ──→ SendQueue [✅ 入队] │
│ │
│ LoginGuard._check_loop │
│ └─ detect_login_state ── direct ──→ xdotool [⚠️ 仅 L1] │
│ │
│ WeChatWatchdog._run │
│ ├─ pgrep/xdpyinfo ── direct ──→ system [✅ 只读] │
│ └─ _autofix_wechat (kill) ── direct ──→ kill [❌ P0 无锁] │
│ │
│ QrCapture │
│ └─ _run (proc.communicate) ── direct ──→ scrot [❌ P0 无超时] │
│ │
└──────────────────────────────────────────────────────────────────────────┘
修改后(目标状态,所有 P0 收口,仅保留只读探测不入队):
┌─ HTTP 路由层 ── 全部经 UIActionScheduler.execute(action, priority) ────┐
│ │
│ routes/send.py │
│ ├─ send_text/file/revoke/forward ── execute(NORMAL) ──→ Scheduler │
│ ├─ detect_login_state × 6 ── direct (中期优化为内联 pgrep) │
│ └─ post_verify_* / capture_baseline ── direct (DB only, 保留) │
│ │
│ routes/contacts.py │
│ ├─ set_remark/add_friend/accept ── execute(NORMAL/HIGH) → Scheduler │
│ └─ detect_login_state × 3 ── direct (中期优化) │
│ │
│ routes/moments.py │
│ ├─ publish/share/delete ── execute(NORMAL) ──→ Scheduler │
│ ├─ like/comment ── execute(LOW) ──→ Scheduler │
│ └─ detect_login_state × 6 ── direct (中期优化) │
│ │
│ routes/login.py │
│ ├─ qr_start (capture_qr_code) ── execute(HIGH) ──→ Scheduler │
│ ├─ logout ── execute(HIGH) ──→ Scheduler │
│ ├─ wechat_restart ── execute(CRITICAL)──→ Scheduler │
│ │ └─ mark_wechat_dead(8.0) before execute │
│ └─ qr_wait (detect_login_state) ── direct (只读,保留) │
│ │
│ routes/screenshot.py │
│ └─ screenshot ── execute(LOW) ──→ Scheduler │
│ │
│ routes/diagnostic.py │
│ └─ autofix_wechat_running ── execute(CRITICAL)──→ Scheduler │
│ └─ mark_wechat_dead(15.0) before execute │
│ │
└──────────────────────────────────────────────────────────────────────────┘
┌─ 后台任务层 ────────────────────────────────────────────────────────────┐
│ │
│ FlowOrchestrator.send_text/send_file │
│ └─ enqueue (向后兼容) ── execute(NORMAL) ──→ Scheduler │
│ │
│ FriendRequestWatcher._handle_request │
│ └─ enqueue (向后兼容) ── execute(HIGH) ──→ Scheduler │
│ │
│ LoginGuard._check_loop │
│ └─ detect_login_state ── direct (只读,保留) │
│ │
│ WeChatWatchdog._autofix_wechat │
│ ├─ scheduler.drain(5.0) ── 等待队列清空 │
│ ├─ scheduler.mark_wechat_dead(15.0) ── 标记 fast-fail 窗口 │
│ └─ pgrep + kill -TERM ── 系统级操作(不入队) │
│ │
│ QrCapture._run │
│ └─ asyncio.wait_for(proc.communicate, 5.0) ── 超时保护 │
│ │
└──────────────────────────────────────────────────────────────────────────┘
1.6 HTTP 端点 → 优先级映射表
| HTTP 端点 | 当前状态 | 目标优先级 | delay_ms 策略 | wait_timeout_ms | 备注 |
|---|---|---|---|---|---|
POST /api/send/text |
✅ enqueue | NORMAL | 同联系人 1000 / 不同 3000 | 15000 | 现有逻辑不变 |
POST /api/send/file |
✅ enqueue | NORMAL | 默认 3000 | 15000 | 现有逻辑不变 |
POST /api/messages/revoke |
✅ enqueue | NORMAL | 默认 3000 | 15000 | 现有逻辑不变 |
POST /api/messages/forward |
✅ enqueue | NORMAL | 默认 3000 | 15000 | 现有逻辑不变 |
POST /api/friends/remark |
✅ enqueue | NORMAL | 默认 3000 | 15000 | 现有逻辑不变 |
POST /api/friends/add |
✅ enqueue | NORMAL | 默认 3000 | 15000 | 现有逻辑不变 |
POST /api/friends/accept |
✅ enqueue | HIGH | 默认 3000 | 30000 | 时效性高 |
POST /api/moments/publish |
✅ enqueue | NORMAL | 默认 3000 | 30000 | 现有逻辑不变 |
POST /api/moments/share |
✅ enqueue | NORMAL | 默认 3000 | 30000 | 现有逻辑不变 |
POST /api/moments/delete |
✅ enqueue | NORMAL | 默认 3000 | 15000 | 现有逻辑不变 |
POST /api/moments/forward |
✅ enqueue | NORMAL | 默认 3000 | 30000 | 现有逻辑不变 |
POST /api/moments/like |
✅ enqueue | LOW | 默认 3000 | 15000 | 可延迟 |
POST /api/moments/comment |
✅ enqueue | LOW | 默认 3000 | 15000 | 可延迟 |
POST /api/screenshot |
❌ direct | LOW | 默认 3000 | 10000 | 阶段 2 收口 |
POST /api/login/qr/start |
❌ direct | HIGH | 默认 3000 | 15000 | 阶段 2 收口 |
POST /api/login/logout |
❌ direct | HIGH | 默认 3000 | 30000 | 阶段 2 收口 |
POST /api/wechat/restart |
❌ direct | CRITICAL | 0(自动) | 60000 | 阶段 2 收口 + fast-fail |
POST /api/diagnostic/autofix/wechat_running |
❌ direct | CRITICAL | 0(自动) | 30000 | 阶段 2 收口 + fast-fail |
GET /api/login/qr/wait |
direct(轮询) | 不入队 | - | - | 只读探测 |
GET /api/status |
内联 pgrep | 不入队 | - | - | 已优化 |
GET /api/diagnostic/run/* |
direct(只读) | 不入队 | - | - | 只读探测 |
二、设计方案
2.1 核心思路
不推翻 SendQueue 重写,而是扩展 + 收口:
- 扩展 SendQueue 为
UIActionScheduler,增加优先级参数 - 收口绕过队列的 6 个调用点,统一入队
- 保留
_ui_lock作为底层单命令互斥(不取消,作为第二道防线) - 新增 metrics 暴露队列深度、等待时长、执行时长
2.2 命名策略
采用渐进式重命名,避免一次性破坏全部调用点:
# messaging/ui_action_scheduler.py(新文件,继承 SendQueue)
class UIActionScheduler(SendQueue):
"""UI 操作统一调度器。
在 SendQueue 基础上增加:
- priority 参数(CRITICAL/HIGH/NORMAL/LOW)
- 优先级队列(asyncio.PriorityQueue 替代 asyncio.Queue)
- metrics 暴露(队列深度、等待时长、执行时长)
"""
async def execute(
self,
action: CoroFactory,
priority: Priority = Priority.NORMAL,
delay_ms: Optional[int] = None,
wait_timeout_ms: Optional[int] = None,
trace_id: str = "",
) -> Any:
"""提交 UI 操作,按优先级调度,串行执行。"""
...
# config.py
_state.send_queue: SendQueue | UIActionScheduler # 渐进式,先保持字段名
向后兼容:
- 保留
SendQueue.enqueue方法签名不变(priority 默认 NORMAL) - 16 个现有
enqueue调用点无需修改 - 新调用点用
execute方法,语义更清晰
2.3 优先级策略
class Priority(enum.IntEnum):
"""UI 操作优先级(数字越小优先级越高)。"""
CRITICAL = 0 # 系统级紧急:wechat_restart、diagnostic_autofix
HIGH = 1 # 用户感知延迟:logout、friend_accept、login_qr_capture
NORMAL = 2 # 常规业务:send_text、send_file、set_remark、add_friend
LOW = 3 # 可延迟:screenshot、moments_like、moments_comment
优先级队列实现:
- 使用
asyncio.PriorityQueue,元素为(priority, seq, coro_factory, future, delay_ms) seq是单调递增序列号,保证同优先级 FIFO(避免 coro_factory 不可比较导致的错误)- worker 出队时按
(priority, seq)排序
优先级反转保护:
- 低优先级操作持锁时,高优先级操作必须等待(无法抢占,xdotool 子进程不可中断)
- 缓解:限制单次操作超时(已有
_Cmd_TIMEOUT_SEC=5.0+ Flow 30s 超时),避免低优先级操作长时间持锁 - 不实现优先级继承(xdotool 子进程无法感知调用方优先级)
2.3.1 delay_ms 与优先级的交互策略
当前 delay_ms 决策逻辑(在 orchestrator.py:212):
- 同联系人发送:
delay_ms=1000(短延时,提升吞吐) - 不同联系人发送:
delay_ms=None→ 用send_delay_ms=3000(默认延时,避免风控) - 路由层
enqueue不传delay_ms→ 用默认send_delay_ms=3000
问题:引入优先级后,CRITICAL/HIGH 操作的延时是否应区别处理?
分析:
- delay_ms 是任务执行后的等待时间,影响下一个任务的开始时机
- 例如:CRITICAL 的 restart 执行后,若 delay_ms=3000,下一个 NORMAL 任务要等 3s
- restart 本身的优先级已经让它排到队首,所以 delay_ms 不影响 restart 自身延迟
策略:
| 优先级 | delay_ms 策略 | 理由 |
|---|---|---|
| CRITICAL | delay_ms=0(不延时) |
restart/autofix 后应立即让后续任务执行,避免无谓等待 |
| HIGH | delay_ms=None(用默认) |
logout/accept_friend 后需要给微信 UI 恢复时间,保持默认延时 |
| NORMAL | 保持现有逻辑(同联系人 1000 / 不同 3000) | send 风控保护,不变 |
| LOW | delay_ms=None(用默认) |
screenshot 后不强制延时,但也不加速,避免连续截图 |
实现:在 scheduler.execute 内根据 priority 自动覆盖 delay_ms:
async def execute(
self,
action: CoroFactory,
priority: Priority = Priority.NORMAL,
delay_ms: Optional[int] = None,
...
) -> Any:
# CRITICAL 任务自动应用 delay_ms=0(除非调用方显式指定)
if priority == Priority.CRITICAL and delay_ms is None:
delay_ms = 0
return await self._enqueue_with_priority(
action, priority, delay_ms, wait_timeout_ms, trace_id
)
注意:调用方显式指定 delay_ms 时优先尊重调用方意图(如批量发送要求固定间隔)。
对现有 enqueue 的影响:无。现有 16 个 enqueue 调用点都走 priority=NORMAL,delay_ms 决策逻辑不变。
2.4 关键决策:watchdog 是否入队
不入队,但加协同机制:
WeChatWatchdog 的 pgrep -x wechat / xdpyinfo / kill -TERM 是系统级操作,不是 UI 操作:
pgrep/xdpyinfo:只读探测,不修改 UI 状态,无需互斥kill -TERM:会杀死微信进程,导致所有进行中的 UI 操作失败
方案:
- watchdog 的
pgrep/xdpyinfo保持现状(不入队,不加锁) - watchdog 的
_autofix_wechat在 kill 前检查scheduler.pending_count():- 若有 pending UI 操作,先
await scheduler.drain(timeout=5.0)等待队列清空 - 超时未清空则强制 kill(记 warning)
- 若有 pending UI 操作,先
- kill 后广播
wechat_killed事件,所有进行中的 Flow 捕获WINDOW_NOT_FOUND后自然失败
2.5 关键决策:QrCapture 如何收口
QrCapture 有两类操作:
capture_qr_code:启动扫码登录流程,多步 UI 操作(activate + screenshot + crop)capture_full_screenshot:单次截图
方案:
capture_qr_code入队(priority=HIGH),与 send_text 串行capture_full_screenshot入队(priority=LOW)- QrCapture 内部增加
_CMD_TIMEOUT_SEC=5.0包裹proc.communicate(),消除无超时风险
2.6 关键决策:LoginGuard 是否入队
不入队,原因:
detect_login_state是只读探测(截图 + 模板匹配),不修改 UI 状态- 5 秒轮询 +
_ui_lock单命令互斥已足够 - 若入队,会被 100 个 send 排队阻塞,登录态检测延迟可达分钟级,影响 friend_watcher 等依赖登录态的任务
但需加保护:
detect_login_state内部已有_ui_lock持锁(经_run),保持现状- 增加监控:若
_ui_lock等待时长 > 2s,记 warning(说明 UI 操作积压)
三、详细设计
3.1 UIActionScheduler 类
# bridge/woc_bridge/messaging/ui_action_scheduler.py
"""UI 操作统一调度器。
所有 xdotool/scrot 子进程调用经此调度器串行执行,避免并发 UI 操作冲突。
在 SendQueue 基础上增加优先级调度与可观测性。
"""
from __future__ import annotations
import asyncio
import enum
import logging
import time
from typing import Any, Awaitable, Callable, Optional
from woc_bridge.messaging.send_queue import SendQueue, CoroFactory
from woc_bridge.models import BridgeError
logger = logging.getLogger("woc-bridge")
class Priority(enum.IntEnum):
"""UI 操作优先级(数字越小优先级越高)。"""
CRITICAL = 0 # 系统级紧急:wechat_restart、diagnostic_autofix
HIGH = 1 # 用户感知延迟:logout、friend_accept、login_qr_capture
NORMAL = 2 # 常规业务:send_text、send_file、set_remark、add_friend
LOW = 3 # 可延迟:screenshot、moments_like、moments_comment
class UIActionScheduler(SendQueue):
"""UI 操作统一调度器。
扩展 SendQueue:
- execute() 方法支持 priority 参数
- 内部用 asyncio.PriorityQueue 替代 asyncio.Queue(覆盖父类 __init__ 创建的 Queue)
- 同优先级 FIFO(通过 seq 序列号保证,避免比较到不可比较的 coro_factory)
- metrics 暴露(pending_count by priority、wait_duration、exec_duration)
向后兼容:
- enqueue() 方法保留,priority 默认 NORMAL
- 现有 16 个 enqueue 调用点无需修改
⚠️ 必须覆盖父类 _run 方法:父类 _run 解包 3 元组 (coro_factory, future, delay_ms),
子类用 5 元组 (priority, seq, coro_factory, future, delay_ms)。
若不覆盖会导致解包失败。
"""
def __init__(
self,
send_delay_ms: int = 3000,
max_calls_per_sec: int = 10,
max_queue_size: int = 100,
) -> None:
super().__init__(send_delay_ms, max_calls_per_sec, max_queue_size)
# 覆盖父类 __init__ 创建的 asyncio.Queue 为 PriorityQueue
# 父类的旧 Queue 对象会被 GC 回收(无其他引用)
self._queue: asyncio.PriorityQueue[
tuple[int, int, CoroFactory, asyncio.Future, Optional[int]]
] = asyncio.PriorityQueue(maxsize=max_queue_size)
self._seq = 0 # 单调递增序列号,保证同优先级 FIFO
# fast-fail 机制:mark_wechat_dead 标记的时间戳
# 在此时间之前,所有非 CRITICAL 任务直接抛 WECHAT_NOT_READY
# 供 watchdog kill / wechat_restart 调用,避免失败风暴
self._wechat_dead_until: float = 0.0
# metrics(简单计数器,供 /api/status 或 Prometheus 暴露)
# asyncio 单线程事件循环,worker 与 HTTP handler 同线程,无需加锁
self._metrics: dict[str, Any] = {
"executed_total": 0,
"executed_by_priority": {p.name: 0 for p in Priority},
"wait_duration_ms_sum": 0.0,
"exec_duration_ms_sum": 0.0,
"timeout_total": 0,
"rate_limited_total": 0,
"fast_fail_total": 0, # 被 mark_wechat_dead 拦截的任务数
}
def mark_wechat_dead(self, duration_sec: float = 10.0) -> None:
"""标记微信已死,期间所有非 CRITICAL 任务直接 fast-fail。
供 watchdog kill / wechat_restart / diagnostic autofix 调用:
- kill 前调 mark_wechat_dead(15.0),覆盖 autostart 拉起窗口
- restart 开始时调 mark_wechat_dead(8.0),覆盖新进程启动窗口
CRITICAL 任务(如 restart 本身)不受影响,确保 restart 能正常执行。
Args:
duration_sec: fast-fail 窗口时长(秒)
"""
self._wechat_dead_until = time.monotonic() + duration_sec
logger.warning(
"[scheduler] mark_wechat_dead %.1fs, %d pending tasks will fast-fail",
duration_sec, self._queue.qsize(),
)
async def execute(
self,
action: CoroFactory,
priority: Priority = Priority.NORMAL,
delay_ms: Optional[int] = None,
wait_timeout_ms: Optional[int] = None,
trace_id: str = "",
) -> Any:
"""提交 UI 操作,按优先级调度,串行执行。
Args:
action: 返回 coroutine 的工厂函数
priority: 优先级(CRITICAL/HIGH/NORMAL/LOW)
delay_ms: 自定义本次延时毫秒;None 时:
- CRITICAL 自动设为 0(执行后不延时)
- 其他优先级用默认 send_delay_ms
wait_timeout_ms: 队列等待超时;None 无限等待
trace_id: 追踪 ID(用于日志关联)
Returns:
coroutine 的实际执行结果
Raises:
BridgeError(RATE_LIMITED): 队列满
BridgeError(TIMEOUT): 等待超时
"""
# CRITICAL 任务自动应用 delay_ms=0(除非调用方显式指定)
# 参见 2.3.1 节 delay_ms 与优先级交互策略
if priority == Priority.CRITICAL and delay_ms is None:
delay_ms = 0
return await self._enqueue_with_priority(
action, priority, delay_ms, wait_timeout_ms, trace_id
)
async def _enqueue_with_priority(
self,
coro_factory: CoroFactory,
priority: Priority,
delay_ms: Optional[int],
wait_timeout_ms: Optional[int],
trace_id: str,
) -> Any:
"""优先级入队(覆盖 SendQueue.enqueue 的内部实现)。
与父类 enqueue 的差异:
- 用 put_nowait 替代 await put(满队时立即抛错,语义更清晰)
- 队列元素从 3 元组扩展为 5 元组(priority, seq, coro_factory, future, delay_ms)
- 增加入队日志(与父类保持一致的可观测性)
"""
loop = asyncio.get_running_loop()
future: asyncio.Future = loop.create_future()
self._seq += 1
item = (int(priority), self._seq, coro_factory, future, delay_ms)
# 入队(带队列满检查,与父类行为一致)
try:
self._queue.put_nowait(item)
except asyncio.QueueFull:
self._metrics["rate_limited_total"] += 1
logger.warning(
"[scheduler] 队列已满 (size=%d/%d),拒绝入队 priority=%s",
self._queue.maxsize, self._queue.maxsize, priority.name,
)
raise BridgeError(
code="RATE_LIMITED",
message=f"UI 调度器队列已满({self._queue.maxsize}),请稍后重试",
details={"retry_after": 3},
)
pending = self._queue.qsize()
logger.info(
"[scheduler] 入队 priority=%s pending=%d trace_id=%s",
priority.name, pending, trace_id or "(none)",
)
# 等待结果(与父类逻辑一致,保留竞争窗口处理)
wait_start = time.monotonic()
if wait_timeout_ms is not None and wait_timeout_ms > 0:
try:
result = await asyncio.wait_for(
future, timeout=wait_timeout_ms / 1000.0
)
except asyncio.TimeoutError:
# 竞争窗口:worker 可能刚好在此时完成并 set_result
if future.done() and not future.cancelled():
logger.info(
"[scheduler] 等待超时但任务刚好完成,取结果 (pending=%d)",
pending,
)
return future.result()
future.cancel()
self._metrics["timeout_total"] += 1
wait_sec = wait_timeout_ms / 1000.0
logger.warning(
"[scheduler] 等待超时 (pending=%d, wait_timeout=%.1fs, priority=%s)",
pending, wait_sec, priority.name,
)
raise BridgeError(
code="TIMEOUT",
message=f"UI 调度器等待超时({pending} 个待处理,已等 {wait_sec:.1f}s)",
details={"retry_after": max(1, int(self.send_delay_ms / 1000))},
)
else:
result = await future
# 记录等待时长
wait_ms = (time.monotonic() - wait_start) * 1000
self._metrics["wait_duration_ms_sum"] += wait_ms
return result
async def enqueue(
self,
coro_factory: CoroFactory,
delay_ms: Optional[int] = None,
wait_timeout_ms: Optional[int] = None,
) -> Any:
"""向后兼容的入队方法(priority 默认 NORMAL)。
现有 16 个 enqueue 调用点无需修改,行为与父类一致。
"""
return await self._enqueue_with_priority(
coro_factory, Priority.NORMAL, delay_ms, wait_timeout_ms, ""
)
async def _run(self) -> None:
"""worker 主循环。
⚠️ 必须覆盖父类 _run:父类解包 3 元组,子类用 5 元组。
与父类 _run 的差异:
- 解包 5 元组 (priority, seq, coro_factory, future, custom_delay_ms)
- 增加 metrics 记录(executed_total / executed_by_priority / exec_duration)
- 保留限流检查(_check_rate_limit 继承父类,不覆盖)
- 保留延时逻辑(delay_ms 优先于 send_delay_ms)
- task_done 调用与父类一致:cancelled 分支和 finally 分支互斥
"""
while True:
try:
priority, seq, coro_factory, future, custom_delay_ms = (
await self._queue.get()
)
except asyncio.CancelledError:
raise
# 调用方已超时取消:跳过执行与延时(与父类一致)
if future.cancelled():
self._queue.task_done()
logger.info(
"[scheduler] 出队任务已取消,跳过 (pending=%d, priority=%s)",
self._queue.qsize(),
Priority(priority).name if 0 <= priority <= 3 else "?",
)
continue
# fast-fail 检查:微信已死期间,非 CRITICAL 任务直接失败
# 参见 6.4 节 失败风暴问题与 fast-fail 机制
if (
self._wechat_dead_until > time.monotonic()
and priority > Priority.CRITICAL # CRITICAL 任务(=0)不受影响
):
logger.info(
"[scheduler] fast-fail (wechat dead, %.1fs remaining, priority=%s)",
self._wechat_dead_until - time.monotonic(),
Priority(priority).name if 0 <= priority <= 3 else "?",
)
if not future.done():
future.set_exception(BridgeError(
code="WECHAT_NOT_READY",
message="微信进程未就绪(重启中),请稍后重试",
details={"retry_after": 5},
))
self._metrics["fast_fail_total"] += 1
self._queue.task_done()
continue # 不延时,立即处理下一个
logger.info(
"[scheduler] 出队,开始处理 (pending=%d, priority=%s)",
self._queue.qsize(),
Priority(priority).name if 0 <= priority <= 3 else "?",
)
executed = False
t_exec = time.perf_counter()
try:
# 继承父类的限流检查
self._check_rate_limit()
self._recent_call_times.append(time.monotonic())
executed = True
logger.info("[scheduler] 开始执行任务 (priority=%s)", Priority(priority).name)
result = await coro_factory()
logger.info(
"[scheduler] 任务执行完成 (%.0fms, priority=%s)",
(time.perf_counter() - t_exec) * 1000,
Priority(priority).name,
)
if not future.done():
future.set_result(result)
except asyncio.CancelledError:
if not future.done():
future.cancel()
raise
except Exception as e:
logger.warning(
"[scheduler] 任务执行抛异常 %s: %s (%.0fms, priority=%s)",
type(e).__name__, e,
(time.perf_counter() - t_exec) * 1000,
Priority(priority).name,
)
if not future.done():
future.set_exception(e)
finally:
self._queue.task_done()
if executed:
# 更新 metrics
exec_ms = (time.perf_counter() - t_exec) * 1000
self._metrics["executed_total"] += 1
self._metrics["exec_duration_ms_sum"] += exec_ms
try:
p = Priority(priority)
self._metrics["executed_by_priority"][p.name] += 1
except ValueError:
pass
delay = (
custom_delay_ms
if custom_delay_ms is not None
else self.send_delay_ms
)
logger.info(
"[scheduler] 延时 %dms 后处理下一个 (priority=%s)",
delay, Priority(priority).name,
)
await asyncio.sleep(delay / 1000.0)
def get_metrics(self) -> dict:
"""返回调度器指标(供 /api/status 或 Prometheus 暴露)。
asyncio 单线程事件循环,与 worker 同线程,无需加锁。
"""
return {
**self._metrics,
"pending_count": self._queue.qsize(),
}
async def drain(self, timeout: float = 5.0) -> bool:
"""等待队列清空(供 watchdog kill 微信前调用)。
通过 _queue.join() 等待所有已入队任务被 task_done。
worker 在 task_done 后可能仍在 sleep(delay),但队列已空,
sleep 结束后阻塞在 get() 上,不会执行新任务。
边界条件:
- worker 已停止时 join 会永远等待(无人调 task_done),靠 timeout 兜底
- drain 返回 True 后、kill 前若有新任务入队,worker 可能开始执行
(概率低,kill 后执行失败被正常捕获)
Args:
timeout: 最大等待秒数
Returns:
True 表示队列已清空,False 表示超时未清空
"""
try:
await asyncio.wait_for(self._queue.join(), timeout=timeout)
return True
except asyncio.TimeoutError:
logger.warning(
"[scheduler] drain timeout %.1fs, %d tasks pending",
timeout, self._queue.qsize(),
)
return False
3.2 优先级分配
| 调用点 | 优先级 | 理由 |
|---|---|---|
routes/login.py wechat_restart |
CRITICAL | 系统级紧急,需尽快执行 |
routes/diagnostic.py autofix_wechat_running |
CRITICAL | 同上 |
routes/login.py logout |
HIGH | 用户主动操作,感知延迟 |
routes/login.py qr_start (capture_qr_code) |
HIGH | 登录流程关键步骤 |
routes/contacts.py accept_friend_request |
HIGH | 好友申请有时效性 |
friend_watcher._handle_request |
HIGH | 同上 |
routes/send.py send_text (flow + legacy) |
NORMAL | 常规业务 |
routes/send.py send_file (flow + legacy) |
NORMAL | 同上 |
routes/send.py revoke / forward |
NORMAL | 同上 |
routes/contacts.py set_remark / add_friend |
NORMAL | 同上 |
routes/moments.py publish / share / delete |
NORMAL | 同上 |
routes/moments.py like / comment |
LOW | 可延迟,不阻塞主流程 |
routes/screenshot.py capture_full_screenshot |
LOW | 可延迟 |
routes/login.py qr_wait (detect_login_state 轮询) |
不入队 | 只读探测,保持现状 |
LoginGuard._check_loop |
不入队 | 只读探测,保持现状 |
WeChatWatchdog._run (pgrep/xdpyinfo) |
不入队 | 系统级探测,保持现状 |
WeChatWatchdog._autofix_wechat (kill) |
不入队但协同 | kill 前调 scheduler.drain() |
3.3 收口方案(6 个 P0 调用点)
3.3.1 routes/login.py - logout
关键:logout 路由在执行 logout 前会先调 detect_login_state(只读探测,不入队)。仅 xdotool.logout() 入队。
# 修改前(routes/login.py:201)
if state_before == LoginState.LOGGED_IN.value:
await xdotool.logout() # 直接调,不入队
# 修改后
if state_before == LoginState.LOGGED_IN.value:
scheduler = _require_send_queue()
await scheduler.execute(
lambda: xdotool.logout(),
priority=Priority.HIGH,
wait_timeout_ms=30000, # logout 多步操作,给 30s
trace_id="logout",
)
说明:detect_login_state 保持不入队(只读探测,与 LoginGuard 一致),仅 xdotool.logout() 入队。
3.3.2 routes/login.py - wechat_restart
# 修改后(routes/login.py:281)
scheduler = _require_send_queue()
new_pid = await scheduler.execute(
lambda: xdotool.restart_wechat(timeout_sec=timeout_sec),
priority=Priority.CRITICAL,
wait_timeout_ms=60000, # restart 可能慢,给 60s
trace_id="wechat_restart",
)
3.3.3 routes/login.py - qr_start
# 修改后(routes/login.py:58-59)
scheduler = _require_send_queue()
async def _qr_capture_flow():
await xdotool._activate_window_fast()
return await qr_capture.capture_qr_code()
qr_data_url = await scheduler.execute(
_qr_capture_flow,
priority=Priority.HIGH,
wait_timeout_ms=15000,
trace_id="qr_start",
)
3.3.4 routes/screenshot.py
# 修改后(routes/screenshot.py:42)
scheduler = _require_send_queue()
png_bytes = await scheduler.execute(
lambda: qr_capture.capture_full_screenshot(),
priority=Priority.LOW,
wait_timeout_ms=10000,
trace_id="screenshot",
)
3.3.5 routes/diagnostic.py - autofix_wechat_running
# 修改后(routes/diagnostic.py:230-249)
scheduler = _require_send_queue()
async def _autofix_flow():
# pkill + check_pid + start_wechat 组合
proc = await asyncio.create_subprocess_exec(
"pkill", "-x", "wechat",
stdout=asyncio.subprocess.DEVNULL,
stderr=asyncio.subprocess.DEVNULL,
)
await proc.wait()
await asyncio.sleep(2) # 等 autostart 拉起
new_pid = await xdotool.check_wechat_pid()
if new_pid is None:
new_pid = await xdotool.start_wechat(timeout_sec=10)
return new_pid
new_pid = await scheduler.execute(
_autofix_flow,
priority=Priority.CRITICAL,
wait_timeout_ms=30000,
trace_id="autofix_wechat_running",
)
3.3.6 QrCapture 超时修复
# ui/qr_capture.py 修改(L48-57)
_CMD_TIMEOUT_SEC = 5.0 # 新增常量
async def _run(self, args: list[str]) -> tuple[int, bytes, bytes]:
"""执行一条命令并返回 (returncode, stdout, stderr)。"""
proc = await asyncio.create_subprocess_exec(
*args,
stdout=asyncio.subprocess.PIPE,
stderr=asyncio.subprocess.PIPE,
env=self._env(),
)
try:
stdout, stderr = await asyncio.wait_for(
proc.communicate(), timeout=_CMD_TIMEOUT_SEC
)
return proc.returncode, stdout, stderr
except asyncio.TimeoutError:
proc.kill()
await proc.wait()
raise
3.3.7 WeChatWatchdog 协同
# ui/watchdog.py 修改
class WeChatWatchdog:
def __init__(
self,
backend: BackendProtocol,
interval: float = 10.0,
fail_threshold: int = 2,
scheduler=None, # 新增可选参数,向后兼容
) -> None:
self.backend = backend
self.interval = interval
self.fail_threshold = fail_threshold
self._scheduler = scheduler # UIActionScheduler 实例(用于 drain)
self._fail_count = 0
self._task: Optional[asyncio.Task] = None
self._stopped = False
async def _autofix_wechat(self) -> None:
logger.warning("[watchdog] autofix: killing wechat for restart")
# kill 前等待 UI 操作队列清空(避免 kill 正在执行的 UI 操作)
if self._scheduler is not None:
drained = await self._scheduler.drain(timeout=5.0)
if not drained:
logger.warning(
"[watchdog] drain timeout, force kill with %d UI ops pending",
self._scheduler.pending_count(),
)
# 原有 kill 逻辑(pgrep + kill -TERM)保持不变
try:
proc = await asyncio.create_subprocess_exec(
"pgrep", "-x", "wechat",
stdout=asyncio.subprocess.PIPE,
stderr=asyncio.subprocess.PIPE,
)
# ... 原有逻辑
app.py 注入:
# app.py:499 修改
_state.watchdog = WeChatWatchdog(
backend=_state.xdotool_backend,
interval=10.0,
scheduler=_state.send_queue, # 新增:注入调度器
)
3.4 注入与生命周期
# app.py _init_state 修改(替换 L383-L387)
from woc_bridge.messaging.ui_action_scheduler import UIActionScheduler, Priority
_state.send_queue = UIActionScheduler( # 替换原 SendQueue
send_delay_ms=cfg.send_delay_ms,
max_calls_per_sec=cfg.max_calls_per_sec,
max_queue_size=cfg.max_queue_size,
)
# _state.send_queue 类型注解保持 SendQueue(多态,UIActionScheduler 是子类)
# 现有 _require_send_queue() 返回 SendQueue 类型,调用方无需修改
# 调用方可用 hasattr(scheduler, 'execute') 判断是否支持优先级
# watchdog 注入 scheduler(app.py:499 修改)
_state.watchdog = WeChatWatchdog(
backend=_state.xdotool_backend,
interval=10.0,
scheduler=_state.send_queue, # 新增:注入调度器用于 drain
)
生命周期不变:
send_queue.start()/send_queue.stop()继承父类,lifespan 中保持原顺序watchdog.start()/watchdog.stop()保持原顺序- shutdown 时
verify_bus.clear()在send_queue.stop()之前(已有逻辑)
3.5 Metrics 暴露
3.5.1 StatusResponse 模型扩展
# bridge/woc_bridge/models/status.py 修改
class StatusResponse(BaseModel):
# ... 现有字段保持不变 ...
send_queue_pending: int = Field(default=0, description="发送队列积压任务数")
# 新增字段(Optional,向后兼容)
ui_scheduler: Optional[dict] = Field(
default=None,
description="UI 调度器指标(仅 UIActionScheduler 实例才有)",
)
3.5.2 routes/status.py 暴露 metrics
# routes/status.py 修改(L96 附近)
send_queue = _state.send_queue
send_queue_pending = send_queue.pending_count() if send_queue else 0
ui_scheduler_metrics = None
if send_queue is not None and hasattr(send_queue, 'get_metrics'):
ui_scheduler_metrics = send_queue.get_metrics()
return StatusResponse(
# ... 现有字段 ...
send_queue_pending=send_queue_pending,
ui_scheduler=ui_scheduler_metrics,
# ...
)
返回示例:
{
"send_queue_pending": 0,
"ui_scheduler": {
"executed_total": 1234,
"executed_by_priority": {
"CRITICAL": 2,
"HIGH": 15,
"NORMAL": 1200,
"LOW": 17
},
"wait_duration_ms_sum": 45678.9,
"exec_duration_ms_sum": 234567.8,
"timeout_total": 3,
"rate_limited_total": 1,
"pending_count": 0
}
}
3.5.3 激活已存在的 Prometheus metrics
发现问题:ui/metrics.py:58 已定义 woc_send_queue_pending Gauge,但从未在任何地方 set 它的值(已存在的 bug)。
# ui/metrics.py 已有定义(无需修改)
send_queue_pending = Gauge(
"woc_send_queue_pending",
"Send queue pending count",
)
# 新增:在 routes/status.py 或 app.py 中激活
# 方案 A:在 status 接口中 set(每次查询时更新)
from woc_bridge.ui.metrics import send_queue_pending as send_queue_pending_gauge
if send_queue is not None:
send_queue_pending_gauge.set(send_queue.pending_count())
# 方案 B:在 UIActionScheduler._run 中 set(每次出队时更新)
# 更实时但会增加 metrics 写入频率
推荐方案 A:在 status 接口中 set,与现有 send_queue_pending 字段同步更新,避免高频写入 Prometheus。
3.5.4 新增 Prometheus metrics(可选,阶段 3)
# ui/metrics.py 新增
from prometheus_client import Counter, Gauge, Histogram
ui_action_executed = Counter(
"woc_ui_action_executed_total",
"UI actions executed total",
["priority"], # CRITICAL / HIGH / NORMAL / LOW
)
ui_action_wait_duration = Histogram(
"woc_ui_action_wait_duration_seconds",
"UI action wait duration (queue wait time)",
buckets=(0.1, 0.5, 1.0, 2.0, 5.0, 10.0, 30.0),
)
ui_action_exec_duration = Histogram(
"woc_ui_action_exec_duration_seconds",
"UI action execution duration",
["priority"],
buckets=(0.5, 1.0, 2.0, 5.0, 10.0, 30.0, 60.0),
)
ui_action_timeout = Counter(
"woc_ui_action_timeout_total",
"UI action wait timeout count",
)
ui_action_rate_limited = Counter(
"woc_ui_action_rate_limited_total",
"UI action rate limited (queue full) count",
)
在 UIActionScheduler._run 和 _enqueue_with_priority 中埋点。
四、迁移路径
阶段 1:基础设施(无破坏性)
目标:引入 UIActionScheduler,向后兼容现有调用
改动:
- 新增
messaging/ui_action_scheduler.py(Priority + UIActionScheduler 类) app.py将SendQueue(...)替换为UIActionScheduler(...)watchdog.py增加scheduler参数(可选,向后兼容)routes/status.py暴露 metrics
验证:
- py_compile 通过
- 现有 16 个
enqueue调用点无需修改,行为不变 /api/status返回ui_scheduler字段
阶段 2:收口 P0 调用点
目标:消除 6 个绕过队列的竞态风险 + QrCapture 超时修复
改动:
routes/login.py:logout / restart / qr_start 改用scheduler.execute- logout 仅
xdotool.logout()入队,detect_login_state保持不入队
- logout 仅
routes/screenshot.py:改用scheduler.execute(priority=LOW)routes/diagnostic.py:autofix 改用scheduler.execute(priority=CRITICAL)ui/qr_capture.py:增加_CMD_TIMEOUT_SEC=5.0包裹proc.communicate()ui/watchdog.py:构造函数新增scheduler参数;_autofix_wechat增加scheduler.drain()调用app.py:WeChatWatchdog(...)注入scheduler=_state.send_queue
验证:
- py_compile 通过
- 并发调用
POST /api/login/logout+POST /api/send/text不再竞态(logout 入队等待) POST /api/screenshot与 send_text 串行执行- watchdog autofix 前等待 UI 队列清空(
drain返回 True 后再 kill) - QrCapture scrot 卡死时 5s 超时(不再无限阻塞)
阶段 3:可观测性增强(可选)
目标:提供 Prometheus metrics 和详细日志
改动:
models/status.py:StatusResponse新增ui_scheduler: Optional[dict]字段routes/status.py:暴露ui_schedulermetrics + 激活已存在的woc_send_queue_pendingGaugeui/metrics.py:新增woc_ui_action_*系列 metrics(Counter/Histogram)UIActionScheduler._run/_enqueue_with_priority:埋点写入 Prometheus metrics- 日志增加
priority和trace_id字段(已在阶段 1 代码中包含)
验证:
GET /api/status返回ui_scheduler字段/metrics端点返回woc_ui_action_executed_total等指标- 日志可按
priority=CRITICAL过滤 - 已存在的
woc_send_queue_pendingGauge 不再为 0(bug 修复)
五、验收标准
AC-01:UIActionScheduler 向后兼容
Given 现有 16 个 enqueue 调用点未修改
When app.py 将 SendQueue(...) 替换为 UIActionScheduler(...)
Then 所有现有功能行为不变,py_compile 通过,单元测试通过
AC-02:优先级调度生效
Given 队列中有 10 个 NORMAL 优先级的 send_text 任务 When 提交一个 CRITICAL 优先级的 wechat_restart Then wechat_restart 优先于剩余 NORMAL 任务执行(可能在当前 NORMAL 任务执行完后立即执行)
AC-03:logout 与 send_text 不再竞态
Given 一个 send_text 正在 send_queue 中执行
When 并发调用 POST /api/login/logout
Then logout 入队等待,send_text 执行完后 logout 才开始执行
AC-04:QrCapture 有超时保护
Given scrot 子进程卡死
When QrCapture._run 执行
Then 5 秒后超时,proc.kill() + await proc.wait(),抛出 TimeoutError
AC-05:watchdog autofix 协同
Given send_queue 中有 5 个 pending 任务
When watchdog 触发 _autofix_wechat
Then 先 scheduler.drain(timeout=5.0) 等待队列清空,超时则强制 kill 并记 warning
AC-06:metrics 暴露
Given 调度器执行了若干操作
When 调用 GET /api/status
Then 返回 ui_scheduler 字段,包含 executed_total / executed_by_priority / pending_count 等指标
5.1 关键场景时序图
5.1.1 修改前:logout 与 send_text 并发竞态
时刻 T0: send_text 正在执行(队列 worker 持有 coro_factory)
┌─────────────────────────────────────────────────────┐
Worker: │ send_text Flow: │
│ activate → click_search → type_query → ... │
│ ↑ 每步持 L1 _ui_lock,步骤间释放 │
└─────────────────────────────────────────────────────┘
时刻 T1: HTTP 请求 POST /api/login/logout 到达
↓ logout 路由直接调 xdotool.logout()(不入队)
↓ logout 尝试 acquire _ui_lock → 等待 send_text 释放
时刻 T2: send_text 执行 click_search_box 完成释放 _ui_lock
↓ logout 抢到 _ui_lock,执行 activate
↓ logout 释放 _ui_lock
↓ send_text 抢到 _ui_lock,执行 type_query
↓ ❌ 焦点已被 logout 改变,type_query 输入到错误位置!
时刻 T3: logout 再次抢 _ui_lock 执行 click_main_menu
↓ ❌ send_text 可能仍在执行,菜单状态不可预测
结果:send_text 发错会话,logout 失败,UI 状态混乱
5.1.2 修改后:logout 与 send_text 串行执行
时刻 T0: send_text 正在执行(worker 持有)
┌─────────────────────────────────────────────────────┐
Worker: │ send_text Flow (NORMAL priority) │
│ activate → click_search → type_query → ... │
└─────────────────────────────────────────────────────┘
时刻 T1: HTTP 请求 POST /api/login/logout 到达
↓ logout 路由调 scheduler.execute(logout, HIGH, 30s)
↓ 入队 (HIGH, seq=N),pending=1
↓ worker 仍在执行 send_text,logout 等待 future
时刻 T2: send_text 完成,worker 出队
↓ PriorityQueue 出队 (HIGH, N)(优先于其他 NORMAL 任务)
↓ 执行 xdotool.logout() 完整流程(多步操作不被穿插)
↓ logout 完成,返回 success
结果:send_text 与 logout 串行,UI 状态正确
5.1.3 watchdog kill 与 drain 协同
时刻 T0: send_queue 中有 5 个 pending send_text
worker 正在执行第 1 个 send_text
时刻 T1: watchdog 检测到 wechat 不响应(fail_count >= 2)
↓ watchdog 调 _autofix_wechat
↓ scheduler.drain(timeout=5.0)
↓ 等待 _queue.join(),worker 继续执行当前任务
时刻 T2: worker 完成当前 send_text(task_done)
↓ _queue.join() 检查未完成任务数
↓ 仍有 4 个 pending(task_done 计数未归零)
↓ worker 出队下一个 send_text 并执行...
↓ ⚠️ 若每个 send_text 耗时 > 1.25s,5s 内无法清空
时刻 T3a (5s 内清空):
↓ drain 返回 True
↓ scheduler.mark_wechat_dead(15.0) ← 启动 fast-fail 窗口
↓ pgrep + kill -TERM wechat
↓ 后续 15s 内新入队任务直接 fast-fail(WECHAT_NOT_READY)
时刻 T3b (5s 超时未清空):
↓ drain 返回 False,记 warning
↓ scheduler.mark_wechat_dead(15.0)
↓ pgrep + kill -TERM wechat
↓ ❌ 仍有 N 个 pending 任务,kill 后每个都会失败一次
时刻 T4: autostart 拉起新 wechat 进程(~5-10s)
↓ mark_wechat_dead 窗口未到期,新任务继续 fast-fail
↓ 窗口到期后,新任务正常执行(wechat 已就绪)
结果:drain 减少失败任务数,mark_wechat_dead 避免失败风暴放大
5.1.4 优先级插队场景
时刻 T0: 队列状态:[send1(NORMAL,seq=1), send2(NORMAL,seq=2), ..., send10(NORMAL,seq=10)]
worker 正在执行 send1
时刻 T1: HTTP 请求 POST /api/wechat/restart 到达
↓ scheduler.execute(restart, CRITICAL, 0ms, 60s)
↓ 入队 (CRITICAL=0, seq=11)
↓ PriorityQueue 重排:[(0,11), (2,2), (2,3), ..., (2,10)]
↓ pending=10(restart 排第 1)
时刻 T2: send1 完成,worker 出队
↓ 出队 (0, 11) → restart 任务
↓ scheduler.mark_wechat_dead(8.0) ← 标记 fast-fail 窗口
↓ 执行 restart_wechat (60s 超时)
时刻 T3: restart 完成(~5s)
↓ 后续 8s 内 (mark_wechat_dead 窗口):
↓ send2-send10 出队 → 检查 _wechat_dead_until → fast-fail
↓ 客户端收到 WECHAT_NOT_READY + retry_after=5
↓ 8s 后窗口到期,新任务正常执行
结果:restart 在 1 个 send 完成后立即执行(不等 send2-send10),
后续 send 通过 fast-fail 快速失败,避免 9 次无谓的 xdotool 超时
六、风险与依赖
6.1 技术风险
| 风险 | 概率 | 影响 | 缓解 |
|---|---|---|---|
| PriorityQueue 元组比较失败 | 低 | worker 崩溃 | coro_factory 不可比较,用 seq 序列号作为第二排序键,(priority, seq) 组合唯一,不会比较到 coro_factory |
父类 _run 与子类 _run 元组结构不一致 |
中 | 父类被误调用时解包失败 | 子类必须覆盖 _run;代码注释明确标注"覆盖父类";单元测试验证子类 _run 被调用 |
父类 __init__ 创建 asyncio.Queue 后被子类 asyncio.PriorityQueue 覆盖 |
低 | 内存短暂多出一个未使用 Queue 对象 | Python 属性遮蔽合法,旧 Queue 被 GC 回收,无副作用 |
| 优先级反转导致 HIGH 操作被 LOW 阻塞 | 中 | logout 等待 screenshot 完成 | 限制单次操作超时(已有 5s/30s),screenshot 入 LOW 但超时 10s |
| drain 阻塞 watchdog | 中 | autofix 延迟 | drain 超时 5s 后强制 kill,不无限等待 |
| drain 返回 True 后 worker 仍在 sleep | 低 | kill 时机略早于 sleep 结束 | 队列已空,worker sleep 结束后阻塞在 get() 上,不会执行新任务;kill 后 worker 执行失败被正常捕获 |
| 收口后 logout 等 HTTP 接口延迟增加 | 中 | 用户感知 | priority=HIGH 保证优先级,wait_timeout_ms 兜底 |
put_nowait 与父类 await put 行为差异 |
低 | 满队时子类立即抛错,父类先检查再 put | 行为等价(父类也有前置 qsize 检查),子类更简洁 |
task_done() 调用次数 |
低 | 计数器异常 | cancelled 分支和 finally 分支互斥(continue 跳过 finally),每个 get() 对应一次 task_done(),与父类一致 |
| metrics 并发读写 | 低 | 数据轻微不准 | asyncio 单线程事件循环,worker 与 HTTP handler 同线程,无需加锁;asyncio.to_thread 释放事件循环时不读写 _metrics |
6.2 依赖
- 依赖现有
SendQueue实现稳定(已验证 16 个调用点) - 依赖
_ui_lock作为底层互斥(保持不变) - 依赖
asyncio.PriorityQueue(Python 3.8+ 标准库,无外部依赖)
6.3 非目标
- 不实现优先级抢占:xdotool 子进程不可中断,低优先级操作开始后必须等其完成
- 不实现优先级继承:xdotool 子进程无法感知调用方优先级
- 不重命名
SendQueue类:保持向后兼容,UIActionScheduler 继承之 - 不修改
_ui_lock:保持作为底层单命令互斥的第二道防线 - 不收口 LoginGuard / MessageStreamer:它们是只读探测,无需入队
6.4 失败风暴问题与 fast-fail 机制
问题场景:当 watchdog kill 微信或 wechat_restart 完成后,send_queue 中可能仍有 N 个已入队的 send_text 任务。这些任务会按顺序执行,每个都因 WINDOW_NOT_FOUND 失败一次(直到 autostart 拉起新进程)。
影响估算:
- 100 个 send_text 任务排队,每个执行失败耗时 ~1s(含 5s xdotool 超时 + 异常处理)
- 微信 autostart 拉起需 ~5-10s
- 失败风暴持续:
min(N, autostart拉起前积压数) × 1s≈ 5-10s 内 5-10 个任务连续失败 - 客户端收到 5-10 个 503 错误,可能触发重试,进一步放大风暴
fast-fail 机制设计(本方案可选实施,建议阶段 2 一并实施):
class UIActionScheduler(SendQueue):
def __init__(self, ...):
super().__init__(...)
# ...
self._wechat_dead_until: float = 0.0 # 时间戳,在此之前所有任务直接 fast-fail
def mark_wechat_dead(self, duration_sec: float = 10.0) -> None:
"""标记微信已死,期间所有新任务直接 fast-fail。
供 watchdog kill / wechat_restart 调用:
- kill 前调 mark_wechat_dead(15.0),覆盖 autostart 拉起窗口
- restart 完成后调 mark_wechat_dead(5.0),覆盖新进程启动窗口
"""
self._wechat_dead_until = time.monotonic() + duration_sec
logger.warning(
"[scheduler] mark_wechat_dead %.1fs, %d pending tasks will fast-fail",
duration_sec, self._queue.qsize(),
)
async def _run(self) -> None:
while True:
priority, seq, coro_factory, future, custom_delay_ms = await self._queue.get()
# ... cancelled 检查 ...
# fast-fail 检查:微信已死期间所有任务直接失败
if self._wechat_dead_until > time.monotonic():
logger.info(
"[scheduler] fast-fail (wechat dead, %.1fs remaining, priority=%s)",
self._wechat_dead_until - time.monotonic(),
Priority(priority).name,
)
if not future.done():
future.set_exception(BridgeError(
code="WECHAT_NOT_READY",
message="微信进程未就绪(重启中),请稍后重试",
details={"retry_after": 5},
))
self._queue.task_done()
continue # 不延时,立即处理下一个
# ... 正常执行流程 ...
调用点:
watchdog._autofix_wechatkill 前调scheduler.mark_wechat_dead(15.0)routes/login.py wechat_restart开始时调scheduler.mark_wechat_dead(8.0)routes/diagnostic.py autofix_wechat_running同上
注意事项:
WECHAT_NOT_READY是新错误码,需在models/errors.py注册 HTTP 503 映射- fast-fail 任务不计入
executed_total,应单独记fast_fail_total指标 - CRITICAL 任务(如 restart 本身)不应被 fast-fail 拦截,需在检查中排除:
if self._wechat_dead_until > time.monotonic() and priority > Priority.CRITICAL:
替代方案对比:
| 方案 | 优点 | 缺点 |
|---|---|---|
| fast-fail(推荐) | 立即拒绝,客户端快速收到错误并退避 | 需新增错误码 + 调用点埋点 |
| drain + cancel pending | 队列清空,无失败任务 | cancel 已入队 future 复杂,可能误取消正在执行的任务 |
| 不处理(保持现状) | 实现简单 | 5-10 个连续失败,客户端可能重试放大风暴 |
七、变更记录
| 版本 | 日期 | 修改人 | 摘要 |
|---|---|---|---|
| v1.0 | 2026-07-17 | - | 初稿,基于 SendQueue 调研设计 UIActionScheduler 统一调度方案 |
| v1.1 | 2026-07-17 | - | 深度复核:修正 task_done 互斥说明、put_nowait 语义、logout 前置 detect_login_state 处理、StatusResponse 模型扩展、激活已存在 Prometheus Gauge、drain 边界条件、metrics 线程安全 |
| v1.2 | 2026-07-18 | - | 深度调研优化:新增 1.4 双层保护模型(L1 _ui_lock vs L2 SendQueue)、1.5 调用点全景图(before/after)、1.6 HTTP 端点→优先级映射表、2.3.1 delay_ms 与优先级交互策略、5.1 关键场景时序图(4 个)、6.4 失败风暴问题与 fast-fail 机制(mark_wechat_dead)、UIActionScheduler 代码新增 mark_wechat_dead 方法与 _run fast-fail 检查、P1.1 detect_login_state 14 处散落调用分析与处理方案 |
八、实施前验证清单
实施前需确认以下事项,避免引入新问题:
8.1 代码复核清单
asyncio.PriorityQueue.put_nowait满时抛asyncio.QueueFull(已确认,与asyncio.Queue行为一致)asyncio.PriorityQueue的task_done()/join()语义与asyncio.Queue一致(已确认,继承关系)- 元组
(int, int, coro_factory, future, delay_ms)中seq单调递增保证唯一,不会比较到coro_factory(已确认) - 父类
__init__创建的asyncio.Queue会被子类asyncio.PriorityQueue覆盖,旧对象无其他引用,GC 回收(已确认) - 父类
_run被子类覆盖,不会调用父类的 3 元组解包(已确认,代码注释标注) _check_rate_limit继承父类,行为不变(已确认,子类不覆盖)send_queue.start()/stop()继承父类,worker task 管理 不变(已确认)pending_count()继承父类,返回self._queue.qsize()(已确认,PriorityQueue 也有 qsize)
8.2 单元测试清单
实施时需编写以下单元测试(参考 verify_bus.py 的测试模式):
# test_ui_action_scheduler.py
async def test_backward_compat_enqueue():
"""enqueue 方法向后兼容,priority 默认 NORMAL。"""
scheduler = UIActionScheduler(send_delay_ms=10, max_calls_per_sec=100)
await scheduler.start()
try:
result = await scheduler.enqueue(lambda: asyncio.sleep(0.01, result="ok"))
assert result == "ok"
metrics = scheduler.get_metrics()
assert metrics["executed_total"] == 1
assert metrics["executed_by_priority"]["NORMAL"] == 1
finally:
await scheduler.stop()
async def test_priority_ordering():
"""CRITICAL 优先于 NORMAL 执行。"""
scheduler = UIActionScheduler(send_delay_ms=0, max_calls_per_sec=100)
await scheduler.start()
try:
# 先入队一个 NORMAL(会立即开始执行)
normal_future = scheduler.enqueue(lambda: asyncio.sleep(0.1, result="normal"))
# 再入队 CRITICAL 和 NORMAL
critical_future = scheduler.execute(
lambda: asyncio.sleep(0.01, result="critical"),
priority=Priority.CRITICAL,
)
normal2_future = scheduler.execute(
lambda: asyncio.sleep(0.01, result="normal2"),
priority=Priority.NORMAL,
)
# CRITICAL 应先于 normal2 完成
critical_result = await critical_future
normal2_result = await normal2_future
assert critical_result == "critical"
assert normal2_result == "normal2"
# 验证执行顺序:CRITICAL 在 normal2 之前
assert scheduler.get_metrics()["executed_by_priority"]["CRITICAL"] == 1
finally:
await scheduler.stop()
async def test_drain_empty_queue():
"""空队列 drain 立即返回 True。"""
scheduler = UIActionScheduler()
result = await scheduler.drain(timeout=1.0)
assert result is True
async def test_drain_with_pending():
"""有 pending 任务时 drain 等待完成。"""
scheduler = UIActionScheduler(send_delay_ms=100)
await scheduler.start()
try:
# 入队一个任务
await scheduler.enqueue(lambda: asyncio.sleep(0.05, result="ok"))
# drain 等待完成
result = await scheduler.drain(timeout=2.0)
assert result is True
finally:
await scheduler.stop()
async def test_drain_timeout():
"""任务执行超过 drain timeout 时返回 False。"""
scheduler = UIActionScheduler(send_delay_ms=0)
await scheduler.start()
try:
# 入队一个长任务
scheduler.enqueue(lambda: asyncio.sleep(1.0, result="slow"))
# drain 短超时
result = await scheduler.drain(timeout=0.1)
assert result is False
finally:
await scheduler.stop()
async def test_queue_full_rate_limited():
"""队列满时抛 RATE_LIMITED。"""
scheduler = UIActionScheduler(send_delay_ms=1000, max_queue_size=2)
await scheduler.start()
try:
# 填满队列(1 个执行中 + 2 个排队)
scheduler.enqueue(lambda: asyncio.sleep(0.5))
scheduler.enqueue(lambda: asyncio.sleep(0.01))
scheduler.enqueue(lambda: asyncio.sleep(0.01))
# 第 4 个应被拒绝
with pytest.raises(BridgeError) as exc_info:
await scheduler.enqueue(lambda: asyncio.sleep(0.01))
assert exc_info.value.code == "RATE_LIMITED"
finally:
await scheduler.stop()
async def test_wait_timeout():
"""等待超时抛 TIMEOUT。"""
scheduler = UIActionScheduler(send_delay_ms=1000, max_calls_per_sec=100)
await scheduler.start()
try:
# 第一个任务慢
scheduler.enqueue(lambda: asyncio.sleep(0.5))
# 第二个任务等待超时
with pytest.raises(BridgeError) as exc_info:
await scheduler.enqueue(
lambda: asyncio.sleep(0.01),
wait_timeout_ms=100, # 100ms 超时
)
assert exc_info.value.code == "TIMEOUT"
finally:
await scheduler.stop()
async def test_metrics_accuracy():
"""metrics 准确记录执行次数和优先级。"""
scheduler = UIActionScheduler(send_delay_ms=0, max_calls_per_sec=100)
await scheduler.start()
try:
await scheduler.execute(lambda: "a", priority=Priority.CRITICAL)
await scheduler.execute(lambda: "b", priority=Priority.HIGH)
await scheduler.execute(lambda: "c", priority=Priority.NORMAL)
await scheduler.execute(lambda: "d", priority=Priority.LOW)
metrics = scheduler.get_metrics()
assert metrics["executed_total"] == 4
assert metrics["executed_by_priority"]["CRITICAL"] == 1
assert metrics["executed_by_priority"]["HIGH"] == 1
assert metrics["executed_by_priority"]["NORMAL"] == 1
assert metrics["executed_by_priority"]["LOW"] == 1
assert metrics["pending_count"] == 0
finally:
await scheduler.stop()
8.3 集成测试清单
阶段 2 实施后需进行集成测试:
- 并发 30 个
POST /api/send/text+ 1 个POST /api/login/logout:logout 在所有 send 之前或之后执行,不穿插 POST /api/screenshot与POST /api/send/text并发:screenshot 串行等待- watchdog autofix 触发时 send_queue 有 pending:先 drain 再 kill
GET /api/status返回ui_scheduler字段且executed_by_priority计数正确/metrics端点woc_send_queue_pendingGauge 不再恒为 0
8.4 回归测试清单
阶段 1 实施后需验证现有功能不退化:
POST /api/send/text正常发送(16 个 enqueue 调用点之一)POST /api/friends/accept正常通过好友(friend_watcher + routes/contacts 共用 enqueue)POST /api/moments/publish正常发朋友圈(routes/moments enqueue)- BatchWorker 群发正常串行(间接经 orchestrator → enqueue)
GET /api/status的send_queue_pending字段仍正确