Add WebDAV backup, statistics, progress and dashboard

- WebDAV uploader: concurrency, request and bandwidth limits, retry with
  backoff, offline catch-up, encrypted manifests, durability tracking
- Automatic unique remote directory claimed with a conditional PUT
- AES-256-GCM client-side encryption with keyed blob names
- /stats, /progress, /progress/stream and a /dashboard page
- Own ctypes inotify binding with a compact watch tree (263 MB -> 76 MB RSS)
- Directory-path ignore patterns and a forget command
- Indexed range queries instead of LIKE for path prefixes

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
This commit is contained in:
retoor
2026-10-08 20:24:04 +02:00
co-authored by Claude Opus 5.5
parent ddf09bea88
commit 7e05e3d26c
18 changed files with 1992 additions and 208 deletions
+6 -2
View File
@@ -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"]
+26 -11
View File
@@ -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: <http://127.0.0.1:9922/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
```
<base_path>/<hostname-slug>-<id8>/
blobs/sha256/ab/cd/<full-hash>.zst # content-addressed, zstd-compressed (optionally encrypted)
manifests/<yyyy>/<mm>/<dd>/<batch-id>.jsonl.zst # append-only version records
meta/format.json # storage format version
blobs/<name[:2]>/<name> # name = HMAC-SHA256(key, sha256); zstd, then AES-256-GCM
manifests/<yyyy>/<mm>/<dd>/<batch-id>.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.
+147 -6
View File
@@ -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)
+170 -2
View File
@@ -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
+41 -6
View File
@@ -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"
+57
View File
@@ -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)
+259
View File
@@ -0,0 +1,259 @@
<!doctype html>
<html lang="en">
<head>
<meta charset="utf-8">
<meta name="viewport" content="width=device-width, initial-scale=1">
<title>versiond</title>
<style>
:root {
--bg: #f6f7f9; --card: #ffffff; --text: #1d2330; --muted: #6b7385; --line: #e3e6ec;
--accent: #2f6fdb; --accent-soft: #dfe9fb; --ok: #1f8a4c; --warn: #b7791f; --bad: #c53030;
--mono: ui-monospace, SFMono-Regular, Menlo, Consolas, monospace;
}
@media (prefers-color-scheme: dark) {
:root {
--bg: #12151b; --card: #1a1e26; --text: #e4e7ee; --muted: #8b93a5; --line: #2a303b;
--accent: #6c9cf0; --accent-soft: #24324a; --ok: #4cc283; --warn: #e0a84a; --bad: #f07171;
}
}
* { box-sizing: border-box; }
body { margin: 0; background: var(--bg); color: var(--text);
font: 14px/1.45 system-ui, -apple-system, "Segoe UI", sans-serif; }
header { display: flex; align-items: center; gap: 12px; flex-wrap: wrap;
padding: 16px 24px; border-bottom: 1px solid var(--line); background: var(--card); }
header h1 { font-size: 18px; margin: 0 8px 0 0; letter-spacing: -0.01em; }
.pill { padding: 2px 10px; border-radius: 999px; font-size: 12px; font-weight: 600;
background: var(--accent-soft); color: var(--accent); }
.pill.ok { color: var(--ok); background: color-mix(in srgb, var(--ok) 14%, transparent); }
.pill.warn { color: var(--warn); background: color-mix(in srgb, var(--warn) 14%, transparent); }
.pill.bad { color: var(--bad); background: color-mix(in srgb, var(--bad) 14%, transparent); }
header .spacer { flex: 1; }
header .meta { color: var(--muted); font-size: 12px; }
main { max-width: 1200px; margin: 0 auto; padding: 20px 16px 40px; display: grid; gap: 16px; }
.tiles { display: grid; grid-template-columns: repeat(auto-fit, minmax(170px, 1fr)); gap: 12px; }
.tile, .card { background: var(--card); border: 1px solid var(--line); border-radius: 10px; }
.tile { padding: 14px 16px; }
.tile .label { color: var(--muted); font-size: 12px; text-transform: uppercase; letter-spacing: 0.04em; }
.tile .value { font-size: 24px; font-weight: 650; margin-top: 4px; font-variant-numeric: tabular-nums; }
.tile .sub { color: var(--muted); font-size: 12px; margin-top: 2px; }
.grid2 { display: grid; grid-template-columns: repeat(auto-fit, minmax(360px, 1fr)); gap: 16px; }
.card { padding: 16px 18px; min-width: 0; }
.card h2 { font-size: 14px; margin: 0 0 12px; display: flex; justify-content: space-between; align-items: center; }
.bar { height: 10px; background: var(--line); border-radius: 999px; overflow: hidden; margin: 8px 0 10px; }
.bar > div { height: 100%; background: var(--accent); border-radius: 999px; transition: width 0.6s ease; }
.kv { display: grid; grid-template-columns: max-content 1fr; gap: 4px 16px; font-size: 13px; }
.kv dt { color: var(--muted); }
.kv dd { margin: 0; font-variant-numeric: tabular-nums; overflow-wrap: anywhere; }
table { width: 100%; border-collapse: collapse; font-size: 13px; }
td, th { text-align: left; padding: 6px 4px; border-bottom: 1px solid var(--line); }
th { color: var(--muted); font-weight: 500; font-size: 12px; }
td.num, th.num { text-align: right; font-variant-numeric: tabular-nums; }
td.path { font-family: var(--mono); font-size: 12px; overflow-wrap: anywhere; }
.empty { color: var(--muted); font-size: 13px; }
.error { color: var(--bad); font-size: 12px; overflow-wrap: anywhere; }
svg.chart { width: 100%; height: 120px; display: block; }
svg.chart rect { fill: var(--accent); }
svg.chart text { fill: var(--muted); font-size: 10px; }
.scan + .scan { margin-top: 14px; padding-top: 14px; border-top: 1px solid var(--line); }
#login { max-width: 420px; margin: 80px auto; }
#login input { width: 100%; padding: 8px 10px; margin: 8px 0; font: inherit; border-radius: 6px;
border: 1px solid var(--line); background: var(--bg); color: var(--text); }
#login button { padding: 8px 14px; border: 0; border-radius: 6px; background: var(--accent); color: #fff; font: inherit; cursor: pointer; }
code { font-family: var(--mono); font-size: 12px; }
</style>
</head>
<body>
<header>
<h1>versiond</h1>
<span id="health" class="pill">…</span>
<span id="remote" class="pill">remote …</span>
<span class="spacer"></span>
<span class="meta" id="updated"></span>
</header>
<div id="login" class="card" hidden>
<h2>API token</h2>
<p class="empty">Run <code>versiond dashboard</code> to open this page signed in, or paste the output of <code>versiond token</code>.</p>
<input id="token" type="password" autocomplete="off" placeholder="token">
<button id="save">Open dashboard</button>
<p id="login-error" class="error"></p>
</div>
<main id="app" hidden>
<section class="tiles">
<div class="tile"><div class="label">Files tracked</div><div class="value" id="t-files">–</div><div class="sub" id="t-files-sub"></div></div>
<div class="tile"><div class="label">Versions</div><div class="value" id="t-versions">–</div><div class="sub" id="t-versions-sub"></div></div>
<div class="tile"><div class="label">Stored locally</div><div class="value" id="t-storage">–</div><div class="sub" id="t-storage-sub"></div></div>
<div class="tile"><div class="label">Backed up</div><div class="value" id="t-durable">–</div><div class="sub" id="t-durable-sub"></div></div>
<div class="tile"><div class="label">Watches</div><div class="value" id="t-watches">–</div><div class="sub" id="t-watches-sub"></div></div>
</section>
<section class="grid2">
<div class="card"><h2>Scans <span class="pill" id="scan-pill">idle</span></h2><div id="scans"></div></div>
<div class="card"><h2>Upload to WebDAV <span class="pill" id="upload-pill">–</span></h2><div id="upload"></div></div>
</section>
<section class="grid2">
<div class="card"><h2>Versions per hour, last 24 h</h2><svg class="chart" id="chart" viewBox="0 0 480 120" preserveAspectRatio="none"></svg></div>
<div class="card"><h2>Monitored directories</h2><div id="roots"></div></div>
</section>
<section class="grid2">
<div class="card"><h2>Most edited files, 7 days</h2><div id="top-files"></div></div>
<div class="card"><h2>Most active projects, 7 days</h2><div id="top-projects"></div></div>
</section>
<section class="grid2">
<div class="card"><h2>Versions by source</h2><div id="sources"></div></div>
<div class="card"><h2>Monitor counters, this session</h2><div id="counters"></div></div>
</section>
</main>
<script>
(() => {
const $ = (id) => document.getElementById(id);
const store = {
get() { try { return localStorage.getItem("versiond-token"); } catch { return null; } },
set(v) { try { localStorage.setItem("versiond-token", v); } catch {} },
clear() { try { localStorage.removeItem("versiond-token"); } catch {} },
};
let token = null;
const fromHash = new URLSearchParams(location.hash.slice(1)).get("token");
if (fromHash) { store.set(fromHash); history.replaceState(null, "", location.pathname); }
token = fromHash || store.get();
const esc = (s) => String(s ?? "").replace(/[&<>"]/g, (c) => ({ "&": "&amp;", "<": "&lt;", ">": "&gt;", '"': "&quot;" }[c]));
const num = (n) => n == null ? "–" : Number(n).toLocaleString();
const bytes = (n) => {
if (n == null) return "–";
const u = ["B", "KB", "MB", "GB", "TB"]; let i = 0;
while (n >= 1024 && i < u.length - 1) { n /= 1024; i++; }
return `${n.toFixed(i ? 1 : 0)} ${u[i]}`;
};
const dur = (s) => {
if (s == null) return "–";
s = Math.round(s);
if (s < 60) return `${s}s`;
if (s < 3600) return `${Math.floor(s / 60)}m ${s % 60}s`;
if (s < 86400) return `${Math.floor(s / 3600)}h ${Math.floor((s % 3600) / 60)}m`;
return `${Math.floor(s / 86400)}d ${Math.floor((s % 86400) / 3600)}h`;
};
const ago = (t) => t ? dur(Date.now() / 1000 - t) + " ago" : "never";
const pillClass = (el, kind, text) => { el.className = "pill " + kind; el.textContent = text; };
async function api(path) {
const r = await fetch(path, { headers: { Authorization: `Bearer ${token}` } });
if (r.status === 401) { store.clear(); showLogin("Token rejected."); throw new Error("unauthorized"); }
if (!r.ok) throw new Error(`${path}: HTTP ${r.status}`);
return r.json();
}
function showLogin(msg) {
$("app").hidden = true; $("login").hidden = false; $("login-error").textContent = msg || "";
}
$("save").onclick = () => { token = $("token").value.trim(); store.set(token); start(); };
$("token").onkeydown = (e) => { if (e.key === "Enter") $("save").click(); };
function table(rows, cols) {
if (!rows.length) return '<p class="empty">Nothing yet.</p>';
const head = cols.map((c) => `<th class="${c.num ? "num" : ""}">${c.label}</th>`).join("");
const body = rows.map((r) => "<tr>" + cols.map((c) =>
`<td class="${c.num ? "num" : c.path ? "path" : ""}">${c.fmt ? c.fmt(r[c.key], r) : esc(r[c.key])}</td>`).join("") + "</tr>").join("");
return `<table><thead><tr>${head}</tr></thead><tbody>${body}</tbody></table>`;
}
function renderProgress(p) {
const health = $("health");
const degraded = p.monitor.roots.some((r) => r.mode !== "inotify");
pillClass(health, degraded ? "warn" : "ok", degraded ? "degraded" : "healthy");
const up = p.upload;
const remoteKind = { online: "ok", connecting: "", offline: "warn", error: "bad", unconfigured: "warn" }[up.state] ?? "";
pillClass($("remote"), remoteKind, "remote " + up.state);
$("updated").textContent = `uptime ${dur(p.uptime_seconds)} · updated ${new Date().toLocaleTimeString()}`;
$("t-watches").textContent = num(p.monitor.watches);
$("t-watches-sub").textContent = `${p.monitor.roots.length} root(s) · ${p.monitor.pending_paths} pending`;
$("t-durable").textContent = up.blobs_total ? `${up.percent}%` : "–";
$("t-durable-sub").textContent = `${num(up.versions_durable)} versions durable · ${num(up.versions_local)} local`;
const active = p.scans.active;
pillClass($("scan-pill"), active.length ? "" : "ok", active.length ? `${active.length} running` : "idle");
const scans = (active.length ? active : p.scans.last).map((s) => `
<div class="scan">
<div><strong>${esc(s.reason)}</strong> <span class="empty">${esc(s.top)}</span></div>
<div class="bar"><div style="width:${s.percent ?? 0}%"></div></div>
<dl class="kv">
<dt>Progress</dt><dd>${s.percent ?? "–"}% · ${num(s.dirs_done)} / ${num(s.dirs_total)} directories</dd>
<dt>Files</dt><dd>${num(s.files_seen)} seen · ${num(s.files_read)} read · ${num(s.versions_stored)} new versions (${bytes(s.bytes_stored)})</dd>
<dt>Speed</dt><dd>${s.files_per_second ?? "–"} files/s</dd>
<dt>${s.finished_at ? "Took" : "Elapsed"}</dt><dd>${dur(s.elapsed_seconds)}${s.eta_seconds != null ? ` · ETA ${dur(s.eta_seconds)}` : ""}</dd>
</dl>
</div>`).join("");
$("scans").innerHTML = scans || '<p class="empty">No scan has run yet in this session.</p>';
pillClass($("upload-pill"), remoteKind, up.state);
if (up.state === "unconfigured") {
$("upload").innerHTML = '<p class="empty">No WebDAV remote configured. Run <code>versiond remote set</code>.</p>';
} else {
$("upload").innerHTML = `
<div class="bar"><div style="width:${up.percent}%"></div></div>
<dl class="kv">
<dt>Uploaded</dt><dd>${up.percent}% · ${num(up.blobs_uploaded)} / ${num(up.blobs_total)} blobs</dd>
<dt>Pending</dt><dd>${num(up.blobs_pending)} blobs · ${bytes(up.bytes_pending)}${up.blobs_failed ? ` · ${num(up.blobs_failed)} retrying` : ""}</dd>
<dt>Speed</dt><dd>${up.blobs_per_second} blobs/s · ${bytes(up.bytes_per_second)}/s</dd>
<dt>ETA</dt><dd>${up.eta_seconds != null ? dur(up.eta_seconds) : up.blobs_pending ? "–" : "up to date"}</dd>
<dt>Directory</dt><dd><code>${esc(up.directory ?? "–")}</code></dd>
<dt>Last upload</dt><dd>${ago(up.last_success_at)}</dd>
<dt>This session</dt><dd>${num(up.session.blobs_uploaded)} blobs · ${bytes(up.session.bytes_uploaded)} · ${num(up.session.manifests_uploaded)} manifests</dd>
</dl>
${up.last_error ? `<p class="error">Last error (${ago(up.last_error_at)}): ${esc(up.last_error)}</p>` : ""}`;
}
$("roots").innerHTML = table(p.monitor.roots, [
{ key: "path", label: "Path", path: true },
{ key: "mode", label: "Mode" },
{ key: "watches", label: "Watches", num: true, fmt: num },
{ key: "baseline_state", label: "Baseline" },
{ key: "last_scan_at", label: "Last scan", fmt: ago },
]);
const counters = Object.entries(p.monitor.counters).sort((a, b) => b[1] - a[1]).map(([k, v]) => ({ k, v }));
$("counters").innerHTML = table(counters, [{ key: "k", label: "Counter" }, { key: "v", label: "Count", num: true, fmt: num }]);
}
function renderStats(s) {
$("t-files").textContent = num(s.files.tracked);
$("t-files-sub").textContent = `${num(s.files.existing)} on disk · ${num(s.projects)} projects`;
$("t-versions").textContent = num(s.versions.total);
$("t-versions-sub").textContent = `${num(s.versions.last_hour)} last hour · ${num(s.versions.last_24h)} last 24 h`;
$("t-storage").textContent = bytes(s.storage.stored_bytes);
$("t-storage-sub").textContent = `${bytes(s.storage.raw_bytes)} raw · ${s.storage.compression_ratio ?? "–"}× compression · ${s.storage.dedup_ratio ?? "–"}× dedup`;
const data = s.versions.per_hour_last_24h, max = Math.max(1, ...data), w = 480 / data.length;
$("chart").innerHTML = data.map((v, i) => {
const h = Math.max(v ? 2 : 0, (v / max) * 100);
return `<rect x="${i * w + 2}" y="${104 - h}" width="${w - 4}" height="${h}" rx="2"><title>${v} versions, ${23 - i}h ago</title></rect>`;
}).join("") + `<text x="2" y="118">−24 h</text><text x="452" y="118">now</text><text x="2" y="10">${num(max)}</text>`;
$("top-files").innerHTML = table(s.top_files_7d, [{ key: "path", label: "File", path: true }, { key: "versions", label: "Versions", num: true, fmt: num }]);
$("top-projects").innerHTML = table(s.top_projects_7d, [{ key: "path", label: "Project", path: true }, { key: "versions", label: "Versions", num: true, fmt: num }]);
const sources = Object.entries(s.versions.by_source).map(([k, v]) => ({ k, v }));
$("sources").innerHTML = table(sources, [{ key: "k", label: "Source" }, { key: "v", label: "Versions", num: true, fmt: num }]);
}
let timers = [];
async function start() {
timers.forEach(clearInterval); timers = [];
if (!token) return showLogin();
try {
const [p, s] = await Promise.all([api("/api/v1/progress"), api("/api/v1/stats")]);
$("login").hidden = true; $("app").hidden = false;
renderProgress(p); renderStats(s);
} catch (e) { if (e.message !== "unauthorized") showLogin(e.message); return; }
timers.push(setInterval(() => api("/api/v1/progress").then(renderProgress).catch(() => {}), 2000));
timers.push(setInterval(() => api("/api/v1/stats").then(renderStats).catch(() => {}), 15000));
}
start();
})();
</script>
</body>
</html>
+23 -1
View File
@@ -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}")
+11
View File
@@ -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):
+54 -6
View File
@@ -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:
+82
View File
@@ -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
+307 -174
View File
@@ -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():
+201
View File
@@ -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/<hmac[:2]>/<hmac>", "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")
+77
View File
@@ -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"],
}
+349
View File
@@ -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]
+162
View File
@@ -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
+12
View File
@@ -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
+8
View File
@@ -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"