Phase 6: remote backup destinations — SFTP & S3 (0.6.0)
- BackupDestination model + backup_destination_service (SFTP via paramiko,
S3-compatible via boto3): upload/list/download/delete/test.
- routers/destinations.py: destinations CRUD (secrets masked, merge-on-update),
test, list/delete remote backups. backups.py: POST /{id}/backup/push and
POST /restore-from (download from a destination + restore, volumes included).
- Frontend: Settings → Backup destinations (SFTP/S3 forms + test); Backup dialog
can push to a destination; Restore dialog can pick a destination + backup.
- deps: paramiko 3.5.0, boto3 1.35.99.
Verified end-to-end against live MinIO + atmoz/sftp: create/test destinations,
push (incl. volumes), list, restore-from to a fresh stack (volume data intact),
delete remote backup.
Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
This commit is contained in:
co-authored by
Claude Opus 4.8
parent
59037f4287
commit
7bd449101d
+3
-1
@@ -18,6 +18,7 @@ from routers import (
|
||||
audit,
|
||||
auth,
|
||||
backups,
|
||||
destinations,
|
||||
editor,
|
||||
images,
|
||||
ports,
|
||||
@@ -49,7 +50,7 @@ async def lifespan(app: FastAPI):
|
||||
update_task.cancel()
|
||||
|
||||
|
||||
app = FastAPI(title="StackPilot", version="0.5.0", lifespan=lifespan)
|
||||
app = FastAPI(title="StackPilot", version="0.6.0", lifespan=lifespan)
|
||||
|
||||
app.add_middleware(
|
||||
CORSMiddleware,
|
||||
@@ -79,6 +80,7 @@ app.include_router(templates.router)
|
||||
app.include_router(audit.router)
|
||||
app.include_router(settings_router.router)
|
||||
app.include_router(backups.router)
|
||||
app.include_router(destinations.router)
|
||||
app.include_router(agents.router)
|
||||
app.include_router(ws.router)
|
||||
|
||||
|
||||
@@ -1,9 +1,13 @@
|
||||
"""SQLModel table models. Importing this package registers all tables."""
|
||||
from models.agent import Agent
|
||||
from models.audit import AuditLog
|
||||
from models.backup_destination import BackupDestination
|
||||
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", "Setting", "Webhook", "Agent"]
|
||||
__all__ = [
|
||||
"User", "Stack", "AuditLog", "Template", "Setting", "Webhook", "Agent",
|
||||
"BackupDestination",
|
||||
]
|
||||
|
||||
@@ -0,0 +1,48 @@
|
||||
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)
|
||||
|
||||
|
||||
DESTINATION_TYPES = ["sftp", "s3"]
|
||||
|
||||
# config keys that hold secrets — masked in API responses.
|
||||
SECRET_KEYS = {"password", "private_key", "secret_key"}
|
||||
|
||||
|
||||
class BackupDestination(SQLModel, table=True):
|
||||
"""A remote target for stack backups (SFTP or S3-compatible)."""
|
||||
|
||||
id: Optional[int] = Field(default=None, primary_key=True)
|
||||
name: str
|
||||
type: str # one of DESTINATION_TYPES
|
||||
config: str = "{}" # JSON-encoded, type-specific (incl. secrets)
|
||||
created_at: datetime = Field(default_factory=_now)
|
||||
|
||||
|
||||
# --- API schemas ---
|
||||
|
||||
|
||||
class DestinationCreate(SQLModel):
|
||||
name: str
|
||||
type: str
|
||||
config: dict
|
||||
|
||||
|
||||
class DestinationUpdate(SQLModel):
|
||||
name: Optional[str] = None
|
||||
config: Optional[dict] = None
|
||||
|
||||
|
||||
class DestinationRead(SQLModel):
|
||||
id: int
|
||||
name: str
|
||||
type: str
|
||||
config: dict # secrets masked
|
||||
created_at: datetime
|
||||
@@ -11,3 +11,5 @@ python-multipart==0.0.20
|
||||
watchdog==6.0.0
|
||||
httpx==0.28.1
|
||||
PyYAML==6.0.2
|
||||
paramiko==3.5.0
|
||||
boto3==1.35.99
|
||||
|
||||
+118
-4
@@ -1,18 +1,26 @@
|
||||
"""Stack backup (incl. volumes) and restore."""
|
||||
from __future__ import annotations
|
||||
|
||||
import asyncio
|
||||
import os
|
||||
import tempfile
|
||||
|
||||
from fastapi import APIRouter, Depends, File, Form, HTTPException, Query, Request, UploadFile
|
||||
from fastapi.responses import FileResponse
|
||||
from pydantic import BaseModel
|
||||
from sqlmodel import Session
|
||||
|
||||
from auth import require_admin
|
||||
from database import get_session
|
||||
from models.backup_destination import BackupDestination
|
||||
from models.stack import Stack
|
||||
from models.user import User
|
||||
from services import audit_service, backup_service, compose_service
|
||||
from services import (
|
||||
audit_service,
|
||||
backup_destination_service as dest_service,
|
||||
backup_service,
|
||||
compose_service,
|
||||
)
|
||||
|
||||
router = APIRouter(prefix="/api/stacks", tags=["backups"])
|
||||
|
||||
@@ -21,6 +29,12 @@ def _ip(request: Request) -> str:
|
||||
return request.client.host if request.client else "unknown"
|
||||
|
||||
|
||||
def _backup_filename(stack_id: str, include_volumes: bool) -> str:
|
||||
date = compose_service.now().strftime("%Y%m%d-%H%M%S")
|
||||
suffix = "full" if include_volumes else "config"
|
||||
return f"backup-{stack_id}-{suffix}-{date}.tar.gz"
|
||||
|
||||
|
||||
@router.get("/{stack_id}/backup")
|
||||
async def backup_stack(
|
||||
stack_id: str,
|
||||
@@ -43,12 +57,10 @@ async def backup_stack(
|
||||
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",
|
||||
filename=_backup_filename(stack_id, include_volumes),
|
||||
)
|
||||
|
||||
|
||||
@@ -94,3 +106,105 @@ async def restore_stack(
|
||||
finally:
|
||||
if os.path.exists(tmp.name):
|
||||
os.unlink(tmp.name)
|
||||
|
||||
|
||||
# --------------------------------------------------------------------------- #
|
||||
# Push to / restore from a remote destination
|
||||
# --------------------------------------------------------------------------- #
|
||||
|
||||
|
||||
class PushBody(BaseModel):
|
||||
destination_id: int
|
||||
include_volumes: bool = True
|
||||
stop_first: bool = True
|
||||
|
||||
|
||||
class RestoreFromBody(BaseModel):
|
||||
destination_id: int
|
||||
name: str
|
||||
target_id: str | None = None
|
||||
overwrite: bool = False
|
||||
restore_volumes: bool = True
|
||||
|
||||
|
||||
def _get_dest(session: Session, dest_id: int) -> BackupDestination:
|
||||
d = session.get(BackupDestination, dest_id)
|
||||
if not d:
|
||||
raise HTTPException(status_code=404, detail=f"Destination {dest_id} not found")
|
||||
return d
|
||||
|
||||
|
||||
@router.post("/{stack_id}/backup/push")
|
||||
async def push_backup(
|
||||
stack_id: str,
|
||||
body: PushBody,
|
||||
request: Request,
|
||||
session: Session = Depends(get_session),
|
||||
user: User = Depends(require_admin),
|
||||
) -> dict:
|
||||
stack = session.get(Stack, stack_id)
|
||||
if not stack:
|
||||
raise HTTPException(status_code=404, detail=f"Stack '{stack_id}' not found")
|
||||
dest = _get_dest(session, body.destination_id)
|
||||
try:
|
||||
path = await backup_service.create_backup(
|
||||
stack_id, stack.name,
|
||||
include_volumes=body.include_volumes, stop_first=body.stop_first,
|
||||
)
|
||||
except backup_service.BackupError as exc:
|
||||
raise HTTPException(status_code=400, detail=str(exc)) from exc
|
||||
|
||||
filename = _backup_filename(stack_id, body.include_volumes)
|
||||
try:
|
||||
remote = await asyncio.to_thread(dest_service.upload, dest, path, filename)
|
||||
except dest_service.DestinationError as exc:
|
||||
raise HTTPException(status_code=502, detail=str(exc)) from exc
|
||||
finally:
|
||||
if os.path.exists(path):
|
||||
os.unlink(path)
|
||||
|
||||
audit_service.record(
|
||||
session, user=user.username, action="stack.backup.push",
|
||||
target=stack_id, detail=f"{dest.name}:{filename}", ip=_ip(request),
|
||||
)
|
||||
return {"ok": True, "destination": dest.name, "name": filename, "remote": remote}
|
||||
|
||||
|
||||
@router.post("/restore-from")
|
||||
async def restore_from_destination(
|
||||
body: RestoreFromBody,
|
||||
request: Request,
|
||||
session: Session = Depends(get_session),
|
||||
user: User = Depends(require_admin),
|
||||
) -> dict:
|
||||
dest = _get_dest(session, body.destination_id)
|
||||
tmp = tempfile.NamedTemporaryFile(delete=False, suffix=".tar.gz")
|
||||
tmp.close()
|
||||
try:
|
||||
try:
|
||||
await asyncio.to_thread(dest_service.download, dest, body.name, tmp.name)
|
||||
except dest_service.DestinationError as exc:
|
||||
raise HTTPException(status_code=502, detail=str(exc)) from exc
|
||||
|
||||
target = compose_service.slugify(body.target_id) if body.target_id else None
|
||||
try:
|
||||
result = backup_service.restore_backup(
|
||||
tmp.name, target_id=target,
|
||||
overwrite=body.overwrite, restore_volumes=body.restore_volumes,
|
||||
)
|
||||
except backup_service.BackupError as exc:
|
||||
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"]
|
||||
if not session.get(Stack, stack_id):
|
||||
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"from {dest.name}:{body.name}", ip=_ip(request),
|
||||
)
|
||||
return result
|
||||
finally:
|
||||
if os.path.exists(tmp.name):
|
||||
os.unlink(tmp.name)
|
||||
|
||||
@@ -0,0 +1,171 @@
|
||||
"""Backup destination management (SFTP / S3-compatible)."""
|
||||
from __future__ import annotations
|
||||
|
||||
import asyncio
|
||||
import json
|
||||
|
||||
from fastapi import APIRouter, Depends, HTTPException, Request
|
||||
from sqlmodel import Session, select
|
||||
|
||||
from auth import get_current_user, require_admin
|
||||
from database import get_session
|
||||
from models.backup_destination import (
|
||||
DESTINATION_TYPES,
|
||||
SECRET_KEYS,
|
||||
BackupDestination,
|
||||
DestinationCreate,
|
||||
DestinationRead,
|
||||
DestinationUpdate,
|
||||
)
|
||||
from models.user import User
|
||||
from services import audit_service, backup_destination_service as dest_service
|
||||
|
||||
router = APIRouter(prefix="/api/backups/destinations", tags=["backups"])
|
||||
|
||||
|
||||
def _ip(request: Request) -> str:
|
||||
return request.client.host if request.client else "unknown"
|
||||
|
||||
|
||||
def _mask(config: dict) -> dict:
|
||||
return {k: ("••••••" if k in SECRET_KEYS and v else v) for k, v in config.items()}
|
||||
|
||||
|
||||
def _to_read(d: BackupDestination) -> DestinationRead:
|
||||
return DestinationRead(
|
||||
id=d.id,
|
||||
name=d.name,
|
||||
type=d.type,
|
||||
config=_mask(dest_service.parse_config(d)),
|
||||
created_at=d.created_at,
|
||||
)
|
||||
|
||||
|
||||
def _get_or_404(session: Session, dest_id: int) -> BackupDestination:
|
||||
d = session.get(BackupDestination, dest_id)
|
||||
if not d:
|
||||
raise HTTPException(status_code=404, detail=f"Destination {dest_id} not found")
|
||||
return d
|
||||
|
||||
|
||||
@router.get("", response_model=list[DestinationRead])
|
||||
def list_destinations(
|
||||
session: Session = Depends(get_session),
|
||||
_user: User = Depends(require_admin),
|
||||
) -> list[DestinationRead]:
|
||||
rows = session.exec(select(BackupDestination).order_by(BackupDestination.id)).all()
|
||||
return [_to_read(d) for d in rows]
|
||||
|
||||
|
||||
@router.post("", response_model=DestinationRead, status_code=201)
|
||||
def create_destination(
|
||||
body: DestinationCreate,
|
||||
request: Request,
|
||||
session: Session = Depends(get_session),
|
||||
user: User = Depends(require_admin),
|
||||
) -> DestinationRead:
|
||||
if body.type not in DESTINATION_TYPES:
|
||||
raise HTTPException(status_code=400, detail=f"Unknown type '{body.type}'")
|
||||
d = BackupDestination(name=body.name, type=body.type, config=json.dumps(body.config))
|
||||
session.add(d)
|
||||
session.commit()
|
||||
session.refresh(d)
|
||||
audit_service.record(
|
||||
session, user=user.username, action="destination.create", target=body.name,
|
||||
detail=body.type, ip=_ip(request),
|
||||
)
|
||||
return _to_read(d)
|
||||
|
||||
|
||||
@router.put("/{dest_id}", response_model=DestinationRead)
|
||||
def update_destination(
|
||||
dest_id: int,
|
||||
body: DestinationUpdate,
|
||||
request: Request,
|
||||
session: Session = Depends(get_session),
|
||||
user: User = Depends(require_admin),
|
||||
) -> DestinationRead:
|
||||
d = _get_or_404(session, dest_id)
|
||||
if body.name is not None:
|
||||
d.name = body.name
|
||||
if body.config is not None:
|
||||
# Merge so masked/blank secrets don't wipe stored ones.
|
||||
existing = dest_service.parse_config(d)
|
||||
for k, v in body.config.items():
|
||||
if k in SECRET_KEYS and (v == "" or v == "••••••"):
|
||||
continue # keep existing secret
|
||||
existing[k] = v
|
||||
d.config = json.dumps(existing)
|
||||
session.add(d)
|
||||
session.commit()
|
||||
session.refresh(d)
|
||||
audit_service.record(
|
||||
session, user=user.username, action="destination.update", target=d.name,
|
||||
ip=_ip(request),
|
||||
)
|
||||
return _to_read(d)
|
||||
|
||||
|
||||
@router.delete("/{dest_id}")
|
||||
def delete_destination(
|
||||
dest_id: int,
|
||||
request: Request,
|
||||
session: Session = Depends(get_session),
|
||||
user: User = Depends(require_admin),
|
||||
) -> dict:
|
||||
d = _get_or_404(session, dest_id)
|
||||
name = d.name
|
||||
session.delete(d)
|
||||
session.commit()
|
||||
audit_service.record(
|
||||
session, user=user.username, action="destination.delete", target=name,
|
||||
ip=_ip(request),
|
||||
)
|
||||
return {"ok": True}
|
||||
|
||||
|
||||
@router.post("/{dest_id}/test")
|
||||
async def test_destination(
|
||||
dest_id: int,
|
||||
session: Session = Depends(get_session),
|
||||
_user: User = Depends(require_admin),
|
||||
) -> dict:
|
||||
d = _get_or_404(session, dest_id)
|
||||
try:
|
||||
await asyncio.to_thread(dest_service.test, d)
|
||||
return {"ok": True}
|
||||
except dest_service.DestinationError as exc:
|
||||
return {"ok": False, "error": str(exc)}
|
||||
|
||||
|
||||
@router.get("/{dest_id}/backups")
|
||||
async def list_destination_backups(
|
||||
dest_id: int,
|
||||
session: Session = Depends(get_session),
|
||||
_user: User = Depends(require_admin),
|
||||
) -> list[dict]:
|
||||
d = _get_or_404(session, dest_id)
|
||||
try:
|
||||
return await asyncio.to_thread(dest_service.list_backups, d)
|
||||
except dest_service.DestinationError as exc:
|
||||
raise HTTPException(status_code=502, detail=str(exc)) from exc
|
||||
|
||||
|
||||
@router.delete("/{dest_id}/backups/{name}")
|
||||
async def delete_destination_backup(
|
||||
dest_id: int,
|
||||
name: str,
|
||||
request: Request,
|
||||
session: Session = Depends(get_session),
|
||||
user: User = Depends(require_admin),
|
||||
) -> dict:
|
||||
d = _get_or_404(session, dest_id)
|
||||
try:
|
||||
await asyncio.to_thread(dest_service.delete, d, name)
|
||||
except dest_service.DestinationError as exc:
|
||||
raise HTTPException(status_code=502, detail=str(exc)) from exc
|
||||
audit_service.record(
|
||||
session, user=user.username, action="destination.backup.delete",
|
||||
target=f"{d.name}/{name}", ip=_ip(request),
|
||||
)
|
||||
return {"ok": True}
|
||||
@@ -0,0 +1,258 @@
|
||||
"""Push/pull stack backups to remote destinations (SFTP or S3-compatible).
|
||||
|
||||
All operations are synchronous (paramiko / boto3); async callers should wrap
|
||||
them with ``asyncio.to_thread``. Destination config is a plain dict parsed from
|
||||
the ``BackupDestination.config`` JSON column.
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
import io
|
||||
import json
|
||||
import logging
|
||||
import os
|
||||
import posixpath
|
||||
import stat
|
||||
from typing import Any
|
||||
|
||||
from models.backup_destination import BackupDestination
|
||||
|
||||
logger = logging.getLogger("stackpilot.backup_dest")
|
||||
|
||||
|
||||
class DestinationError(Exception):
|
||||
pass
|
||||
|
||||
|
||||
def parse_config(dest: BackupDestination) -> dict:
|
||||
try:
|
||||
return json.loads(dest.config or "{}")
|
||||
except json.JSONDecodeError:
|
||||
return {}
|
||||
|
||||
|
||||
# --------------------------------------------------------------------------- #
|
||||
# SFTP (paramiko)
|
||||
# --------------------------------------------------------------------------- #
|
||||
|
||||
|
||||
def _sftp_connect(cfg: dict):
|
||||
import paramiko
|
||||
|
||||
host = cfg.get("host")
|
||||
if not host:
|
||||
raise DestinationError("SFTP host is required")
|
||||
port = int(cfg.get("port") or 22)
|
||||
username = cfg.get("username")
|
||||
transport = paramiko.Transport((host, port))
|
||||
try:
|
||||
pkey = None
|
||||
if cfg.get("private_key"):
|
||||
pkey = _load_key(cfg["private_key"])
|
||||
transport.connect(username=username, password=cfg.get("password") or None, pkey=pkey)
|
||||
except Exception as exc: # noqa: BLE001
|
||||
transport.close()
|
||||
raise DestinationError(f"SFTP connection failed: {exc}") from exc
|
||||
return paramiko.SFTPClient.from_transport(transport), transport
|
||||
|
||||
|
||||
def _load_key(key_str: str):
|
||||
import paramiko
|
||||
|
||||
for cls in (paramiko.Ed25519Key, paramiko.ECDSAKey, paramiko.RSAKey):
|
||||
try:
|
||||
return cls.from_private_key(io.StringIO(key_str))
|
||||
except Exception: # noqa: BLE001
|
||||
continue
|
||||
raise DestinationError("Could not parse SFTP private key")
|
||||
|
||||
|
||||
def _sftp_makedirs(sftp, path: str) -> None:
|
||||
if not path or path in (".", "/"):
|
||||
return
|
||||
parts = path.strip("/").split("/")
|
||||
cur = "/" if path.startswith("/") else ""
|
||||
for p in parts:
|
||||
cur = posixpath.join(cur, p) if cur else p
|
||||
try:
|
||||
sftp.stat(cur)
|
||||
except IOError:
|
||||
sftp.mkdir(cur)
|
||||
|
||||
|
||||
def _sftp_upload(cfg: dict, local_path: str, filename: str) -> str:
|
||||
sftp, transport = _sftp_connect(cfg)
|
||||
try:
|
||||
base = cfg.get("path") or "."
|
||||
if base not in (".", ""):
|
||||
_sftp_makedirs(sftp, base)
|
||||
remote = posixpath.join(base, filename) if base not in (".", "") else filename
|
||||
sftp.put(local_path, remote)
|
||||
return remote
|
||||
finally:
|
||||
sftp.close()
|
||||
transport.close()
|
||||
|
||||
|
||||
def _sftp_list(cfg: dict) -> list[dict]:
|
||||
sftp, transport = _sftp_connect(cfg)
|
||||
try:
|
||||
base = cfg.get("path") or "."
|
||||
out = []
|
||||
try:
|
||||
entries = sftp.listdir_attr(base)
|
||||
except IOError:
|
||||
return []
|
||||
for e in entries:
|
||||
if stat.S_ISDIR(e.st_mode):
|
||||
continue
|
||||
if not e.filename.endswith(".tar.gz"):
|
||||
continue
|
||||
out.append({"name": e.filename, "size": e.st_size, "modified": e.st_mtime})
|
||||
return sorted(out, key=lambda x: x["modified"] or 0, reverse=True)
|
||||
finally:
|
||||
sftp.close()
|
||||
transport.close()
|
||||
|
||||
|
||||
def _sftp_download(cfg: dict, name: str, local_path: str) -> None:
|
||||
sftp, transport = _sftp_connect(cfg)
|
||||
try:
|
||||
base = cfg.get("path") or "."
|
||||
remote = posixpath.join(base, name) if base not in (".", "") else name
|
||||
sftp.get(remote, local_path)
|
||||
finally:
|
||||
sftp.close()
|
||||
transport.close()
|
||||
|
||||
|
||||
def _sftp_delete(cfg: dict, name: str) -> None:
|
||||
sftp, transport = _sftp_connect(cfg)
|
||||
try:
|
||||
base = cfg.get("path") or "."
|
||||
remote = posixpath.join(base, name) if base not in (".", "") else name
|
||||
sftp.remove(remote)
|
||||
finally:
|
||||
sftp.close()
|
||||
transport.close()
|
||||
|
||||
|
||||
# --------------------------------------------------------------------------- #
|
||||
# S3-compatible (boto3)
|
||||
# --------------------------------------------------------------------------- #
|
||||
|
||||
|
||||
def _s3_client(cfg: dict):
|
||||
import boto3
|
||||
|
||||
bucket = cfg.get("bucket")
|
||||
if not bucket:
|
||||
raise DestinationError("S3 bucket is required")
|
||||
return boto3.client(
|
||||
"s3",
|
||||
endpoint_url=cfg.get("endpoint_url") or None,
|
||||
region_name=cfg.get("region") or None,
|
||||
aws_access_key_id=cfg.get("access_key") or None,
|
||||
aws_secret_access_key=cfg.get("secret_key") or None,
|
||||
)
|
||||
|
||||
|
||||
def _s3_key(cfg: dict, filename: str) -> str:
|
||||
prefix = (cfg.get("prefix") or "").strip("/")
|
||||
return f"{prefix}/{filename}" if prefix else filename
|
||||
|
||||
|
||||
def _s3_upload(cfg: dict, local_path: str, filename: str) -> str:
|
||||
client = _s3_client(cfg)
|
||||
key = _s3_key(cfg, filename)
|
||||
try:
|
||||
client.upload_file(local_path, cfg["bucket"], key)
|
||||
except Exception as exc: # noqa: BLE001
|
||||
raise DestinationError(f"S3 upload failed: {exc}") from exc
|
||||
return key
|
||||
|
||||
|
||||
def _s3_list(cfg: dict) -> list[dict]:
|
||||
client = _s3_client(cfg)
|
||||
prefix = (cfg.get("prefix") or "").strip("/")
|
||||
kwargs: dict[str, Any] = {"Bucket": cfg["bucket"]}
|
||||
if prefix:
|
||||
kwargs["Prefix"] = prefix + "/"
|
||||
try:
|
||||
resp = client.list_objects_v2(**kwargs)
|
||||
except Exception as exc: # noqa: BLE001
|
||||
raise DestinationError(f"S3 list failed: {exc}") from exc
|
||||
out = []
|
||||
for obj in resp.get("Contents", []):
|
||||
name = obj["Key"].split("/")[-1]
|
||||
if not name.endswith(".tar.gz"):
|
||||
continue
|
||||
out.append(
|
||||
{
|
||||
"name": name,
|
||||
"size": obj.get("Size", 0),
|
||||
"modified": obj["LastModified"].timestamp() if obj.get("LastModified") else None,
|
||||
}
|
||||
)
|
||||
return sorted(out, key=lambda x: x["modified"] or 0, reverse=True)
|
||||
|
||||
|
||||
def _s3_download(cfg: dict, name: str, local_path: str) -> None:
|
||||
client = _s3_client(cfg)
|
||||
try:
|
||||
client.download_file(cfg["bucket"], _s3_key(cfg, name), local_path)
|
||||
except Exception as exc: # noqa: BLE001
|
||||
raise DestinationError(f"S3 download failed: {exc}") from exc
|
||||
|
||||
|
||||
def _s3_delete(cfg: dict, name: str) -> None:
|
||||
client = _s3_client(cfg)
|
||||
client.delete_object(Bucket=cfg["bucket"], Key=_s3_key(cfg, name))
|
||||
|
||||
|
||||
# --------------------------------------------------------------------------- #
|
||||
# Dispatch
|
||||
# --------------------------------------------------------------------------- #
|
||||
|
||||
|
||||
def upload(dest: BackupDestination, local_path: str, filename: str) -> str:
|
||||
cfg = parse_config(dest)
|
||||
if dest.type == "sftp":
|
||||
return _sftp_upload(cfg, local_path, filename)
|
||||
if dest.type == "s3":
|
||||
return _s3_upload(cfg, local_path, filename)
|
||||
raise DestinationError(f"Unknown destination type '{dest.type}'")
|
||||
|
||||
|
||||
def list_backups(dest: BackupDestination) -> list[dict]:
|
||||
cfg = parse_config(dest)
|
||||
if dest.type == "sftp":
|
||||
return _sftp_list(cfg)
|
||||
if dest.type == "s3":
|
||||
return _s3_list(cfg)
|
||||
raise DestinationError(f"Unknown destination type '{dest.type}'")
|
||||
|
||||
|
||||
def download(dest: BackupDestination, name: str, local_path: str) -> None:
|
||||
cfg = parse_config(dest)
|
||||
if dest.type == "sftp":
|
||||
_sftp_download(cfg, name, local_path)
|
||||
elif dest.type == "s3":
|
||||
_s3_download(cfg, name, local_path)
|
||||
else:
|
||||
raise DestinationError(f"Unknown destination type '{dest.type}'")
|
||||
|
||||
|
||||
def delete(dest: BackupDestination, name: str) -> None:
|
||||
cfg = parse_config(dest)
|
||||
if dest.type == "sftp":
|
||||
_sftp_delete(cfg, name)
|
||||
elif dest.type == "s3":
|
||||
_s3_delete(cfg, name)
|
||||
else:
|
||||
raise DestinationError(f"Unknown destination type '{dest.type}'")
|
||||
|
||||
|
||||
def test(dest: BackupDestination) -> bool:
|
||||
"""Connectivity check — lists the target (cheap, validates auth + path)."""
|
||||
list_backups(dest)
|
||||
return True
|
||||
Reference in New Issue
Block a user