Files
stackpilot/backend/services/backup_destination_service.py
T
menzeljandClaude Opus 5 adfd77a983 Fix NFS uploads against root_squash exports (0.40.2)
Uploading a backup extracted the tar straight into the NFS mount via
put_archive, and the daemon chowns every entry while extracting — an export
with root_squash refuses that ("failed to Lchown ... for UID 0, GID 0:
operation not permitted"), so the upload died with a docker 500 even though
plain writes to the share work (which is why the destination test passed).

The helper container now unpacks into its own filesystem and copies the file
into the mount with cat, which never chowns. Restores hit the same wall when a
volume or bind folder lives on a squashed mount, so import_path/import_volume
fall back to a copy-through-staging when (and only when) the failure is a chown
denial — local restores keep preserving ownership. NFS file names are validated
against the same safe charset as the subdir parts, since both are interpolated
into the helper's shell commands.

Verified against a real root_squash NFS export: test/upload/list/download/delete
round trip, byte-identical download, restore into an NFS-backed volume via the
fallback, and ownership still preserved (1000:1000, 0600) on a local volume.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-08-16 18:53:25 +00:00

494 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:
import re
# Same safe charset as the subdir parts: the name is interpolated into the
# helper container's shell commands.
if not name or not re.fullmatch(r"[A-Za-z0-9._-]+", name) 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, command: list[str] | str = "true"):
"""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=command,
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])
# Unpack into the container's own filesystem, then copy the file across:
# extracting straight into the NFS mount makes the daemon chown the file,
# which a root_squash export refuses ("failed to Lchown ... for UID 0").
container = _nfs_helper(
volume, ["sh", "-c", f"cat '/tmp/{filename}' > '{target}/{filename}'"]
)
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("/tmp", tmp)
container.start()
status = container.wait(timeout=3600).get("StatusCode", 1)
if status != 0:
err = (container.logs(stdout=True, stderr=True) or b"").decode("utf-8", "replace")
raise DestinationError(f"NFS upload failed: {err.strip() or f'exit {status}'}")
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