diff --git a/pyproject.toml b/pyproject.toml index 1b240af..705b246 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -11,11 +11,12 @@ requires-python = ">=3.12" dependencies = [ "fastapi>=0.110", "uvicorn>=0.29", - "asyncinotify>=4.0", + "httpx>=0.27", + "cryptography>=42", ] [project.optional-dependencies] -dev = ["pytest>=8", "httpx>=0.27"] +dev = ["pytest>=8"] [project.scripts] versiond = "versiond.cli:main" @@ -23,5 +24,8 @@ versiond = "versiond.cli:main" [tool.setuptools.packages.find] where = ["src"] +[tool.setuptools.package-data] +versiond = ["dashboard.html"] + [tool.pytest.ini_options] testpaths = ["tests"] diff --git a/readme.md b/readme.md index 82df0c4..c5364c5 100644 --- a/readme.md +++ b/readme.md @@ -10,10 +10,11 @@ | Milestone | State | |---|---| -| M1 – Core (monitor, filters, index, spool, history, diff, restore, systemd install) | **Implemented** (`src/versiond/`) | -| M2 – WebDAV remote, unique remote directory, encryption | Not started. Versions are stored locally (`durability = local`) | -| M3 – Retention, purge, GC | Not started (pinning exists) | -| M4 – Agent skill, metrics, reindex/adopt | Not started | +| M1 – Core (monitor, filters, index, spool, history, diff, restore, systemd install) | **Implemented** | +| M2 – WebDAV remote, unique remote directory, encryption, manifests | **Implemented** | +| Statistics, progress (incl. live stream) and web dashboard | **Implemented** | +| M3 – Retention, thinning, remote GC | Not started (pinning and `forget` exist) | +| M4 – Agent skill, Prometheus metrics, reindex from remote | Not started | ```bash pipx install -e . # installs the `versiond` command @@ -24,6 +25,12 @@ versiond history ~/projects/app/main.py versiond diff ~/projects/app/main.py # last change versiond restore '~/projects/app/src/*' --as-of 2026-10-08T14:00 # dry run versiond restore '~/projects/app/src/*' --as-of 2026-10-08T14:00 --execute +versiond remote set --url https://u123456.your-storagebox.de --user u123456 # prompts for the password +versiond key export # store the backup key somewhere safe, off this machine +versiond progress --follow # scan + upload progress, rate and ETA +versiond stats # files, versions, storage, activity +versiond dashboard # opens the web dashboard, already signed in +versiond forget ~/projects/app/data --execute # drop history of something you now ignore ``` API docs: . API calls need `Authorization: Bearer $(versiond token)`. @@ -66,7 +73,7 @@ The primary use case is a safety net for workflows where files are rewritten fre | Web framework | FastAPI + Uvicorn | | Bind address | `127.0.0.1:9922` (fixed default, overridable for testing only) | | Metadata store | SQLite (WAL mode) | -| File monitoring | Linux inotify via `asyncinotify` (pure ctypes, asyncio-native, no dependencies) | +| File monitoring | Linux inotify via a built-in ctypes binding (no dependencies), read on the asyncio loop | | Remote storage | WebDAV (RFC 4918) | | API docs | Swagger UI at `/docs`, ReDoc at `/redoc`, schema at `/openapi.json` | @@ -181,11 +188,14 @@ Monitored roots do not need their own remote directories: blobs are shared (cont ``` /-/ - blobs/sha256/ab/cd/.zst # content-addressed, zstd-compressed (optionally encrypted) - manifests///
/.jsonl.zst # append-only version records - meta/format.json # storage format version + blobs// # name = HMAC-SHA256(key, sha256); zstd, then AES-256-GCM + manifests///
/.jsonl.zst.enc # append-only version + rename records, encrypted + meta/owner.json # installation id, hostname, key id (plaintext) + meta/format.json # storage format version (plaintext) ``` +One directory level under `blobs/` (256 prefixes) keeps the number of `MKCOL` requests small. + Blobs are immutable and deduplicated by SHA-256 of the plaintext content. Manifests are append-only journals. Together they make it possible to fully rebuild the index (`POST /admin/reindex`) after local data loss. ### 4.2 Ingestion @@ -198,7 +208,7 @@ The user registers one or more **roots** (`POST /roots {"path": "~/projects"}`); | Option | Verdict | |---|---| -| **Raw inotify (`asyncinotify`)** | **Chosen.** One kernel watch per *directory* (not per file), costing about 1 KB of kernel memory each. Events are delivered on one file descriptor read directly by the asyncio loop: no threads, no polling, ~0 % CPU when idle. Pure ctypes, no dependencies. | +| **Raw inotify (own ctypes binding)** | **Chosen.** One kernel watch per *directory* (not per file), costing about 1 KB of kernel memory each. Events are delivered on one file descriptor read directly by the asyncio loop (`add_reader`): no threads, no polling, ~0 % CPU when idle. Watched directories are kept as a tree of `(wd, parent, name)` nodes, so a directory rename is one pointer change and memory stays at ~100 bytes per watch. `asyncinotify` was tried first; its per-watch `Path` objects cost ~260 MB RSS at 136k watches versus ~90 MB with the tree. | | `watchfiles` (Rust `notify`) | Good library, but its recursive mode adds a watch to *every* directory, including `node_modules`/`.venv`; filters only drop events afterwards. It can use up the watch limit on large trees, and it pulls in a compiled dependency. | | `watchdog` | Heavier (emitter threads, snapshot objects), same one-watch-per-directory cost, cross-platform layer we don't need. | | fanotify | Unprivileged use (kernel ≥ 5.13) does not allow filesystem- or mount-wide marks, so it still needs one mark per directory, with a more complex API. No benefit without root. | @@ -286,7 +296,7 @@ All upload behaviour is configurable (`[upload]` section and `PATCH /config/uplo | `batch_manifest_interval_s` | 30 | Manifest flush interval | | `spool_max_mb` | 1024 | When full, ingest returns `507 Insufficient Storage` | -Uploads are idempotent (content-addressed; `HEAD` before `PUT`). A version counts as **durable** only after its blob and manifest entry are both stored remotely. The API shows this state per version (`pending` / `durable`). +Uploads are idempotent (content-addressed names, so a repeated `PUT` writes identical data). No `HEAD` is sent before `PUT`; that would double the requests for the common case where the blob is new. A version counts as **durable** only after its blob and manifest entry are both stored remotely. The API shows this state per version (`pending` / `durable`). When the remote is unreachable, the service keeps working from the spool and catches up when the connection returns. @@ -373,6 +383,11 @@ All endpoints are JSON, under `/api/v1` (the docs, schema and skill endpoints ar | POST | `/versions/{id}/pin` / `DELETE` | Pin / unpin | | POST | `/purge` | Run a purge (criteria + `dry_run`) | | DELETE | `/versions/{id}`, `/files?path=`, `/projects/{id}` | Explicit deletion | +| GET | `/stats` | Files, versions (per source, per hour for 24 h), storage, dedup/compression, top files and projects | +| GET | `/progress` | Active and last scans (percent, files/s, ETA), upload queue (percent, rate, ETA, errors), monitor counters | +| GET | `/progress/stream` | Server-sent events with `/progress` every `interval` seconds | +| GET | `/dashboard` | Web dashboard (token via `#token=` fragment or prompt) | +| POST | `/forget` | Delete all history below a path (`dry_run` by default) | | POST | `/admin/gc` | Garbage-collect remote blobs | | POST | `/admin/reindex` | Rebuild the index from remote manifests | | GET | `/admin/queue` | Upload queue status | @@ -405,7 +420,7 @@ Indexes on `(file_id, captured_at)`, `files(path)`, `versions(blob_sha256)`. Sch - **Loopback only.** The server refuses to start if configured to bind to a non-loopback address. - **Local authentication.** Other local users and processes (including browsers, via DNS rebinding) can reach `127.0.0.1`. Every request therefore needs a bearer token stored in `~/.config/versiond/credentials` (`0600`). The `Host` header is checked against `127.0.0.1:9922`/`localhost:9922`, and CORS is disabled. -- **Secrets in backed-up files.** `.env` and similar files contain credentials and are sent off-machine. Client-side encryption is therefore **on by default**: blobs and manifests are encrypted with XChaCha20-Poly1305 using a key held locally (with an export/recovery procedure documented at first run). Content-addressing uses a keyed hash (HMAC-SHA-256) so the remote cannot confirm guesses of known content. +- **Secrets in backed-up files.** `.env` and similar files contain credentials and are sent off-machine. Client-side encryption is therefore **always on**: blobs and manifests are encrypted with AES-256-GCM (random 96-bit nonce) using a 64-byte key in `~/.config/versiond/backup.key` (`0600`). Remote blob names are HMAC-SHA-256 of the content hash, so the server cannot confirm guesses of known content. **The key is the only way to read the remote copy**: `versiond key export` prints it, and it must be stored off the machine. - **Credential storage.** WebDAV secrets are kept in the `0600` credentials file or, when available, the Secret Service/keyring. They never appear in logs or API responses. - **Restore path safety.** See §4.7. - **Audit log.** All configuration changes, restores and purges are recorded with timestamp and client identity. diff --git a/src/versiond/api.py b/src/versiond/api.py index 29470ce..38aee69 100644 --- a/src/versiond/api.py +++ b/src/versiond/api.py @@ -2,9 +2,11 @@ from __future__ import annotations +import asyncio import base64 import binascii import fnmatch +import json import logging import os import secrets @@ -15,18 +17,24 @@ from datetime import datetime from typing import Any, Literal from fastapi import APIRouter, Depends, FastAPI, Query, Request -from fastapi.responses import HTMLResponse, JSONResponse, PlainTextResponse, Response +from fastapi.responses import HTMLResponse, JSONResponse, PlainTextResponse, Response, StreamingResponse from pydantic import BaseModel, Field +from importlib import resources + from . import __version__ from . import diff as difflib_ +from . import stats as stats_ from .config import Config +from .crypto import KeyRing from .db import Database from .filters import Filters -from .ingest import Coalescer, Rejected, Repository, TokenBucket +from .ingest import Coalescer, Rejected, Repository, TokenBucket, tree_range from .monitor import Monitor, RootError +from .remote import RemoteError, RemoteSettings from .restore import Criteria, Restorer, RestoreError from .store import BlobStore +from .uploader import Uploader log = logging.getLogger(__name__) @@ -40,6 +48,8 @@ class Problem(Exception): REJECT_STATUS = {"file-too-large": 413} ROOT_STATUS = {"not-found": 404, "not-a-directory": 400, "root-overlap": 409, "watch-limit": 409} +REMOTE_STATUS = {"remote-auth": 400, "remote-unreachable": 502, "remote-http": 502, + "remote-full": 507, "key-mismatch": 409, "claim-failed": 409} RESTORE_STATUS = {"not-found": 404, "unsafe-path": 400, "criteria-too-broad": 400, "plan-expired": 410, "plan-used": 409, "conflict": 409} @@ -110,6 +120,20 @@ class VersionRestoreIn(BaseModel): conflict: Literal["overwrite", "skip", "rename"] = "overwrite" +class RemoteIn(BaseModel): + url: str = Field(..., description="WebDAV server URL, e.g. https://user.your-storagebox.de") + username: str + password: str | None = Field(None, description="omit to keep the stored password") + base_path: str = "/versioned/" + verify_tls: bool = True + timeout_seconds: float = 30.0 + + +class ForgetIn(BaseModel): + path: str = Field(..., description="file or directory whose history is deleted") + dry_run: bool = True + + class RestoreIn(BaseModel): as_of: str | float | None = None paths: list[str] = [] @@ -128,7 +152,7 @@ class RestoreIn(BaseModel): # services class Services: - def __init__(self, cfg: Config): + def __init__(self, cfg: Config, remote_transport=None): self.cfg = cfg self.filters = Filters.from_config(cfg) self.db = Database(cfg.paths.index_file) @@ -147,21 +171,26 @@ class Services: max_watch_fraction=float(cfg.get("monitor", "max_watch_fraction")), ) self.restorer = Restorer(self.db, self.repo, self.blobs) + self.keys = KeyRing.load_or_create(cfg.key_file) + self.uploader = Uploader(cfg, self.db, self.blobs, self.keys, transport=remote_transport) + self.repo.on_commit = self.uploader.wake rate = float(cfg.get("limits", "ingest_requests_per_second")) self.ingest_limit = TokenBucket(rate, 1.0, time.monotonic) self.token = cfg.api_token() self.started_at = time.time() -def create_app(cfg: Config | None = None, allowed_hosts: set[str] | None = None) -> FastAPI: +def create_app(cfg: Config | None = None, allowed_hosts: set[str] | None = None, + remote_transport=None) -> FastAPI: cfg = cfg or Config.load() hosts = allowed_hosts or {f"127.0.0.1:{cfg.port}", f"localhost:{cfg.port}"} @asynccontextmanager async def lifespan(app: FastAPI): - services = Services(cfg) + services = Services(cfg, remote_transport) app.state.services = services await services.monitor.start() + await services.uploader.start() sd_notify("READY=1") log.info("versiond %s ready on %s:%s", __version__, cfg.host, cfg.port) try: @@ -169,6 +198,7 @@ def create_app(cfg: Config | None = None, allowed_hosts: set[str] | None = None) finally: sd_notify("STOPPING=1") await services.monitor.stop() + await services.uploader.stop() flushed = services.coalescer.flush() log.info("shutdown: flushed %d pending versions", flushed) services.db.close() @@ -199,6 +229,10 @@ def create_app(cfg: Config | None = None, allowed_hosts: set[str] | None = None) async def on_root_error(_: Request, exc: RootError): return _problem(ROOT_STATUS.get(exc.code, 400), exc.code, exc.detail) + @app.exception_handler(RemoteError) + async def on_remote_error(_: Request, exc: RemoteError): + return _problem(REMOTE_STATUS.get(exc.code, 502), exc.code, exc.detail) + @app.exception_handler(RestoreError) async def on_restore_error(_: Request, exc: RestoreError): return _problem(RESTORE_STATUS.get(exc.code, 400), exc.code, exc.detail) @@ -225,9 +259,14 @@ def create_app(cfg: Config | None = None, allowed_hosts: set[str] | None = None) "roots": len(roots), "degraded_roots": degraded, "pending_paths": sum(1 for p in s.coalescer.pending.values() if p.content is not None), - "watches": len(s.monitor.watches), + "watches": s.monitor.watch_total, + "remote": s.uploader.state, } + @app.get("/dashboard", include_in_schema=False) + async def dashboard() -> HTMLResponse: + return HTMLResponse(resources.files("versiond").joinpath("dashboard.html").read_text()) + api = APIRouter(prefix="/api/v1", dependencies=[Depends(authenticate)]) # config & filters @@ -247,6 +286,83 @@ def create_app(cfg: Config | None = None, allowed_hosts: set[str] | None = None) out.append({"path": path, "accepted": reason is None, "reason": reason}) return out + # remote + + def _settings(s: Services, body: RemoteIn) -> RemoteSettings: + password = body.password if body.password is not None else s.cfg.credential("remote_password") + if not password: + raise Problem(400, "missing-password", "password is required") + return RemoteSettings(url=body.url, username=body.username, password=password, + base_path=body.base_path, verify_tls=body.verify_tls, + timeout_seconds=body.timeout_seconds) + + @api.get("/config/remote", tags=["remote"]) + async def get_remote(s: Services = Depends(services)) -> dict[str, Any]: + return s.uploader.public_settings() + + @api.put("/config/remote", tags=["remote"]) + async def set_remote(body: RemoteIn, s: Services = Depends(services)) -> dict[str, Any]: + return await s.uploader.configure(_settings(s, body)) + + @api.post("/config/remote/test", tags=["remote"]) + async def test_remote(body: RemoteIn, s: Services = Depends(services)) -> dict[str, Any]: + await s.uploader.test(_settings(s, body)) + return {"ok": True} + + # statistics and progress + + def _progress(s: Services) -> dict[str, Any]: + roots = [] + for root in s.db.roots(): + roots.append({ + "id": root["id"], "path": root["path"], + "mode": s.monitor.root_mode(root["id"], root["path"]), + "watches": s.monitor.watch_counts[root["id"]], + "baseline_state": root["baseline_state"], "last_scan_at": root["last_scan_at"], + }) + return { + "time": time.time(), + "uptime_seconds": round(time.time() - s.started_at), + "monitor": { + "watches": s.monitor.watch_total, + "roots": roots, + "pending_paths": sum(1 for p in s.coalescer.pending.values() if p.content is not None), + "counters": dict(s.monitor.counters), + }, + "scans": { + "active": [p.as_dict() for p in s.monitor.scans.values()], + "last": [p.as_dict() for p in s.monitor.last_scans.values()], + }, + "upload": s.uploader.progress(), + } + + @api.get("/stats", tags=["status"]) + async def get_stats(s: Services = Depends(services)) -> dict[str, Any]: + return stats_.collect(s.db) + + @api.get("/progress", tags=["status"]) + async def get_progress(s: Services = Depends(services)) -> dict[str, Any]: + return _progress(s) + + @api.get("/progress/stream", tags=["status"]) + async def stream_progress( + request: Request, + interval: float = Query(1.0, ge=0.25, le=60), + max_events: int = Query(0, ge=0, description="stop after this many events (0 = until disconnect)"), + s: Services = Depends(services), + ): + """Server-sent events: one `progress` event per interval.""" + async def events(): + sent = 0 + while not await request.is_disconnected(): + yield f"event: progress\ndata: {json.dumps(_progress(s))}\n\n" + sent += 1 + if max_events and sent >= max_events: + return + await asyncio.sleep(interval) + return StreamingResponse(events(), media_type="text/event-stream", + headers={"Cache-Control": "no-store"}) + # snapshots def _snapshot(s: Services, body: SnapshotIn) -> dict[str, Any]: @@ -492,6 +608,31 @@ def create_app(cfg: Config | None = None, allowed_hosts: set[str] | None = None) return {"from": from_label, "to": to_label, "identical": not hunk_list, "hunks": hunk_list} return PlainTextResponse(difflib_.unified(hunk_list, from_label, to_label)) + @api.post("/forget", tags=["versions"]) + async def forget(body: ForgetIn, s: Services = Depends(services)) -> dict[str, Any]: + """Delete all history for a file or directory tree (e.g. after adding an ignore pattern).""" + path = os.path.abspath(os.path.expanduser(body.path)) + lo, hi = tree_range(path) + counts = s.db.one( + "SELECT count(*) AS files, (SELECT count(*) FROM versions v JOIN files f ON f.id = v.file_id " + "WHERE f.path = ? OR (f.path >= ? AND f.path < ?)) AS versions " + "FROM files WHERE path = ? OR (path >= ? AND path < ?)", + (path, lo, hi, path, lo, hi), + ) + if body.dry_run: + return {"path": path, "dry_run": True, **counts} + files = blobs = 0 + while chunk := s.repo.forget_chunk(path): + files += chunk + await asyncio.sleep(0) + while chunk := s.repo.drop_orphan_blobs(): + blobs += chunk + await asyncio.sleep(0) + projects = s.repo.drop_empty_projects() + s.db.audit("forget", path=path, files=files, blobs=blobs) + return {"path": path, "dry_run": False, "files_removed": files, + "versions_removed": counts["versions"], "blobs_removed": blobs, "projects_removed": projects} + # bulk restore @api.post("/restores", tags=["restore"], status_code=201) diff --git a/src/versiond/cli.py b/src/versiond/cli.py index 4d7db3f..c379b8d 100644 --- a/src/versiond/cli.py +++ b/src/versiond/cli.py @@ -3,6 +3,7 @@ from __future__ import annotations import argparse +import getpass import json import logging import os @@ -134,7 +135,7 @@ class Client: self.token = self.cfg.api_token() def request(self, method: str, path: str, body: Any = None, query: dict | None = None, - raw: bool = False) -> Any: + raw: bool = False, timeout: float = 60) -> Any: if query: path += "?" + urllib.parse.urlencode({k: v for k, v in query.items() if v is not None}) data = json.dumps(body).encode() if body is not None else None @@ -143,7 +144,7 @@ class Client: if data is not None: req.add_header("Content-Type", "application/json") try: - with urllib.request.urlopen(req, timeout=60) as resp: + with urllib.request.urlopen(req, timeout=timeout) as resp: payload = resp.read() except urllib.error.HTTPError as exc: detail = exc.read().decode(errors="replace") @@ -261,6 +262,141 @@ def cmd_restore(args: argparse.Namespace) -> int: return 0 +def _bytes(n: float | None) -> str: + if n is None: + return "-" + for unit in ("B", "KB", "MB", "GB"): + if n < 1024: + return f"{n:.0f} {unit}" if unit == "B" else f"{n:.1f} {unit}" + n /= 1024 + return f"{n:.1f} TB" + + +def _dur(seconds: float | None) -> str: + if seconds is None: + return "-" + seconds = int(seconds) + if seconds < 3600: + return f"{seconds // 60}m{seconds % 60:02d}s" + return f"{seconds // 3600}h{(seconds % 3600) // 60:02d}m" + + +def cmd_remote(args: argparse.Namespace) -> int: + c = Client() + if args.action == "show": + print(json.dumps(c.request("GET", "/api/v1/config/remote"), indent=2)) + return 0 + current = c.request("GET", "/api/v1/config/remote") + url = args.url or current["url"] + user = args.user or current["username"] + if not url or not user: + raise SystemExit("give --url and --user") + if args.password_stdin: + password = sys.stdin.readline().rstrip("\n") + elif current["password_set"] and not args.new_password: + password = None + else: + password = getpass.getpass(f"WebDAV password for {user}: ") + body = {"url": url, "username": user, "password": password, + "base_path": args.base_path or current["base_path"]} + if args.action == "test": + c.request("POST", "/api/v1/config/remote/test", body) + print("connection ok") + return 0 + result = c.request("PUT", "/api/v1/config/remote", body) + verb = "adopted existing" if result.get("adopted") else "claimed new" + print(f"remote configured; {verb} directory {result['directory']} on {result['url']}") + print(f"encryption: {result['encryption']}, key id {result['key_id']}") + print("IMPORTANT: back up your key with `versiond key export`; without it the remote copy cannot be read.") + return 0 + + +def cmd_forget(args: argparse.Namespace) -> int: + result = Client().request("POST", "/api/v1/forget", {"path": _abs(args.path), "dry_run": not args.execute}, + timeout=3600) + if result["dry_run"]: + print(f"would delete {result['files']:,} files / {result['versions']:,} versions under {result['path']}") + print("add --execute to delete") + else: + print(f"deleted {result['files_removed']:,} files, {result['versions_removed']:,} versions, " + f"{result['blobs_removed']:,} blobs, {result['projects_removed']} projects") + return 0 + + +def cmd_key(args: argparse.Namespace) -> int: + from .crypto import KeyRing + + keys = KeyRing.load_or_create(Config.load().key_file) + print(keys.export()) + print(f"# key id {keys.key_id} - store this somewhere safe and off this machine", file=sys.stderr) + return 0 + + +def cmd_stats(args: argparse.Namespace) -> int: + s = Client().request("GET", "/api/v1/stats") + if args.json: + print(json.dumps(s, indent=2)) + return 0 + st = s["storage"] + print(f"files {s['files']['tracked']:,} tracked, {s['files']['existing']:,} on disk, {s['projects']:,} projects") + v = s["versions"] + print(f"versions {v['total']:,} total, {v['last_hour']:,} last hour, {v['last_24h']:,} last 24h, {v['last_7d']:,} last 7d") + print(f"storage {_bytes(st['stored_bytes'])} stored ({_bytes(st['raw_bytes'])} raw, " + f"{st['compression_ratio']}x compression, {st['dedup_ratio']}x dedup)") + print("sources " + ", ".join(f"{k}={n:,}" for k, n in v["by_source"].items())) + print("per hour " + " ".join(str(n) for n in v["per_hour_last_24h"])) + if s["top_files_7d"]: + print("most edited (7d):") + for row in s["top_files_7d"]: + print(f" {row['versions']:>5} {row['path']}") + return 0 + + +def _print_progress(p: dict) -> None: + for scan in p["scans"]["active"]: + print(f"scan {scan['reason']} {scan['top']}: {scan['percent']}% " + f"({scan['dirs_done']:,}/{scan['dirs_total']:,} dirs, {scan['files_seen']:,} files, " + f"{scan['versions_stored']:,} new) eta {_dur(scan['eta_seconds'])}") + if not p["scans"]["active"]: + print("scan idle") + u = p["upload"] + if u["state"] == "unconfigured": + print("upload not configured (versiond remote set)") + else: + print(f"upload {u['state']}: {u['percent']}% ({u['blobs_uploaded']:,}/{u['blobs_total']:,} blobs, " + f"{u['blobs_pending']:,} pending {_bytes(u['bytes_pending'])}) " + f"{u['blobs_per_second']}/s {_bytes(u['bytes_per_second'])}/s eta {_dur(u['eta_seconds'])}") + print(f"durable {u['versions_durable']:,} versions backed up, {u['versions_local']:,} local only") + if u["last_error"]: + print(f"error {u['last_error']}") + m = p["monitor"] + print(f"monitor {m['watches']:,} watches, {m['pending_paths']} pending, {m['counters'].get('events', 0):,} events") + + +def cmd_progress(args: argparse.Namespace) -> int: + c = Client() + while True: + p = c.request("GET", "/api/v1/progress") + if args.json: + print(json.dumps(p)) + else: + if args.follow: + print("\033[2J\033[H", end="") + _print_progress(p) + if not args.follow: + return 0 + time.sleep(args.interval) + + +def cmd_dashboard(args: argparse.Namespace) -> int: + cfg = Config.load() + url = f"http://{cfg.host}:{cfg.port}/dashboard#token={cfg.api_token()}" + print(url) + if not args.no_open and shutil.which("xdg-open"): + subprocess.Popen(["xdg-open", url], stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL) + return 0 + + def build_parser() -> argparse.ArgumentParser: parser = argparse.ArgumentParser(prog="versiond", description="Local file versioning service") parser.add_argument("--version", action="version", version=f"versiond {__version__}") @@ -317,6 +453,38 @@ def build_parser() -> argparse.ArgumentParser: p.add_argument("--conflict", choices=["overwrite", "skip", "rename", "fail"], default="overwrite") p.add_argument("--execute", action="store_true") p.set_defaults(func=cmd_restore) + + p = sub.add_parser("remote", help="configure the WebDAV backup target") + p.add_argument("action", choices=["set", "show", "test"]) + p.add_argument("--url", help="e.g. https://u123456.your-storagebox.de") + p.add_argument("--user") + p.add_argument("--base-path", help="parent directory on the server (default /versioned/)") + p.add_argument("--password-stdin", action="store_true", help="read the password from stdin") + p.add_argument("--new-password", action="store_true", help="prompt even if a password is stored") + p.set_defaults(func=cmd_remote) + + p = sub.add_parser("forget", help="delete all history for a file or directory (dry run unless --execute)") + p.add_argument("path") + p.add_argument("--execute", action="store_true") + p.set_defaults(func=cmd_forget) + + p = sub.add_parser("key", help="backup encryption key") + p.add_argument("action", choices=["export"]) + p.set_defaults(func=cmd_key) + + p = sub.add_parser("stats", help="statistics: files, versions, storage, activity") + p.add_argument("--json", action="store_true") + p.set_defaults(func=cmd_stats) + + p = sub.add_parser("progress", help="scan and upload progress") + p.add_argument("-f", "--follow", action="store_true") + p.add_argument("--interval", type=float, default=2.0) + p.add_argument("--json", action="store_true") + p.set_defaults(func=cmd_progress) + + p = sub.add_parser("dashboard", help="open the web dashboard") + p.add_argument("--no-open", action="store_true") + p.set_defaults(func=cmd_dashboard) return parser diff --git a/src/versiond/config.py b/src/versiond/config.py index 7650c8a..80810fc 100644 --- a/src/versiond/config.py +++ b/src/versiond/config.py @@ -30,6 +30,21 @@ DEFAULTS: dict[str, Any] = { "patterns": [], "extra_dotfiles": [], }, + "remote": { + "url": "", + "username": "", + "base_path": "/versioned/", + "verify_tls": True, + "timeout_seconds": 30.0, + "directory": "", + }, + "upload": { + "concurrency": 4, + "max_requests_per_second": 10.0, + "max_bandwidth_kbps": 0, + "retry_max_backoff_seconds": 300.0, + "manifest_interval_seconds": 30.0, + }, "monitor": { "poll_interval_seconds": 60.0, "reconcile_interval_seconds": 86_400.0, @@ -154,14 +169,34 @@ class Config: def port(self) -> int: return int(self.data["server"]["port"]) + def save(self) -> None: + _write_private(self.paths.config_file, dump_toml(self.data)) + + def update(self, section: str, values: dict[str, Any]) -> None: + self.data[section].update(values) + self.save() + + def _credentials(self) -> dict[str, Any]: + if self.paths.credentials_file.exists(): + return tomllib.loads(self.paths.credentials_file.read_text()) + return {} + + def credential(self, key: str) -> str: + return self._credentials().get(key, "") + + def set_credential(self, key: str, value: str) -> None: + creds = self._credentials() + creds[key] = value + _write_private(self.paths.credentials_file, dump_toml(creds)) + def api_token(self) -> str: """Return the local API token, creating it on first use.""" - creds: dict[str, Any] = {} - if self.paths.credentials_file.exists(): - creds = tomllib.loads(self.paths.credentials_file.read_text()) - token = creds.get("api_token") + token = self.credential("api_token") if not token: token = secrets.token_urlsafe(32) - creds["api_token"] = token - _write_private(self.paths.credentials_file, dump_toml(creds)) + self.set_credential("api_token", token) return token + + @property + def key_file(self): + return self.paths.config_dir / "backup.key" diff --git a/src/versiond/crypto.py b/src/versiond/crypto.py new file mode 100644 index 0000000..71f9056 --- /dev/null +++ b/src/versiond/crypto.py @@ -0,0 +1,57 @@ +"""Client-side encryption for everything sent to the remote.""" + +from __future__ import annotations + +import base64 +import hashlib +import hmac +import os +from pathlib import Path + +from cryptography.hazmat.primitives.ciphers.aead import AESGCM + +MAGIC = b"VD1" +NONCE_BYTES = 12 +AAD = b"versiond" + + +class KeyRing: + """A 64-byte master key: 32 bytes for AES-256-GCM, 32 bytes for keyed naming.""" + + def __init__(self, key: bytes): + if len(key) != 64: + raise ValueError("backup key must be 64 bytes") + self._aead = AESGCM(key[:32]) + self._mac_key = key[32:] + self.key = key + + @classmethod + def load_or_create(cls, path: Path) -> "KeyRing": + if path.exists(): + return cls(base64.b64decode(path.read_text().strip())) + key = os.urandom(64) + fd = os.open(path, os.O_WRONLY | os.O_CREAT | os.O_EXCL, 0o600) + with os.fdopen(fd, "w") as fh: + fh.write(base64.b64encode(key).decode() + "\n") + return cls(key) + + def export(self) -> str: + return base64.b64encode(self.key).decode() + + @property + def key_id(self) -> str: + return hmac.new(self._mac_key, b"versiond-key-id", hashlib.sha256).hexdigest()[:16] + + def blob_name(self, digest: str) -> str: + """Remote name for a blob; keyed so the server cannot confirm guessed content.""" + return hmac.new(self._mac_key, digest.encode(), hashlib.sha256).hexdigest() + + def encrypt(self, data: bytes) -> bytes: + nonce = os.urandom(NONCE_BYTES) + return MAGIC + nonce + self._aead.encrypt(nonce, data, AAD) + + def decrypt(self, data: bytes) -> bytes: + if not data.startswith(MAGIC): + raise ValueError("not a versiond encrypted object") + nonce = data[len(MAGIC): len(MAGIC) + NONCE_BYTES] + return self._aead.decrypt(nonce, data[len(MAGIC) + NONCE_BYTES:], AAD) diff --git a/src/versiond/dashboard.html b/src/versiond/dashboard.html new file mode 100644 index 0000000..a9ca0f6 --- /dev/null +++ b/src/versiond/dashboard.html @@ -0,0 +1,259 @@ + + + + + +versiond + + + +
+

versiond

+ … + remote … + + +
+ + + +
+
+
Files tracked
–
+
Versions
–
+
Stored locally
–
+
Backed up
–
+
Watches
–
+
+ +
+

Scans idle

+

Upload to WebDAV –

+
+ +
+

Versions per hour, last 24 h

+

Monitored directories

+
+ +
+

Most edited files, 7 days

+

Most active projects, 7 days

+
+ +
+

Versions by source

+

Monitor counters, this session

+
+
+ + + + diff --git a/src/versiond/db.py b/src/versiond/db.py index fa08041..f36264b 100644 --- a/src/versiond/db.py +++ b/src/versiond/db.py @@ -8,7 +8,18 @@ import time from pathlib import Path from typing import Any, Iterable -SCHEMA_VERSION = 1 +SCHEMA_VERSION = 2 + +# Columns added after the first release: (table, column, definition). +MIGRATIONS = [ + ("blobs", "remote_state", "TEXT NOT NULL DEFAULT 'pending'"), + ("blobs", "attempts", "INTEGER NOT NULL DEFAULT 0"), + ("blobs", "next_attempt_at", "REAL NOT NULL DEFAULT 0"), + ("blobs", "last_error", "TEXT"), + ("blobs", "uploaded_at", "REAL"), + ("versions", "manifest_batch", "TEXT"), + ("renames", "manifest_batch", "TEXT"), +] SCHEMA = """ CREATE TABLE IF NOT EXISTS roots ( @@ -85,6 +96,12 @@ CREATE INDEX IF NOT EXISTS files_root ON files(root_id); CREATE INDEX IF NOT EXISTS files_project ON files(project_id); """ +INDEXES = """ +CREATE INDEX IF NOT EXISTS blobs_remote ON blobs(remote_state, next_attempt_at); +CREATE INDEX IF NOT EXISTS versions_durability ON versions(durability); +CREATE INDEX IF NOT EXISTS versions_time ON versions(captured_at); +""" + def row_dict(row: sqlite3.Row | None) -> dict[str, Any] | None: return dict(row) if row is not None else None @@ -98,6 +115,11 @@ class Database: self.conn.execute("PRAGMA synchronous=NORMAL") self.conn.execute("PRAGMA foreign_keys=ON") self.conn.executescript(SCHEMA) + for table, column, definition in MIGRATIONS: + columns = {r[1] for r in self.conn.execute(f"PRAGMA table_info({table})")} + if column not in columns: + self.conn.execute(f"ALTER TABLE {table} ADD COLUMN {column} {definition}") + self.conn.executescript(INDEXES) current = self.conn.execute("PRAGMA user_version").fetchone()[0] if current < SCHEMA_VERSION: self.conn.execute(f"PRAGMA user_version={SCHEMA_VERSION}") diff --git a/src/versiond/filters.py b/src/versiond/filters.py index 81a1d03..df63451 100644 --- a/src/versiond/filters.py +++ b/src/versiond/filters.py @@ -70,6 +70,13 @@ class Filters: or name.endswith(IGNORED_DIR_SUFFIXES) ) + def dir_path_ignored(self, rel_posix: str) -> bool: + """True if a directory path (relative to its root) matches a user pattern.""" + return any(fnmatch.fnmatch(rel_posix, p.rstrip("/")) for p in self.patterns) + + def skip_dir(self, name: str, rel_posix: str) -> bool: + return self.dir_ignored(name) or bool(self.patterns) and self.dir_path_ignored(rel_posix) + def _dotfile_allowed(self, name: str) -> bool: return ( name in ALLOWED_DOTFILES @@ -85,6 +92,10 @@ class Filters: for part in parts[:-1]: if self.dir_ignored(part): return "ignored-directory" + if self.patterns: + for depth in range(1, len(parts)): + if self.dir_path_ignored("/".join(parts[:depth])): + return "ignored-pattern" name = parts[-1] if name.startswith("."): if not self._dotfile_allowed(name): diff --git a/src/versiond/ingest.py b/src/versiond/ingest.py index d2ea389..2322ad8 100644 --- a/src/versiond/ingest.py +++ b/src/versiond/ingest.py @@ -54,6 +54,7 @@ class Repository: self.filters = filters self._roots: list[dict[str, Any]] = [] self._project_cache: dict[str, tuple[str, str]] = {} + self.on_commit: Callable[[], None] | None = None self.reload_roots() def reload_roots(self) -> None: @@ -203,6 +204,8 @@ class Repository: (now, file["project_id"]), ) log.debug("committed %s (%s, %s)", path, source, reason) + if self.on_commit is not None: + self.on_commit() return {"status": "committed", "version_id": cur.lastrowid, "path": path} def mark_deleted(self, path: str) -> None: @@ -210,8 +213,7 @@ class Repository: def mark_deleted_tree(self, directory: str) -> None: self.db.execute( - "UPDATE files SET exists_on_disk = 0 WHERE path LIKE ? ESCAPE '\\'", - (_like_prefix(directory),), + "UPDATE files SET exists_on_disk = 0 WHERE path >= ? AND path < ?", tree_range(directory) ) def rename(self, old: str, new: str) -> bool: @@ -229,7 +231,7 @@ class Repository: def rename_tree(self, old_dir: str, new_dir: str) -> int: rows = self.db.all( - "SELECT id, path FROM files WHERE path LIKE ? ESCAPE '\\'", (_like_prefix(old_dir),) + "SELECT id, path FROM files WHERE path >= ? AND path < ?", tree_range(old_dir) ) moved = 0 with self.db.transaction(): @@ -247,9 +249,55 @@ class Repository: return moved -def _like_prefix(directory: str) -> str: - escaped = directory.rstrip("/").replace("\\", "\\\\").replace("%", "\\%").replace("_", "\\_") - return escaped + "/%" + def forget_chunk(self, path: str, limit: int = 2000) -> int: + """Delete history for path (a file or a directory tree), one chunk at a time.""" + ids = [r["id"] for r in self.db.all( + "SELECT id FROM files WHERE path = ? OR (path >= ? AND path < ?) LIMIT ?", + (path, *tree_range(path), limit), + )] + if not ids: + return 0 + marks = ",".join("?" * len(ids)) + with self.db.transaction(): + self.db.execute(f"DELETE FROM versions WHERE file_id IN ({marks})", ids) + self.db.execute(f"DELETE FROM renames WHERE file_id IN ({marks})", ids) + self.db.execute(f"DELETE FROM files WHERE id IN ({marks})", ids) + return len(ids) + + def drop_orphan_blobs(self, limit: int = 2000) -> int: + """Remove blobs no version references any more (rows and local spool files).""" + rows = self.db.all( + "SELECT sha256 FROM blobs b WHERE NOT EXISTS " + "(SELECT 1 FROM versions v WHERE v.blob_sha256 = b.sha256) LIMIT ?", + (limit,), + ) + if not rows: + return 0 + digests = [r["sha256"] for r in rows] + with self.db.transaction(): + self.db.execute(f"DELETE FROM blobs WHERE sha256 IN ({','.join('?' * len(digests))})", digests) + for digest in digests: + try: + self.blobs.path_for(digest).unlink() + except FileNotFoundError: + pass + return len(digests) + + def drop_empty_projects(self) -> int: + cur = self.db.execute( + "DELETE FROM projects WHERE NOT EXISTS (SELECT 1 FROM files f WHERE f.project_id = projects.id)" + ) + self._project_cache.clear() + return cur.rowcount + + +def tree_range(directory: str) -> tuple[str, str]: + """Bounds selecting every path strictly below directory with an index range scan. + + '0' is the character right after '/', so [dir/, dir0) covers exactly 'dir/...'. + """ + base = directory.rstrip("/") + return base + "/", base + "0" class TokenBucket: diff --git a/src/versiond/inotify.py b/src/versiond/inotify.py new file mode 100644 index 0000000..06448dc --- /dev/null +++ b/src/versiond/inotify.py @@ -0,0 +1,82 @@ +"""Minimal Linux inotify binding (ctypes, no dependencies).""" + +from __future__ import annotations + +import ctypes +import ctypes.util +import os +import struct + +ACCESS = 0x00000001 +MODIFY = 0x00000002 +ATTRIB = 0x00000004 +CLOSE_WRITE = 0x00000008 +MOVED_FROM = 0x00000040 +MOVED_TO = 0x00000080 +CREATE = 0x00000100 +DELETE = 0x00000200 +DELETE_SELF = 0x00000400 +MOVE_SELF = 0x00000800 +UNMOUNT = 0x00002000 +Q_OVERFLOW = 0x00004000 +IGNORED = 0x00008000 +ONLYDIR = 0x01000000 +DONT_FOLLOW = 0x02000000 +EXCL_UNLINK = 0x04000000 +ISDIR = 0x40000000 + +_EVENT = struct.Struct("iIII") +_READ_SIZE = 256 * 1024 + +_libc = ctypes.CDLL(ctypes.util.find_library("c") or "libc.so.6", use_errno=True) +_libc.inotify_init1.argtypes = [ctypes.c_int] +_libc.inotify_init1.restype = ctypes.c_int +_libc.inotify_add_watch.argtypes = [ctypes.c_int, ctypes.c_char_p, ctypes.c_uint32] +_libc.inotify_add_watch.restype = ctypes.c_int +_libc.inotify_rm_watch.argtypes = [ctypes.c_int, ctypes.c_int] +_libc.inotify_rm_watch.restype = ctypes.c_int + + +def _error(path: str | None = None) -> OSError: + code = ctypes.get_errno() + return OSError(code, os.strerror(code), path) + + +class Inotify: + def __init__(self) -> None: + self.fd = _libc.inotify_init1(os.O_NONBLOCK | os.O_CLOEXEC) + if self.fd < 0: + raise _error() + + def fileno(self) -> int: + return self.fd + + def add_watch(self, path: str, mask: int) -> int: + wd = _libc.inotify_add_watch(self.fd, os.fsencode(path), mask) + if wd < 0: + raise _error(path) + return wd + + def rm_watch(self, wd: int) -> None: + _libc.inotify_rm_watch(self.fd, wd) + + def read(self) -> list[tuple[int, int, int, str | None]]: + """Return all queued events as (wd, mask, cookie, name) without blocking.""" + events: list[tuple[int, int, int, str | None]] = [] + while True: + try: + data = os.read(self.fd, _READ_SIZE) + except BlockingIOError: + return events + offset = 0 + while offset < len(data): + wd, mask, cookie, length = _EVENT.unpack_from(data, offset) + offset += _EVENT.size + raw = data[offset: offset + length].rstrip(b"\0") + offset += length + events.append((wd, mask, cookie, os.fsdecode(raw) if raw else None)) + + def close(self) -> None: + if self.fd >= 0: + os.close(self.fd) + self.fd = -1 diff --git a/src/versiond/monitor.py b/src/versiond/monitor.py index a00f5e9..f7e4767 100644 --- a/src/versiond/monitor.py +++ b/src/versiond/monitor.py @@ -7,28 +7,21 @@ import errno import logging import os import time +from collections import Counter +from dataclasses import asdict, dataclass, field from pathlib import Path from typing import Any, AsyncIterator -from asyncinotify import Inotify, Mask, Watch - +from . import inotify as ino from .db import Database from .filters import Filters -from .ingest import Coalescer, Rejected, Repository, is_within +from .ingest import Coalescer, Rejected, Repository, is_within, tree_range log = logging.getLogger(__name__) DIR_MASK = ( - Mask.CLOSE_WRITE - | Mask.MOVED_TO - | Mask.MOVED_FROM - | Mask.CREATE - | Mask.DELETE - | Mask.DELETE_SELF - | Mask.MOVE_SELF - | Mask.ONLYDIR - | Mask.DONT_FOLLOW - | Mask.EXCL_UNLINK + ino.CLOSE_WRITE | ino.MOVED_TO | ino.MOVED_FROM | ino.CREATE | ino.DELETE + | ino.DELETE_SELF | ino.MOVE_SELF | ino.ONLYDIR | ino.DONT_FOLLOW | ino.EXCL_UNLINK ) NETWORK_FS = frozenset({ @@ -93,6 +86,83 @@ def needs_polling(path: str) -> bool: return fstype in NETWORK_FS or fstype.startswith("fuse") +class _Dir: + """A watched directory. Only the name is stored; paths are rebuilt from the parent chain.""" + + __slots__ = ("wd", "parent", "name", "children") + + def __init__(self, wd: int, parent: "_Dir | None", name: str): + self.wd = wd + self.parent = parent + self.name = name + self.children: dict[str, _Dir] | None = None + + def path(self) -> str: + parts = [] + node: _Dir | None = self + while node is not None: + parts.append(node.name) + node = node.parent + return "/".join(reversed(parts)) + + def top(self) -> "_Dir": + node = self + while node.parent is not None: + node = node.parent + return node + + def attach(self, parent: "_Dir", name: str) -> None: + self.parent, self.name = parent, name + if parent.children is None: + parent.children = {} + parent.children[name] = self + + def detach(self) -> None: + if self.parent is not None and self.parent.children: + self.parent.children.pop(self.name, None) + if not self.parent.children: + self.parent.children = None + self.parent = None + + def walk(self): + stack = [self] + while stack: + node = stack.pop() + yield node + if node.children: + stack.extend(node.children.values()) + + +@dataclass +class ScanProgress: + root_id: int + reason: str + top: str + started_at: float = field(default_factory=time.time) + dirs_total: int = 0 + dirs_done: int = 0 + files_seen: int = 0 + files_read: int = 0 + versions_stored: int = 0 + bytes_stored: int = 0 + finished_at: float | None = None + + def as_dict(self) -> dict[str, Any]: + data = asdict(self) + elapsed = (self.finished_at or time.time()) - self.started_at + data["elapsed_seconds"] = round(elapsed, 1) + data["percent"] = ( + round(100 * self.dirs_done / self.dirs_total, 1) if self.dirs_total else None + ) + rate = self.dirs_done / elapsed if elapsed > 0 else 0 + remaining = max(0, self.dirs_total - self.dirs_done) + data["eta_seconds"] = ( + round(remaining / rate) if rate > 0 and self.finished_at is None and self.dirs_total else None + ) + data["files_per_second"] = round(self.files_seen / elapsed, 1) if elapsed > 0 else None + return data + + class Monitor: def __init__( self, @@ -111,21 +181,30 @@ class Monitor: self.poll_interval = poll_interval self.reconcile_interval = reconcile_interval self.max_watch_fraction = max_watch_fraction - self.inotify: Inotify | None = None - self.watches: dict[str, Watch] = {} - self.watch_paths: dict[Watch, str] = {} + self.inotify: ino.Inotify | None = None + self.by_wd: dict[int, _Dir] = {} + self.tops: dict[str, _Dir] = {} + self.top_ids: dict[str, int] = {} + self.watch_counts: Counter[int] = Counter() self.polled: dict[int, set[str]] = {} self.missing: set[int] = set() - self.moves: dict[int, tuple[str, bool, asyncio.TimerHandle]] = {} + self.moves: dict[int, tuple[_Dir, str, bool, asyncio.TimerHandle]] = {} self.locks: dict[int, asyncio.Lock] = {} self.tasks: set[asyncio.Task] = set() + self.scans: dict[int, ScanProgress] = {} + self.last_scans: dict[int, ScanProgress] = {} + self.counters: Counter[str] = Counter() self._last_full_reconcile = time.monotonic() + @property + def watch_total(self) -> int: + return len(self.by_wd) + # lifecycle async def start(self) -> None: - self.inotify = Inotify() - self._spawn(self._reader()) + self.inotify = ino.Inotify() + asyncio.get_running_loop().add_reader(self.inotify.fileno(), self._on_readable) for root in self.db.roots(): await self._activate(root) self._spawn(self._periodic()) @@ -134,14 +213,15 @@ class Monitor: for task in list(self.tasks): task.cancel() await asyncio.gather(*self.tasks, return_exceptions=True) - for _, _, handle in self.moves.values(): + for *_, handle in self.moves.values(): handle.cancel() self.moves.clear() if self.inotify is not None: + asyncio.get_running_loop().remove_reader(self.inotify.fileno()) self.inotify.close() self.inotify = None - self.watches.clear() - self.watch_paths.clear() + self.by_wd.clear() + self.tops.clear() def _spawn(self, coro) -> asyncio.Task: task = asyncio.get_running_loop().create_task(coro) @@ -159,7 +239,10 @@ class Monitor: def _lock(self, root_id: int) -> asyncio.Lock: return self.locks.setdefault(root_id, asyncio.Lock()) - async def walk_dirs(self, top: str) -> AsyncIterator[str]: + def _skip(self, path: str, name: str, root_path: str) -> bool: + return self.filters.skip_dir(name, path[len(root_path) + 1:]) + + async def walk_dirs(self, top: str, root_path: str) -> AsyncIterator[str]: """Yield top and every non-ignored directory below it (symlinks not followed).""" stack = [top] seen = 0 @@ -169,7 +252,7 @@ class Monitor: try: with os.scandir(current) as it: for entry in it: - if entry.is_dir(follow_symlinks=False) and not self.filters.dir_ignored(entry.name): + if entry.is_dir(follow_symlinks=False) and not self._skip(entry.path, entry.name, root_path): stack.append(entry.path) except OSError: continue @@ -178,7 +261,7 @@ class Monitor: await asyncio.sleep(0) async def count_dirs(self, top: str) -> int: - return sum([1 async for _ in self.walk_dirs(top)]) + return sum([1 async for _ in self.walk_dirs(top, top)]) async def add_root(self, raw_path: str, force: bool = False) -> dict[str, Any]: path = os.path.realpath(os.path.expanduser(raw_path)) @@ -189,8 +272,7 @@ class Monitor: raise RootError("root-overlap", f"{path} is already inside root {root['path']}") if is_within(root["path"], path): raise RootError("root-overlap", f"{path} contains existing root {root['path']}; remove it first") - polling = needs_polling(path) - if not polling and not force: + if not needs_polling(path) and not force: needed = await self.count_dirs(path) limit = max_user_watches() available = limit - watches_in_use() @@ -200,28 +282,40 @@ class Monitor: f"{path} needs {needed} inotify watches but only {available} of {limit} are free; " "choose a smaller directory, raise fs.inotify.max_user_watches, or pass force=true", ) - cur = self.db.execute( - "INSERT INTO roots(path, added_at) VALUES (?, ?)", (path, time.time()) + cur = self.db.execute("INSERT INTO roots(path, added_at) VALUES (?, ?)", (path, time.time())) + root_id = cur.lastrowid + # Files and projects kept from an earlier root at this location are re-adopted. + lo, hi = tree_range(path) + self.db.execute("UPDATE files SET root_id = ? WHERE root_id IS NULL AND path >= ? AND path < ?", + (root_id, lo, hi)) + self.db.execute( + "UPDATE projects SET root_id = ? WHERE root_id IS NULL AND (path = ? OR (path >= ? AND path < ?))", + (root_id, path, lo, hi), ) self.db.audit("root.add", path=path) self.repo.reload_roots() - root = self.db.root(cur.lastrowid) + root = self.db.root(root_id) await self._activate(root) - return self.db.root(root["id"]) + return self.db.root(root_id) async def remove_root(self, root_id: int) -> None: root = self.db.root(root_id) if root is None: raise RootError("not-found", f"root {root_id} does not exist") - self._unwatch_tree(root["path"]) + top = self.tops.get(root["path"]) + if top is not None: + self._unwatch(top) + self.top_ids.pop(root["path"], None) self.polled.pop(root_id, None) self.missing.discard(root_id) + self.scans.pop(root_id, None) self.db.execute("DELETE FROM roots WHERE id = ?", (root_id,)) self.db.audit("root.remove", path=root["path"]) self.repo.reload_roots() async def _activate(self, root: dict[str, Any]) -> None: root_id, path = root["id"], root["path"] + self.top_ids[path] = root_id if not os.path.isdir(path): self.missing.add(root_id) self.db.update_root(root_id, mode="missing") @@ -232,188 +326,209 @@ class Monitor: if needs_polling(path): self.polled[root_id].add(path) else: - await self._watch_tree(root_id, path) + await self._watch_tree(root_id, path, None, path, path) self._refresh_root_status(root_id) reason = "reconcile" if root["baseline_state"] == "done" else "baseline" self._spawn(self.reconcile(root_id, reason)) + def root_mode(self, root_id: int, path: str) -> str: + if root_id in self.missing: + return "missing" + polled = self.polled.get(root_id, ()) + if path in polled: + return "polling" + return "degraded" if polled else "inotify" + def _refresh_root_status(self, root_id: int) -> None: root = self.db.root(root_id) if root is None: return - if root_id in self.missing: - mode = "missing" - elif root["path"] in self.polled.get(root_id, ()): - mode = "polling" - elif self.polled.get(root_id): - mode = "degraded" - else: - mode = "inotify" self.db.update_root( root_id, - mode=mode, - watch_count=sum(1 for p in self.watches if is_within(p, root["path"])), + mode=self.root_mode(root_id, root["path"]), + watch_count=self.watch_counts[root_id], polled_dirs=len(self.polled.get(root_id, ())), ) # watches - def _add_watch(self, path: str) -> bool: - if path in self.watches or self.inotify is None: - return True - try: - watch = self.inotify.add_watch(path, DIR_MASK) - except OSError as exc: - if exc.errno == errno.ENOSPC: - raise - log.debug("cannot watch %s: %s", path, exc) - return False - self.watches[path] = watch - self.watch_paths[watch] = path - return True - - async def _watch_tree(self, root_id: int, top: str) -> None: - skip: list[str] = [] - async for directory in self.walk_dirs(top): - if any(is_within(directory, s) for s in skip): - continue + async def _watch_tree(self, root_id: int, top_path: str, parent: _Dir | None, name: str, + root_path: str) -> None: + """Add watches for a directory and every non-ignored directory below it.""" + assert self.inotify is not None + stack: list[tuple[str, _Dir | None, str]] = [(top_path, parent, name)] + added = 0 + while stack: + path, parent_node, node_name = stack.pop() try: - self._add_watch(directory) + wd = self.inotify.add_watch(path, DIR_MASK) + except OSError as exc: + if exc.errno == errno.ENOSPC: + log.warning("inotify watch limit reached; polling %s instead", path) + self.polled.setdefault(root_id, set()).add(path) + self.counters["watch_limit_hits"] += 1 + continue + if wd in self.by_wd: + continue # same directory reached twice (bind mount); keep the first + node = _Dir(wd, None, node_name) + if parent_node is None: + self.tops[path] = node + else: + node.attach(parent_node, node_name) + self.by_wd[wd] = node + self.watch_counts[root_id] += 1 + try: + with os.scandir(path) as it: + for entry in it: + if entry.is_dir(follow_symlinks=False) and not self._skip(entry.path, entry.name, root_path): + stack.append((entry.path, node, entry.name)) except OSError: - log.warning("inotify watch limit reached; polling %s instead", directory) - self.polled.setdefault(root_id, set()).add(directory) - skip.append(directory) + pass + added += 1 + if added % YIELD_EVERY == 0: + await asyncio.sleep(0) - def _unwatch_tree(self, top: str) -> None: - for path in [p for p in self.watches if is_within(p, top)]: - watch = self.watches.pop(path) - self.watch_paths.pop(watch, None) - if self.inotify is not None: - try: - self.inotify.rm_watch(watch) - except (OSError, ValueError, KeyError): - pass + def _unwatch(self, node: _Dir) -> None: + top = node.top() + root_id = self.top_ids.get(top.name) + removed = 0 + for child in list(node.walk()): + if self.by_wd.pop(child.wd, None) is not None: + removed += 1 + if self.inotify is not None: + self.inotify.rm_watch(child.wd) + if node.parent is None: + self.tops.pop(node.name, None) + node.detach() + if root_id is not None: + self.watch_counts[root_id] -= removed - def _move_watches(self, old: str, new: str) -> None: - for path in [p for p in self.watches if is_within(p, old)]: - watch = self.watches.pop(path) - moved = new + path[len(old):] - self.watches[moved] = watch - self.watch_paths[watch] = moved + def _root_id_of(self, node: _Dir) -> int | None: + return self.top_ids.get(node.top().name) # events - async def _reader(self) -> None: - assert self.inotify is not None - async for event in self.inotify: + def _on_readable(self) -> None: + if self.inotify is None: + return + for wd, mask, cookie, name in self.inotify.read(): + self.counters["events"] += 1 try: - self._handle(event) + self._handle(wd, mask, cookie, name) except Exception: log.exception("failed handling inotify event") - def _root_id_for(self, path: str) -> int | None: - root = self.repo.root_for(path) - return root["id"] if root else None - - def _handle(self, event) -> None: - mask = event.mask - if Mask.Q_OVERFLOW in mask: + def _handle(self, wd: int, mask: int, cookie: int, name: str | None) -> None: + if mask & ino.Q_OVERFLOW: + self.counters["queue_overflows"] += 1 log.warning("inotify queue overflow; reconciling all roots") for root in self.db.roots(): self._spawn(self.reconcile(root["id"], "overflow")) return - watch = event.watch - if Mask.IGNORED in mask: - path = self.watch_paths.pop(watch, None) if watch is not None else None - if path is not None and self.watches.get(path) is watch: - del self.watches[path] + node = self.by_wd.get(wd) + if node is None: return - directory = self.watch_paths.get(watch) if watch is not None else None - if directory is None: + if mask & ino.IGNORED: + self.by_wd.pop(wd, None) + root_id = self._root_id_of(node) + if root_id is not None: + self.watch_counts[root_id] -= 1 + node.detach() + return + if name is None: + if mask & (ino.DELETE_SELF | ino.MOVE_SELF) and node.parent is None: + root_id = self.top_ids.get(node.name) + log.warning("root %s was removed or moved", node.name) + self._unwatch(node) + if root_id is not None: + self.missing.add(root_id) + self._refresh_root_status(root_id) return - if event.name is None: - if mask & (Mask.DELETE_SELF | Mask.MOVE_SELF): - root = self.repo.root_for(directory) - if root and root["path"] == directory: - log.warning("root %s was removed or moved", directory) - self._unwatch_tree(directory) - self.missing.add(root["id"]) - self._refresh_root_status(root["id"]) + is_dir = bool(mask & ino.ISDIR) + if mask & ino.MOVED_FROM: + handle = asyncio.get_running_loop().call_later(MOVE_PAIR_TIMEOUT, self._moved_out, cookie) + self.moves[cookie] = (node, name, is_dir, handle) return - - path = os.path.join(directory, str(event.name)) - is_dir = Mask.ISDIR in mask - - if Mask.MOVED_FROM in mask: - handle = asyncio.get_running_loop().call_later( - MOVE_PAIR_TIMEOUT, self._moved_out, event.cookie - ) - self.moves[event.cookie] = (path, is_dir, handle) - return - - if Mask.MOVED_TO in mask: - origin = self.moves.pop(event.cookie, None) + if mask & ino.MOVED_TO: + origin = self.moves.pop(cookie, None) if origin is not None: - origin[2].cancel() - self._moved_in(origin[0] if origin else None, path, is_dir) + origin[3].cancel() + self._moved_in(origin, node, name, is_dir) return + path = node.path() + "/" + name if is_dir: - if Mask.CREATE in mask: - self._new_directory(path) - elif Mask.DELETE in mask: - self._unwatch_tree(path) + if mask & ino.CREATE: + self._new_directory(node, name) + elif mask & ino.DELETE: + child = node.children.get(name) if node.children else None + if child is not None: + self._unwatch(child) self.repo.mark_deleted_tree(path) return - - if Mask.CLOSE_WRITE in mask: + if mask & ino.CLOSE_WRITE: self.capture(path, "monitor", "modified") - elif Mask.DELETE in mask: + elif mask & ino.DELETE: + self.counters["deletes"] += 1 self.repo.mark_deleted(path) - def _new_directory(self, path: str) -> None: - root_id = self._root_id_for(path) - if root_id is None or any(self.filters.dir_ignored(p) for p in self.repo.relative(path).parts): + def _new_directory(self, parent: _Dir, name: str) -> None: + root_path = parent.top().name + if self._skip(parent.path() + "/" + name, name, root_path): return - self._spawn(self._watch_and_scan(root_id, path)) + root_id = self._root_id_of(parent) + if root_id is None: + return + self._spawn(self._watch_and_scan(root_id, parent, name)) - async def _watch_and_scan(self, root_id: int, path: str) -> None: + async def _watch_and_scan(self, root_id: int, parent: _Dir, name: str) -> None: # Watch first, then scan, so files created before the watch existed are not missed. - await self._watch_tree(root_id, path) - self._refresh_root_status(root_id) + path = parent.path() + "/" + name + await self._watch_tree(root_id, path, parent, name, parent.top().name) await self.reconcile(root_id, "new-directory", subtree=path) def _moved_out(self, cookie: int) -> None: origin = self.moves.pop(cookie, None) if origin is None: return - path, is_dir, _ = origin + parent, name, is_dir, _ = origin + path = parent.path() + "/" + name if is_dir: - self._unwatch_tree(path) + child = parent.children.get(name) if parent.children else None + if child is not None: + self._unwatch(child) self.repo.mark_deleted_tree(path) else: self.repo.mark_deleted(path) - def _moved_in(self, old: str | None, new: str, is_dir: bool) -> None: - new_ok = self._root_id_for(new) is not None and not any( - self.filters.dir_ignored(p) for p in self.repo.relative(new).parts[: None if is_dir else -1] - ) + def _moved_in(self, origin, parent: _Dir, name: str, is_dir: bool) -> None: + new_path = parent.path() + "/" + name + old_path = origin[0].path() + "/" + origin[1] if origin else None if is_dir: - if old is not None and new_ok and old in self.watches: - self._move_watches(old, new) - self.repo.rename_tree(old, new) + new_ok = not self._skip(new_path, name, parent.top().name) + child = None + if origin is not None and origin[0].children: + child = origin[0].children.get(origin[1]) + if child is not None and new_ok and self._root_id_of(origin[0]) == self._root_id_of(parent): + child.detach() + child.attach(parent, name) + self.repo.rename_tree(old_path, new_path) + self.counters["renames"] += 1 return - if old is not None: - self._unwatch_tree(old) - self.repo.mark_deleted_tree(old) + if child is not None: + self._unwatch(child) + if old_path is not None: + self.repo.mark_deleted_tree(old_path) if new_ok: - self._new_directory(new) + self._new_directory(parent, name) return - if old is not None and not self.repo.rename(old, new): - self.repo.mark_deleted(old) - if new_ok: - self.capture(new, "monitor", "renamed" if old else "moved-in") + if old_path is not None: + if self.repo.rename(old_path, new_path): + self.counters["renames"] += 1 + else: + self.repo.mark_deleted(old_path) + self.capture(new_path, "monitor", "renamed" if old_path else "moved-in") def capture(self, path: str, source: str, reason: str) -> dict[str, Any] | None: try: @@ -421,14 +536,13 @@ class Monitor: content, st = self.repo.read_file(path) self.repo.check_content(content) except Rejected as exc: - log.debug("skip %s: %s", path, exc.reason) + self.counters[f"skipped:{exc.reason}"] += 1 return None - except FileNotFoundError: + except OSError: return None - except OSError as exc: - log.debug("cannot read %s: %s", path, exc) - return None - return self.coalescer.submit(path, content, source, reason, st) + result = self.coalescer.submit(path, content, source, reason, st) + self.counters[f"capture:{result['status']}"] += 1 + return result # scanning @@ -438,12 +552,15 @@ class Monitor: if root is None or root_id in self.missing: return 0 top = subtree or root["path"] - committed = 0 - seen: set[str] = set() async with self._lock(root_id): + progress = ScanProgress(root_id, reason, top) + if subtree is None: + progress.dirs_total = self.watch_counts[root_id] + len(self.polled.get(root_id, ())) + self.scans[root_id] = progress if reason == "baseline": self.db.update_root(root_id, baseline_state="running") - async for directory in self.walk_dirs(top): + async for directory in self.walk_dirs(top, root["path"]): + progress.dirs_done += 1 try: entries = list(os.scandir(directory)) except OSError: @@ -460,7 +577,7 @@ class Monitor: continue if self.filters.size_reason(st.st_size): continue - seen.add(path) + progress.files_seen += 1 if path in self.coalescer.pending: continue row = self.db.file_by_path(path) @@ -477,24 +594,37 @@ class Monitor: self.repo.check_content(content) except (Rejected, OSError): continue + progress.files_read += 1 result = self.repo.commit(path, content, "scan", reason, fst) if result["status"] == "committed": - committed += 1 - await asyncio.sleep(0) + progress.versions_stored += 1 + progress.bytes_stored += len(content) + if progress.dirs_done % 20 == 0: + await asyncio.sleep(0) + if progress.dirs_total < progress.dirs_done: + progress.dirs_total = progress.dirs_done + # Anything indexed as existing but no longer a file on disk is marked deleted. rows = self.db.all( - "SELECT path FROM files WHERE root_id = ? AND exists_on_disk = 1", (root_id,) + "SELECT path FROM files WHERE root_id = ? AND exists_on_disk = 1 " + "AND (path = ? OR (path >= ? AND path < ?))", + (root_id, top, *tree_range(top)), ) - for row in rows: - if is_within(row["path"], top) and row["path"] not in seen: - if not os.path.isfile(row["path"]): - self.repo.mark_deleted(row["path"]) + for i, row in enumerate(rows): + if not os.path.isfile(row["path"]): + self.repo.mark_deleted(row["path"]) + if i % 1000 == 0: + await asyncio.sleep(0) fields: dict[str, Any] = {"last_scan_at": time.time()} if reason == "baseline": - fields.update(baseline_state="done", baseline_files=committed) + fields.update(baseline_state="done", baseline_files=progress.versions_stored) self.db.update_root(root_id, **fields) - if committed: - log.info("%s scan of %s stored %d versions", reason, top, committed) - return committed + progress.finished_at = time.time() + if subtree is None: + self.scans.pop(root_id, None) + self.last_scans[root_id] = progress + if progress.versions_stored: + log.info("%s scan of %s stored %d versions", reason, top, progress.versions_stored) + return progress.versions_stored async def _periodic(self) -> None: while True: @@ -509,6 +639,9 @@ class Monitor: for root_id, dirs in list(self.polled.items()): for directory in list(dirs): await self.reconcile(root_id, "poll", subtree=directory) + for root in self.db.roots(): + if root["id"] not in self.missing: + self._refresh_root_status(root["id"]) if time.monotonic() - self._last_full_reconcile >= self.reconcile_interval: self._last_full_reconcile = time.monotonic() for root in self.db.roots(): diff --git a/src/versiond/remote.py b/src/versiond/remote.py new file mode 100644 index 0000000..69d9bb4 --- /dev/null +++ b/src/versiond/remote.py @@ -0,0 +1,201 @@ +"""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 +from dataclasses import dataclass +from typing import Any + +import httpx + +INSTALL_APP_ID = b"versiond:4c1b9e0f7a2d4e8b" +FORMAT_VERSION = 1 + + +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) + msg = INSTALL_APP_ID + b":" + str(os.getuid()).encode() + return hmac.new(bytes.fromhex(machine_id), 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, key_id: str) -> 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. + """ + 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(), + "key_id": key_id, + } + 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": "aes-256-gcm", "compression": "zstd", + "blob_layout": "blobs//", "key_id": key_id} + 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: + if owner.get("key_id") not in (None, key_id): + raise RemoteError( + "key-mismatch", + f"{directory} was written with a different backup key; restore that key " + "(~/.config/versiond/backup.key) or remove the directory", + ) + 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") diff --git a/src/versiond/stats.py b/src/versiond/stats.py new file mode 100644 index 0000000..19b3b9e --- /dev/null +++ b/src/versiond/stats.py @@ -0,0 +1,77 @@ +"""Aggregate statistics over the index.""" + +from __future__ import annotations + +import time +from typing import Any + +from .db import Database + + +def collect(db: Database) -> dict[str, Any]: + now = time.time() + files = db.one( + "SELECT count(*) AS tracked, coalesce(sum(exists_on_disk), 0) AS existing FROM files" + ) + versions = db.one( + "SELECT count(*) AS n, coalesce(sum(b.size), 0) AS logical_bytes " + "FROM versions v JOIN blobs b ON b.sha256 = v.blob_sha256" + ) + blobs = db.one( + "SELECT count(*) AS n, coalesce(sum(size), 0) AS raw_bytes, " + "coalesce(sum(stored_size), 0) AS stored_bytes FROM blobs" + ) + windows = {} + for label, seconds in (("last_hour", 3600), ("last_24h", 86_400), ("last_7d", 604_800)): + windows[label] = db.one( + "SELECT count(*) AS n FROM versions WHERE captured_at >= ?", (now - seconds,) + )["n"] + hour_start = now - (now % 3600) + hourly = [0] * 24 + for row in db.all( + "SELECT CAST((? - captured_at) / 3600 AS INTEGER) AS ago, count(*) AS n " + "FROM versions WHERE captured_at >= ? GROUP BY ago", + (hour_start + 3600, hour_start - 23 * 3600), + ): + if 0 <= row["ago"] < 24: + hourly[23 - row["ago"]] = row["n"] + week_ago = now - 604_800 + return { + "generated_at": now, + "files": { + "tracked": files["tracked"], + "existing": files["existing"], + "deleted": files["tracked"] - files["existing"], + }, + "versions": { + "total": versions["n"], + **windows, + "logical_bytes": versions["logical_bytes"], + "by_source": {r["source"]: r["n"] for r in db.all( + "SELECT source, count(*) AS n FROM versions GROUP BY source ORDER BY n DESC")}, + "by_reason": {r["reason"]: r["n"] for r in db.all( + "SELECT reason, count(*) AS n FROM versions GROUP BY reason ORDER BY n DESC LIMIT 20")}, + "per_hour_last_24h": hourly, + }, + "storage": { + "unique_blobs": blobs["n"], + "raw_bytes": blobs["raw_bytes"], + "stored_bytes": blobs["stored_bytes"], + "dedup_ratio": round(versions["logical_bytes"] / blobs["raw_bytes"], 2) if blobs["raw_bytes"] else None, + "compression_ratio": round(blobs["raw_bytes"] / blobs["stored_bytes"], 2) if blobs["stored_bytes"] else None, + }, + "top_files_7d": db.all( + "SELECT f.path, count(*) AS versions FROM versions v JOIN files f ON f.id = v.file_id " + "WHERE v.captured_at >= ? AND v.reason NOT IN ('baseline', 'reconcile') " + "GROUP BY v.file_id ORDER BY versions DESC LIMIT 10", + (week_ago,), + ), + "top_projects_7d": db.all( + "SELECT p.path, count(*) AS versions FROM versions v JOIN files f ON f.id = v.file_id " + "JOIN projects p ON p.id = f.project_id " + "WHERE v.captured_at >= ? AND v.reason NOT IN ('baseline', 'reconcile') " + "GROUP BY p.id ORDER BY versions DESC LIMIT 10", + (week_ago,), + ), + "projects": db.one("SELECT count(*) AS n FROM projects")["n"], + } diff --git a/src/versiond/uploader.py b/src/versiond/uploader.py new file mode 100644 index 0000000..b776450 --- /dev/null +++ b/src/versiond/uploader.py @@ -0,0 +1,349 @@ +"""Background upload of blobs and manifests to the WebDAV remote.""" + +from __future__ import annotations + +import asyncio +import json +import logging +import random +import secrets +import time +from collections import deque +from typing import Any + +import httpx + +from .config import Config +from .crypto import KeyRing +from .db import Database +from .remote import RemoteError, RemoteSettings, WebDAV, check_connection, claim_directory, join +from .store import BlobStore, _codec + +log = logging.getLogger(__name__) + +BATCH_SIZE = 256 +MANIFEST_MAX_RECORDS = 5000 +IDLE_SLEEP = 2.0 +OFFLINE_SLEEP = 30.0 +RATE_WINDOW = 60.0 + + +class Throttle: + """Async token bucket for requests per second and bytes per second.""" + + def __init__(self, per_second: float): + self.rate = per_second + self.tokens = per_second + self.stamp = time.monotonic() + + async def take(self, amount: float = 1.0) -> None: + if self.rate <= 0: + return + now = time.monotonic() + self.tokens = min(self.rate, self.tokens + (now - self.stamp) * self.rate) + self.stamp = now + self.tokens -= amount + if self.tokens < 0: + await asyncio.sleep(-self.tokens / self.rate) + + +class Uploader: + def __init__(self, cfg: Config, db: Database, blobs: BlobStore, keys: KeyRing, + transport: httpx.AsyncBaseTransport | None = None): + self.cfg = cfg + self.db = db + self.blobs = blobs + self.keys = keys + self.transport = transport + self.dav: WebDAV | None = None + self.task: asyncio.Task | None = None + self.state = "unconfigured" + self.directory = "" + self.counters: dict[str, Any] = { + "blobs_uploaded": 0, "bytes_uploaded": 0, "upload_failures": 0, + "manifests_uploaded": 0, "versions_made_durable": 0, + } + self.last_error: str | None = None + self.last_error_at: float | None = None + self.last_success_at: float | None = None + self.consecutive_failures = 0 + self.recent: deque[tuple[float, int]] = deque() + self._last_manifest = 0.0 + self._wake = asyncio.Event() + self.offline_sleep = OFFLINE_SLEEP + + # settings + + def settings(self) -> RemoteSettings: + r = self.cfg.data["remote"] + return RemoteSettings( + url=r["url"], username=r["username"], password=self.cfg.credential("remote_password"), + base_path=r["base_path"], verify_tls=bool(r["verify_tls"]), + timeout_seconds=float(r["timeout_seconds"]), directory=r["directory"], + ) + + def _upload_cfg(self, key: str) -> float: + return float(self.cfg.data["upload"][key]) + + def public_settings(self) -> dict[str, Any]: + s = self.settings() + return { + "configured": s.configured, "url": s.url, "username": s.username, + "password_set": bool(s.password), "base_path": s.base_path, "directory": s.directory, + "verify_tls": s.verify_tls, "timeout_seconds": s.timeout_seconds, + "encryption": "aes-256-gcm", "key_id": self.keys.key_id, "state": self.state, + } + + async def test(self, settings: RemoteSettings) -> None: + dav = WebDAV(settings, transport=self.transport) + try: + await check_connection(dav) + finally: + await dav.close() + + async def configure(self, settings: RemoteSettings) -> dict[str, Any]: + """Validate, claim a unique directory, persist, and (re)start uploading.""" + dav = WebDAV(settings, transport=self.transport) + try: + await check_connection(dav) + claim = await claim_directory(dav, self.keys.key_id) + finally: + await dav.close() + await self.stop() + self.cfg.update("remote", { + "url": settings.url, "username": settings.username, "base_path": settings.base_path, + "verify_tls": settings.verify_tls, "timeout_seconds": settings.timeout_seconds, + "directory": claim["directory"], + }) + self.cfg.set_credential("remote_password", settings.password) + self.db.audit("remote.configure", url=settings.url, directory=claim["directory"], + adopted=claim["adopted"]) + await self.start() + return {**self.public_settings(), "adopted": claim["adopted"]} + + # lifecycle + + async def start(self) -> None: + settings = self.settings() + if not settings.configured: + self.state = "unconfigured" + return + concurrency = int(self._upload_cfg("concurrency")) + self.dav = WebDAV(settings, concurrency=concurrency, transport=self.transport) + self.directory = settings.directory + self.state = "connecting" + self.task = asyncio.get_running_loop().create_task(self._run()) + + async def stop(self) -> None: + if self.task is not None: + self.task.cancel() + await asyncio.gather(self.task, return_exceptions=True) + self.task = None + if self.dav is not None: + await self.dav.close() + self.dav = None + + def wake(self) -> None: + self._wake.set() + + # main loop + + def _failed(self, exc: Exception) -> None: + self.last_error = str(exc) + self.last_error_at = time.time() + self.consecutive_failures += 1 + + async def _ensure_directory(self) -> None: + assert self.dav is not None + if not self.directory: + claim = await claim_directory(self.dav, self.keys.key_id) + self.directory = claim["directory"] + self.cfg.update("remote", {"directory": self.directory}) + await self.dav.makedirs(join(self.directory, "blobs")) + await self.dav.makedirs(join(self.directory, "manifests")) + + async def _run(self) -> None: + while True: + try: + await self._ensure_directory() + self.state = "online" + self.consecutive_failures = 0 + while True: + did_work = await self._upload_batch() + if time.monotonic() - self._last_manifest >= self._upload_cfg("manifest_interval_seconds") or not did_work: + did_work = await self._write_manifest() or did_work + if not did_work: + self._wake.clear() + try: + await asyncio.wait_for(self._wake.wait(), IDLE_SLEEP) + except asyncio.TimeoutError: + pass + except asyncio.CancelledError: + raise + except RemoteError as exc: + self._failed(exc) + self.state = "error" if exc.code in ("remote-auth", "key-mismatch", "remote-full") else "offline" + log.warning("remote %s: %s; retrying in %ss", self.state, exc.detail, self.offline_sleep) + await asyncio.sleep(self.offline_sleep) + except Exception as exc: + self._failed(exc) + self.state = "offline" + log.exception("uploader failed; retrying in %ss", self.offline_sleep) + await asyncio.sleep(self.offline_sleep) + + async def _upload_batch(self) -> bool: + rows = self.db.all( + "SELECT sha256, stored_size, attempts FROM blobs WHERE remote_state IN ('pending', 'failed') " + "AND next_attempt_at <= ? ORDER BY created_at LIMIT ?", + (time.time(), BATCH_SIZE), + ) + if not rows: + return False + concurrency = max(1, int(self._upload_cfg("concurrency"))) + requests = Throttle(self._upload_cfg("max_requests_per_second")) + bandwidth = Throttle(self._upload_cfg("max_bandwidth_kbps") * 1024) + queue: asyncio.Queue = asyncio.Queue() + for row in rows: + queue.put_nowait(row) + fatal: list[RemoteError] = [] + + async def worker() -> None: + while not queue.empty() and not fatal: + row = queue.get_nowait() + await requests.take() + try: + await self._upload_blob(row, bandwidth) + except RemoteError as exc: + if exc.code in ("remote-unreachable", "remote-auth", "remote-full"): + fatal.append(exc) + + await asyncio.gather(*(worker() for _ in range(concurrency))) + if fatal: + raise fatal[0] + return True + + async def _upload_blob(self, row: dict[str, Any], bandwidth: Throttle) -> None: + assert self.dav is not None + digest = row["sha256"] + try: + stored = self.blobs.path_for(digest).read_bytes() + except FileNotFoundError: + self.db.execute("UPDATE blobs SET remote_state = 'missing-local', last_error = ? WHERE sha256 = ?", + ("local blob file is missing", digest)) + return + name = self.keys.blob_name(digest) + prefix = join(self.directory, "blobs", name[:2]) + payload = self.keys.encrypt(stored) + await bandwidth.take(len(payload)) + try: + await self.dav.mkcol(prefix) + await self.dav.put(join(prefix, name), payload) + except RemoteError as exc: + attempts = row["attempts"] + 1 + backoff = min(self._upload_cfg("retry_max_backoff_seconds"), 2 ** min(attempts, 16)) + backoff *= random.uniform(0.8, 1.2) + self.db.execute( + "UPDATE blobs SET remote_state = 'failed', attempts = ?, next_attempt_at = ?, " + "last_error = ? WHERE sha256 = ?", + (attempts, time.time() + backoff, exc.detail, digest), + ) + self.counters["upload_failures"] += 1 + self._failed(exc) + raise + now = time.time() + self.db.execute( + "UPDATE blobs SET remote_state = 'uploaded', uploaded_at = ?, attempts = attempts + 1, " + "last_error = NULL WHERE sha256 = ?", + (now, digest), + ) + self.counters["blobs_uploaded"] += 1 + self.counters["bytes_uploaded"] += len(payload) + self.recent.append((time.monotonic(), len(payload))) + self.last_success_at = now + self.consecutive_failures = 0 + + async def _write_manifest(self) -> bool: + assert self.dav is not None + self._last_manifest = time.monotonic() + versions = self.db.all( + "SELECT v.id, f.path, v.blob_sha256 AS sha256, b.size, v.captured_at, v.source, v.reason, v.pinned " + "FROM versions v JOIN files f ON f.id = v.file_id JOIN blobs b ON b.sha256 = v.blob_sha256 " + "WHERE v.durability = 'local' AND b.remote_state = 'uploaded' ORDER BY v.id LIMIT ?", + (MANIFEST_MAX_RECORDS,), + ) + renames = self.db.all( + "SELECT id, file_id, old_path, new_path, at FROM renames WHERE manifest_batch IS NULL " + "ORDER BY id LIMIT ?", + (MANIFEST_MAX_RECORDS,), + ) + if not versions and not renames: + return False + stamp = time.gmtime() + batch_id = time.strftime("%Y%m%dT%H%M%SZ", stamp) + "-" + secrets.token_hex(4) + lines = [json.dumps({"type": "version", **v, "blob": self.keys.blob_name(v["sha256"])}) for v in versions] + lines += [json.dumps({"type": "rename", **r}) for r in renames] + payload = self.keys.encrypt(_codec.compress(("\n".join(lines) + "\n").encode())) + folder = join(self.directory, "manifests", time.strftime("%Y/%m/%d", stamp)) + await self.dav.makedirs(folder) + await self.dav.put(join(folder, f"{batch_id}.jsonl.zst.enc"), payload) + with self.db.transaction(): + for chunk in _chunks([v["id"] for v in versions], 500): + self.db.execute( + f"UPDATE versions SET durability = 'durable', manifest_batch = ? " + f"WHERE id IN ({','.join('?' * len(chunk))})", + (batch_id, *chunk), + ) + for chunk in _chunks([r["id"] for r in renames], 500): + self.db.execute( + f"UPDATE renames SET manifest_batch = ? WHERE id IN ({','.join('?' * len(chunk))})", + (batch_id, *chunk), + ) + self.counters["manifests_uploaded"] += 1 + self.counters["versions_made_durable"] += len(versions) + log.info("manifest %s: %d versions, %d renames durable", batch_id, len(versions), len(renames)) + return True + + # progress + + def progress(self) -> dict[str, Any]: + now = time.monotonic() + while self.recent and now - self.recent[0][0] > RATE_WINDOW: + self.recent.popleft() + window = min(RATE_WINDOW, max(1.0, now - self.recent[0][0])) if self.recent else RATE_WINDOW + blobs_per_sec = len(self.recent) / window + bytes_per_sec = sum(b for _, b in self.recent) / window + totals = {r["remote_state"]: r for r in self.db.all( + "SELECT remote_state, count(*) AS n, coalesce(sum(stored_size), 0) AS bytes FROM blobs GROUP BY remote_state" + )} + pending_n = sum(t["n"] for s, t in totals.items() if s in ("pending", "failed")) + pending_bytes = sum(t["bytes"] for s, t in totals.items() if s in ("pending", "failed")) + uploaded = totals.get("uploaded", {"n": 0, "bytes": 0}) + all_n = sum(t["n"] for t in totals.values()) + durability = {r["durability"]: r["n"] for r in self.db.all( + "SELECT durability, count(*) AS n FROM versions GROUP BY durability" + )} + return { + "state": self.state, + "directory": self.directory or None, + "blobs_total": all_n, + "blobs_uploaded": uploaded["n"], + "blobs_pending": pending_n, + "blobs_failed": totals.get("failed", {"n": 0})["n"], + "bytes_pending": pending_bytes, + "percent": round(100 * uploaded["n"] / all_n, 1) if all_n else 100.0, + "blobs_per_second": round(blobs_per_sec, 2), + "bytes_per_second": round(bytes_per_sec), + "eta_seconds": round(pending_n / blobs_per_sec) if blobs_per_sec > 0 and pending_n else None, + "versions_durable": durability.get("durable", 0), + "versions_local": durability.get("local", 0), + "last_success_at": self.last_success_at, + "last_error": self.last_error, + "last_error_at": self.last_error_at, + "consecutive_failures": self.consecutive_failures, + "session": dict(self.counters), + } + + +def _chunks(items: list, size: int): + for i in range(0, len(items), size): + yield items[i: i + size] diff --git a/tests/test_remote.py b/tests/test_remote.py new file mode 100644 index 0000000..268922a --- /dev/null +++ b/tests/test_remote.py @@ -0,0 +1,162 @@ +"""WebDAV backup against an in-memory WebDAV server.""" + +import base64 +import json +import time + +import httpx +import pytest +from fastapi.testclient import TestClient + +from versiond.api import create_app +from versiond.config import Config, Paths +from versiond.crypto import KeyRing +from versiond.store import _codec + +from test_service import CONFIG, history, wait_for + + +class FakeDAV: + """Just enough RFC 4918 for the client: PROPFIND, MKCOL, PUT, GET, DELETE.""" + + def __init__(self, user="u1", password="secret"): + self.auth = "Basic " + base64.b64encode(f"{user}:{password}".encode()).decode() + self.dirs = {"/"} + self.files: dict[str, bytes] = {} + self.down = False + self.requests = 0 + + @staticmethod + def parent(path): + return path.rsplit("/", 1)[0] or "/" + + def handler(self, request: httpx.Request) -> httpx.Response: + self.requests += 1 + if self.down: + raise httpx.ConnectError("connection refused", request=request) + if request.headers.get("authorization") != self.auth: + return httpx.Response(401) + path = request.url.path.rstrip("/") or "/" + method = request.method + if method == "PROPFIND": + return httpx.Response(207 if path in self.dirs or path in self.files else 404) + if method == "MKCOL": + if path in self.dirs: + return httpx.Response(405) + if self.parent(path) not in self.dirs: + return httpx.Response(409) + self.dirs.add(path) + return httpx.Response(201) + if method == "PUT": + if self.parent(path) not in self.dirs: + return httpx.Response(409) + if request.headers.get("if-none-match") == "*" and path in self.files: + return httpx.Response(412) + self.files[path] = request.read() + return httpx.Response(201) + if method == "GET": + return httpx.Response(200, content=self.files[path]) if path in self.files else httpx.Response(404) + if method == "DELETE": + self.files.pop(path, None) + return httpx.Response(204) + return httpx.Response(405) + + +@pytest.fixture +def remote_env(tmp_path, monkeypatch): + monkeypatch.setenv("VERSIOND_HOME", str(tmp_path / "home")) + paths = Paths.resolve() + paths.ensure() + paths.config_file.write_text(CONFIG + "\n[upload]\nmanifest_interval_seconds = 0.2\n") + cfg = Config.load(paths) + dav = FakeDAV() + project = tmp_path / "proj" + project.mkdir() + (project / "main.py").write_text("print('hello')\n") + (project / ".env").write_text("SECRET=1\n") + app = create_app(cfg, remote_transport=httpx.MockTransport(dav.handler)) + with TestClient(app, base_url="http://127.0.0.1:9922") as client: + client.headers["Authorization"] = f"Bearer {cfg.api_token()}" + app.state.services.uploader.offline_sleep = 0.3 + yield client, dav, project, cfg + + +REMOTE = {"url": "https://dav.example", "username": "u1", "password": "secret", "base_path": "/versioned/"} + + +def progress(client): + return client.get("/api/v1/progress").json() + + +def test_backup_to_webdav(remote_env): + client, dav, project, cfg = remote_env + + bad = client.put("/api/v1/config/remote", json={**REMOTE, "password": "wrong"}) + assert bad.status_code == 400 and "remote-auth" in bad.text + + r = client.put("/api/v1/config/remote", json=REMOTE) + assert r.status_code == 200, r.text + remote = r.json() + directory = remote["directory"] + assert directory.startswith("/versioned/") and remote["adopted"] is False + assert "password" not in remote and remote["password_set"] is True + owner = json.loads(dav.files[directory + "/meta/owner.json"]) + assert owner["key_id"] == remote["key_id"] + assert "secret" not in cfg.paths.config_file.read_text() # password only in credentials + + client.post("/api/v1/roots", json={"path": str(project)}) + wait_for(lambda: progress(client)["upload"]["versions_local"] == 0 + and progress(client)["upload"]["versions_durable"] == 2, timeout=10) + + keys = KeyRing.load_or_create(cfg.key_file) + blob_paths = [p for p in dav.files if p.startswith(directory + "/blobs/")] + contents = {_codec.decompress(keys.decrypt(dav.files[p])) for p in blob_paths} + assert contents == {b"print('hello')\n", b"SECRET=1\n"} + assert all(b"SECRET" not in dav.files[p] for p in blob_paths) # encrypted at rest + + manifests = [p for p in dav.files if p.startswith(directory + "/manifests/")] + records = [json.loads(line) for p in manifests + for line in _codec.decompress(keys.decrypt(dav.files[p])).decode().splitlines()] + assert {r["path"] for r in records if r["type"] == "version"} == {str(project / "main.py"), str(project / ".env")} + + # reconfiguring the same machine adopts its own directory + again = client.put("/api/v1/config/remote", json={**REMOTE, "password": None}) + assert again.json()["directory"] == directory and again.json()["adopted"] is True + + # offline: changes wait locally, then upload when the server is back + dav.down = True + (project / "main.py").write_text("print('offline edit')\n") + wait_for(lambda: progress(client)["upload"]["state"] == "offline", timeout=10) + assert progress(client)["upload"]["versions_local"] == 1 + dav.down = False + wait_for(lambda: progress(client)["upload"]["versions_local"] == 0, timeout=15) + assert progress(client)["upload"]["state"] == "online" + + +def test_stats_progress_dashboard(remote_env): + client, _, project, _ = remote_env + root = client.post("/api/v1/roots", json={"path": str(project)}).json() + wait_for(lambda: client.get(f"/api/v1/roots/{root['id']}").json()["baseline_state"] == "done") + (project / "main.py").write_text("v2\n") + wait_for(lambda: len(history(client, project / "main.py")) == 2) + + stats = client.get("/api/v1/stats").json() + assert stats["files"]["tracked"] == 2 + assert stats["versions"]["total"] == 3 + assert stats["versions"]["by_source"] == {"scan": 2, "monitor": 1} + assert sum(stats["versions"]["per_hour_last_24h"]) == 3 + assert stats["top_files_7d"][0]["path"] == str(project / "main.py") + + p = progress(client) + assert p["upload"]["state"] == "unconfigured" + assert p["scans"]["last"][0]["percent"] == 100.0 + assert p["monitor"]["counters"]["capture:committed"] == 1 + + page = client.get("/dashboard", headers={"Authorization": ""}) + assert page.status_code == 200 and "versiond" in page.text + + with client.stream("GET", "/api/v1/progress/stream", params={"interval": 0.25, "max_events": 1}) as stream: + for line in stream.iter_lines(): + if line.startswith("data: "): + assert "upload" in json.loads(line[6:]) + break diff --git a/tests/test_service.py b/tests/test_service.py index a56e1c4..bccec19 100644 --- a/tests/test_service.py +++ b/tests/test_service.py @@ -174,3 +174,15 @@ def test_restore_refuses_unsafe_paths(env): os.symlink("/etc", project / "escape") r = client.post(f"/api/v1/versions/{vid}/restore", json={"target_path": str(project / "escape" / "x")}) assert r.status_code == 400 + + +def test_forget(env): + client, project = env + root = client.post("/api/v1/roots", json={"path": str(project)}).json() + wait_for(lambda: client.get(f"/api/v1/roots/{root['id']}").json()["baseline_state"] == "done") + dry = client.post("/api/v1/forget", json={"path": str(project / "src")}).json() + assert dry["dry_run"] and dry["files"] == 1 and dry["versions"] == 1 + done = client.post("/api/v1/forget", json={"path": str(project / "src"), "dry_run": False}).json() + assert done["files_removed"] == 1 and done["blobs_removed"] == 1 + assert history(client, project / "src" / "app.py") == [] + assert history(client, project / "pyproject.toml") # outside the forgotten tree diff --git a/tests/test_units.py b/tests/test_units.py index 7a4d7eb..b7963c9 100644 --- a/tests/test_units.py +++ b/tests/test_units.py @@ -40,6 +40,14 @@ def test_filter_patterns(): assert f.path_reason(PurePath("src/key.txt")) is None +def test_directory_patterns(): + f = Filters(patterns=("app/data/uploads",)) + assert f.skip_dir("uploads", "app/data/uploads") + assert not f.skip_dir("uploads", "other/uploads") + assert f.path_reason(PurePath("app/data/uploads/a/b.txt")) == "ignored-pattern" + assert f.path_reason(PurePath("app/data/config.json")) is None + + def test_diff_unified(): a = b"one\ntwo\nthree\n" b = b"one\n2\nthree\nfour\n"