- No client-side encryption: plaintext content-addressed remote (SHA-256 names), no backup key, no cryptography dependency (server disk encryption is the trust model). - Filters back user data: images, archives, databases, PDFs accepted; 10 MiB cap; binary sniffing removed; temp names hardened (~$, #..#, .temp). - SQLite zero-error policy: backup-API snapshots + integrity_check, journal folding, locked/corrupt loud skips, verified restores. - Recovery: reindex from manifests, remote adopt, on-demand blob fetch, 5-day retention + thinning, GC, date-guarded remote purge, metrics. - Scheduler with WebDAV quota signal and 70% pressure backstop (floor kept). - 37 tests incl. live-monitor capture safety and DB safety.
324 lines
11 KiB
Python
324 lines
11 KiB
Python
"""WebDAV client and the automatic unique remote directory."""
|
|
|
|
from __future__ import annotations
|
|
|
|
import hashlib
|
|
import hmac
|
|
import json
|
|
import os
|
|
import re
|
|
import secrets
|
|
import socket
|
|
import time
|
|
import xml.etree.ElementTree as ET
|
|
from dataclasses import dataclass
|
|
from typing import Any
|
|
from urllib.parse import unquote, urlparse
|
|
|
|
import httpx
|
|
|
|
INSTALL_APP_ID = b"versiond:4c1b9e0f7a2d4e8b"
|
|
FORMAT_VERSION = 2
|
|
|
|
|
|
class RemoteError(Exception):
|
|
def __init__(self, code: str, detail: str):
|
|
super().__init__(detail)
|
|
self.code = code
|
|
self.detail = detail
|
|
|
|
|
|
@dataclass
|
|
class RemoteSettings:
|
|
url: str
|
|
username: str
|
|
password: str
|
|
base_path: str = "/versioned/"
|
|
verify_tls: bool = True
|
|
timeout_seconds: float = 30.0
|
|
directory: str = ""
|
|
|
|
@property
|
|
def configured(self) -> bool:
|
|
return bool(self.url and self.username and self.password)
|
|
|
|
|
|
def installation_id() -> str:
|
|
"""Stable per-user, per-machine id; the raw machine-id never leaves the host.
|
|
|
|
Same construction as systemd's sd_id128_get_machine_app_specific():
|
|
HMAC-SHA256 keyed by the machine id.
|
|
"""
|
|
try:
|
|
machine_id = open("/etc/machine-id").read().strip()
|
|
except OSError:
|
|
machine_id = ""
|
|
if not machine_id:
|
|
return secrets.token_hex(8)
|
|
try:
|
|
key = bytes.fromhex(machine_id)
|
|
except ValueError:
|
|
return secrets.token_hex(8)
|
|
msg = INSTALL_APP_ID + b":" + str(os.getuid()).encode()
|
|
return hmac.new(key, msg, hashlib.sha256).hexdigest()[:16]
|
|
|
|
|
|
def hostname_slug() -> str:
|
|
slug = re.sub(r"[^a-z0-9-]+", "-", socket.gethostname().lower()).strip("-")
|
|
return slug[:40] or "host"
|
|
|
|
|
|
def join(*parts: str) -> str:
|
|
"""Join remote path segments into '/a/b/c' (no trailing slash)."""
|
|
cleaned = [p.strip("/") for p in parts if p and p.strip("/")]
|
|
return "/" + "/".join(cleaned)
|
|
|
|
|
|
class WebDAV:
|
|
def __init__(self, settings: RemoteSettings, concurrency: int = 4,
|
|
transport: httpx.AsyncBaseTransport | None = None):
|
|
self.settings = settings
|
|
self.client = httpx.AsyncClient(
|
|
base_url=settings.url.rstrip("/"),
|
|
auth=(settings.username, settings.password),
|
|
verify=settings.verify_tls,
|
|
timeout=settings.timeout_seconds,
|
|
limits=httpx.Limits(max_connections=concurrency, max_keepalive_connections=concurrency),
|
|
transport=transport,
|
|
headers={"User-Agent": "versiond"},
|
|
)
|
|
self._known_dirs: set[str] = set()
|
|
|
|
async def close(self) -> None:
|
|
await self.client.aclose()
|
|
|
|
@staticmethod
|
|
def _check(resp: httpx.Response, *ok: int) -> httpx.Response:
|
|
if resp.status_code in ok:
|
|
return resp
|
|
if resp.status_code == 401:
|
|
raise RemoteError("remote-auth", "WebDAV server rejected the username or password")
|
|
if resp.status_code == 507:
|
|
raise RemoteError("remote-full", "WebDAV server is out of storage")
|
|
raise RemoteError("remote-http", f"{resp.request.method} {resp.request.url.path} -> HTTP {resp.status_code}")
|
|
|
|
async def request(self, method: str, path: str, **kwargs) -> httpx.Response:
|
|
try:
|
|
return await self.client.request(method, path, **kwargs)
|
|
except httpx.HTTPError as exc:
|
|
raise RemoteError("remote-unreachable", f"{method} {path}: {exc.__class__.__name__}: {exc}") from exc
|
|
|
|
async def exists(self, path: str) -> bool:
|
|
resp = await self.request("PROPFIND", path, headers={"Depth": "0"})
|
|
if resp.status_code == 404:
|
|
return False
|
|
self._check(resp, 207, 200)
|
|
return True
|
|
|
|
async def mkcol(self, path: str) -> None:
|
|
if path in self._known_dirs:
|
|
return
|
|
resp = await self.request("MKCOL", path + "/")
|
|
# 405: already exists. 301 is sent by some servers for existing collections.
|
|
self._check(resp, 201, 405, 301)
|
|
self._known_dirs.add(path)
|
|
|
|
async def makedirs(self, path: str) -> None:
|
|
current = ""
|
|
for part in path.strip("/").split("/"):
|
|
current += "/" + part
|
|
await self.mkcol(current)
|
|
|
|
async def put(self, path: str, data: bytes, if_none_match: bool = False) -> bool:
|
|
"""Upload; returns False if if_none_match was set and the object already exists."""
|
|
headers = {"Content-Type": "application/octet-stream"}
|
|
if if_none_match:
|
|
headers["If-None-Match"] = "*"
|
|
resp = await self.request("PUT", path, content=data, headers=headers)
|
|
if if_none_match and resp.status_code == 412:
|
|
return False
|
|
self._check(resp, 200, 201, 204)
|
|
return True
|
|
|
|
async def get(self, path: str) -> bytes | None:
|
|
resp = await self.request("GET", path)
|
|
if resp.status_code == 404:
|
|
return None
|
|
self._check(resp, 200)
|
|
return resp.content
|
|
|
|
async def delete(self, path: str) -> None:
|
|
resp = await self.request("DELETE", path)
|
|
self._check(resp, 200, 204, 404)
|
|
|
|
|
|
async def check_connection(dav: WebDAV) -> None:
|
|
"""Fail with a clear RemoteError unless we can list, create and delete."""
|
|
base = join(dav.settings.base_path)
|
|
await dav.makedirs(base)
|
|
probe = join(base, f".versiond-probe-{secrets.token_hex(4)}")
|
|
await dav.put(probe, b"probe")
|
|
await dav.delete(probe)
|
|
|
|
|
|
async def claim_directory(dav: WebDAV) -> dict[str, Any]:
|
|
"""Pick and claim this installation's unique directory below base_path.
|
|
|
|
Returns {"directory": ..., "adopted": bool}. Never takes over a directory that
|
|
belongs to a different installation id. Remote data is stored unencrypted
|
|
(relying on server disk encryption); blobs are plain content-addressed
|
|
objects named by their SHA-256.
|
|
"""
|
|
base = join(dav.settings.base_path)
|
|
await dav.makedirs(base)
|
|
install_id = installation_id()
|
|
name = f"{hostname_slug()}-{install_id[:8]}"
|
|
for _ in range(5):
|
|
directory = join(base, name)
|
|
owner_path = join(directory, "meta", "owner.json")
|
|
existing = await dav.get(owner_path)
|
|
if existing is None:
|
|
for sub in ("meta", "blobs", "manifests"):
|
|
await dav.makedirs(join(directory, sub))
|
|
owner = {
|
|
"installation_id": install_id,
|
|
"hostname": socket.gethostname(),
|
|
"user": os.environ.get("USER", ""),
|
|
"created_at": time.time(),
|
|
}
|
|
if not await dav.put(owner_path, json.dumps(owner, indent=2).encode(), if_none_match=True):
|
|
continue # someone created it between our GET and PUT; re-read
|
|
fmt = {"format": FORMAT_VERSION, "encryption": "none", "compression": "zstd-or-zlib",
|
|
"blob_layout": "blobs/<sha256[:2]>/<sha256>"}
|
|
await dav.put(join(directory, "meta", "format.json"), json.dumps(fmt, indent=2).encode())
|
|
return {"directory": directory, "adopted": False}
|
|
try:
|
|
owner = json.loads(existing)
|
|
except ValueError:
|
|
owner = {}
|
|
if owner.get("installation_id") == install_id:
|
|
return {"directory": directory, "adopted": True}
|
|
name = f"{hostname_slug()}-{install_id[:8]}-{secrets.token_hex(2)}"
|
|
raise RemoteError("claim-failed", "could not claim a unique remote directory")
|
|
|
|
|
|
# remote listing (read-only PROPFIND) for the purge endpoint
|
|
|
|
DAV_NS = "DAV:"
|
|
_PROPFIND_BODY = (
|
|
'<?xml version="1.0" encoding="utf-8"?>'
|
|
'<d:propfind xmlns:d="DAV:"><d:prop>'
|
|
"<d:resourcetype/><d:getcontentlength/>"
|
|
"</d:prop></d:propfind>"
|
|
).encode()
|
|
|
|
_QUOTA_BODY = (
|
|
'<?xml version="1.0" encoding="utf-8"?>'
|
|
'<d:propfind xmlns:d="DAV:"><d:prop>'
|
|
"<d:quota-used-bytes/><d:quota-available-bytes/>"
|
|
"</d:prop></d:propfind>"
|
|
).encode()
|
|
|
|
|
|
@dataclass
|
|
class RemoteEntry:
|
|
path: str
|
|
is_dir: bool
|
|
size: int | None = None
|
|
|
|
|
|
def _href_to_path(href: str) -> str:
|
|
path = unquote(urlparse(href).path or href)
|
|
return path.rstrip("/") or "/"
|
|
|
|
|
|
async def list_dir(dav: "WebDAV", path: str) -> list[RemoteEntry]:
|
|
"""Immediate children of a collection via PROPFIND Depth:1 (read-only)."""
|
|
resp = await dav.request(
|
|
"PROPFIND", path,
|
|
headers={"Depth": "1", "Content-Type": "application/xml"},
|
|
content=_PROPFIND_BODY,
|
|
)
|
|
if resp.status_code == 404:
|
|
return []
|
|
WebDAV._check(resp, 207)
|
|
try:
|
|
root = ET.fromstring(resp.content)
|
|
except ET.ParseError as exc:
|
|
raise RemoteError("remote-http", f"PROPFIND {path} returned unparseable XML") from exc
|
|
entries: list[RemoteEntry] = []
|
|
for response in root.findall(f"{{{DAV_NS}}}response"):
|
|
href = response.findtext(f"{{{DAV_NS}}}href")
|
|
if not href:
|
|
continue
|
|
child = _href_to_path(href)
|
|
if child == path.rstrip("/") or child == path:
|
|
continue # the collection itself, not a child
|
|
is_dir = False
|
|
size: int | None = None
|
|
for propstat in response.findall(f"{{{DAV_NS}}}propstat"):
|
|
prop = propstat.find(f"{{{DAV_NS}}}prop")
|
|
if prop is None:
|
|
continue
|
|
restype = prop.find(f"{{{DAV_NS}}}resourcetype")
|
|
if restype is not None and restype.find(f"{{{DAV_NS}}}collection") is not None:
|
|
is_dir = True
|
|
length = prop.findtext(f"{{{DAV_NS}}}getcontentlength")
|
|
if length is not None:
|
|
try:
|
|
size = int(length)
|
|
except ValueError:
|
|
size = None
|
|
entries.append(RemoteEntry(child, is_dir, size))
|
|
return entries
|
|
|
|
|
|
async def list_tree(dav: "WebDAV", top: str) -> tuple[list[RemoteEntry], list[RemoteEntry]]:
|
|
"""All files and directories strictly below top. Returns (files, dirs)."""
|
|
files: list[RemoteEntry] = []
|
|
dirs: list[RemoteEntry] = []
|
|
stack = [top]
|
|
while stack:
|
|
current = stack.pop()
|
|
for entry in await list_dir(dav, current):
|
|
if entry.is_dir:
|
|
dirs.append(entry)
|
|
stack.append(entry.path)
|
|
else:
|
|
files.append(entry)
|
|
return files, dirs
|
|
|
|
|
|
async def quota(dav: "WebDAV", path: str) -> tuple[int | None, int | None]:
|
|
"""RFC 4331 quota properties. Returns (used_bytes, available_bytes or None).
|
|
|
|
(None, None) when the server exposes no quota information.
|
|
"""
|
|
resp = await dav.request(
|
|
"PROPFIND", path,
|
|
headers={"Depth": "0", "Content-Type": "application/xml"},
|
|
content=_QUOTA_BODY,
|
|
)
|
|
if resp.status_code == 404:
|
|
return None, None
|
|
WebDAV._check(resp, 207)
|
|
try:
|
|
root = ET.fromstring(resp.content)
|
|
except ET.ParseError:
|
|
return None, None
|
|
used = available = None
|
|
for prop in root.findall(f".//{{{DAV_NS}}}prop"):
|
|
used_text = prop.findtext(f"{{{DAV_NS}}}quota-used-bytes")
|
|
avail_text = prop.findtext(f"{{{DAV_NS}}}quota-available-bytes")
|
|
if used_text is not None:
|
|
try:
|
|
used = int(used_text)
|
|
except ValueError:
|
|
used = None
|
|
if avail_text is not None:
|
|
try:
|
|
available = int(avail_text)
|
|
except ValueError:
|
|
available = None
|
|
return used, available
|