WechatOnCloud/bridge/woc_bridge/messaging/send_queue.py

159 lines
6.0 KiB
Python
Raw Normal View History

"""发送串行化队列。
所有 xdotool 操作经此队列串行执行避免并发 UI 操作冲突
内部维护最近 1 秒调用时间戳用于限流
"""
from __future__ import annotations
import asyncio
import logging
import time
from typing import Any, Awaitable, Callable
from woc_bridge.models import BridgeError
logger = logging.getLogger("woc-bridge")
# 工厂类型:返回一个待执行的 coroutine
CoroFactory = Callable[[], Awaitable[Any]]
class SendQueue:
"""串行化发送队列 + 限流。
通过 asyncio.Queue 串行执行所有发送任务执行间隔可配置
默认 800ms单实例每秒调用上限可配置默认 10
"""
def __init__(
self,
send_delay_ms: int = 800,
max_calls_per_sec: int = 10,
) -> None:
"""初始化队列配置。
Args:
send_delay_ms: 两次发送之间的最小间隔毫秒
max_calls_per_sec: 每秒最大调用次数
"""
self.send_delay_ms = send_delay_ms
self.max_calls_per_sec = max_calls_per_sec
self._queue: asyncio.Queue[tuple[CoroFactory, asyncio.Future]] = asyncio.Queue()
self._worker: asyncio.Task | None = None
# 最近 1 秒内的调用时间戳
self._recent_call_times: list[float] = []
async def start(self) -> None:
"""启动 worker task。"""
if self._worker is None or self._worker.done():
self._worker = asyncio.create_task(self._run())
async def stop(self) -> None:
"""取消 worker。"""
if self._worker is not None and not self._worker.done():
self._worker.cancel()
try:
await self._worker
except asyncio.CancelledError:
pass
self._worker = None
async def enqueue(self, coro_factory: CoroFactory) -> Any:
"""将一个返回 coroutine 的工厂入队,等待执行结果。
Args:
coro_factory: 调用后返回 coroutine 的工厂函数
Returns:
coroutine 的执行结果
Raises:
BridgeError: 限流命中时立即抛 RATE_LIMITED
任务执行抛出的异常会透传给调用方
"""
loop = asyncio.get_running_loop()
future: asyncio.Future = loop.create_future()
await self._queue.put((coro_factory, future))
logger.info("send_queue: 入队 (pending=%d)", self._queue.qsize())
return await future
def pending_count(self) -> int:
"""返回当前队列中待执行任务数(供 /api/status 暴露给客户端做退避决策)。"""
return self._queue.qsize()
def _check_rate_limit(self) -> None:
"""检查限流。
清理 1 秒前的时间戳若当前已满 max_calls_per_sec 则抛
BridgeError(RATE_LIMITED)并在 details 中携带 retry_after 秒数
供上层设置 Retry-After 响应头
"""
now = time.monotonic()
# 清理 1 秒前的时间戳
self._recent_call_times = [t for t in self._recent_call_times if now - t < 1.0]
if len(self._recent_call_times) >= self.max_calls_per_sec:
# 计算建议等待秒数:最早一次调用距窗口边界还差多久
oldest = self._recent_call_times[0]
retry_after = max(1, int(1.0 - (now - oldest)) + 1)
logger.warning(
"send_queue: 限流命中,拒绝执行 (recent=%d/%d, retry_after=%ds)",
len(self._recent_call_times), self.max_calls_per_sec, retry_after,
)
raise BridgeError(
code="RATE_LIMITED",
message=f"发送限流:每秒最多 {self.max_calls_per_sec}",
details={"retry_after": retry_after},
)
async def _run(self) -> None:
"""worker 主循环。
循环取出任务执行执行前检查限流超限则失败该任务
执行前记录开始时间戳避免长任务导致 1 秒窗口内超限
执行后 sleep send_delay_ms/1000
"""
while True:
coro_factory, future = await self._queue.get()
logger.info("send_queue: 出队,开始处理 (pending=%d)", self._queue.qsize())
# 标记任务是否真正开始执行(用于决定 finally 是否延时)
executed = False
t_exec = time.perf_counter()
try:
# 执行前检查限流
self._check_rate_limit()
# 记录开始时间戳(限流窗口基于开始时刻,避免长任务后窗口偏移)
self._recent_call_times.append(time.monotonic())
executed = True
# 执行任务
logger.info("send_queue: 开始执行任务")
result = await coro_factory()
logger.info(
"send_queue: 任务执行完成 (%.0fms)",
(time.perf_counter() - t_exec) * 1000,
)
if not future.done():
future.set_result(result)
except asyncio.CancelledError:
# worker 被取消时,把取消传播给等待的调用方
if not future.done():
future.cancel()
raise
except Exception as e:
logger.warning(
"send_queue: 任务执行抛异常 %s: %s (%.0fms)",
type(e).__name__, e, (time.perf_counter() - t_exec) * 1000,
)
if not future.done():
future.set_exception(e)
finally:
self._queue.task_done()
# 仅在任务真正执行过时延时,限流失败的任务不延时
if executed:
logger.info(
"send_queue: 延时 %dms 后处理下一个",
self.send_delay_ms,
)
await asyncio.sleep(self.send_delay_ms / 1000.0)