广播日志到所有 WebSocket 连接
(self, task_uuid: str, log_message: str)
| 114 | _log_queues[task_uuid].append(log_message) |
| 115 | |
| 116 | async def _broadcast_log(self, task_uuid: str, log_message: str): |
| 117 | """广播日志到所有 WebSocket 连接""" |
| 118 | with _ws_lock: |
| 119 | connections = _ws_connections.get(task_uuid, []).copy() |
| 120 | # 注意:不在这里更新 sent_index,因为日志已经通过 add_log 添加到队列 |
| 121 | # sent_index 应该只在 get_unsent_logs 或发送历史日志时更新 |
| 122 | # 这样可以避免竞态条件 |
| 123 | |
| 124 | for ws in connections: |
| 125 | try: |
| 126 | await ws.send_json({ |
| 127 | "type": "log", |
| 128 | "task_uuid": task_uuid, |
| 129 | "message": log_message, |
| 130 | "timestamp": utcnow_naive().isoformat() |
| 131 | }) |
| 132 | # 发送成功后更新 sent_index |
| 133 | with _ws_lock: |
| 134 | ws_id = id(ws) |
| 135 | if task_uuid in _ws_sent_index and ws_id in _ws_sent_index[task_uuid]: |
| 136 | _ws_sent_index[task_uuid][ws_id] += 1 |
| 137 | except Exception as e: |
| 138 | logger.warning(f"WebSocket 发送失败: {e}") |
| 139 | |
| 140 | async def broadcast_status(self, task_uuid: str, status: str, **kwargs): |
| 141 | """广播任务状态更新""" |
no test coverage detected