Phase 4: backups w/ volumes, notifications, settings & users, audit page (0.4.0)
- Backup/restore: per-stack tar.gz incl. named-volume snapshots (helper container), upload restore with rename/overwrite/conflict detection. - Notifications: ntfy/Discord/Slack/Gotify/generic webhooks, per-event subscriptions; wired into the update checker and stack lifecycle. - Settings page: update-check interval, webhook CRUD + test, user management (with last-admin safeguards). - Audit log page (searchable, paginated). - Mobile-responsive sidebar/layout. Multi-host agents and remote backup destinations (SFTP/S3) deferred. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
This commit is contained in:
co-authored by
Claude Opus 4.8
parent
22d9864436
commit
8d19b09abd
@@ -32,6 +32,9 @@ class Settings(BaseSettings):
|
||||
DOCKER_SOCKET: str = "/var/run/docker.sock"
|
||||
HOST_PROC_PATH: str = "/host_proc"
|
||||
|
||||
# Throwaway image used to read/write named-volume contents during backup.
|
||||
BACKUP_HELPER_IMAGE: str = "alpine:latest"
|
||||
|
||||
# Host browser sandbox roots
|
||||
ALLOWED_BROWSE_ROOTS: Annotated[list[str], NoDecode] = [
|
||||
"/", "/mnt", "/media", "/srv", "/opt",
|
||||
|
||||
+5
-1
@@ -16,9 +16,11 @@ from docker_client import DockerError
|
||||
from routers import (
|
||||
audit,
|
||||
auth,
|
||||
backups,
|
||||
editor,
|
||||
images,
|
||||
ports,
|
||||
settings as settings_router,
|
||||
stacks,
|
||||
system,
|
||||
templates,
|
||||
@@ -46,7 +48,7 @@ async def lifespan(app: FastAPI):
|
||||
update_task.cancel()
|
||||
|
||||
|
||||
app = FastAPI(title="StackPilot", version="0.3.0", lifespan=lifespan)
|
||||
app = FastAPI(title="StackPilot", version="0.4.0", lifespan=lifespan)
|
||||
|
||||
app.add_middleware(
|
||||
CORSMiddleware,
|
||||
@@ -74,6 +76,8 @@ app.include_router(images.router)
|
||||
app.include_router(ports.router)
|
||||
app.include_router(templates.router)
|
||||
app.include_router(audit.router)
|
||||
app.include_router(settings_router.router)
|
||||
app.include_router(backups.router)
|
||||
app.include_router(ws.router)
|
||||
|
||||
|
||||
|
||||
@@ -1,7 +1,8 @@
|
||||
"""SQLModel table models. Importing this package registers all tables."""
|
||||
from models.audit import AuditLog
|
||||
from models.setting import Setting, Webhook
|
||||
from models.stack import Stack
|
||||
from models.template import Template
|
||||
from models.user import User
|
||||
|
||||
__all__ = ["User", "Stack", "AuditLog", "Template"]
|
||||
__all__ = ["User", "Stack", "AuditLog", "Template", "Setting", "Webhook"]
|
||||
|
||||
@@ -0,0 +1,85 @@
|
||||
from __future__ import annotations
|
||||
|
||||
from datetime import datetime, timezone
|
||||
from typing import Optional
|
||||
|
||||
from sqlmodel import Field, SQLModel
|
||||
|
||||
|
||||
def _now() -> datetime:
|
||||
return datetime.now(timezone.utc)
|
||||
|
||||
|
||||
# Notification event names ----------------------------------------------------
|
||||
EVENT_UPDATE_AVAILABLE = "update_available"
|
||||
EVENT_STACK_START = "stack_start"
|
||||
EVENT_STACK_STOP = "stack_stop"
|
||||
EVENT_STACK_ERROR = "stack_error"
|
||||
EVENT_PULL_FAILED = "pull_failed"
|
||||
|
||||
ALL_EVENTS = [
|
||||
EVENT_UPDATE_AVAILABLE,
|
||||
EVENT_STACK_START,
|
||||
EVENT_STACK_STOP,
|
||||
EVENT_STACK_ERROR,
|
||||
EVENT_PULL_FAILED,
|
||||
]
|
||||
|
||||
WEBHOOK_TYPES = ["ntfy", "discord", "slack", "gotify", "generic"]
|
||||
|
||||
|
||||
class Setting(SQLModel, table=True):
|
||||
"""Simple key/value store for runtime-tunable settings (JSON-encoded value)."""
|
||||
|
||||
key: str = Field(primary_key=True)
|
||||
value: str # JSON-encoded
|
||||
|
||||
|
||||
class Webhook(SQLModel, table=True):
|
||||
id: Optional[int] = Field(default=None, primary_key=True)
|
||||
name: str
|
||||
url: str
|
||||
type: str = "generic" # one of WEBHOOK_TYPES
|
||||
events: str = ",".join(ALL_EVENTS) # comma-separated subscribed events
|
||||
enabled: bool = True
|
||||
created_at: datetime = Field(default_factory=_now)
|
||||
|
||||
|
||||
# --- API schemas ---
|
||||
|
||||
|
||||
class WebhookCreate(SQLModel):
|
||||
name: str
|
||||
url: str
|
||||
type: str = "generic"
|
||||
events: list[str] = ALL_EVENTS
|
||||
enabled: bool = True
|
||||
|
||||
|
||||
class WebhookUpdate(SQLModel):
|
||||
name: Optional[str] = None
|
||||
url: Optional[str] = None
|
||||
type: Optional[str] = None
|
||||
events: Optional[list[str]] = None
|
||||
enabled: Optional[bool] = None
|
||||
|
||||
|
||||
class WebhookRead(SQLModel):
|
||||
id: int
|
||||
name: str
|
||||
url: str
|
||||
type: str
|
||||
events: list[str]
|
||||
enabled: bool
|
||||
created_at: datetime
|
||||
|
||||
|
||||
class SettingsRead(SQLModel):
|
||||
update_check_interval_minutes: int
|
||||
env_webhook_count: int
|
||||
available_events: list[str]
|
||||
webhook_types: list[str]
|
||||
|
||||
|
||||
class SettingsUpdate(SQLModel):
|
||||
update_check_interval_minutes: Optional[int] = None
|
||||
@@ -35,6 +35,12 @@ class UserCreate(SQLModel):
|
||||
role: str = "admin"
|
||||
|
||||
|
||||
class UserUpdate(SQLModel):
|
||||
password: Optional[str] = None
|
||||
role: Optional[str] = None
|
||||
is_active: Optional[bool] = None
|
||||
|
||||
|
||||
class LoginRequest(SQLModel):
|
||||
username: str
|
||||
password: str
|
||||
|
||||
+111
-1
@@ -5,7 +5,7 @@ import time
|
||||
from collections import defaultdict, deque
|
||||
|
||||
from fastapi import APIRouter, Depends, HTTPException, Request, status
|
||||
from sqlmodel import Session
|
||||
from sqlmodel import Session, select
|
||||
|
||||
import auth as auth_mod
|
||||
from database import get_session
|
||||
@@ -16,6 +16,7 @@ from models.user import (
|
||||
User,
|
||||
UserCreate,
|
||||
UserRead,
|
||||
UserUpdate,
|
||||
)
|
||||
from services import audit_service
|
||||
|
||||
@@ -109,3 +110,112 @@ def refresh(
|
||||
@router.get("/me", response_model=UserRead)
|
||||
def me(user: User = Depends(auth_mod.get_current_user)) -> User:
|
||||
return user
|
||||
|
||||
|
||||
# --------------------------------------------------------------------------- #
|
||||
# User management (admin only)
|
||||
# --------------------------------------------------------------------------- #
|
||||
|
||||
|
||||
def _ip(request: Request) -> str:
|
||||
return request.client.host if request.client else "unknown"
|
||||
|
||||
|
||||
@router.get("/users", response_model=list[UserRead])
|
||||
def list_users(
|
||||
session: Session = Depends(get_session),
|
||||
_admin: User = Depends(auth_mod.require_admin),
|
||||
) -> list[User]:
|
||||
return session.exec(select(User).order_by(User.id)).all()
|
||||
|
||||
|
||||
@router.post("/users", response_model=UserRead, status_code=201)
|
||||
def create_user(
|
||||
body: UserCreate,
|
||||
request: Request,
|
||||
session: Session = Depends(get_session),
|
||||
admin: User = Depends(auth_mod.require_admin),
|
||||
) -> User:
|
||||
if not body.username.strip() or not body.password:
|
||||
raise HTTPException(status_code=400, detail="Username and password required")
|
||||
if auth_mod.get_user(session, body.username):
|
||||
raise HTTPException(status_code=409, detail="Username already exists")
|
||||
role = body.role if body.role in ("admin", "user") else "user"
|
||||
user = User(
|
||||
username=body.username,
|
||||
hashed_password=auth_mod.hash_password(body.password),
|
||||
role=role,
|
||||
)
|
||||
session.add(user)
|
||||
session.commit()
|
||||
session.refresh(user)
|
||||
audit_service.record(
|
||||
session, user=admin.username, action="user.create", target=user.username,
|
||||
detail=f"role={role}", ip=_ip(request),
|
||||
)
|
||||
return user
|
||||
|
||||
|
||||
@router.patch("/users/{user_id}", response_model=UserRead)
|
||||
def update_user(
|
||||
user_id: int,
|
||||
body: UserUpdate,
|
||||
request: Request,
|
||||
session: Session = Depends(get_session),
|
||||
admin: User = Depends(auth_mod.require_admin),
|
||||
) -> User:
|
||||
user = session.get(User, user_id)
|
||||
if not user:
|
||||
raise HTTPException(status_code=404, detail="User not found")
|
||||
# Guard against locking yourself out / demoting the last admin.
|
||||
demoting = (body.role is not None and body.role != "admin") or body.is_active is False
|
||||
if user.role == "admin" and demoting:
|
||||
other_admins = session.exec(
|
||||
select(User).where(User.role == "admin", User.is_active == True, User.id != user_id) # noqa: E712
|
||||
).first()
|
||||
if not other_admins:
|
||||
raise HTTPException(status_code=400, detail="Cannot demote or disable the last active admin")
|
||||
if body.password:
|
||||
user.hashed_password = auth_mod.hash_password(body.password)
|
||||
if body.role is not None:
|
||||
if body.role not in ("admin", "user"):
|
||||
raise HTTPException(status_code=400, detail="Invalid role")
|
||||
user.role = body.role
|
||||
if body.is_active is not None:
|
||||
user.is_active = body.is_active
|
||||
session.add(user)
|
||||
session.commit()
|
||||
session.refresh(user)
|
||||
audit_service.record(
|
||||
session, user=admin.username, action="user.update", target=user.username,
|
||||
ip=_ip(request),
|
||||
)
|
||||
return user
|
||||
|
||||
|
||||
@router.delete("/users/{user_id}")
|
||||
def delete_user(
|
||||
user_id: int,
|
||||
request: Request,
|
||||
session: Session = Depends(get_session),
|
||||
admin: User = Depends(auth_mod.require_admin),
|
||||
) -> dict:
|
||||
user = session.get(User, user_id)
|
||||
if not user:
|
||||
raise HTTPException(status_code=404, detail="User not found")
|
||||
if user.id == admin.id:
|
||||
raise HTTPException(status_code=400, detail="You cannot delete your own account")
|
||||
if user.role == "admin":
|
||||
other_admins = session.exec(
|
||||
select(User).where(User.role == "admin", User.is_active == True, User.id != user_id) # noqa: E712
|
||||
).first()
|
||||
if not other_admins:
|
||||
raise HTTPException(status_code=400, detail="Cannot delete the last active admin")
|
||||
username = user.username
|
||||
session.delete(user)
|
||||
session.commit()
|
||||
audit_service.record(
|
||||
session, user=admin.username, action="user.delete", target=username,
|
||||
ip=_ip(request),
|
||||
)
|
||||
return {"ok": True}
|
||||
|
||||
@@ -0,0 +1,96 @@
|
||||
"""Stack backup (incl. volumes) and restore."""
|
||||
from __future__ import annotations
|
||||
|
||||
import os
|
||||
import tempfile
|
||||
|
||||
from fastapi import APIRouter, Depends, File, Form, HTTPException, Query, Request, UploadFile
|
||||
from fastapi.responses import FileResponse
|
||||
from sqlmodel import Session
|
||||
|
||||
from auth import require_admin
|
||||
from database import get_session
|
||||
from models.stack import Stack
|
||||
from models.user import User
|
||||
from services import audit_service, backup_service, compose_service
|
||||
|
||||
router = APIRouter(prefix="/api/stacks", tags=["backups"])
|
||||
|
||||
|
||||
def _ip(request: Request) -> str:
|
||||
return request.client.host if request.client else "unknown"
|
||||
|
||||
|
||||
@router.get("/{stack_id}/backup")
|
||||
async def backup_stack(
|
||||
stack_id: str,
|
||||
request: Request,
|
||||
include_volumes: bool = Query(True),
|
||||
stop_first: bool = Query(True),
|
||||
session: Session = Depends(get_session),
|
||||
user: User = Depends(require_admin),
|
||||
):
|
||||
stack = session.get(Stack, stack_id)
|
||||
if not stack:
|
||||
raise HTTPException(status_code=404, detail=f"Stack '{stack_id}' not found")
|
||||
try:
|
||||
path = await backup_service.create_backup(
|
||||
stack_id, stack.name, include_volumes=include_volumes, stop_first=stop_first,
|
||||
)
|
||||
except backup_service.BackupError as exc:
|
||||
raise HTTPException(status_code=400, detail=str(exc)) from exc
|
||||
audit_service.record(
|
||||
session, user=user.username, action="stack.backup", target=stack_id,
|
||||
detail=f"volumes={include_volumes}", ip=_ip(request),
|
||||
)
|
||||
date = compose_service.now().strftime("%Y%m%d-%H%M%S")
|
||||
suffix = "full" if include_volumes else "config"
|
||||
return FileResponse(
|
||||
path,
|
||||
media_type="application/gzip",
|
||||
filename=f"backup-{stack_id}-{suffix}-{date}.tar.gz",
|
||||
)
|
||||
|
||||
|
||||
@router.post("/restore")
|
||||
async def restore_stack(
|
||||
request: Request,
|
||||
file: UploadFile = File(...),
|
||||
target_id: str | None = Form(None),
|
||||
overwrite: bool = Form(False),
|
||||
restore_volumes: bool = Form(True),
|
||||
session: Session = Depends(get_session),
|
||||
user: User = Depends(require_admin),
|
||||
) -> dict:
|
||||
tmp = tempfile.NamedTemporaryFile(delete=False, suffix=".tar.gz")
|
||||
try:
|
||||
while chunk := await file.read(1024 * 1024):
|
||||
tmp.write(chunk)
|
||||
tmp.close()
|
||||
|
||||
target = compose_service.slugify(target_id) if target_id else None
|
||||
try:
|
||||
result = backup_service.restore_backup(
|
||||
tmp.name,
|
||||
target_id=target,
|
||||
overwrite=overwrite,
|
||||
restore_volumes=restore_volumes,
|
||||
)
|
||||
except backup_service.BackupError as exc:
|
||||
# 409 for the "already exists" conflict, 400 for malformed backups.
|
||||
code = 409 if "already exists" in str(exc) else 400
|
||||
raise HTTPException(status_code=code, detail=str(exc)) from exc
|
||||
|
||||
stack_id = result["stack_id"]
|
||||
stack = session.get(Stack, stack_id)
|
||||
if not stack:
|
||||
session.add(Stack(id=stack_id, name=result.get("name", stack_id)))
|
||||
session.commit()
|
||||
audit_service.record(
|
||||
session, user=user.username, action="stack.restore", target=stack_id,
|
||||
detail=f"volumes={result['volumes_restored']}", ip=_ip(request),
|
||||
)
|
||||
return result
|
||||
finally:
|
||||
if os.path.exists(tmp.name):
|
||||
os.unlink(tmp.name)
|
||||
@@ -0,0 +1,192 @@
|
||||
"""Application settings: update interval + notification webhooks."""
|
||||
from __future__ import annotations
|
||||
|
||||
from fastapi import APIRouter, Depends, HTTPException, Request
|
||||
from sqlmodel import Session, select
|
||||
|
||||
from auth import get_current_user, require_admin
|
||||
from config import settings as env_settings
|
||||
from database import get_session
|
||||
from models.setting import (
|
||||
ALL_EVENTS,
|
||||
WEBHOOK_TYPES,
|
||||
SettingsRead,
|
||||
SettingsUpdate,
|
||||
Webhook,
|
||||
WebhookCreate,
|
||||
WebhookRead,
|
||||
WebhookUpdate,
|
||||
)
|
||||
from models.user import User
|
||||
from services import audit_service, notify_service, settings_service
|
||||
|
||||
router = APIRouter(prefix="/api/settings", tags=["settings"])
|
||||
|
||||
|
||||
def _ip(request: Request) -> str:
|
||||
return request.client.host if request.client else "unknown"
|
||||
|
||||
|
||||
def _to_read(wh: Webhook) -> WebhookRead:
|
||||
return WebhookRead(
|
||||
id=wh.id,
|
||||
name=wh.name,
|
||||
url=wh.url,
|
||||
type=wh.type,
|
||||
events=[e.strip() for e in (wh.events or "").split(",") if e.strip()],
|
||||
enabled=wh.enabled,
|
||||
created_at=wh.created_at,
|
||||
)
|
||||
|
||||
|
||||
# --------------------------------------------------------------------------- #
|
||||
# General settings
|
||||
# --------------------------------------------------------------------------- #
|
||||
|
||||
|
||||
@router.get("", response_model=SettingsRead)
|
||||
def get_settings(
|
||||
session: Session = Depends(get_session),
|
||||
_user: User = Depends(get_current_user),
|
||||
) -> SettingsRead:
|
||||
return SettingsRead(
|
||||
update_check_interval_minutes=settings_service.get_update_interval(session),
|
||||
env_webhook_count=len(env_settings.NOTIFY_WEBHOOKS),
|
||||
available_events=ALL_EVENTS,
|
||||
webhook_types=WEBHOOK_TYPES,
|
||||
)
|
||||
|
||||
|
||||
@router.put("", response_model=SettingsRead)
|
||||
def update_settings(
|
||||
body: SettingsUpdate,
|
||||
request: Request,
|
||||
session: Session = Depends(get_session),
|
||||
user: User = Depends(require_admin),
|
||||
) -> SettingsRead:
|
||||
if body.update_check_interval_minutes is not None:
|
||||
if body.update_check_interval_minutes < 5:
|
||||
raise HTTPException(status_code=400, detail="Interval must be at least 5 minutes")
|
||||
settings_service.set_value(
|
||||
session,
|
||||
settings_service.KEY_UPDATE_INTERVAL,
|
||||
body.update_check_interval_minutes,
|
||||
)
|
||||
audit_service.record(
|
||||
session, user=user.username, action="settings.update",
|
||||
target="update_interval", detail=str(body.update_check_interval_minutes),
|
||||
ip=_ip(request),
|
||||
)
|
||||
return get_settings(session, user)
|
||||
|
||||
|
||||
# --------------------------------------------------------------------------- #
|
||||
# Webhooks
|
||||
# --------------------------------------------------------------------------- #
|
||||
|
||||
|
||||
def _validate(wtype: str, events: list[str]) -> None:
|
||||
if wtype not in WEBHOOK_TYPES:
|
||||
raise HTTPException(status_code=400, detail=f"Unknown webhook type '{wtype}'")
|
||||
bad = [e for e in events if e not in ALL_EVENTS]
|
||||
if bad:
|
||||
raise HTTPException(status_code=400, detail=f"Unknown event(s): {', '.join(bad)}")
|
||||
|
||||
|
||||
@router.get("/webhooks", response_model=list[WebhookRead])
|
||||
def list_webhooks(
|
||||
session: Session = Depends(get_session),
|
||||
_user: User = Depends(require_admin),
|
||||
) -> list[WebhookRead]:
|
||||
rows = session.exec(select(Webhook).order_by(Webhook.id)).all()
|
||||
return [_to_read(w) for w in rows]
|
||||
|
||||
|
||||
@router.post("/webhooks", response_model=WebhookRead, status_code=201)
|
||||
def create_webhook(
|
||||
body: WebhookCreate,
|
||||
request: Request,
|
||||
session: Session = Depends(get_session),
|
||||
user: User = Depends(require_admin),
|
||||
) -> WebhookRead:
|
||||
_validate(body.type, body.events)
|
||||
wh = Webhook(
|
||||
name=body.name,
|
||||
url=body.url,
|
||||
type=body.type,
|
||||
events=",".join(body.events),
|
||||
enabled=body.enabled,
|
||||
)
|
||||
session.add(wh)
|
||||
session.commit()
|
||||
session.refresh(wh)
|
||||
audit_service.record(
|
||||
session, user=user.username, action="webhook.create", target=str(wh.id),
|
||||
detail=body.name, ip=_ip(request),
|
||||
)
|
||||
return _to_read(wh)
|
||||
|
||||
|
||||
@router.put("/webhooks/{webhook_id}", response_model=WebhookRead)
|
||||
def update_webhook(
|
||||
webhook_id: int,
|
||||
body: WebhookUpdate,
|
||||
request: Request,
|
||||
session: Session = Depends(get_session),
|
||||
user: User = Depends(require_admin),
|
||||
) -> WebhookRead:
|
||||
wh = session.get(Webhook, webhook_id)
|
||||
if not wh:
|
||||
raise HTTPException(status_code=404, detail="Webhook not found")
|
||||
if body.type is not None or body.events is not None:
|
||||
_validate(body.type or wh.type, body.events if body.events is not None else _to_read(wh).events)
|
||||
if body.name is not None:
|
||||
wh.name = body.name
|
||||
if body.url is not None:
|
||||
wh.url = body.url
|
||||
if body.type is not None:
|
||||
wh.type = body.type
|
||||
if body.events is not None:
|
||||
wh.events = ",".join(body.events)
|
||||
if body.enabled is not None:
|
||||
wh.enabled = body.enabled
|
||||
session.add(wh)
|
||||
session.commit()
|
||||
session.refresh(wh)
|
||||
audit_service.record(
|
||||
session, user=user.username, action="webhook.update", target=str(wh.id),
|
||||
ip=_ip(request),
|
||||
)
|
||||
return _to_read(wh)
|
||||
|
||||
|
||||
@router.delete("/webhooks/{webhook_id}")
|
||||
def delete_webhook(
|
||||
webhook_id: int,
|
||||
request: Request,
|
||||
session: Session = Depends(get_session),
|
||||
user: User = Depends(require_admin),
|
||||
) -> dict:
|
||||
wh = session.get(Webhook, webhook_id)
|
||||
if not wh:
|
||||
raise HTTPException(status_code=404, detail="Webhook not found")
|
||||
session.delete(wh)
|
||||
session.commit()
|
||||
audit_service.record(
|
||||
session, user=user.username, action="webhook.delete", target=str(webhook_id),
|
||||
ip=_ip(request),
|
||||
)
|
||||
return {"ok": True}
|
||||
|
||||
|
||||
@router.post("/webhooks/{webhook_id}/test")
|
||||
async def test_webhook(
|
||||
webhook_id: int,
|
||||
session: Session = Depends(get_session),
|
||||
_user: User = Depends(require_admin),
|
||||
) -> dict:
|
||||
wh = session.get(Webhook, webhook_id)
|
||||
if not wh:
|
||||
raise HTTPException(status_code=404, detail="Webhook not found")
|
||||
ok = await notify_service.test_webhook(wh.type, wh.url)
|
||||
return {"ok": ok}
|
||||
@@ -19,8 +19,14 @@ from models.stack import (
|
||||
StackCreate,
|
||||
StackUpdate,
|
||||
)
|
||||
from models.setting import (
|
||||
EVENT_PULL_FAILED,
|
||||
EVENT_STACK_ERROR,
|
||||
EVENT_STACK_START,
|
||||
EVENT_STACK_STOP,
|
||||
)
|
||||
from models.user import User
|
||||
from services import audit_service, compose_service
|
||||
from services import audit_service, compose_service, notify_service
|
||||
from services.convert_service import convert_docker_run
|
||||
|
||||
router = APIRouter(prefix="/api/stacks", tags=["stacks"])
|
||||
@@ -220,6 +226,35 @@ def clone_stack(
|
||||
# --------------------------------------------------------------------------- #
|
||||
|
||||
|
||||
# which lifecycle actions emit a notification on success
|
||||
_START_ACTIONS = {"start", "restart", "update"}
|
||||
_STOP_ACTIONS = {"stop", "down"}
|
||||
|
||||
|
||||
async def _notify_lifecycle(action_name: str, stack_id: str, ok: bool, detail: str, session) -> None:
|
||||
try:
|
||||
if not ok:
|
||||
event = EVENT_PULL_FAILED if action_name in ("pull", "update") else EVENT_STACK_ERROR
|
||||
await notify_service.notify(
|
||||
event,
|
||||
f"Stack '{stack_id}' {action_name} failed",
|
||||
detail or f"compose {action_name} returned a non-zero exit code.",
|
||||
session,
|
||||
)
|
||||
elif action_name in _START_ACTIONS:
|
||||
await notify_service.notify(
|
||||
EVENT_STACK_START, f"Stack '{stack_id}' started",
|
||||
f"compose {action_name} completed successfully.", session,
|
||||
)
|
||||
elif action_name in _STOP_ACTIONS:
|
||||
await notify_service.notify(
|
||||
EVENT_STACK_STOP, f"Stack '{stack_id}' stopped",
|
||||
f"compose {action_name} completed successfully.", session,
|
||||
)
|
||||
except Exception: # noqa: BLE001 - notifications are best-effort
|
||||
pass
|
||||
|
||||
|
||||
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)
|
||||
@@ -227,12 +262,15 @@ async def _lifecycle(action_fn, action_name, stack_id, request, session, user):
|
||||
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):
|
||||
ok = result.get("returncode") in (0, None)
|
||||
stderr = result.get("stderr", "").strip()[-2000:]
|
||||
await _notify_lifecycle(action_name, stack_id, ok, stderr, session)
|
||||
if not ok:
|
||||
raise HTTPException(
|
||||
status_code=500,
|
||||
detail={
|
||||
"error": f"compose {action_name} failed",
|
||||
"detail": result.get("stderr", "").strip()[-2000:],
|
||||
"detail": stderr,
|
||||
},
|
||||
)
|
||||
return result
|
||||
|
||||
@@ -0,0 +1,268 @@
|
||||
"""Stack backup & restore, including named-volume contents.
|
||||
|
||||
A backup is a single ``.tar.gz`` with this layout::
|
||||
|
||||
manifest.json metadata + volume/bind inventory
|
||||
compose/... the full stack directory (compose file, .env, ...)
|
||||
volumes/<full>.tar raw contents of each compose-managed named volume
|
||||
|
||||
Named-volume contents are read/written through a throwaway helper container
|
||||
(``BACKUP_HELPER_IMAGE``) with the volume bind-mounted — this is the portable
|
||||
way to snapshot a volume regardless of its driver/mountpoint.
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
import asyncio
|
||||
import io
|
||||
import json
|
||||
import logging
|
||||
import os
|
||||
import shutil
|
||||
import tarfile
|
||||
import tempfile
|
||||
from typing import Optional
|
||||
|
||||
from config import settings
|
||||
from docker_client import DockerError, get_client, safe_call
|
||||
from services import compose_service
|
||||
|
||||
logger = logging.getLogger("stackpilot.backup")
|
||||
|
||||
COMPOSE_PROJECT_LABEL = "com.docker.compose.project"
|
||||
COMPOSE_VOLUME_LABEL = "com.docker.compose.volume"
|
||||
MANIFEST_NAME = "manifest.json"
|
||||
BACKUP_FORMAT_VERSION = 1
|
||||
|
||||
|
||||
class BackupError(Exception):
|
||||
pass
|
||||
|
||||
|
||||
# --------------------------------------------------------------------------- #
|
||||
# Helper container for volume I/O
|
||||
# --------------------------------------------------------------------------- #
|
||||
|
||||
|
||||
def _ensure_helper_image(client) -> None:
|
||||
image = settings.BACKUP_HELPER_IMAGE
|
||||
try:
|
||||
safe_call(client.images.get, image)
|
||||
except DockerError:
|
||||
logger.info("Pulling backup helper image %s", image)
|
||||
safe_call(client.images.pull, image)
|
||||
|
||||
|
||||
def _export_volume(full_name: str) -> bytes:
|
||||
client = get_client()
|
||||
_ensure_helper_image(client)
|
||||
container = safe_call(
|
||||
client.containers.create,
|
||||
settings.BACKUP_HELPER_IMAGE,
|
||||
command="true",
|
||||
volumes={full_name: {"bind": "/v", "mode": "ro"}},
|
||||
)
|
||||
try:
|
||||
# "/v/." copies the *contents* of the volume (no leading "v/" prefix),
|
||||
# so restore can extract straight back into the volume root.
|
||||
bits, _ = container.get_archive("/v/.")
|
||||
buf = io.BytesIO()
|
||||
for chunk in bits:
|
||||
buf.write(chunk)
|
||||
return buf.getvalue()
|
||||
finally:
|
||||
try:
|
||||
container.remove(force=True)
|
||||
except Exception: # noqa: BLE001
|
||||
pass
|
||||
|
||||
|
||||
def _restore_volume(full_name: str, labels: dict, tar_bytes: bytes) -> None:
|
||||
client = get_client()
|
||||
_ensure_helper_image(client)
|
||||
try:
|
||||
safe_call(client.volumes.get, full_name)
|
||||
except DockerError:
|
||||
safe_call(client.volumes.create, name=full_name, labels=labels or {})
|
||||
container = safe_call(
|
||||
client.containers.create,
|
||||
settings.BACKUP_HELPER_IMAGE,
|
||||
command="true",
|
||||
volumes={full_name: {"bind": "/v", "mode": "rw"}},
|
||||
)
|
||||
try:
|
||||
container.put_archive("/v", tar_bytes)
|
||||
finally:
|
||||
try:
|
||||
container.remove(force=True)
|
||||
except Exception: # noqa: BLE001
|
||||
pass
|
||||
|
||||
|
||||
def _compose_volumes(stack_id: str) -> list[dict]:
|
||||
"""Return [{full, short, labels}] for compose-managed named volumes."""
|
||||
try:
|
||||
client = get_client()
|
||||
vols = safe_call(
|
||||
client.volumes.list,
|
||||
filters={"label": f"{COMPOSE_PROJECT_LABEL}={stack_id}"},
|
||||
)
|
||||
except DockerError:
|
||||
return []
|
||||
out = []
|
||||
for v in vols:
|
||||
labels = v.attrs.get("Labels") or {}
|
||||
out.append(
|
||||
{
|
||||
"full": v.name,
|
||||
"short": labels.get(COMPOSE_VOLUME_LABEL, v.name),
|
||||
"labels": labels,
|
||||
}
|
||||
)
|
||||
return out
|
||||
|
||||
|
||||
# --------------------------------------------------------------------------- #
|
||||
# Backup
|
||||
# --------------------------------------------------------------------------- #
|
||||
|
||||
|
||||
async def create_backup(
|
||||
stack_id: str,
|
||||
name: str,
|
||||
include_volumes: bool = True,
|
||||
stop_first: bool = True,
|
||||
) -> str:
|
||||
"""Create a backup tar.gz and return its path on disk."""
|
||||
directory = compose_service.stack_dir(stack_id)
|
||||
if not os.path.isdir(directory):
|
||||
raise BackupError("Stack directory missing")
|
||||
|
||||
volumes = _compose_volumes(stack_id) if include_volumes else []
|
||||
|
||||
# For a consistent volume snapshot, stop the stack first.
|
||||
stopped = False
|
||||
if include_volumes and stop_first and volumes:
|
||||
try:
|
||||
await compose_service.stop(stack_id)
|
||||
stopped = True
|
||||
except Exception as exc: # noqa: BLE001
|
||||
logger.warning("Could not stop %s before backup: %s", stack_id, exc)
|
||||
|
||||
try:
|
||||
manifest = {
|
||||
"format_version": BACKUP_FORMAT_VERSION,
|
||||
"stack_id": stack_id,
|
||||
"name": name,
|
||||
"created_at": compose_service.now().isoformat(),
|
||||
"include_volumes": include_volumes,
|
||||
"volumes": [{"full": v["full"], "short": v["short"], "labels": v["labels"]} for v in volumes],
|
||||
}
|
||||
|
||||
tmp = tempfile.NamedTemporaryFile(delete=False, suffix=".tar.gz")
|
||||
tmp.close()
|
||||
with tarfile.open(tmp.name, "w:gz") as tar:
|
||||
# manifest
|
||||
data = json.dumps(manifest, indent=2).encode("utf-8")
|
||||
info = tarfile.TarInfo(MANIFEST_NAME)
|
||||
info.size = len(data)
|
||||
tar.addfile(info, io.BytesIO(data))
|
||||
# stack directory
|
||||
tar.add(directory, arcname="compose")
|
||||
# volume contents
|
||||
for v in volumes:
|
||||
vbytes = await asyncio.to_thread(_export_volume, v["full"])
|
||||
info = tarfile.TarInfo(f"volumes/{v['full']}.tar")
|
||||
info.size = len(vbytes)
|
||||
tar.addfile(info, io.BytesIO(vbytes))
|
||||
return tmp.name
|
||||
finally:
|
||||
if stopped:
|
||||
try:
|
||||
await compose_service.up(stack_id)
|
||||
except Exception as exc: # noqa: BLE001
|
||||
logger.warning("Could not restart %s after backup: %s", stack_id, exc)
|
||||
|
||||
|
||||
# --------------------------------------------------------------------------- #
|
||||
# Restore
|
||||
# --------------------------------------------------------------------------- #
|
||||
|
||||
|
||||
def read_manifest(tar_path: str) -> dict:
|
||||
with tarfile.open(tar_path, "r:gz") as tar:
|
||||
member = tar.getmember(MANIFEST_NAME)
|
||||
fh = tar.extractfile(member)
|
||||
if fh is None:
|
||||
raise BackupError("Backup is missing its manifest")
|
||||
return json.loads(fh.read().decode("utf-8"))
|
||||
|
||||
|
||||
def _safe_extract_compose(tar: tarfile.TarFile, dest_dir: str) -> None:
|
||||
"""Extract the ``compose/`` subtree into dest_dir, guarding path traversal."""
|
||||
os.makedirs(dest_dir, exist_ok=True)
|
||||
for member in tar.getmembers():
|
||||
if not member.name.startswith("compose/"):
|
||||
continue
|
||||
rel = member.name[len("compose/") :]
|
||||
if not rel:
|
||||
continue
|
||||
target = os.path.normpath(os.path.join(dest_dir, rel))
|
||||
if not target.startswith(os.path.abspath(dest_dir) + os.sep) and target != os.path.abspath(dest_dir):
|
||||
raise BackupError(f"Refusing unsafe path in backup: {member.name}")
|
||||
if member.isdir():
|
||||
os.makedirs(target, exist_ok=True)
|
||||
elif member.isreg():
|
||||
os.makedirs(os.path.dirname(target), exist_ok=True)
|
||||
src = tar.extractfile(member)
|
||||
if src is not None:
|
||||
with open(target, "wb") as out:
|
||||
shutil.copyfileobj(src, out)
|
||||
|
||||
|
||||
def restore_backup(
|
||||
tar_path: str,
|
||||
target_id: Optional[str] = None,
|
||||
overwrite: bool = False,
|
||||
restore_volumes: bool = True,
|
||||
) -> dict:
|
||||
"""Restore a backup. Returns {stack_id, name, volumes_restored}."""
|
||||
manifest = read_manifest(tar_path)
|
||||
stack_id = target_id or manifest.get("stack_id")
|
||||
if not stack_id:
|
||||
raise BackupError("Backup manifest has no stack id")
|
||||
|
||||
directory = compose_service.stack_dir(stack_id)
|
||||
exists = os.path.isdir(directory)
|
||||
if exists and not overwrite:
|
||||
raise BackupError(f"Stack '{stack_id}' already exists")
|
||||
|
||||
with tarfile.open(tar_path, "r:gz") as tar:
|
||||
if exists:
|
||||
shutil.rmtree(directory)
|
||||
_safe_extract_compose(tar, directory)
|
||||
|
||||
volumes_restored = 0
|
||||
if restore_volumes:
|
||||
for v in manifest.get("volumes", []):
|
||||
member_name = f"volumes/{v['full']}.tar"
|
||||
try:
|
||||
member = tar.getmember(member_name)
|
||||
except KeyError:
|
||||
continue
|
||||
fh = tar.extractfile(member)
|
||||
if fh is None:
|
||||
continue
|
||||
# Re-target volume labels to the (possibly new) stack id.
|
||||
labels = dict(v.get("labels") or {})
|
||||
labels[COMPOSE_PROJECT_LABEL] = stack_id
|
||||
full = v["full"]
|
||||
if target_id and manifest.get("stack_id") and full.startswith(manifest["stack_id"] + "_"):
|
||||
full = stack_id + full[len(manifest["stack_id"]):]
|
||||
_restore_volume(full, labels, fh.read())
|
||||
volumes_restored += 1
|
||||
|
||||
return {
|
||||
"stack_id": stack_id,
|
||||
"name": manifest.get("name", stack_id),
|
||||
"volumes_restored": volumes_restored,
|
||||
}
|
||||
@@ -0,0 +1,133 @@
|
||||
"""Outbound notification webhooks.
|
||||
|
||||
Webhooks are configured two ways:
|
||||
* DB-managed (the ``Webhook`` table) — per-webhook type + event subscriptions,
|
||||
editable from the Settings page.
|
||||
* Env ``NOTIFY_WEBHOOKS`` — a comma-separated list of generic JSON endpoints
|
||||
that receive every event (kept for backward compatibility / GitOps setups).
|
||||
|
||||
Supported types: ntfy, discord, slack, gotify, generic (JSON POST).
|
||||
All delivery is best-effort: failures are logged, never raised to the caller.
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
import logging
|
||||
from typing import Optional
|
||||
|
||||
import httpx
|
||||
from sqlmodel import Session, select
|
||||
|
||||
from config import settings as env_settings
|
||||
from database import engine
|
||||
from models.setting import Webhook
|
||||
|
||||
logger = logging.getLogger("stackpilot.notify")
|
||||
|
||||
_TIMEOUT = 10.0
|
||||
|
||||
|
||||
# --------------------------------------------------------------------------- #
|
||||
# Payload formatting per webhook type
|
||||
# --------------------------------------------------------------------------- #
|
||||
|
||||
|
||||
def _build_request(wtype: str, url: str, event: str, title: str, message: str):
|
||||
"""Return (method-kwargs) for httpx.post for the given webhook type."""
|
||||
if wtype == "ntfy":
|
||||
return {
|
||||
"url": url,
|
||||
"content": message.encode("utf-8"),
|
||||
"headers": {"Title": title, "Tags": _ntfy_tag(event)},
|
||||
}
|
||||
if wtype == "discord":
|
||||
return {"url": url, "json": {"content": f"**{title}**\n{message}"}}
|
||||
if wtype == "slack":
|
||||
return {"url": url, "json": {"text": f"*{title}*\n{message}"}}
|
||||
if wtype == "gotify":
|
||||
priority = 8 if event in ("stack_error", "pull_failed") else 5
|
||||
return {
|
||||
"url": url,
|
||||
"json": {"title": title, "message": message, "priority": priority},
|
||||
}
|
||||
# generic
|
||||
return {
|
||||
"url": url,
|
||||
"json": {"event": event, "title": title, "message": message},
|
||||
}
|
||||
|
||||
|
||||
def _ntfy_tag(event: str) -> str:
|
||||
return {
|
||||
"update_available": "arrow_up",
|
||||
"stack_start": "white_check_mark",
|
||||
"stack_stop": "stop_button",
|
||||
"stack_error": "rotating_light",
|
||||
"pull_failed": "warning",
|
||||
}.get(event, "bell")
|
||||
|
||||
|
||||
async def _deliver(client: httpx.AsyncClient, wtype: str, url: str, event: str, title: str, message: str) -> bool:
|
||||
kwargs = _build_request(wtype, url, event, title, message)
|
||||
target = kwargs.pop("url")
|
||||
try:
|
||||
resp = await client.post(target, timeout=_TIMEOUT, **kwargs)
|
||||
resp.raise_for_status()
|
||||
return True
|
||||
except httpx.HTTPError as exc:
|
||||
logger.warning("Webhook delivery failed (%s): %s", wtype, exc)
|
||||
return False
|
||||
|
||||
|
||||
# --------------------------------------------------------------------------- #
|
||||
# Public API
|
||||
# --------------------------------------------------------------------------- #
|
||||
|
||||
|
||||
def _targets_for_event(session: Session, event: str) -> list[tuple[str, str]]:
|
||||
"""Return [(type, url)] of all destinations subscribed to ``event``."""
|
||||
targets: list[tuple[str, str]] = []
|
||||
for wh in session.exec(select(Webhook)).all():
|
||||
if not wh.enabled:
|
||||
continue
|
||||
subscribed = [e.strip() for e in (wh.events or "").split(",") if e.strip()]
|
||||
if event in subscribed:
|
||||
targets.append((wh.type, wh.url))
|
||||
# Env-configured generic endpoints receive everything.
|
||||
for url in env_settings.NOTIFY_WEBHOOKS:
|
||||
targets.append(("generic", url))
|
||||
return targets
|
||||
|
||||
|
||||
async def notify(
|
||||
event: str,
|
||||
title: str,
|
||||
message: str,
|
||||
session: Optional[Session] = None,
|
||||
) -> int:
|
||||
"""Fan out ``event`` to all subscribed webhooks. Returns delivered count."""
|
||||
if session is None:
|
||||
with Session(engine) as own:
|
||||
return await notify(event, title, message, own)
|
||||
|
||||
targets = _targets_for_event(session, event)
|
||||
if not targets:
|
||||
return 0
|
||||
delivered = 0
|
||||
async with httpx.AsyncClient(follow_redirects=True) as client:
|
||||
for wtype, url in targets:
|
||||
if await _deliver(client, wtype, url, event, title, message):
|
||||
delivered += 1
|
||||
return delivered
|
||||
|
||||
|
||||
async def test_webhook(wtype: str, url: str) -> bool:
|
||||
"""Send a one-off test notification to a single destination."""
|
||||
async with httpx.AsyncClient(follow_redirects=True) as client:
|
||||
return await _deliver(
|
||||
client,
|
||||
wtype,
|
||||
url,
|
||||
"update_available",
|
||||
"StackPilot test notification",
|
||||
"If you can read this, your webhook is configured correctly. 🚀",
|
||||
)
|
||||
@@ -0,0 +1,45 @@
|
||||
"""Runtime settings stored in the DB (key/value), with env fallbacks."""
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
from typing import Any, Optional
|
||||
|
||||
from sqlmodel import Session, select
|
||||
|
||||
from config import settings as env_settings
|
||||
from database import engine
|
||||
from models.setting import Setting
|
||||
|
||||
KEY_UPDATE_INTERVAL = "update_check_interval_minutes"
|
||||
|
||||
|
||||
def get(session: Session, key: str, default: Any = None) -> Any:
|
||||
row = session.get(Setting, key)
|
||||
if row is None:
|
||||
return default
|
||||
try:
|
||||
return json.loads(row.value)
|
||||
except json.JSONDecodeError:
|
||||
return default
|
||||
|
||||
|
||||
def set_value(session: Session, key: str, value: Any) -> None:
|
||||
row = session.get(Setting, key)
|
||||
encoded = json.dumps(value)
|
||||
if row is None:
|
||||
session.add(Setting(key=key, value=encoded))
|
||||
else:
|
||||
row.value = encoded
|
||||
session.add(row)
|
||||
session.commit()
|
||||
|
||||
|
||||
def get_update_interval(session: Optional[Session] = None) -> int:
|
||||
"""Effective update-check interval in minutes (DB override or env default)."""
|
||||
if session is None:
|
||||
with Session(engine) as own:
|
||||
return get_update_interval(own)
|
||||
val = get(session, KEY_UPDATE_INTERVAL)
|
||||
if isinstance(val, int) and val > 0:
|
||||
return val
|
||||
return env_settings.UPDATE_CHECK_INTERVAL_MINUTES
|
||||
@@ -16,6 +16,8 @@ import httpx
|
||||
|
||||
from config import settings
|
||||
from docker_client import DockerError, get_client, safe_call
|
||||
from models.setting import EVENT_UPDATE_AVAILABLE
|
||||
from services import notify_service, settings_service
|
||||
|
||||
logger = logging.getLogger("stackpilot.update")
|
||||
|
||||
@@ -45,6 +47,10 @@ class UpdateStatus:
|
||||
# image ref -> UpdateStatus
|
||||
_CACHE: dict[str, UpdateStatus] = {}
|
||||
|
||||
# images we've already sent an "update available" notification for, so the
|
||||
# background loop doesn't re-notify on every cycle.
|
||||
_NOTIFIED: set[str] = set()
|
||||
|
||||
|
||||
# --------------------------------------------------------------------------- #
|
||||
# Image reference parsing
|
||||
@@ -164,6 +170,18 @@ async def check_image(image: str) -> UpdateStatus:
|
||||
error=error,
|
||||
)
|
||||
_CACHE[image] = status
|
||||
if update_available and image not in _NOTIFIED:
|
||||
_NOTIFIED.add(image)
|
||||
try:
|
||||
await notify_service.notify(
|
||||
EVENT_UPDATE_AVAILABLE,
|
||||
"Image update available",
|
||||
f"A newer image is available for {image}.",
|
||||
)
|
||||
except Exception as exc: # noqa: BLE001 - notifications are best-effort
|
||||
logger.debug("update notify failed for %s: %s", image, exc)
|
||||
elif not update_available:
|
||||
_NOTIFIED.discard(image)
|
||||
return status
|
||||
|
||||
|
||||
@@ -192,7 +210,6 @@ def get_cache() -> dict[str, dict]:
|
||||
|
||||
|
||||
async def background_loop():
|
||||
interval = max(settings.UPDATE_CHECK_INTERVAL_MINUTES, 5) * 60
|
||||
# initial delay so startup isn't blocked
|
||||
await asyncio.sleep(30)
|
||||
while True:
|
||||
@@ -201,4 +218,6 @@ async def background_loop():
|
||||
logger.info("Image update check complete (%d images)", len(_CACHE))
|
||||
except Exception as exc: # noqa: BLE001
|
||||
logger.warning("Image update check failed: %s", exc)
|
||||
# Re-read the interval each cycle so Settings changes take effect.
|
||||
interval = max(settings_service.get_update_interval(), 5) * 60
|
||||
await asyncio.sleep(interval)
|
||||
|
||||
Reference in New Issue
Block a user