F7 — Nothing stopped two compose operations landing on the same stack. There was a busy flag, but is_busy() was only ever read to colour the status column; no lifecycle handler consulted it before acting. Two tabs, or auto-update picking up a stack somebody had just clicked, both ran pull + up -d against the same project and raced over recreating containers. Lifecycle calls, the two deploy WebSockets and the auto-update pass now take a real lock; a second caller gets 409 (or an error frame and close 4409) and auto-update skips and retries next cycle. The lock is a row rather than a set in one worker's memory, so it holds across workers and across a restart, and it carries an expiry — a worker killed mid-deploy would otherwise strand the stack with no fix short of editing the database. F10 — /api/stacks/stats sampled every running container on every call, one blocking daemon request each, and both the dashboard and the stacks list poll it every five seconds. Two tabs on a 40-container host meant a sustained ~16 samples a second. Cached for 4s behind a lock so concurrent callers share one sweep, the same shape dashboard_service already used for its fleet aggregate. F11 — Three module dicts assumed exactly one uvicorn worker without saying so and were lost on restart. The busy set is the lock above. The image update cache is now mirrored to SQLite, so a restart shows the badges immediately instead of blanking them for up to an hour, and the already-notified marks come back with them rather than re-announcing the same updates. The login rate limiter is a table, so it cannot be cleared by getting the process to restart and no longer multiplies by the worker count. The constraint that shaped this: compose_service and update_service are shared with the agent, which has no database. Neither may import one. So the lock is a separate service the central app enforces at its own entry points, and update persistence is an opt-in callback the central app registers in its lifespan — the agent registers nothing and behaves exactly as before. A test asserts update_service never imports the database, since that is the kind of thing a later change breaks silently. Both new nets were checked by reverting the fix: dropping the lock from _lifecycle fails six tests, removing the stats cache fails the one that names the behaviour. Also wires up cache pruning in the same sweep — without it both the dict and the table grew one entry per image tag ever run, for the life of the install. 31 new tests (729 total). Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01Dk43rmEeRfYi5wsLDbmfyG
605 lines
22 KiB
Python
605 lines
22 KiB
Python
"""WebSocket endpoints for real-time log streaming and Docker events."""
|
|
from __future__ import annotations
|
|
|
|
import asyncio
|
|
import json
|
|
import logging
|
|
import urllib.parse
|
|
|
|
import contextlib
|
|
|
|
import websockets
|
|
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.agent import Agent
|
|
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
|
|
compose_service.mark_busy(stack_id)
|
|
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:
|
|
compose_service.clear_busy(stack_id)
|
|
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/agent-logs/{agent_id}/{stack_id}")
|
|
async def ws_agent_logs(
|
|
websocket: WebSocket,
|
|
agent_id: int,
|
|
stack_id: str,
|
|
token: str | None = Query(default=None),
|
|
):
|
|
"""Proxy live compose logs from a remote agent through to the browser."""
|
|
await websocket.accept()
|
|
if not await _authorize(websocket, token):
|
|
return
|
|
|
|
with Session(engine) as session:
|
|
agent = session.get(Agent, agent_id)
|
|
if not agent:
|
|
await websocket.send_text(json.dumps({"type": "error", "detail": "agent not found"}))
|
|
await websocket.close()
|
|
return
|
|
|
|
base = agent.url.rstrip("/")
|
|
ws_url = ("wss://" + base[8:] if base.startswith("https://")
|
|
else "ws://" + base[7:] if base.startswith("http://")
|
|
else "ws://" + base)
|
|
# URL-encode the token: agent tokens may contain base64 chars (+ / =) that
|
|
# would otherwise be mangled in the query string and rejected as 4401.
|
|
ws_url += f"/agent/ws/logs/{stack_id}?token={urllib.parse.quote(agent.token, safe='')}"
|
|
|
|
async def _err(detail: str) -> None:
|
|
with contextlib.suppress(Exception):
|
|
await websocket.send_text(json.dumps({"type": "error", "detail": detail}))
|
|
|
|
# Connect to the agent. Surface connection problems (agent down, wrong URL,
|
|
# an outdated agent that lacks /agent/ws/logs, TLS issues) instead of
|
|
# silently dropping the socket.
|
|
try:
|
|
upstream = await websockets.connect(ws_url, open_timeout=10, ping_interval=20)
|
|
except websockets.InvalidStatus as exc:
|
|
code = getattr(getattr(exc, "response", None), "status_code", None)
|
|
hint = " — the agent may be running an old version without live-log support; update it." if code == 404 else ""
|
|
logger.warning("Agent log proxy: handshake to %s failed (%s)", agent.name, code)
|
|
await _err(f"Agent '{agent.name}' rejected the log stream (HTTP {code}){hint}")
|
|
with contextlib.suppress(Exception):
|
|
await websocket.close()
|
|
return
|
|
except Exception as exc: # noqa: BLE001
|
|
logger.warning("Agent log proxy: cannot reach %s at %s: %s", agent.name, agent.url, exc)
|
|
await _err(f"Could not connect to agent '{agent.name}' at {agent.url}: {exc}")
|
|
with contextlib.suppress(Exception):
|
|
await websocket.close()
|
|
return
|
|
|
|
try:
|
|
async for message in upstream:
|
|
await websocket.send_text(
|
|
message if isinstance(message, str) else message.decode("utf-8", "replace")
|
|
)
|
|
except WebSocketDisconnect:
|
|
pass
|
|
except websockets.ConnectionClosed as exc:
|
|
# Abnormal upstream close (e.g. 4401 bad token, or agent-side error).
|
|
if exc.code not in (1000, 1001):
|
|
await _err(f"Agent log stream closed unexpectedly (code {exc.code}).")
|
|
except Exception as exc: # noqa: BLE001
|
|
logger.warning("Agent log proxy: stream error from %s: %s", agent.name, exc)
|
|
await _err(str(exc))
|
|
finally:
|
|
with contextlib.suppress(Exception):
|
|
await upstream.close()
|
|
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
|
|
compose_service.mark_busy(stack_id)
|
|
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:
|
|
compose_service.clear_busy(stack_id)
|
|
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/agent-deploy/{agent_id}/{stack_id}")
|
|
async def ws_agent_deploy(
|
|
websocket: WebSocket,
|
|
agent_id: int,
|
|
stack_id: str,
|
|
token: str | None = Query(default=None),
|
|
):
|
|
"""Proxy a remote agent's `compose up` deploy stream through to the browser,
|
|
then record the same audit entry as the REST agent lifecycle endpoint."""
|
|
await websocket.accept()
|
|
if not await _authorize(websocket, token):
|
|
return
|
|
username = decode_token(token, "access").get("sub", "unknown") if token else "unknown"
|
|
|
|
with Session(engine) as session:
|
|
agent = session.get(Agent, agent_id)
|
|
if not agent:
|
|
await websocket.send_text(json.dumps({"type": "error", "detail": "agent not found"}))
|
|
await websocket.close()
|
|
return
|
|
|
|
base = agent.url.rstrip("/")
|
|
ws_url = ("wss://" + base[8:] if base.startswith("https://")
|
|
else "ws://" + base[7:] if base.startswith("http://")
|
|
else "ws://" + base)
|
|
ws_url += f"/agent/ws/deploy/{stack_id}?token={urllib.parse.quote(agent.token, safe='')}"
|
|
|
|
async def _err(detail: str) -> None:
|
|
with contextlib.suppress(Exception):
|
|
await websocket.send_text(json.dumps({"type": "error", "detail": detail}))
|
|
|
|
try:
|
|
upstream = await websockets.connect(ws_url, open_timeout=10, ping_interval=20)
|
|
except websockets.InvalidStatus as exc:
|
|
code = getattr(getattr(exc, "response", None), "status_code", None)
|
|
hint = " — the agent may be running an old version without deploy-console support; update it." if code == 404 else ""
|
|
logger.warning("Agent deploy proxy: handshake to %s failed (%s)", agent.name, code)
|
|
await _err(f"Agent '{agent.name}' rejected the deploy stream (HTTP {code}){hint}")
|
|
with contextlib.suppress(Exception):
|
|
await websocket.close()
|
|
return
|
|
except Exception as exc: # noqa: BLE001
|
|
logger.warning("Agent deploy proxy: cannot reach %s at %s: %s", agent.name, agent.url, exc)
|
|
await _err(f"Could not connect to agent '{agent.name}' at {agent.url}: {exc}")
|
|
with contextlib.suppress(Exception):
|
|
await websocket.close()
|
|
return
|
|
|
|
rc: int | None = None
|
|
try:
|
|
async for message in upstream:
|
|
text = message if isinstance(message, str) else message.decode("utf-8", "replace")
|
|
with contextlib.suppress(Exception):
|
|
msg = json.loads(text)
|
|
if msg.get("type") == "done":
|
|
rc = msg.get("returncode")
|
|
await websocket.send_text(text)
|
|
except WebSocketDisconnect:
|
|
pass
|
|
except websockets.ConnectionClosed as exc:
|
|
if exc.code not in (1000, 1001):
|
|
await _err(f"Agent deploy stream closed unexpectedly (code {exc.code}).")
|
|
except Exception as exc: # noqa: BLE001
|
|
logger.warning("Agent deploy proxy: stream error from %s: %s", agent.name, exc)
|
|
await _err(str(exc))
|
|
finally:
|
|
with contextlib.suppress(Exception):
|
|
await upstream.close()
|
|
with contextlib.suppress(Exception):
|
|
await websocket.close()
|
|
|
|
with contextlib.suppress(Exception):
|
|
with Session(engine) as session:
|
|
audit_service.record(
|
|
session, user=username, action="agent.stack.start",
|
|
target=f"{agent.name}/{stack_id}", detail=f"rc={rc} (deploy console)", ip="ws",
|
|
)
|
|
|
|
|
|
@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()
|
|
|
|
|
|
@router.websocket("/ws/agent-exec/{agent_id}/{container_id}")
|
|
async def ws_agent_exec(
|
|
websocket: WebSocket,
|
|
agent_id: int,
|
|
container_id: str,
|
|
token: str | None = Query(default=None),
|
|
cmd: str | None = Query(default=None),
|
|
):
|
|
"""Proxy an interactive exec session to a remote agent (admin only).
|
|
|
|
Unlike the log/deploy proxies this forwards in BOTH directions so keystrokes
|
|
reach the container and its output streams back."""
|
|
await websocket.accept()
|
|
if not await _authorize_admin(websocket, token):
|
|
return
|
|
|
|
with Session(engine) as session:
|
|
agent = session.get(Agent, agent_id)
|
|
if not agent:
|
|
await websocket.send_text(json.dumps({"type": "error", "detail": "agent not found"}))
|
|
await websocket.close()
|
|
return
|
|
|
|
base = agent.url.rstrip("/")
|
|
ws_url = ("wss://" + base[8:] if base.startswith("https://")
|
|
else "ws://" + base[7:] if base.startswith("http://")
|
|
else "ws://" + base)
|
|
ws_url += f"/agent/ws/exec/{container_id}?token={urllib.parse.quote(agent.token, safe='')}"
|
|
if cmd:
|
|
ws_url += f"&cmd={urllib.parse.quote(cmd, safe='')}"
|
|
|
|
async def _err(detail: str) -> None:
|
|
with contextlib.suppress(Exception):
|
|
await websocket.send_text(json.dumps({"type": "error", "detail": detail}))
|
|
|
|
try:
|
|
upstream = await websockets.connect(ws_url, open_timeout=10, ping_interval=20)
|
|
except websockets.InvalidStatus as exc:
|
|
code = getattr(getattr(exc, "response", None), "status_code", None)
|
|
hint = " — the agent may be running an old version without terminal support; update it." if code == 404 else ""
|
|
logger.warning("Agent exec proxy: handshake to %s failed (%s)", agent.name, code)
|
|
await _err(f"Agent '{agent.name}' rejected the terminal (HTTP {code}){hint}")
|
|
with contextlib.suppress(Exception):
|
|
await websocket.close()
|
|
return
|
|
except Exception as exc: # noqa: BLE001
|
|
logger.warning("Agent exec proxy: cannot reach %s at %s: %s", agent.name, agent.url, exc)
|
|
await _err(f"Could not connect to agent '{agent.name}' at {agent.url}: {exc}")
|
|
with contextlib.suppress(Exception):
|
|
await websocket.close()
|
|
return
|
|
|
|
with contextlib.suppress(Exception):
|
|
with Session(engine) as session:
|
|
username = decode_token(token, "access").get("sub", "unknown") if token else "unknown"
|
|
audit_service.record(
|
|
session, user=username, action="agent.container.exec",
|
|
target=f"{agent.name}/{container_id[:12]}", ip="ws",
|
|
)
|
|
|
|
async def browser_to_agent() -> None:
|
|
try:
|
|
while True:
|
|
msg = await websocket.receive_text()
|
|
await upstream.send(msg)
|
|
except (WebSocketDisconnect, websockets.ConnectionClosed):
|
|
pass
|
|
|
|
async def agent_to_browser() -> None:
|
|
try:
|
|
async for message in upstream:
|
|
await websocket.send_text(
|
|
message if isinstance(message, str) else message.decode("utf-8", "replace")
|
|
)
|
|
except (WebSocketDisconnect, websockets.ConnectionClosed):
|
|
pass
|
|
|
|
b2a = asyncio.create_task(browser_to_agent())
|
|
a2b = asyncio.create_task(agent_to_browser())
|
|
done, pending = await asyncio.wait({b2a, a2b}, return_when=asyncio.FIRST_COMPLETED)
|
|
for task in pending:
|
|
task.cancel()
|
|
with contextlib.suppress(Exception):
|
|
await asyncio.gather(*pending, return_exceptions=True)
|
|
with contextlib.suppress(Exception):
|
|
await upstream.close()
|
|
with contextlib.suppress(Exception):
|
|
await websocket.close()
|