广播批量任务状态
(self, batch_id: str)
| 305 | logger.warning(f"广播批量状态失败: {e}") |
| 306 | |
| 307 | async def _broadcast_batch_status(self, batch_id: str): |
| 308 | """广播批量任务状态""" |
| 309 | with _ws_lock: |
| 310 | connections = _ws_connections.get(f"batch_{batch_id}", []).copy() |
| 311 | |
| 312 | status = _batch_status.get(batch_id, {}) |
| 313 | |
| 314 | for ws in connections: |
| 315 | try: |
| 316 | await ws.send_json({ |
| 317 | "type": "status", |
| 318 | "batch_id": batch_id, |
| 319 | "timestamp": utcnow_naive().isoformat(), |
| 320 | **status |
| 321 | }) |
| 322 | except Exception as e: |
| 323 | logger.warning(f"WebSocket 发送批量状态失败: {e}") |
| 324 | |
| 325 | def get_batch_status(self, batch_id: str) -> Optional[dict]: |
| 326 | """获取批量任务状态""" |
no test coverage detected