FastAPI service on 127.0.0.1:9922 that monitors user-chosen directories with pruned inotify watches, stores versions in a local content-addressed spool indexed in SQLite, coalesces bursts (first + last), and offers history, diff and dry-run-first bulk restore. Includes a CLI with a systemd user unit installer. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
516 lines
19 KiB
Python
516 lines
19 KiB
Python
"""Directory monitoring with inotify, pruned watches and reconciliation scans."""
|
|
|
|
from __future__ import annotations
|
|
|
|
import asyncio
|
|
import errno
|
|
import logging
|
|
import os
|
|
import time
|
|
from pathlib import Path
|
|
from typing import Any, AsyncIterator
|
|
|
|
from asyncinotify import Inotify, Mask, Watch
|
|
|
|
from .db import Database
|
|
from .filters import Filters
|
|
from .ingest import Coalescer, Rejected, Repository, is_within
|
|
|
|
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
|
|
)
|
|
|
|
NETWORK_FS = frozenset({
|
|
"nfs", "nfs4", "cifs", "smb3", "smbfs", "sshfs", "9p", "ceph", "glusterfs",
|
|
"afs", "davfs", "lustre", "gpfs", "virtiofs",
|
|
})
|
|
|
|
MOVE_PAIR_TIMEOUT = 0.5
|
|
YIELD_EVERY = 200
|
|
|
|
|
|
class RootError(Exception):
|
|
def __init__(self, code: str, detail: str):
|
|
super().__init__(detail)
|
|
self.code = code
|
|
self.detail = detail
|
|
|
|
|
|
def max_user_watches() -> int:
|
|
try:
|
|
return int(Path("/proc/sys/fs/inotify/max_user_watches").read_text())
|
|
except (OSError, ValueError):
|
|
return 8192
|
|
|
|
|
|
def watches_in_use() -> int:
|
|
"""Best-effort count of inotify watches held by this user's processes."""
|
|
uid = os.getuid()
|
|
total = 0
|
|
for proc in Path("/proc").iterdir():
|
|
if not proc.name.isdigit():
|
|
continue
|
|
try:
|
|
if proc.stat().st_uid != uid:
|
|
continue
|
|
for fdinfo in (proc / "fdinfo").iterdir():
|
|
with open(fdinfo) as fh:
|
|
total += sum(1 for line in fh if line.startswith("inotify wd:"))
|
|
except OSError:
|
|
continue
|
|
return total
|
|
|
|
|
|
def filesystem_type(path: str) -> str:
|
|
best, fstype = "", ""
|
|
try:
|
|
with open("/proc/self/mounts") as fh:
|
|
for line in fh:
|
|
fields = line.split()
|
|
if len(fields) < 3:
|
|
continue
|
|
mountpoint = fields[1].replace("\\040", " ")
|
|
if is_within(path, mountpoint) and len(mountpoint) > len(best):
|
|
best, fstype = mountpoint, fields[2]
|
|
except OSError:
|
|
pass
|
|
return fstype
|
|
|
|
|
|
def needs_polling(path: str) -> bool:
|
|
fstype = filesystem_type(path)
|
|
return fstype in NETWORK_FS or fstype.startswith("fuse")
|
|
|
|
|
|
class Monitor:
|
|
def __init__(
|
|
self,
|
|
db: Database,
|
|
repo: Repository,
|
|
coalescer: Coalescer,
|
|
filters: Filters,
|
|
poll_interval: float = 60.0,
|
|
reconcile_interval: float = 86_400.0,
|
|
max_watch_fraction: float = 0.5,
|
|
):
|
|
self.db = db
|
|
self.repo = repo
|
|
self.coalescer = coalescer
|
|
self.filters = filters
|
|
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.polled: dict[int, set[str]] = {}
|
|
self.missing: set[int] = set()
|
|
self.moves: dict[int, tuple[str, bool, asyncio.TimerHandle]] = {}
|
|
self.locks: dict[int, asyncio.Lock] = {}
|
|
self.tasks: set[asyncio.Task] = set()
|
|
self._last_full_reconcile = time.monotonic()
|
|
|
|
# lifecycle
|
|
|
|
async def start(self) -> None:
|
|
self.inotify = Inotify()
|
|
self._spawn(self._reader())
|
|
for root in self.db.roots():
|
|
await self._activate(root)
|
|
self._spawn(self._periodic())
|
|
|
|
async def stop(self) -> None:
|
|
for task in list(self.tasks):
|
|
task.cancel()
|
|
await asyncio.gather(*self.tasks, return_exceptions=True)
|
|
for _, _, handle in self.moves.values():
|
|
handle.cancel()
|
|
self.moves.clear()
|
|
if self.inotify is not None:
|
|
self.inotify.close()
|
|
self.inotify = None
|
|
self.watches.clear()
|
|
self.watch_paths.clear()
|
|
|
|
def _spawn(self, coro) -> asyncio.Task:
|
|
task = asyncio.get_running_loop().create_task(coro)
|
|
self.tasks.add(task)
|
|
task.add_done_callback(self._task_done)
|
|
return task
|
|
|
|
def _task_done(self, task: asyncio.Task) -> None:
|
|
self.tasks.discard(task)
|
|
if not task.cancelled() and task.exception():
|
|
log.error("monitor task failed", exc_info=task.exception())
|
|
|
|
# roots
|
|
|
|
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]:
|
|
"""Yield top and every non-ignored directory below it (symlinks not followed)."""
|
|
stack = [top]
|
|
seen = 0
|
|
while stack:
|
|
current = stack.pop()
|
|
yield current
|
|
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):
|
|
stack.append(entry.path)
|
|
except OSError:
|
|
continue
|
|
seen += 1
|
|
if seen % YIELD_EVERY == 0:
|
|
await asyncio.sleep(0)
|
|
|
|
async def count_dirs(self, top: str) -> int:
|
|
return sum([1 async for _ in self.walk_dirs(top)])
|
|
|
|
async def add_root(self, raw_path: str, force: bool = False) -> dict[str, Any]:
|
|
path = os.path.realpath(os.path.expanduser(raw_path))
|
|
if not os.path.isdir(path):
|
|
raise RootError("not-a-directory", f"{path} is not a directory")
|
|
for root in self.db.roots():
|
|
if is_within(path, root["path"]):
|
|
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:
|
|
needed = await self.count_dirs(path)
|
|
limit = max_user_watches()
|
|
available = limit - watches_in_use()
|
|
if needed > available * self.max_watch_fraction:
|
|
raise RootError(
|
|
"watch-limit",
|
|
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())
|
|
)
|
|
self.db.audit("root.add", path=path)
|
|
self.repo.reload_roots()
|
|
root = self.db.root(cur.lastrowid)
|
|
await self._activate(root)
|
|
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"])
|
|
self.polled.pop(root_id, None)
|
|
self.missing.discard(root_id)
|
|
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"]
|
|
if not os.path.isdir(path):
|
|
self.missing.add(root_id)
|
|
self.db.update_root(root_id, mode="missing")
|
|
log.warning("root %s is missing; waiting for it to reappear", path)
|
|
return
|
|
self.missing.discard(root_id)
|
|
self.polled[root_id] = set()
|
|
if needs_polling(path):
|
|
self.polled[root_id].add(path)
|
|
else:
|
|
await self._watch_tree(root_id, path)
|
|
self._refresh_root_status(root_id)
|
|
reason = "reconcile" if root["baseline_state"] == "done" else "baseline"
|
|
self._spawn(self.reconcile(root_id, reason))
|
|
|
|
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"])),
|
|
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
|
|
try:
|
|
self._add_watch(directory)
|
|
except OSError:
|
|
log.warning("inotify watch limit reached; polling %s instead", directory)
|
|
self.polled.setdefault(root_id, set()).add(directory)
|
|
skip.append(directory)
|
|
|
|
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 _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
|
|
|
|
# events
|
|
|
|
async def _reader(self) -> None:
|
|
assert self.inotify is not None
|
|
async for event in self.inotify:
|
|
try:
|
|
self._handle(event)
|
|
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:
|
|
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]
|
|
return
|
|
directory = self.watch_paths.get(watch) if watch is not None else None
|
|
if directory is None:
|
|
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"])
|
|
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 origin is not None:
|
|
origin[2].cancel()
|
|
self._moved_in(origin[0] if origin else None, path, is_dir)
|
|
return
|
|
|
|
if is_dir:
|
|
if Mask.CREATE in mask:
|
|
self._new_directory(path)
|
|
elif Mask.DELETE in mask:
|
|
self._unwatch_tree(path)
|
|
self.repo.mark_deleted_tree(path)
|
|
return
|
|
|
|
if Mask.CLOSE_WRITE in mask:
|
|
self.capture(path, "monitor", "modified")
|
|
elif Mask.DELETE in mask:
|
|
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):
|
|
return
|
|
self._spawn(self._watch_and_scan(root_id, path))
|
|
|
|
async def _watch_and_scan(self, root_id: int, path: 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)
|
|
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
|
|
if is_dir:
|
|
self._unwatch_tree(path)
|
|
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]
|
|
)
|
|
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)
|
|
return
|
|
if old is not None:
|
|
self._unwatch_tree(old)
|
|
self.repo.mark_deleted_tree(old)
|
|
if new_ok:
|
|
self._new_directory(new)
|
|
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")
|
|
|
|
def capture(self, path: str, source: str, reason: str) -> dict[str, Any] | None:
|
|
try:
|
|
self.repo.check_path(path)
|
|
content, st = self.repo.read_file(path)
|
|
self.repo.check_content(content)
|
|
except Rejected as exc:
|
|
log.debug("skip %s: %s", path, exc.reason)
|
|
return None
|
|
except FileNotFoundError:
|
|
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)
|
|
|
|
# scanning
|
|
|
|
async def reconcile(self, root_id: int, reason: str, subtree: str | None = None) -> int:
|
|
"""Compare disk with the index; store changed files, mark vanished ones deleted."""
|
|
root = self.db.root(root_id)
|
|
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):
|
|
if reason == "baseline":
|
|
self.db.update_root(root_id, baseline_state="running")
|
|
async for directory in self.walk_dirs(top):
|
|
try:
|
|
entries = list(os.scandir(directory))
|
|
except OSError:
|
|
continue
|
|
for entry in entries:
|
|
try:
|
|
if not entry.is_file(follow_symlinks=False):
|
|
continue
|
|
st = entry.stat(follow_symlinks=False)
|
|
except OSError:
|
|
continue
|
|
path = entry.path
|
|
if self.filters.path_reason(self.repo.relative(path, root)):
|
|
continue
|
|
if self.filters.size_reason(st.st_size):
|
|
continue
|
|
seen.add(path)
|
|
if path in self.coalescer.pending:
|
|
continue
|
|
row = self.db.file_by_path(path)
|
|
if (
|
|
row
|
|
and row["exists_on_disk"]
|
|
and row["inode"] == st.st_ino
|
|
and row["mtime_ns"] == st.st_mtime_ns
|
|
and row["size"] == st.st_size
|
|
):
|
|
continue
|
|
try:
|
|
content, fst = self.repo.read_file(path)
|
|
self.repo.check_content(content)
|
|
except (Rejected, OSError):
|
|
continue
|
|
result = self.repo.commit(path, content, "scan", reason, fst)
|
|
if result["status"] == "committed":
|
|
committed += 1
|
|
await asyncio.sleep(0)
|
|
rows = self.db.all(
|
|
"SELECT path FROM files WHERE root_id = ? AND exists_on_disk = 1", (root_id,)
|
|
)
|
|
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"])
|
|
fields: dict[str, Any] = {"last_scan_at": time.time()}
|
|
if reason == "baseline":
|
|
fields.update(baseline_state="done", baseline_files=committed)
|
|
self.db.update_root(root_id, **fields)
|
|
if committed:
|
|
log.info("%s scan of %s stored %d versions", reason, top, committed)
|
|
return committed
|
|
|
|
async def _periodic(self) -> None:
|
|
while True:
|
|
await asyncio.sleep(self.poll_interval)
|
|
for root_id in list(self.missing):
|
|
root = self.db.root(root_id)
|
|
if root is None:
|
|
self.missing.discard(root_id)
|
|
elif os.path.isdir(root["path"]):
|
|
log.info("root %s is back", root["path"])
|
|
await self._activate(root)
|
|
for root_id, dirs in list(self.polled.items()):
|
|
for directory in list(dirs):
|
|
await self.reconcile(root_id, "poll", subtree=directory)
|
|
if time.monotonic() - self._last_full_reconcile >= self.reconcile_interval:
|
|
self._last_full_reconcile = time.monotonic()
|
|
for root in self.db.roots():
|
|
await self.reconcile(root["id"], "reconcile")
|