"""Directory monitoring with inotify, pruned watches and reconciliation scans.""" from __future__ import annotations import asyncio 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 . import inotify as ino from .db import Database from .filters import Filters from .ingest import Coalescer, Rejected, Repository, is_within, tree_range log = logging.getLogger(__name__) DIR_MASK = ( 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({ "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 _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, 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: 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[_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 = 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()) 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: asyncio.get_running_loop().remove_reader(self.inotify.fileno()) self.inotify.close() self.inotify = None self.by_wd.clear() self.tops.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()) 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 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._skip(entry.path, entry.name, root_path): 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, 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") if not needs_polling(path) 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())) 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(root_id) 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") 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") 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, 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 self.db.update_root( root_id, mode=self.root_mode(root_id, root["path"]), watch_count=self.watch_counts[root_id], polled_dirs=len(self.polled.get(root_id, ())), ) # watches 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: 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: pass added += 1 if added % YIELD_EVERY == 0: await asyncio.sleep(0) 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 _root_id_of(self, node: _Dir) -> int | None: return self.top_ids.get(node.top().name) # events 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(wd, mask, cookie, name) except Exception: log.exception("failed handling inotify event") 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 node = self.by_wd.get(wd) if node is None: return 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 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 if mask & ino.MOVED_TO: origin = self.moves.pop(cookie, None) if origin is not None: origin[3].cancel() self._moved_in(origin, node, name, is_dir) return path = node.path() + "/" + name if is_dir: 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 & ino.CLOSE_WRITE: self.capture(path, "monitor", "modified") elif mask & ino.DELETE: self.counters["deletes"] += 1 self.repo.mark_deleted(path) def _new_directory(self, parent: _Dir, name: str) -> None: root_path = parent.top().name if self._skip(parent.path() + "/" + name, name, root_path): return 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, parent: _Dir, name: str) -> None: # Watch first, then scan, so files created before the watch existed are not missed. 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 parent, name, is_dir, _ = origin path = parent.path() + "/" + name if is_dir: 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, 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: 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 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(parent, name) return 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: self.repo.check_path(path) content, st = self.repo.read_file(path) self.repo.check_content(content) except Rejected as exc: self.counters[f"skipped:{exc.reason}"] += 1 return None except OSError: return None result = self.coalescer.submit(path, content, source, reason, st) self.counters[f"capture:{result['status']}"] += 1 return result # 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"] 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, root["path"]): progress.dirs_done += 1 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 progress.files_seen += 1 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 progress.files_read += 1 result = self.repo.commit(path, content, "scan", reason, fst) if result["status"] == "committed": 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 " "AND (path = ? OR (path >= ? AND path < ?))", (root_id, top, *tree_range(top)), ) 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=progress.versions_stored) self.db.update_root(root_id, **fields) 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: 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) 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(): await self.reconcile(root["id"], "reconcile")