WechatOnCloud/bridge/woc_bridge/routes/login.py
Kris ed2ee086de refactor: 重构UI操作调度,引入优先级调度与监控
主要变更:
1. 新增UIActionScheduler实现优先级UI任务调度,支持CRITICAL/HIGH/NORMAL/LOW四级优先级
2. 替换原有SendQueue为UIActionScheduler,统一所有UI操作的调度逻辑
3. 为截图、登录、重启等API添加调度器封装,支持超时、限流与fast-fail机制
4. 新增Prometheus监控指标,统计UI任务执行、等待耗时、超时、限流与快速失败情况
5. 为QrCapture添加命令执行超时保护,避免X11操作卡死
6. 扩展StatusResponse与状态接口,暴露调度器指标
7. 新增完整的单元测试覆盖调度器功能
8. 为看门狗注入调度器,实现kill前队列清空与快速失败窗口
2026-07-18 03:43:28 +08:00

354 lines
14 KiB
Python
Raw Permalink Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

from __future__ import annotations
import asyncio
import logging
import time
from fastapi import APIRouter
from woc_bridge.config import (
_require_xdotool,
_require_qr_capture,
_require_db_reader,
_require_send_queue,
)
from woc_bridge.db.coordinator import with_db_retry
from woc_bridge.messaging.ui_action_scheduler import Priority
from woc_bridge.models import (
QrLoginStartResult,
QrLoginWaitResult,
LogoutResponse,
RestartResponse,
BridgeError,
LoginState,
)
logger = logging.getLogger("woc-bridge")
router = APIRouter()
@router.post("/api/login/qr/start", response_model=QrLoginStartResult)
async def login_qr_start() -> QrLoginStartResult:
"""启动扫码登录:截取二维码区域并返回 data URL。
channels LoginAdapter 的扫码入口。调用本接口后,调用方应把 qr_data_url
渲染给用户扫描,然后立即调 GET /api/login/qr/wait 轮询登录结果。
Returns:
QrLoginStartResult含 qr_data_urldata:image/png;base64,... /
message / connected=false
Raises:
BridgeError(WINDOW_NOT_FOUND): 微信窗口未找到无法截图HTTP 503
其他异常由全局兜底处理器返回 BRIDGE_INTERNAL_ERROR
Notes:
- 激活微信窗口后再截图,确保二维码可见(不被其他窗口遮挡)
- 二维码有时效(微信约 60s 刷新),调用方应在 qr/wait 超时后
重新调本接口获取新二维码
- 若微信已登录,本接口仍会返回二维码截图(调用方应先调
/api/status 判断 login_state避免重复登录
"""
logger.info("login/qr/start: 收到请求")
xdotool = _require_xdotool()
qr_capture = _require_qr_capture()
scheduler = _require_send_queue()
# 激活微信窗口(非阻塞,避免 VNC 无人操作时 --sync 死等)
t0 = time.perf_counter()
logger.info("login/qr/start: 激活窗口 + 截图开始")
async def _qr_capture_flow():
await xdotool._activate_window_fast()
return await qr_capture.capture_qr_code()
try:
if hasattr(scheduler, "execute"):
qr_data_url = await scheduler.execute(
_qr_capture_flow,
priority=Priority.HIGH,
wait_timeout_ms=15000,
trace_id="qr_start",
)
else:
# legacy 回退:直接调用
await xdotool._activate_window_fast()
qr_data_url = await qr_capture.capture_qr_code()
except BridgeError as e:
logger.warning(
"login/qr/start: 失败 BridgeError code=%s msg=%s (%.0fms)",
e.code, e.message, (time.perf_counter() - t0) * 1000,
)
raise
logger.info(
"login/qr/start: 激活窗口 + 截图完成 (%.0fms)",
(time.perf_counter() - t0) * 1000,
)
return QrLoginStartResult(
qr_data_url=qr_data_url,
message="请使用微信扫描二维码",
connected=False,
)
# ---------------------------------------------------------------------------
# 路由GET /api/login/qr/wait
# ---------------------------------------------------------------------------
@router.get("/api/login/qr/wait", response_model=QrLoginWaitResult)
@with_db_retry
async def login_qr_wait(timeout: int = 30) -> QrLoginWaitResult:
"""轮询登录态,等待扫码登录完成。
每 2 秒调一次 detect_login_state直到 logged_in / not_running / 超时。
长轮询模式(响应在登录成功或超时后才返回),调用方无需在 client 端做
轮询间隔控制,直接发请求阻塞等待即可。
Args:
timeout: 最大等待秒数,默认 30上限 120防止长时间占用连接
Returns:
QrLoginWaitResult含 connected / message / qr_data_url成功时为空
/ credentials成功时含 wxid + nickname否则为 None
Raises:
BridgeError(INVALID_PARAMS): timeout < 1HTTP 400
Notes:
- 微信进程未运行not_running时立即返回不继续等待继续等无意义
- 登录成功时尝试从 DB 读 wxid/nicknameDB 不可达时返回空字符串
credentials 不为 None但 wxid 为空)
- 超时返回 connected=falsemessage="等待扫码超时",调用方应重新调
POST /api/login/qr/start 获取新二维码(旧二维码可能已过期)
"""
if timeout < 1:
raise BridgeError(code="INVALID_PARAMS", message="timeout 必须 >= 1")
if timeout > 120:
timeout = 120
xdotool = _require_xdotool()
db_reader = _require_db_reader()
deadline = time.monotonic() + timeout
while time.monotonic() < deadline:
state = await xdotool.detect_login_state()
if state == "logged_in":
# 登录成功:从 DB 读取 wxid / nickname
try:
self_info = await asyncio.to_thread(db_reader.get_self_info)
wxid = self_info.get("wxid", "")
nickname = self_info.get("nickname", "")
except Exception:
wxid = ""
nickname = ""
logger.info("login/qr/wait: → 登录成功 wxid=%s nickname=%s", wxid, nickname)
return QrLoginWaitResult(
connected=True,
message="登录成功",
qr_data_url="",
credentials={"wxid": wxid, "nickname": nickname},
)
if state == "not_running":
# 微信进程未运行,继续等无意义,直接返回
logger.warning("login/qr/wait: → 微信进程未运行")
return QrLoginWaitResult(
connected=False,
message="微信进程未运行,无法扫码登录",
qr_data_url="",
credentials=None,
)
# not_logged_in / logging_in继续等待
await asyncio.sleep(2)
# 超时
logger.warning("login/qr/wait: → 等待扫码超时(%ds", timeout)
return QrLoginWaitResult(
connected=False,
message="等待扫码超时",
qr_data_url="",
credentials=None,
)
# ---------------------------------------------------------------------------
# 路由POST /api/login/logout
# ---------------------------------------------------------------------------
@router.post("/api/login/logout", response_model=LogoutResponse)
async def login_logout() -> LogoutResponse:
"""退出微信登录channels LifecycleAdapter 调用)。
通过 UI 操作退出登录,避免直接 kill 进程导致登录态文件未刷盘。
幂等:未登录时直接返回成功(不抛错),方便调用方无需先查状态再调用。
流程(仅 logged_in 时执行):
1. detect_login_state 确认登录态
2. xdotool.logout() 执行 UI 操作:激活窗口 → 点击主菜单 → 方向键导航
到「退出登录」→ 回车 → 确认对话框(具体坐标/次数为估算值,需实测调优)
Returns:
LogoutResponsesuccess=true / message"已退出登录""当前未登录,无需退出"
Raises:
BridgeError(LOGOUT_FAILED): 窗口未找到或 UI 操作失败HTTP 500
Notes:
- 幂等not_running / not_logged_in 状态下都返回 success=true
- logout() 的 UI 坐标/方向键次数为估算值spec 已记录需实测调优
- 退出后微信进程仍在运行(只是登录态变 not_logged_inautostart
不会拉起新进程;如需重新登录走 /api/login/qr/start
- 不删除任何文件登录态由微信自身管理bridge 不破坏数据
"""
logger.info("login/logout: 收到请求")
t_total = time.perf_counter()
xdotool = _require_xdotool()
# 记录调用前的登录态(用于返回 message
t0 = time.perf_counter()
state_before = await xdotool.detect_login_state()
logger.info(
"login/logout: 登录态检测 → %s (%.0fms)",
state_before, (time.perf_counter() - t0) * 1000,
)
# 若已登录,执行 UI 退出操作
if state_before == LoginState.LOGGED_IN.value:
t1 = time.perf_counter()
logger.info("login/logout: logout 调用开始")
scheduler = _require_send_queue()
try:
if hasattr(scheduler, "execute"):
await scheduler.execute(
lambda: xdotool.logout(),
priority=Priority.HIGH,
wait_timeout_ms=30000,
trace_id="logout",
)
else:
await xdotool.logout()
except BridgeError as e:
logger.warning(
"login/logout: 失败 BridgeError code=%s msg=%s (%.0fms)",
e.code, e.message, (time.perf_counter() - t1) * 1000,
)
raise
except Exception as e:
logger.warning(
"login/logout: 失败 %s: %s (%.0fms)",
type(e).__name__, e, (time.perf_counter() - t1) * 1000,
)
raise BridgeError(
code="LOGOUT_FAILED",
message=f"退出登录失败: {e}",
)
logger.info(
"login/logout: logout 调用完成 (%.0fms)",
(time.perf_counter() - t1) * 1000,
)
logger.info(
"login/logout: ✓ (总耗时 %.0fms)",
(time.perf_counter() - t_total) * 1000,
)
return LogoutResponse(success=True, message="已退出登录")
# 未登录,幂等返回
logger.info("login/logout: 当前未登录,无需登出")
return LogoutResponse(success=True, message="当前未登录,无需退出")
# ---------------------------------------------------------------------------
# 路由POST /api/wechat/restart
# ---------------------------------------------------------------------------
@router.post("/api/wechat/restart", response_model=RestartResponse)
async def wechat_restart() -> RestartResponse:
"""重启微信进程(不破坏登录态,数据卷保留)。
流程:
1. pgrep -x wechat 获取当前所有 PID
2. 对每个 PID 发 SIGTERM不强制 SIGKILL给微信优雅退出机会避免 DB 写入未刷盘)
3. 轮询 pgrep -x wechat等待新 PID 出现autostart 2 秒后拉起)
4. 30 秒内未出现新进程抛 RESTART_TIMEOUT
与 docker stop/start 的区别:
- 本接口只重启微信进程不重启容器X 会话/VNC 连接保持,速度更快(~5s
- docker stop 会杀整个容器VNC 断连,需重新连接
Returns:
RestartResponsesuccess=true / message"微信已重启""微信已启动"
/ pid新进程 PID
Raises:
BridgeError(RESTART_TIMEOUT): 30 秒内未检测到新进程HTTP 408
BridgeError(BRIDGE_INTERNAL_ERROR): pgrep 等命令异常HTTP 500
Notes:
- 数据卷保留,登录态不丢(除非 SIGTERM 期间微信主动退出登录)
- 不提供 stop/start 接口autostart watchdog 会立即拉起被 stop
的微信进程stop 没有意义;启动由 autostart 管理
- 若 was_running=false重启前未运行message="微信已启动"
- 调用方应在收到 success=true 后调 /api/status 确认 login_state
恢复到 logged_inautostart 拉起后微信会自动尝试恢复登录)
"""
timeout_sec = 30
logger.info("wechat/restart: 收到请求 timeout=%ds", timeout_sec)
t_total = time.perf_counter()
xdotool = _require_xdotool()
# 记录调用前是否在运行
t0 = time.perf_counter()
was_running = await xdotool.is_wechat_running()
logger.info(
"wechat/restart: 进程检测 → running=%s (%.0fms)",
was_running, (time.perf_counter() - t0) * 1000,
)
t1 = time.perf_counter()
logger.info("wechat/restart: restart 调用开始")
scheduler = _require_send_queue()
# 启动 fast-fail 窗口:覆盖新进程启动期(~5-10s避免 restart 期间
# 队列内积压的 NORMAL 任务对新窗口执行 xdotool 失败形成风暴
if hasattr(scheduler, "mark_wechat_dead"):
try:
scheduler.mark_wechat_dead(8.0)
except Exception as exc:
logger.warning("wechat/restart: mark_wechat_dead failed: %s", exc)
try:
if hasattr(scheduler, "execute"):
new_pid = await scheduler.execute(
lambda: xdotool.restart_wechat(timeout_sec=timeout_sec),
priority=Priority.CRITICAL,
wait_timeout_ms=60000,
trace_id="wechat_restart",
)
else:
new_pid = await xdotool.restart_wechat(timeout_sec=timeout_sec)
except BridgeError as e:
logger.warning(
"wechat/restart: 失败 BridgeError code=%s msg=%s (%.0fms)",
e.code, e.message, (time.perf_counter() - t1) * 1000,
)
raise
except Exception as e:
# 非 timeout 的意外错误(如 pgrep 不可用)归为内部错误,
# 不误报为 RESTART_TIMEOUT调用方会按 timeout 语义重试,无意义)
logger.warning(
"wechat/restart: 失败 %s: %s (%.0fms)",
type(e).__name__, e, (time.perf_counter() - t1) * 1000,
)
raise BridgeError(
code="BRIDGE_INTERNAL_ERROR",
message=f"重启微信失败: {e}",
)
logger.info(
"wechat/restart: restart 调用完成 (%.0fms)",
(time.perf_counter() - t1) * 1000,
)
message = "微信已重启" if was_running else "微信已启动"
logger.info(
"wechat/restart: ✓ (总耗时 %.0fms)",
(time.perf_counter() - t_total) * 1000,
)
return RestartResponse(success=True, message=message, pid=new_pid)