80 lines
2.9 KiB
Python
80 lines
2.9 KiB
Python
# --- routers/logs.py ---
|
||
import asyncio
|
||
import logging
|
||
from fastapi import APIRouter, WebSocket, WebSocketDisconnect
|
||
from logger import ColoredFormatter, register_websocket_handler
|
||
|
||
router = APIRouter()
|
||
|
||
class LogBroadcaster:
|
||
"""管理所有活躍 WebSocket 連線的廣播器 (Connection Manager)"""
|
||
def __init__(self):
|
||
self.active_connections: set[WebSocket] = set()
|
||
self.loop: asyncio.AbstractEventLoop = None
|
||
|
||
async def connect(self, websocket: WebSocket):
|
||
await websocket.accept()
|
||
self.active_connections.add(websocket)
|
||
|
||
def disconnect(self, websocket: WebSocket):
|
||
self.active_connections.discard(websocket)
|
||
|
||
async def broadcast(self, message: str):
|
||
"""非同步推播日誌給所有連線中的客戶端"""
|
||
if not self.active_connections:
|
||
return
|
||
|
||
# 併發發送,並透過 return_exceptions=True 確保單一連線異常不影響其他客戶端
|
||
tasks = [self._safe_send(conn, message) for conn in self.active_connections]
|
||
await asyncio.gather(*tasks, return_exceptions=True)
|
||
|
||
async def _safe_send(self, websocket: WebSocket, message: str):
|
||
try:
|
||
await websocket.send_text(message)
|
||
except Exception:
|
||
self.disconnect(websocket)
|
||
|
||
# 建立全域廣播器實例
|
||
log_broadcaster = LogBroadcaster()
|
||
|
||
class WebSocketLogHandler(logging.Handler):
|
||
"""自訂 Log Handler,攔截系統日誌並安全地派發至非同步廣播器"""
|
||
def __init__(self, broadcaster: LogBroadcaster):
|
||
super().__init__()
|
||
self.broadcaster = broadcaster
|
||
# 沿用系統 ColoredFormatter,完美保留 ANSI 色碼
|
||
self.setFormatter(ColoredFormatter())
|
||
|
||
def emit(self, record):
|
||
try:
|
||
msg = self.format(record)
|
||
loop = self.broadcaster.loop
|
||
# 確保在 Event Loop 處於執行狀態時,安全地跨執行緒派發任務
|
||
if loop and loop.is_running():
|
||
loop.call_soon_threadsafe(
|
||
lambda: asyncio.create_task(self.broadcaster.broadcast(msg))
|
||
)
|
||
except Exception:
|
||
self.handleError(record)
|
||
|
||
# 建立 Handler 實例並註冊至所有受管控的 Logger
|
||
websocket_log_handler = WebSocketLogHandler(log_broadcaster)
|
||
register_websocket_handler(websocket_log_handler)
|
||
|
||
@router.websocket("/ws/logs")
|
||
async def websocket_logs(websocket: WebSocket):
|
||
"""WebSocket 實時日誌串流端點"""
|
||
# 若尚未綁定 Event Loop,則於首次連線時動態綁定
|
||
if not log_broadcaster.loop:
|
||
log_broadcaster.loop = asyncio.get_running_loop()
|
||
|
||
await log_broadcaster.connect(websocket)
|
||
try:
|
||
# 保持連線,監聽客戶端斷線狀態
|
||
while True:
|
||
await websocket.receive_text()
|
||
except WebSocketDisconnect:
|
||
log_broadcaster.disconnect(websocket)
|
||
except Exception:
|
||
log_broadcaster.disconnect(websocket)
|