_ensure_helper_image moved to stack_assets_service, but backup_destination_service imports it lazily inside _nfs_run/_nfs_helper, so nothing failed at import time — every NFS destination operation raised ImportError at runtime instead. The helper is now a public ensure_helper_image() and the NFS helpers import it from its new home. Verified: every services/ and routers/ module imports, and both NFS helper paths run through to a Docker call instead of ImportError. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
480 lines
16 KiB
Python
480 lines
16 KiB
Python
"""Push/pull stack backups to remote destinations (SFTP, S3-compatible or NFS).
|
|
|
|
All operations are synchronous (paramiko / boto3 / docker); 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
|
|
import tarfile
|
|
import tempfile
|
|
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))
|
|
|
|
|
|
# --------------------------------------------------------------------------- #
|
|
# NFS — the Docker daemon mounts the export as a named volume; file I/O runs
|
|
# through a throwaway helper container (same pattern as volume backups), so
|
|
# the backend itself needs no mount privileges.
|
|
# --------------------------------------------------------------------------- #
|
|
|
|
_NFS_VOLUME_PREFIX = "stackpilot-nfs-dest-"
|
|
|
|
|
|
def _nfs_check_name(name: str) -> None:
|
|
if not name or "/" in name or name in (".", "..") or name.startswith(".."):
|
|
raise DestinationError(f"Invalid backup file name '{name}'")
|
|
|
|
|
|
def _nfs_subdir(cfg: dict) -> str:
|
|
"""Sanitized relative directory inside the export ('' = export root).
|
|
|
|
Parts are restricted to a safe charset because the path is interpolated
|
|
into helper-container shell commands.
|
|
"""
|
|
import re
|
|
|
|
raw = (cfg.get("subdir") or "").strip().strip("/")
|
|
if not raw:
|
|
return ""
|
|
parts = [p for p in raw.split("/") if p]
|
|
for p in parts:
|
|
if p == ".." or not re.fullmatch(r"[A-Za-z0-9._-]+", p):
|
|
raise DestinationError(
|
|
"Subdirectory may only contain letters, digits, '.', '_' and '-'"
|
|
)
|
|
return "/".join(parts)
|
|
|
|
|
|
def _nfs_volume(dest: BackupDestination, cfg: dict) -> str:
|
|
"""Ensure the named volume describing this NFS mount exists; recreate it
|
|
when the destination's server/path/options changed (opts are immutable)."""
|
|
from docker_client import DockerError, get_client, safe_call
|
|
|
|
server = (cfg.get("server") or "").strip()
|
|
path = (cfg.get("path") or "").strip()
|
|
if not server:
|
|
raise DestinationError("NFS server is required")
|
|
if not path.startswith("/"):
|
|
raise DestinationError("NFS export path must be absolute (start with /)")
|
|
options = (cfg.get("options") or "rw").strip().strip(",")
|
|
driver_opts = {"type": "nfs", "o": f"addr={server},{options}", "device": f":{path}"}
|
|
|
|
name = f"{_NFS_VOLUME_PREFIX}{dest.id}"
|
|
client = get_client()
|
|
try:
|
|
vol = safe_call(client.volumes.get, name)
|
|
if (vol.attrs.get("Options") or {}) != driver_opts:
|
|
safe_call(vol.remove)
|
|
raise DockerError("recreate", "options changed")
|
|
except DockerError:
|
|
safe_call(
|
|
client.volumes.create,
|
|
name=name,
|
|
driver="local",
|
|
driver_opts=driver_opts,
|
|
labels={"stackpilot.nfs-destination": str(dest.id)},
|
|
)
|
|
return name
|
|
|
|
|
|
def _nfs_target(cfg: dict) -> str:
|
|
sub = _nfs_subdir(cfg)
|
|
return f"/nfs/{sub}" if sub else "/nfs"
|
|
|
|
|
|
def _nfs_run(volume: str, command: list[str]) -> str:
|
|
"""Run a helper container with the NFS volume at /nfs; return stdout."""
|
|
import docker.errors
|
|
|
|
from config import settings
|
|
from docker_client import DockerError, get_client
|
|
from services.stack_assets_service import ensure_helper_image
|
|
|
|
client = get_client()
|
|
ensure_helper_image(client)
|
|
try:
|
|
out = client.containers.run(
|
|
settings.BACKUP_HELPER_IMAGE,
|
|
command,
|
|
volumes={volume: {"bind": "/nfs", "mode": "rw"}},
|
|
remove=True,
|
|
)
|
|
return (out or b"").decode("utf-8", "replace")
|
|
except docker.errors.ContainerError as exc:
|
|
stderr = (exc.stderr or b"").decode("utf-8", "replace").strip()
|
|
raise DestinationError(f"NFS operation failed: {stderr or exc}") from exc
|
|
except (docker.errors.APIError, DockerError) as exc:
|
|
# Mount errors surface here (unreachable server, bad export, ...).
|
|
raise DestinationError(f"NFS mount failed: {exc}") from exc
|
|
|
|
|
|
def _nfs_helper(volume: str):
|
|
"""A created (not started) helper container for archive I/O on /nfs."""
|
|
import docker.errors
|
|
|
|
from config import settings
|
|
from docker_client import DockerError, get_client, safe_call
|
|
from services.stack_assets_service import ensure_helper_image
|
|
|
|
client = get_client()
|
|
ensure_helper_image(client)
|
|
try:
|
|
return safe_call(
|
|
client.containers.create,
|
|
settings.BACKUP_HELPER_IMAGE,
|
|
command="true",
|
|
volumes={volume: {"bind": "/nfs", "mode": "rw"}},
|
|
)
|
|
except (docker.errors.APIError, DockerError) as exc:
|
|
raise DestinationError(f"NFS mount failed: {exc}") from exc
|
|
|
|
|
|
def _nfs_upload(dest: BackupDestination, cfg: dict, local_path: str, filename: str) -> str:
|
|
import docker.errors
|
|
|
|
_nfs_check_name(filename)
|
|
volume = _nfs_volume(dest, cfg)
|
|
target = _nfs_target(cfg)
|
|
# Creates the subdir if needed AND fails early with a clear mount error.
|
|
_nfs_run(volume, ["mkdir", "-p", target])
|
|
container = _nfs_helper(volume)
|
|
try:
|
|
with tempfile.TemporaryFile() as tmp:
|
|
with tarfile.open(fileobj=tmp, mode="w") as tar:
|
|
tar.add(local_path, arcname=filename)
|
|
tmp.seek(0)
|
|
container.put_archive(target, tmp)
|
|
except docker.errors.APIError as exc:
|
|
raise DestinationError(f"NFS upload failed: {exc}") from exc
|
|
finally:
|
|
try:
|
|
container.remove(force=True)
|
|
except Exception: # noqa: BLE001
|
|
pass
|
|
sub = _nfs_subdir(cfg)
|
|
return posixpath.join(sub, filename) if sub else filename
|
|
|
|
|
|
def _nfs_list(dest: BackupDestination, cfg: dict) -> list[dict]:
|
|
volume = _nfs_volume(dest, cfg)
|
|
target = _nfs_target(cfg)
|
|
out = _nfs_run(
|
|
volume,
|
|
["sh", "-c", f"cd {target} 2>/dev/null && stat -c '%n|%s|%Y' *.tar.gz 2>/dev/null; true"],
|
|
)
|
|
entries = []
|
|
for line in out.splitlines():
|
|
parts = line.strip().split("|")
|
|
if len(parts) != 3 or parts[0] == "*.tar.gz":
|
|
continue
|
|
try:
|
|
entries.append({"name": parts[0], "size": int(parts[1]), "modified": int(parts[2])})
|
|
except ValueError:
|
|
continue
|
|
return sorted(entries, key=lambda x: x["modified"] or 0, reverse=True)
|
|
|
|
|
|
def _nfs_download(dest: BackupDestination, cfg: dict, name: str, local_path: str) -> None:
|
|
import docker.errors
|
|
|
|
_nfs_check_name(name)
|
|
volume = _nfs_volume(dest, cfg)
|
|
target = _nfs_target(cfg)
|
|
container = _nfs_helper(volume)
|
|
try:
|
|
bits, _ = container.get_archive(f"{target}/{name}")
|
|
with tempfile.TemporaryFile() as tmp:
|
|
for chunk in bits:
|
|
tmp.write(chunk)
|
|
tmp.seek(0)
|
|
with tarfile.open(fileobj=tmp) as tar:
|
|
member = next((m for m in tar.getmembers() if m.isreg()), None)
|
|
fh = tar.extractfile(member) if member else None
|
|
if fh is None:
|
|
raise DestinationError(f"'{name}' not found on NFS destination")
|
|
with open(local_path, "wb") as out:
|
|
while chunk := fh.read(1024 * 1024):
|
|
out.write(chunk)
|
|
except docker.errors.APIError as exc:
|
|
raise DestinationError(f"NFS download failed (does '{name}' exist?): {exc}") from exc
|
|
finally:
|
|
try:
|
|
container.remove(force=True)
|
|
except Exception: # noqa: BLE001
|
|
pass
|
|
|
|
|
|
def _nfs_delete(dest: BackupDestination, cfg: dict, name: str) -> None:
|
|
_nfs_check_name(name)
|
|
volume = _nfs_volume(dest, cfg)
|
|
_nfs_run(volume, ["rm", "-f", f"{_nfs_target(cfg)}/{name}"])
|
|
|
|
|
|
def _nfs_test(dest: BackupDestination, cfg: dict) -> bool:
|
|
volume = _nfs_volume(dest, cfg)
|
|
target = _nfs_target(cfg)
|
|
_nfs_run(
|
|
volume,
|
|
["sh", "-c", f"mkdir -p {target} && touch {target}/.stackpilot-test && rm -f {target}/.stackpilot-test"],
|
|
)
|
|
return True
|
|
|
|
|
|
# --------------------------------------------------------------------------- #
|
|
# 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)
|
|
if dest.type == "nfs":
|
|
return _nfs_upload(dest, 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)
|
|
if dest.type == "nfs":
|
|
return _nfs_list(dest, 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)
|
|
elif dest.type == "nfs":
|
|
_nfs_download(dest, 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)
|
|
elif dest.type == "nfs":
|
|
_nfs_delete(dest, cfg, name)
|
|
else:
|
|
raise DestinationError(f"Unknown destination type '{dest.type}'")
|
|
|
|
|
|
def test(dest: BackupDestination) -> bool:
|
|
"""Connectivity check — validates reachability, auth and write access."""
|
|
if dest.type == "nfs":
|
|
return _nfs_test(dest, parse_config(dest))
|
|
list_backups(dest)
|
|
return True
|