WechatOnCloud/doc/优化方案/05-UIActionScheduler统一调度方案.md
Kris 67aab58a00 docs: 新增PRD规范与WechatOnCloud改造、多应用桥接框架文档
新增三份文档:
1. 产品需求文档(PRD)编写规范
2. WechatOnCloud容器化微信改造需求方案
3. 多应用桥接框架整体设计方案
2026-07-18 16:04:25 +08:00

76 KiB
Raw Permalink Blame History

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 / L8526 处) 每次 HTTP 请求 每次持锁 ~200-500msscreenshot+模板匹配100 并发 send 时累计 12-30s 额外延迟
routes/moments.py L151 / L255 / L346 / L447 / L545 / L6776 处) 每次 HTTP 请求 同上
routes/contacts.py L236 / L313 / L4333 处) 每次 HTTP 请求 同上
routes/login.py L117 / L1902 处) 每次 HTTP 请求 同上
routes/diagnostic.py L1141 处) 诊断触发 频率低,影响小
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"的优化路径是可行且被项目采纳的。

处理方案本方案不强制收口,但记录优化路径

  1. 短期(本方案不实施):保持现状。detect_login_state 是单次截图+模板匹配,持锁时间可控(< 1s。在 send_queue 任务执行期间(通常 1-5s最多 1-2 次 detect_login_state 抢占 _ui_lock影响有限。
  2. 中期(独立优化项):参照 routes/status.py 的优化模式,将路由前置 detect_login_state 改为内联 pgrep + find_window 判定(不持 _ui_lock仅在判定为 "logged_in" 后才入队执行真正的 UI 操作。
  3. 长期(架构演进)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 秒滑窗限流 + 队列满拒绝 + 等待超时

可直接复用为基础,但需扩展:

  1. 增加 priority 参数(当前队列元素是 tuple[CoroFactory, Future, Optional[int]],第三段是 custom_delay_ms,无 priority 槽位)
  2. 重命名以反映通用语义(向后兼容保留别名)
  3. 收口绕过队列的调用点

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_locksend_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 重写,而是扩展 + 收口

  1. 扩展 SendQueueUIActionScheduler,增加优先级参数
  2. 收口绕过队列的 6 个调用点,统一入队
  3. 保留 _ui_lock 作为底层单命令互斥(不取消,作为第二道防线)
  4. 新增 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=NORMALdelay_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
  • kill 后广播 wechat_killed 事件,所有进行中的 Flow 捕获 WINDOW_NOT_FOUND 后自然失败

2.5 关键决策QrCapture 如何收口

QrCapture 有两类操作:

  1. capture_qr_code:启动扫码登录流程,多步 UI 操作activate + screenshot + crop
  2. 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 注入 schedulerapp.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向后兼容现有调用

改动

  1. 新增 messaging/ui_action_scheduler.pyPriority + UIActionScheduler 类)
  2. app.pySendQueue(...) 替换为 UIActionScheduler(...)
  3. watchdog.py 增加 scheduler 参数(可选,向后兼容)
  4. routes/status.py 暴露 metrics

验证

  • py_compile 通过
  • 现有 16 个 enqueue 调用点无需修改,行为不变
  • /api/status 返回 ui_scheduler 字段

阶段 2收口 P0 调用点

目标:消除 6 个绕过队列的竞态风险 + QrCapture 超时修复

改动

  1. routes/login.pylogout / restart / qr_start 改用 scheduler.execute
    • logout 仅 xdotool.logout() 入队,detect_login_state 保持不入队
  2. routes/screenshot.py:改用 scheduler.execute(priority=LOW)
  3. routes/diagnostic.pyautofix 改用 scheduler.execute(priority=CRITICAL)
  4. ui/qr_capture.py:增加 _CMD_TIMEOUT_SEC=5.0 包裹 proc.communicate()
  5. ui/watchdog.py:构造函数新增 scheduler 参数;_autofix_wechat 增加 scheduler.drain() 调用
  6. app.pyWeChatWatchdog(...) 注入 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 和详细日志

改动

  1. models/status.pyStatusResponse 新增 ui_scheduler: Optional[dict] 字段
  2. routes/status.py:暴露 ui_scheduler metrics + 激活已存在的 woc_send_queue_pending Gauge
  3. ui/metrics.py:新增 woc_ui_action_* 系列 metricsCounter/Histogram
  4. UIActionScheduler._run / _enqueue_with_priority:埋点写入 Prometheus metrics
  5. 日志增加 prioritytrace_id 字段(已在阶段 1 代码中包含)

验证

  • GET /api/status 返回 ui_scheduler 字段
  • /metrics 端点返回 woc_ui_action_executed_total 等指标
  • 日志可按 priority=CRITICAL 过滤
  • 已存在的 woc_send_queue_pending Gauge 不再为 0bug 修复)

五、验收标准

AC-01UIActionScheduler 向后兼容

Given 现有 16 个 enqueue 调用点未修改 When app.pySendQueue(...) 替换为 UIActionScheduler(...) Then 所有现有功能行为不变py_compile 通过,单元测试通过

AC-02优先级调度生效

Given 队列中有 10 个 NORMAL 优先级的 send_text 任务 When 提交一个 CRITICAL 优先级的 wechat_restart Then wechat_restart 优先于剩余 NORMAL 任务执行(可能在当前 NORMAL 任务执行完后立即执行)

AC-03logout 与 send_text 不再竞态

Given 一个 send_text 正在 send_queue 中执行 When 并发调用 POST /api/login/logout Then logout 入队等待send_text 执行完后 logout 才开始执行

AC-04QrCapture 有超时保护

Given scrot 子进程卡死 When QrCapture._run 执行 Then 5 秒后超时,proc.kill() + await proc.wait(),抛出 TimeoutError

AC-05watchdog autofix 协同

Given send_queue 中有 5 个 pending 任务 When watchdog 触发 _autofix_wechat Thenscheduler.drain(timeout=5.0) 等待队列清空,超时则强制 kill 并记 warning

AC-06metrics 暴露

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_textlogout 等待 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_texttask_done
         ↓ _queue.join() 检查未完成任务数
         ↓ 仍有 4 个 pendingtask_done 计数未归零)
         ↓ worker 出队下一个 send_text 并执行...
         ↓ ⚠️ 若每个 send_text 耗时 > 1.25s5s 内无法清空

时刻 T3a (5s 内清空):
         ↓ drain 返回 True
         ↓ scheduler.mark_wechat_dead(15.0)  ← 启动 fast-fail 窗口
         ↓ pgrep + kill -TERM wechat
         ↓ 后续 15s 内新入队任务直接 fast-failWECHAT_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=10restart 排第 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/30sscreenshot 入 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.PriorityQueuePython 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_wechat kill 前调 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.PriorityQueuetask_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/logoutlogout 在所有 send 之前或之后执行,不穿插
  • POST /api/screenshotPOST /api/send/text 并发screenshot 串行等待
  • watchdog autofix 触发时 send_queue 有 pending先 drain 再 kill
  • GET /api/status 返回 ui_scheduler 字段且 executed_by_priority 计数正确
  • /metrics 端点 woc_send_queue_pending Gauge 不再恒为 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/statussend_queue_pending 字段仍正确