"""WebSocket endpoints for real-time log streaming and Docker events.""" from __future__ import annotations import asyncio import json from fastapi import APIRouter, Query, WebSocket, WebSocketDisconnect from jose import JWTError from auth import decode_token from services import compose_service router = APIRouter(tags=["ws"]) async def _authorize(websocket: WebSocket, token: str | None) -> bool: """Validate the JWT supplied as a query param. Closes socket on failure.""" if not token: await websocket.close(code=4401) return False try: decode_token(token, "access") except (JWTError, Exception): # noqa: BLE001 await websocket.close(code=4401) return False return True async def _stream_logs(websocket: WebSocket, stack_id: str, service: str | None): """Stream `docker compose logs -f` output to the client.""" args = ["logs", "--no-color", "--tail", "200", "--timestamps", "-f"] if service: args.append(service) try: async for line in compose_service.stream_compose(stack_id, args): await websocket.send_text( json.dumps( { "type": "log", "stack_id": stack_id, "service": service, "line": line, } ) ) except WebSocketDisconnect: raise except Exception as exc: # noqa: BLE001 await websocket.send_text( json.dumps({"type": "error", "detail": str(exc)}) ) @router.websocket("/ws/logs/{stack_id}") async def ws_stack_logs( websocket: WebSocket, stack_id: str, token: str | None = Query(default=None), ): await websocket.accept() if not await _authorize(websocket, token): return try: await _stream_logs(websocket, stack_id, None) except WebSocketDisconnect: pass @router.websocket("/ws/logs/{stack_id}/{service}") async def ws_service_logs( websocket: WebSocket, stack_id: str, service: str, token: str | None = Query(default=None), ): await websocket.accept() if not await _authorize(websocket, token): return try: await _stream_logs(websocket, stack_id, service) except WebSocketDisconnect: pass @router.websocket("/ws/events") async def ws_events( websocket: WebSocket, token: str | None = Query(default=None), ): """Stream global Docker events (decoded subset).""" await websocket.accept() if not await _authorize(websocket, token): return from docker_client import get_client loop = asyncio.get_event_loop() queue: asyncio.Queue = asyncio.Queue() stop = asyncio.Event() def reader(): try: client = get_client() for event in client.events(decode=True): if stop.is_set(): break loop.call_soon_threadsafe(queue.put_nowait, event) except Exception: # noqa: BLE001 pass task = loop.run_in_executor(None, reader) try: while True: event = await queue.get() actor = event.get("Actor", {}) or {} attrs = actor.get("Attributes", {}) or {} await websocket.send_text( json.dumps( { "type": "event", "action": event.get("Action"), "container": attrs.get("name"), "stack": attrs.get("com.docker.compose.project"), } ) ) except WebSocketDisconnect: pass finally: stop.set() task.cancel()