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)
|