WechatOnCloud/doc/优化方案/05-UIActionScheduler统一调度方案.md

1646 lines
76 KiB
Markdown
Raw Permalink Normal View 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](../../bridge/woc_bridge/routes/login.py#L201) |
| restart kill 微信与并发 UI 操作冲突 | restart 不入队 | xdotool 对空窗口操作、TimeoutError 雪崩 | [routes/login.py:281](../../bridge/woc_bridge/routes/login.py#L281) |
| diagnostic autofix pkill 无锁 | 直接 `pkill -x wechat` | 与并发 UI 操作冲突高风险 | [routes/diagnostic.py:232-239](../../bridge/woc_bridge/routes/diagnostic.py#L232) |
| QrCapture 无锁无超时 | `proc.communicate()` 无 timeout | 与 send_text 截图竞争 X11、可能无限阻塞 | [ui/qr_capture.py:48-57](../../bridge/woc_bridge/ui/qr_capture.py#L48) |
| WeChatWatchdog autofix 无锁 | pgrep+kill 不持 `_ui_lock` | kill 微信时 UI 操作中途失败 | [ui/watchdog.py:99-191](../../bridge/woc_bridge/ui/watchdog.py#L99) |
#### 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. **扩展 SendQueue**`UIActionScheduler`,增加优先级参数
2. **收口绕过队列的 6 个调用点**,统一入队
3. **保留 `_ui_lock`** 作为底层单命令互斥(不取消,作为第二道防线)
4. **新增 metrics** 暴露队列深度、等待时长、执行时长
### 2.2 命名策略
采用**渐进式重命名**,避免一次性破坏全部调用点:
```python
# 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 优先级策略
```python
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
```python
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
- 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 类
```python
# 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()` 入队。
```python
# 修改前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
```python
# 修改后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
```python
# 修改后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
```python
# 修改后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
```python
# 修改后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 超时修复
```python
# 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 协同
```python
# 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 注入**
```python
# app.py:499 修改
_state.watchdog = WeChatWatchdog(
backend=_state.xdotool_backend,
interval=10.0,
scheduler=_state.send_queue, # 新增:注入调度器
)
```
### 3.4 注入与生命周期
```python
# 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 模型扩展
```python
# 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
```python
# 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,
# ...
)
```
返回示例:
```json
{
"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
```python
# 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
```python
# 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.py`Priority + UIActionScheduler 类)
2. `app.py``SendQueue(...)` 替换为 `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.py`logout / restart / qr_start 改用 `scheduler.execute`
- logout 仅 `xdotool.logout()` 入队,`detect_login_state` 保持不入队
2. `routes/screenshot.py`:改用 `scheduler.execute(priority=LOW)`
3. `routes/diagnostic.py`autofix 改用 `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.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 和详细日志
**改动**
1. `models/status.py``StatusResponse` 新增 `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. 日志增加 `priority``trace_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.py``SendQueue(...)` 替换为 `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`
**Then** 先 `scheduler.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.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 一并实施**
```python
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 拦截,需在检查中排除:
```python
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` 的测试模式):
```python
# 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_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/status``send_queue_pending` 字段仍正确