Initial commit: StackPilot Phase 1 (Core)
Self-hosted Docker Compose manager. - Backend: FastAPI + docker-py + SQLite (JWT auth, file-first stacks, lifecycle, live status, WebSocket logs, docker-run converter, audit log) - Frontend: React + Vite + Tailwind (login/setup, dashboard, stacks, stack detail, Monaco editor, dark/light theme) - Deployment: docker-compose.yml, Dockerfiles, nginx reverse proxy Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
This commit is contained in:
@@ -0,0 +1,29 @@
|
||||
"""Audit log query endpoint."""
|
||||
from __future__ import annotations
|
||||
|
||||
from typing import Optional
|
||||
|
||||
from fastapi import APIRouter, Depends, Query
|
||||
from sqlmodel import Session, select
|
||||
|
||||
from auth import get_current_user
|
||||
from database import get_session
|
||||
from models.audit import AuditLog
|
||||
from models.user import User
|
||||
|
||||
router = APIRouter(prefix="/api/audit", tags=["audit"])
|
||||
|
||||
|
||||
@router.get("")
|
||||
def list_audit(
|
||||
limit: int = Query(100, le=500),
|
||||
offset: int = 0,
|
||||
stack_id: Optional[str] = None,
|
||||
session: Session = Depends(get_session),
|
||||
_user: User = Depends(get_current_user),
|
||||
) -> list[AuditLog]:
|
||||
stmt = select(AuditLog).order_by(AuditLog.timestamp.desc())
|
||||
if stack_id:
|
||||
stmt = stmt.where(AuditLog.target == stack_id)
|
||||
stmt = stmt.offset(offset).limit(limit)
|
||||
return session.exec(stmt).all()
|
||||
@@ -0,0 +1,111 @@
|
||||
"""Authentication routes + first-launch setup wizard."""
|
||||
from __future__ import annotations
|
||||
|
||||
import time
|
||||
from collections import defaultdict, deque
|
||||
|
||||
from fastapi import APIRouter, Depends, HTTPException, Request, status
|
||||
from sqlmodel import Session
|
||||
|
||||
import auth as auth_mod
|
||||
from database import get_session
|
||||
from models.user import (
|
||||
LoginRequest,
|
||||
RefreshRequest,
|
||||
TokenPair,
|
||||
User,
|
||||
UserCreate,
|
||||
UserRead,
|
||||
)
|
||||
from services import audit_service
|
||||
|
||||
router = APIRouter(prefix="/api/auth", tags=["auth"])
|
||||
|
||||
# Simple in-memory rate limiter for login (max 10 / minute / IP).
|
||||
_LOGIN_HITS: dict[str, deque] = defaultdict(deque)
|
||||
_RATE_LIMIT = 10
|
||||
_RATE_WINDOW = 60.0
|
||||
|
||||
|
||||
def _check_rate_limit(ip: str) -> None:
|
||||
now = time.monotonic()
|
||||
hits = _LOGIN_HITS[ip]
|
||||
while hits and now - hits[0] > _RATE_WINDOW:
|
||||
hits.popleft()
|
||||
if len(hits) >= _RATE_LIMIT:
|
||||
raise HTTPException(
|
||||
status_code=status.HTTP_429_TOO_MANY_REQUESTS,
|
||||
detail="Too many login attempts, slow down.",
|
||||
)
|
||||
hits.append(now)
|
||||
|
||||
|
||||
def _tokens_for(user: User) -> TokenPair:
|
||||
return TokenPair(
|
||||
access_token=auth_mod.create_access_token(user),
|
||||
refresh_token=auth_mod.create_refresh_token(user),
|
||||
)
|
||||
|
||||
|
||||
@router.get("/needs-setup")
|
||||
def needs_setup(session: Session = Depends(get_session)) -> dict:
|
||||
"""First-launch wizard check: True if no users exist yet."""
|
||||
return {"needs_setup": not auth_mod.users_exist(session)}
|
||||
|
||||
|
||||
@router.post("/setup", response_model=TokenPair)
|
||||
def setup(
|
||||
body: UserCreate, session: Session = Depends(get_session)
|
||||
) -> TokenPair:
|
||||
if auth_mod.users_exist(session):
|
||||
raise HTTPException(status_code=400, detail="Setup already completed")
|
||||
user = User(
|
||||
username=body.username,
|
||||
hashed_password=auth_mod.hash_password(body.password),
|
||||
role="admin",
|
||||
)
|
||||
session.add(user)
|
||||
session.commit()
|
||||
session.refresh(user)
|
||||
audit_service.record(
|
||||
session, user=user.username, action="user.setup", target=user.username
|
||||
)
|
||||
return _tokens_for(user)
|
||||
|
||||
|
||||
@router.post("/login", response_model=TokenPair)
|
||||
def login(
|
||||
body: LoginRequest,
|
||||
request: Request,
|
||||
session: Session = Depends(get_session),
|
||||
) -> TokenPair:
|
||||
ip = request.client.host if request.client else "unknown"
|
||||
_check_rate_limit(ip)
|
||||
user = auth_mod.authenticate(session, body.username, body.password)
|
||||
if not user:
|
||||
raise HTTPException(
|
||||
status_code=status.HTTP_401_UNAUTHORIZED,
|
||||
detail="Incorrect username or password",
|
||||
)
|
||||
audit_service.record(
|
||||
session, user=user.username, action="auth.login", target=user.username, ip=ip
|
||||
)
|
||||
return _tokens_for(user)
|
||||
|
||||
|
||||
@router.post("/refresh", response_model=TokenPair)
|
||||
def refresh(
|
||||
body: RefreshRequest, session: Session = Depends(get_session)
|
||||
) -> TokenPair:
|
||||
payload = auth_mod.decode_token(body.refresh_token, "refresh")
|
||||
user = auth_mod.get_user(session, payload.get("sub", ""))
|
||||
if not user or not user.is_active:
|
||||
raise HTTPException(
|
||||
status_code=status.HTTP_401_UNAUTHORIZED, detail="Invalid refresh token"
|
||||
)
|
||||
return _tokens_for(user)
|
||||
|
||||
|
||||
@router.get("/me", response_model=UserRead)
|
||||
def me(user: User = Depends(auth_mod.get_current_user)) -> User:
|
||||
return user
|
||||
@@ -0,0 +1,334 @@
|
||||
"""Stack CRUD + lifecycle endpoints."""
|
||||
from __future__ import annotations
|
||||
|
||||
import os
|
||||
from dataclasses import asdict
|
||||
|
||||
from fastapi import APIRouter, Depends, HTTPException, Query, Request
|
||||
from fastapi.responses import FileResponse
|
||||
from sqlmodel import Session, select
|
||||
|
||||
from auth import get_current_user, require_admin
|
||||
from database import get_session
|
||||
from docker_client import DockerError
|
||||
from models.stack import (
|
||||
ConvertRequest,
|
||||
ConvertResponse,
|
||||
Stack,
|
||||
StackCloneRequest,
|
||||
StackCreate,
|
||||
StackUpdate,
|
||||
)
|
||||
from models.user import User
|
||||
from services import audit_service, compose_service
|
||||
from services.convert_service import convert_docker_run
|
||||
|
||||
router = APIRouter(prefix="/api/stacks", tags=["stacks"])
|
||||
|
||||
|
||||
# --------------------------------------------------------------------------- #
|
||||
# helpers
|
||||
# --------------------------------------------------------------------------- #
|
||||
|
||||
|
||||
def _client_ip(request: Request) -> str:
|
||||
return request.client.host if request.client else "unknown"
|
||||
|
||||
|
||||
def sync_discovered_stacks(session: Session) -> None:
|
||||
"""Register any on-disk stacks not yet in the database."""
|
||||
known = {s.id for s in session.exec(select(Stack)).all()}
|
||||
for stack_id in compose_service.discover_stacks():
|
||||
if stack_id not in known:
|
||||
stack = Stack(id=stack_id, name=stack_id)
|
||||
session.add(stack)
|
||||
session.commit()
|
||||
|
||||
|
||||
def _get_stack_or_404(session: Session, stack_id: str) -> Stack:
|
||||
stack = session.get(Stack, stack_id)
|
||||
if not stack:
|
||||
raise HTTPException(status_code=404, detail=f"Stack '{stack_id}' not found")
|
||||
return stack
|
||||
|
||||
|
||||
def _stack_summary(stack: Stack) -> dict:
|
||||
try:
|
||||
containers = compose_service.containers_for_stack(stack.id)
|
||||
status = compose_service.compute_status(stack.id)
|
||||
except DockerError:
|
||||
containers = []
|
||||
status = "unknown"
|
||||
return {
|
||||
"id": stack.id,
|
||||
"name": stack.name,
|
||||
"description": stack.description,
|
||||
"status": status,
|
||||
"service_count": len(containers),
|
||||
"running_count": sum(1 for c in containers if c.state == "running"),
|
||||
"created_at": stack.created_at,
|
||||
"updated_at": stack.updated_at,
|
||||
}
|
||||
|
||||
|
||||
# --------------------------------------------------------------------------- #
|
||||
# CRUD
|
||||
# --------------------------------------------------------------------------- #
|
||||
|
||||
|
||||
@router.get("")
|
||||
def list_stacks(
|
||||
session: Session = Depends(get_session),
|
||||
_user: User = Depends(get_current_user),
|
||||
) -> list[dict]:
|
||||
sync_discovered_stacks(session)
|
||||
stacks = session.exec(select(Stack)).all()
|
||||
return [_stack_summary(s) for s in stacks]
|
||||
|
||||
|
||||
@router.post("", status_code=201)
|
||||
def create_stack(
|
||||
body: StackCreate,
|
||||
request: Request,
|
||||
session: Session = Depends(get_session),
|
||||
user: User = Depends(require_admin),
|
||||
) -> dict:
|
||||
stack_id = compose_service.slugify(body.name)
|
||||
if session.get(Stack, stack_id) or os.path.isdir(compose_service.stack_dir(stack_id)):
|
||||
raise HTTPException(status_code=409, detail=f"Stack '{stack_id}' already exists")
|
||||
compose_service.write_compose(stack_id, body.yaml or "services:\n")
|
||||
if body.env:
|
||||
compose_service.write_env(stack_id, body.env)
|
||||
stack = Stack(id=stack_id, name=body.name, description=body.description)
|
||||
session.add(stack)
|
||||
session.commit()
|
||||
session.refresh(stack)
|
||||
audit_service.record(
|
||||
session, user=user.username, action="stack.create", target=stack_id,
|
||||
ip=_client_ip(request),
|
||||
)
|
||||
return _stack_summary(stack)
|
||||
|
||||
|
||||
@router.get("/{stack_id}")
|
||||
def get_stack(
|
||||
stack_id: str,
|
||||
session: Session = Depends(get_session),
|
||||
_user: User = Depends(get_current_user),
|
||||
) -> dict:
|
||||
stack = _get_stack_or_404(session, stack_id)
|
||||
try:
|
||||
containers = [asdict(c) for c in compose_service.containers_for_stack(stack_id)]
|
||||
status = compose_service.compute_status(stack_id)
|
||||
except DockerError as exc:
|
||||
containers = []
|
||||
status = "unknown"
|
||||
return {
|
||||
"id": stack.id,
|
||||
"name": stack.name,
|
||||
"description": stack.description,
|
||||
"status": status,
|
||||
"yaml": compose_service.read_compose(stack_id),
|
||||
"env": compose_service.read_env(stack_id),
|
||||
"containers": containers,
|
||||
"created_at": stack.created_at,
|
||||
"updated_at": stack.updated_at,
|
||||
}
|
||||
|
||||
|
||||
@router.put("/{stack_id}")
|
||||
def update_stack(
|
||||
stack_id: str,
|
||||
body: StackUpdate,
|
||||
request: Request,
|
||||
session: Session = Depends(get_session),
|
||||
user: User = Depends(require_admin),
|
||||
) -> dict:
|
||||
stack = _get_stack_or_404(session, stack_id)
|
||||
if body.yaml is not None:
|
||||
compose_service.write_compose(stack_id, body.yaml)
|
||||
if body.env is not None:
|
||||
compose_service.write_env(stack_id, body.env)
|
||||
if body.name is not None:
|
||||
stack.name = body.name
|
||||
if body.description is not None:
|
||||
stack.description = body.description
|
||||
stack.updated_at = compose_service.now()
|
||||
session.add(stack)
|
||||
session.commit()
|
||||
session.refresh(stack)
|
||||
audit_service.record(
|
||||
session, user=user.username, action="stack.update", target=stack_id,
|
||||
ip=_client_ip(request),
|
||||
)
|
||||
return _stack_summary(stack)
|
||||
|
||||
|
||||
@router.delete("/{stack_id}")
|
||||
async def delete_stack(
|
||||
stack_id: str,
|
||||
request: Request,
|
||||
delete_files: bool = Query(True),
|
||||
session: Session = Depends(get_session),
|
||||
user: User = Depends(require_admin),
|
||||
) -> dict:
|
||||
stack = _get_stack_or_404(session, stack_id)
|
||||
try:
|
||||
await compose_service.down(stack_id)
|
||||
except Exception: # noqa: BLE001 - best-effort teardown
|
||||
pass
|
||||
if delete_files:
|
||||
compose_service.delete_stack_files(stack_id)
|
||||
session.delete(stack)
|
||||
session.commit()
|
||||
audit_service.record(
|
||||
session, user=user.username, action="stack.delete", target=stack_id,
|
||||
detail=f"delete_files={delete_files}", ip=_client_ip(request),
|
||||
)
|
||||
return {"ok": True}
|
||||
|
||||
|
||||
@router.post("/{stack_id}/clone")
|
||||
def clone_stack(
|
||||
stack_id: str,
|
||||
body: StackCloneRequest,
|
||||
request: Request,
|
||||
session: Session = Depends(get_session),
|
||||
user: User = Depends(require_admin),
|
||||
) -> dict:
|
||||
_get_stack_or_404(session, stack_id)
|
||||
new_id = compose_service.slugify(body.name)
|
||||
if session.get(Stack, new_id):
|
||||
raise HTTPException(status_code=409, detail=f"Stack '{new_id}' already exists")
|
||||
try:
|
||||
compose_service.clone_stack_files(stack_id, new_id)
|
||||
except compose_service.StackFileError as exc:
|
||||
raise HTTPException(status_code=409, detail=str(exc)) from exc
|
||||
stack = Stack(id=new_id, name=body.name)
|
||||
session.add(stack)
|
||||
session.commit()
|
||||
session.refresh(stack)
|
||||
audit_service.record(
|
||||
session, user=user.username, action="stack.clone",
|
||||
target=new_id, detail=f"from {stack_id}", ip=_client_ip(request),
|
||||
)
|
||||
return _stack_summary(stack)
|
||||
|
||||
|
||||
# --------------------------------------------------------------------------- #
|
||||
# lifecycle
|
||||
# --------------------------------------------------------------------------- #
|
||||
|
||||
|
||||
async def _lifecycle(action_fn, action_name, stack_id, request, session, user):
|
||||
_get_stack_or_404(session, stack_id)
|
||||
result = await action_fn(stack_id)
|
||||
audit_service.record(
|
||||
session, user=user.username, action=f"stack.{action_name}", target=stack_id,
|
||||
detail=f"rc={result.get('returncode')}", ip=_client_ip(request),
|
||||
)
|
||||
if result.get("returncode") not in (0, None):
|
||||
raise HTTPException(
|
||||
status_code=500,
|
||||
detail={
|
||||
"error": f"compose {action_name} failed",
|
||||
"detail": result.get("stderr", "").strip()[-2000:],
|
||||
},
|
||||
)
|
||||
return result
|
||||
|
||||
|
||||
@router.post("/{stack_id}/start")
|
||||
async def start_stack(stack_id: str, request: Request, session: Session = Depends(get_session), user: User = Depends(require_admin)):
|
||||
return await _lifecycle(compose_service.up, "start", stack_id, request, session, user)
|
||||
|
||||
|
||||
@router.post("/{stack_id}/stop")
|
||||
async def stop_stack(stack_id: str, request: Request, session: Session = Depends(get_session), user: User = Depends(require_admin)):
|
||||
return await _lifecycle(compose_service.stop, "stop", stack_id, request, session, user)
|
||||
|
||||
|
||||
@router.post("/{stack_id}/restart")
|
||||
async def restart_stack(stack_id: str, request: Request, session: Session = Depends(get_session), user: User = Depends(require_admin)):
|
||||
return await _lifecycle(compose_service.restart, "restart", stack_id, request, session, user)
|
||||
|
||||
|
||||
@router.post("/{stack_id}/pull")
|
||||
async def pull_stack(stack_id: str, request: Request, session: Session = Depends(get_session), user: User = Depends(require_admin)):
|
||||
return await _lifecycle(compose_service.pull, "pull", stack_id, request, session, user)
|
||||
|
||||
|
||||
@router.post("/{stack_id}/update")
|
||||
async def update_stack_images(stack_id: str, request: Request, session: Session = Depends(get_session), user: User = Depends(require_admin)):
|
||||
return await _lifecycle(compose_service.update, "update", stack_id, request, session, user)
|
||||
|
||||
|
||||
@router.post("/{stack_id}/down")
|
||||
async def down_stack(stack_id: str, request: Request, session: Session = Depends(get_session), user: User = Depends(require_admin)):
|
||||
return await _lifecycle(compose_service.down, "down", stack_id, request, session, user)
|
||||
|
||||
|
||||
# --------------------------------------------------------------------------- #
|
||||
# logs / export / convert
|
||||
# --------------------------------------------------------------------------- #
|
||||
|
||||
|
||||
@router.get("/{stack_id}/logs")
|
||||
async def stack_logs(
|
||||
stack_id: str,
|
||||
tail: int = Query(200, le=2000),
|
||||
session: Session = Depends(get_session),
|
||||
_user: User = Depends(get_current_user),
|
||||
) -> dict:
|
||||
_get_stack_or_404(session, stack_id)
|
||||
result = await compose_service.logs(stack_id, tail=tail)
|
||||
return {"logs": result.get("stdout", "") + result.get("stderr", "")}
|
||||
|
||||
|
||||
@router.get("/{stack_id}/services/{service}/logs")
|
||||
async def service_logs(
|
||||
stack_id: str,
|
||||
service: str,
|
||||
tail: int = Query(200, le=2000),
|
||||
session: Session = Depends(get_session),
|
||||
_user: User = Depends(get_current_user),
|
||||
) -> dict:
|
||||
_get_stack_or_404(session, stack_id)
|
||||
result = await compose_service.logs(stack_id, service=service, tail=tail)
|
||||
return {"logs": result.get("stdout", "") + result.get("stderr", "")}
|
||||
|
||||
|
||||
@router.get("/{stack_id}/export")
|
||||
def export_stack(
|
||||
stack_id: str,
|
||||
session: Session = Depends(get_session),
|
||||
_user: User = Depends(get_current_user),
|
||||
):
|
||||
import io
|
||||
import tarfile
|
||||
import tempfile
|
||||
|
||||
stack = _get_stack_or_404(session, stack_id)
|
||||
directory = compose_service.stack_dir(stack_id)
|
||||
if not os.path.isdir(directory):
|
||||
raise HTTPException(status_code=404, detail="Stack directory missing")
|
||||
tmp = tempfile.NamedTemporaryFile(delete=False, suffix=".tar.gz")
|
||||
with tarfile.open(tmp.name, "w:gz") as tar:
|
||||
tar.add(directory, arcname=stack_id)
|
||||
date = compose_service.now().strftime("%Y%m%d")
|
||||
return FileResponse(
|
||||
tmp.name,
|
||||
media_type="application/gzip",
|
||||
filename=f"stack-{stack_id}-{date}.tar.gz",
|
||||
)
|
||||
|
||||
|
||||
@router.post("/convert", response_model=ConvertResponse)
|
||||
def convert(
|
||||
body: ConvertRequest,
|
||||
_user: User = Depends(get_current_user),
|
||||
) -> ConvertResponse:
|
||||
try:
|
||||
return ConvertResponse(yaml=convert_docker_run(body.command))
|
||||
except ValueError as exc:
|
||||
raise HTTPException(status_code=400, detail=str(exc)) from exc
|
||||
@@ -0,0 +1,92 @@
|
||||
"""Host / Docker system information."""
|
||||
from __future__ import annotations
|
||||
|
||||
import os
|
||||
import shutil
|
||||
|
||||
from fastapi import APIRouter, Depends
|
||||
|
||||
from auth import get_current_user
|
||||
from config import settings
|
||||
from docker_client import DockerError, get_client, safe_call
|
||||
from models.user import User
|
||||
|
||||
router = APIRouter(prefix="/api/system", tags=["system"])
|
||||
|
||||
|
||||
def _read_proc(path: str) -> str:
|
||||
full = os.path.join(settings.HOST_PROC_PATH, path)
|
||||
if not os.path.isfile(full):
|
||||
full = os.path.join("/proc", path)
|
||||
try:
|
||||
with open(full, "r", encoding="utf-8") as fh:
|
||||
return fh.read()
|
||||
except OSError:
|
||||
return ""
|
||||
|
||||
|
||||
def _mem_info() -> dict:
|
||||
info = {}
|
||||
for line in _read_proc("meminfo").splitlines():
|
||||
parts = line.split(":")
|
||||
if len(parts) == 2:
|
||||
key = parts[0].strip()
|
||||
val = parts[1].strip().split()[0]
|
||||
try:
|
||||
info[key] = int(val) * 1024 # kB -> bytes
|
||||
except ValueError:
|
||||
pass
|
||||
total = info.get("MemTotal", 0)
|
||||
available = info.get("MemAvailable", info.get("MemFree", 0))
|
||||
return {"total": total, "available": available, "used": max(total - available, 0)}
|
||||
|
||||
|
||||
def _uptime() -> float:
|
||||
raw = _read_proc("uptime")
|
||||
try:
|
||||
return float(raw.split()[0])
|
||||
except (IndexError, ValueError):
|
||||
return 0.0
|
||||
|
||||
|
||||
def _cpu_count() -> int:
|
||||
return os.cpu_count() or 0
|
||||
|
||||
|
||||
def _disk_usage() -> dict:
|
||||
try:
|
||||
usage = shutil.disk_usage(settings.DATA_DIR)
|
||||
return {"total": usage.total, "used": usage.used, "free": usage.free}
|
||||
except OSError:
|
||||
return {"total": 0, "used": 0, "free": 0}
|
||||
|
||||
|
||||
@router.get("/info")
|
||||
def system_info(_user: User = Depends(get_current_user)) -> dict:
|
||||
docker_version = ""
|
||||
host_os = ""
|
||||
containers_running = 0
|
||||
containers_total = 0
|
||||
try:
|
||||
client = get_client()
|
||||
version = safe_call(client.version)
|
||||
docker_version = version.get("Version", "")
|
||||
info = safe_call(client.info)
|
||||
host_os = info.get("OperatingSystem", "")
|
||||
containers_running = info.get("ContainersRunning", 0)
|
||||
containers_total = info.get("Containers", 0)
|
||||
except DockerError as exc:
|
||||
docker_version = f"unavailable ({exc.error})"
|
||||
|
||||
return {
|
||||
"docker_version": docker_version,
|
||||
"host_os": host_os,
|
||||
"hostname": os.uname().nodename,
|
||||
"cpu_cores": _cpu_count(),
|
||||
"ram": _mem_info(),
|
||||
"disk": _disk_usage(),
|
||||
"uptime_seconds": _uptime(),
|
||||
"containers_running": containers_running,
|
||||
"containers_total": containers_total,
|
||||
"gpus": [], # populated in Phase 2
|
||||
}
|
||||
@@ -0,0 +1,130 @@
|
||||
"""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()
|
||||
Reference in New Issue
Block a user