"""WebSocket endpoints for real-time log streaming and Docker events.""" from __future__ import annotations import asyncio import json import logging import contextlib from fastapi import APIRouter, Query, WebSocket, WebSocketDisconnect from jose import JWTError from sqlmodel import Session from auth import decode_token, resolve_token_user from database import engine from models.setting import EVENT_PULL_FAILED, EVENT_STACK_ERROR, EVENT_STACK_START from services import ( audit_service, compose_service, exec_service, notify_service, stack_lock_service, update_service, ) logger = logging.getLogger("stackpilot.ws") router = APIRouter(tags=["ws"]) def _user_for(token: str | None): """The live user behind a socket's token, or None. Resolves against the database rather than reading the role straight off the JWT: a socket can outlive a demotion, a disabled account or a password reset, and the exec endpoint below is root-equivalent on the host. Same check the HTTP routes make. """ if not token: return None try: payload = decode_token(token, "access") except (JWTError, Exception): # noqa: BLE001 return None with Session(engine) as session: user = resolve_token_user(session, payload) if user: session.expunge(user) return user async def _authorize(websocket: WebSocket, token: str | None) -> bool: """Validate the JWT supplied as a query param. Closes socket on failure.""" if _user_for(token) is None: await websocket.close(code=4401) return False return True async def _authorize_admin(websocket: WebSocket, token: str | None) -> bool: """Like _authorize but also requires the admin role (exec is root-equivalent). Closes 4401 on a missing/invalid token, 4403 on a valid non-admin token.""" user = _user_for(token) if user is None: await websocket.close(code=4401) return False if user.role != "admin": await websocket.close(code=4403) 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/deploy/{stack_id}") async def ws_deploy( websocket: WebSocket, stack_id: str, token: str | None = Query(default=None), ): """Run `docker compose up -d` and stream its output (image pulls, container creation) to the browser so the user sees deploy progress live. Records the same audit entry and notification as the REST `/start` endpoint.""" await websocket.accept() if not await _authorize(websocket, token): return username = decode_token(token, "access").get("sub", "unknown") if token else "unknown" rc: int | None = None disconnected = False # Same guard the REST lifecycle uses — the deploy console runs the very # same `compose up`, so it has to queue behind an in-flight operation # rather than race it. lock_session = Session(engine) try: stack_lock_service.acquire(lock_session, stack_id, "start", username) except stack_lock_service.StackBusy as exc: lock_session.close() with contextlib.suppress(Exception): await websocket.send_text(json.dumps({"type": "error", "detail": str(exc)})) await websocket.close(code=4409) return try: async for kind, payload in compose_service.stream_up(stack_id): if kind == "log": await websocket.send_text(json.dumps({"type": "log", "line": payload})) else: rc = payload await websocket.send_text(json.dumps({"type": "done", "returncode": rc})) except WebSocketDisconnect: # Client navigated away; the compose subprocess keeps running so the # deploy still completes in the background. disconnected = True except Exception as exc: # noqa: BLE001 with contextlib.suppress(Exception): await websocket.send_text(json.dumps({"type": "error", "detail": str(exc)})) finally: with contextlib.suppress(Exception): stack_lock_service.release(lock_session, stack_id) lock_session.close() ok = rc in (0, None) try: with Session(engine) as session: audit_service.record( session, user=username, action="stack.start", target=stack_id, detail=f"rc={rc} (deploy console)", ip="ws", ) if ok: await notify_service.notify( EVENT_STACK_START, f"Stack '{stack_id}' started", "compose up completed successfully.", session, ) else: await notify_service.notify( EVENT_STACK_ERROR, f"Stack '{stack_id}' start failed", "compose up returned a non-zero exit code.", session, ) except Exception: # noqa: BLE001 - audit/notify are best-effort pass if not disconnected: with contextlib.suppress(Exception): await websocket.close() @router.websocket("/ws/update/{stack_id}") async def ws_update( websocket: WebSocket, stack_id: str, token: str | None = Query(default=None), ): """Run `docker compose pull && up -d` and stream its output, so the stacks list can render real update progress. Same audit/notify contract as the REST `/update` endpoint, which stays for non-interactive callers.""" await websocket.accept() if not await _authorize_admin(websocket, token): return username = decode_token(token, "access").get("sub", "unknown") if token else "unknown" rc: int | None = None disconnected = False lock_session = Session(engine) try: stack_lock_service.acquire(lock_session, stack_id, "update", username) except stack_lock_service.StackBusy as exc: lock_session.close() with contextlib.suppress(Exception): await websocket.send_text(json.dumps({"type": "error", "detail": str(exc)})) await websocket.close(code=4409) return try: async for kind, payload in compose_service.stream_update(stack_id): if kind == "log": await websocket.send_text(json.dumps({"type": "log", "line": payload})) else: rc = payload await websocket.send_text(json.dumps({"type": "done", "returncode": rc})) except WebSocketDisconnect: # Client navigated away; compose keeps running so the update finishes. disconnected = True except Exception as exc: # noqa: BLE001 with contextlib.suppress(Exception): await websocket.send_text(json.dumps({"type": "error", "detail": str(exc)})) finally: with contextlib.suppress(Exception): stack_lock_service.release(lock_session, stack_id) lock_session.close() ok = rc in (0, None) try: with Session(engine) as session: audit_service.record( session, user=username, action="stack.update", target=stack_id, detail=f"rc={rc} (update stream)", ip="ws", ) if ok: await notify_service.notify( EVENT_STACK_START, f"Stack '{stack_id}' updated", "compose pull + up completed successfully.", session, ) else: await notify_service.notify( EVENT_PULL_FAILED, f"Stack '{stack_id}' update failed", "compose pull/up returned a non-zero exit code.", session, ) except Exception: # noqa: BLE001 - audit/notify are best-effort pass if not disconnected: with contextlib.suppress(Exception): await websocket.close() if ok: update_service.refresh_stack_local(stack_id) @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() @router.websocket("/ws/exec/{container_id}") async def ws_exec( websocket: WebSocket, container_id: str, token: str | None = Query(default=None), cmd: str | None = Query(default=None), ): """Interactive shell into a compose-managed container (admin only).""" await websocket.accept() if not await _authorize_admin(websocket, token): return username = decode_token(token, "access").get("sub", "unknown") if token else "unknown" shell = cmd or exec_service.DEFAULT_SHELL try: exec_id = exec_service.create_exec(container_id, [shell]) holder, raw = exec_service.start_exec(exec_id) except Exception as exc: # noqa: BLE001 with contextlib.suppress(Exception): await websocket.send_text(json.dumps({"type": "error", "detail": str(exc)})) with contextlib.suppress(Exception): await websocket.close() return with contextlib.suppress(Exception): with Session(engine) as session: audit_service.record( session, user=username, action="container.exec", target=container_id[:12], detail=shell, ip="ws", ) try: await exec_service.pump_exec(websocket, exec_id, holder, raw) except WebSocketDisconnect: pass finally: with contextlib.suppress(Exception): await websocket.close()