"""WebSocket 日志流""" import asyncio import json from datetime import datetime from typing import Dict, AsyncGenerator class LogStreamer: """日志流管理器""" def __init__(self): self._queues: Dict[str, asyncio.Queue] = {} self._subscribers: Dict[str, list] = {} def create_queue(self, task_id: str) -> asyncio.Queue: """创建任务的日志队列""" queue = asyncio.Queue(maxsize=1000) self._queues[task_id] = queue return queue async def emit(self, task_id: str, message: str, level: str = "info"): """发送日志消息""" if task_id not in self._queues: return log_entry = { "timestamp": datetime.now().isoformat(), "level": level, "message": message, } try: self._queues[task_id].put_nowait(log_entry) except asyncio.QueueFull: # 队列满了,丢弃最旧的消息 try: self._queues[task_id].get_nowait() self._queues[task_id].put_nowait(log_entry) except asyncio.QueueEmpty: pass async def emit_step(self, task_id: str, step: str): """发送步骤标记""" await self.emit(task_id, f"[{step}]", level="step") async def emit_error(self, task_id: str, error: str): """发送错误消息""" await self.emit(task_id, error, level="error") async def emit_warning(self, task_id: str, warning: str): """发送警告消息""" await self.emit(task_id, warning, level="warn") async def subscribe(self, task_id: str) -> AsyncGenerator[dict, None]: """订阅任务日志""" queue = self._queues.get(task_id) if not queue: return while True: try: msg = await asyncio.wait_for(queue.get(), timeout=30) if msg is None: # 结束信号 break yield msg except asyncio.TimeoutError: # 发送心跳 yield {"timestamp": datetime.now().isoformat(), "level": "heartbeat", "message": ""} except Exception: break def complete(self, task_id: str): """标记任务日志结束""" if task_id in self._queues: try: self._queues[task_id].put_nowait(None) except asyncio.QueueFull: pass def cleanup(self, task_id: str): """清理任务日志队列""" self._queues.pop(task_id, None) # 全局单例 log_streamer = LogStreamer()