diff options
Diffstat (limited to 'packages/meshbay-node/src/meshbay_node/indexer')
4 files changed, 330 insertions, 50 deletions
diff --git a/packages/meshbay-node/src/meshbay_node/indexer/__init__.py b/packages/meshbay-node/src/meshbay_node/indexer/__init__.py index bb6231b..c92c730 100644 --- a/packages/meshbay-node/src/meshbay_node/indexer/__init__.py +++ b/packages/meshbay-node/src/meshbay_node/indexer/__init__.py @@ -1,5 +1,6 @@ """Directory indexer and Mesh Group Index.""" from .indexer import DirectoryIndexer from .group_index import GroupIndex +from .cache import IndexCache -__all__ = ["DirectoryIndexer", "GroupIndex"] +__all__ = ["DirectoryIndexer", "GroupIndex", "IndexCache"] diff --git a/packages/meshbay-node/src/meshbay_node/indexer/cache.py b/packages/meshbay-node/src/meshbay_node/indexer/cache.py new file mode 100644 index 0000000..c26ddf1 --- /dev/null +++ b/packages/meshbay-node/src/meshbay_node/indexer/cache.py @@ -0,0 +1,95 @@ +""" +MeshBay Node — persistent (path, size, mtime) -> hash cache, one per group. + +Without this, every node restart re-reads and re-hashes every file in every +root, even when nothing changed — measured at 23 minutes for a 114 GB library +on a USB hard drive. This cache lets a scan skip the read entirely for a file +whose size and mtime still match what was hashed last time. + +It is a path-keyed accelerator only. The GroupIndex itself stays keyed by +content hash (see indexer.py's note on why two identical files are one +entry) — this cache never changes that, it only avoids recomputing a hash +that has not changed. +""" + +import logging +from dataclasses import dataclass +from pathlib import Path + +import aiosqlite + +log = logging.getLogger(__name__) + +_SCHEMA = """ +CREATE TABLE IF NOT EXISTS files ( + path TEXT PRIMARY KEY, + mtime REAL NOT NULL, + size INTEGER NOT NULL, + hash TEXT NOT NULL, + type TEXT NOT NULL, + added_at INTEGER NOT NULL +); +""" + + +@dataclass +class CachedEntry: + hash: str + type: str + added_at: int + + +class IndexCache: + """Async SQLite (path, size, mtime) -> hash cache for one group.""" + + def __init__(self, db_path: Path): + self._db_path = db_path + self._db: aiosqlite.Connection | None = None + + async def open(self) -> None: + self._db_path.parent.mkdir(parents=True, exist_ok=True) + self._db = await aiosqlite.connect(str(self._db_path)) + await self._db.executescript(_SCHEMA) + await self._db.commit() + + async def close(self) -> None: + if self._db: + await self._db.close() + self._db = None + + async def __aenter__(self): + await self.open() + return self + + async def __aexit__(self, *_): + await self.close() + + async def lookup(self, path: str, size: int, mtime: float) -> CachedEntry | None: + """ + A cache hit requires an EXACT match on both size and mtime. A mtime + touched without a content change is a false negative (an unnecessary + rehash) — accepted, since the alternative (trusting a stale hash) is + a silent wrong answer instead of an occasional wasted read. + """ + async with self._db.execute( + "SELECT hash, type, added_at FROM files " + "WHERE path = ? AND size = ? AND mtime = ?", + (path, size, mtime)) as cur: + row = await cur.fetchone() + return CachedEntry(hash=row[0], type=row[1], added_at=row[2]) if row else None + + async def put(self, path: str, size: int, mtime: float, hash: str, + type: str, added_at: int) -> None: + """ + Written only once a file has been hashed in full — never partway + through — so a crash mid-hash leaves no stale/partial row behind: the + next scan simply finds no cache entry and hashes the file again. + """ + await self._db.execute( + "INSERT INTO files (path, mtime, size, hash, type, added_at) " + "VALUES (?, ?, ?, ?, ?, ?) " + "ON CONFLICT(path) DO UPDATE SET " + "mtime = excluded.mtime, size = excluded.size, hash = excluded.hash, " + "type = excluded.type, added_at = excluded.added_at", + (path, mtime, size, hash, type, added_at)) + await self._db.commit() diff --git a/packages/meshbay-node/src/meshbay_node/indexer/group_index.py b/packages/meshbay-node/src/meshbay_node/indexer/group_index.py index 5dbdc5d..ec98667 100644 --- a/packages/meshbay-node/src/meshbay_node/indexer/group_index.py +++ b/packages/meshbay-node/src/meshbay_node/indexer/group_index.py @@ -86,6 +86,23 @@ class GroupIndex: def count(self) -> int: return len(self._entries) + def entries_by_id(self) -> dict: + """A snapshot copy, for diff() to compare a later version against — + see daemon.py _on_index_change, the only caller.""" + return dict(self._entries) + + @classmethod + def _snapshot(cls, group_id: str, sk_node: Ed25519PrivateKey, gek: bytes | None, + version: int, entries_by_id: dict) -> "GroupIndex": + """ + A lightweight stand-in for diff()'s `previous` argument — never + serialized or sent anywhere, just a comparison point built from an + earlier entries_by_id() snapshot rather than a live GroupIndex. + """ + idx = cls(group_id=group_id, sk_node=sk_node, gek=gek, version=version) + idx._entries = dict(entries_by_id) + return idx + # ── Serialisation ───────────────────────────────────────────────────────── def serialize(self) -> bytes: diff --git a/packages/meshbay-node/src/meshbay_node/indexer/indexer.py b/packages/meshbay-node/src/meshbay_node/indexer/indexer.py index 482b556..b479f7e 100644 --- a/packages/meshbay-node/src/meshbay_node/indexer/indexer.py +++ b/packages/meshbay-node/src/meshbay_node/indexer/indexer.py @@ -26,6 +26,7 @@ import asyncio import logging import time from concurrent.futures import ThreadPoolExecutor +from dataclasses import dataclass from pathlib import Path from typing import Callable, Awaitable @@ -36,6 +37,7 @@ from watchdog.observers import Observer from meshbay_common.paths import fold, find_fold_collisions, long_path from meshbay_common.protocol import IndexEntry +from meshbay_node.indexer.cache import IndexCache from meshbay_node.indexer.group_index import GroupIndex from meshbay_node.roots import Root, RootSet @@ -75,6 +77,27 @@ def _is_indexable(path: Path) -> bool: _HASH_CHUNK = 8 * 1024 * 1024 # 8 MB streaming hash chunks +@dataclass +class IndexProgress: + """ + A snapshot of "is this indexer mid-scan right now", for the status shown + to the operator (Create Group wizard, adding a directory) and pushed to + connected members (a presence dot, never anything more specific — see + daemon.py/webrtc_server.py). Reset per _scan_root() call rather than + accumulated across a group's roots: the consumers that matter always + watch exactly one root being scanned. + + Mutated only from the asyncio loop thread (the hashing itself runs in an + executor thread, but never touches this), so no lock is needed. + """ + scanning: bool = False + scanned_bytes: int = 0 + total_bytes: int = 0 + # Basename only, deliberately not the full path — enough to show progress + # without broadcasting the operator's directory structure. + current_dir: str = "" + + def _virtual_dir(root: Root, file_path: Path) -> str: """ The directory a file appears in, as members see it: `"Films/2024"`. @@ -87,6 +110,17 @@ def _virtual_dir(root: Root, file_path: Path) -> str: return root.name if str(rel) == "." else f"{root.name}/{rel.as_posix()}" +def _walk_root(root: Root) -> list[Path]: + """ + Blocking directory walk — always run via an executor, never awaited + directly in the asyncio loop. A tree of tens of thousands of files (or + one on a slow network share) can take seconds; run inline, that stalls + every other thing the daemon is doing — WebRTC sessions, chat, the admin + UI — for as long as it takes. + """ + return [p for p in root.path.rglob("*") if p.is_file()] + + def _scan_file(root: Root, file_path: Path) -> IndexEntry | None: """Compute IndexEntry for a file. Blocking — run in executor. Uses streaming blake3 so arbitrarily large files (ISOs, VM images, etc.) @@ -136,8 +170,17 @@ class DirectoryIndexer: # How often to re-check which roots are readable and reconcile the index # against what is actually on disk. Not a poll for changes — a backstop for # the events the OS did not deliver, and the way a re-plugged drive is - # noticed. - RECONCILE_SECS = 60.0 + # noticed. Watchdog already covers the common case in real time, so this + # does not need to run often to do its job; it backs off further still + # (see _reconcile_loop) when nothing has changed for a while, and a + # per-group operator setting (roster.py SETTING_RECONCILE_INTERVAL) can + # override the starting point. + DEFAULT_RECONCILE_SECS = 600.0 # 10 min + RECONCILE_BACKOFF_CAP = 7200.0 # 2 h — never sleeps longer than this + # How long to wait after the last event on a given path before acting on + # it — several writes to the same file in quick succession (a slow copy + # in several passes) collapse into one hash instead of one per write. + DEFAULT_DEBOUNCE_SECS = 2.0 def __init__( self, @@ -146,20 +189,43 @@ class DirectoryIndexer: sk_node: Ed25519PrivateKey, gek: bytes | None, on_change: Callable[["DirectoryIndexer"], Awaitable[None]] | None = None, + cache: IndexCache | None = None, + reconcile_secs: float = DEFAULT_RECONCILE_SECS, + debounce_secs: float = DEFAULT_DEBOUNCE_SECS, ): self.roots = roots self.group_id = group_id self.sk_node = sk_node self.gek = gek self.on_change = on_change + self.reconcile_secs = reconcile_secs + self.debounce_secs = debounce_secs + # Current backoff delay — starts at reconcile_secs, doubles on every + # tick that finds nothing changed (up to RECONCILE_BACKOFF_CAP), and + # resets the moment something real happens (a change, or a peer + # connecting — see note_activity()). + self._reconcile_delay = reconcile_secs + # Path -> (size, mtime, hash) accelerator, so a restart does not have + # to re-read a file it already hashed last time (see cache.py). None + # in tests that do not care about it — every hash is then a miss. + self._cache = cache self._index = GroupIndex(group_id=group_id, sk_node=sk_node, gek=gek) self._index.roots = roots.describe() - self._executor = ThreadPoolExecutor(max_workers=2, thread_name_prefix="indexer") + # One worker, deliberately, not a real pool: hashing two files at once + # buys nothing here and can cost a lot. It only offloads the blocking + # read+hash off the asyncio loop; it was never used for concurrency — + # every call site awaits one run_in_executor before starting the next + # (see _hash_or_cached below) — and measured on a spinning USB drive, + # two interleaved multi-GB reads would seek-thrash against each other + # rather than go faster. Left at 1 so the number does not promise a + # concurrency this code never provided. + self._executor = ThreadPoolExecutor(max_workers=1, thread_name_prefix="indexer") self._observer: Observer | None = None self._loop: asyncio.AbstractEventLoop | None = None self._reconciler: asyncio.Task | None = None self._pending_timers: dict[str, asyncio.TimerHandle] = {} + self.progress = IndexProgress() @property def index(self) -> GroupIndex: @@ -204,21 +270,77 @@ class DirectoryIndexer: async def _scan_root(self, root: Root) -> int: log.info("Scanning %s (root %r) ...", root.path, root.name) - loop = asyncio.get_event_loop() count = 0 + loop = asyncio.get_event_loop() try: - files = [p for p in root.path.rglob("*") if p.is_file()] + files = await loop.run_in_executor(self._executor, _walk_root, root) except OSError as e: log.warning("Cannot scan root %r: %s", root.name, e) return 0 - for file_path in files: - entry = await loop.run_in_executor( - self._executor, _scan_file, root, file_path) - if entry: - self._index.add_entry(entry) - count += 1 + + # Sizes up front, off the same listing that already walked the tree — + # the progress bar's denominator, not a second pass over the disk. + sized: list[tuple[Path, int]] = [] + for p in files: + try: + sized.append((p, p.stat().st_size)) + except OSError: + continue + + self.progress.scanning = True + self.progress.scanned_bytes = 0 + self.progress.total_bytes = sum(size for _, size in sized) + self.progress.current_dir = "" + try: + for file_path, size in sized: + self.progress.current_dir = file_path.parent.name + entry = await self._hash_or_cached(root, file_path) + self.progress.scanned_bytes += size + if entry: + self._index.add_entry(entry) + count += 1 + finally: + # Must run even if a hash/IO error propagates out of the loop + # above — an indexing state that never turns back off is worse + # than the scan itself failing. + self.progress.scanning = False + self.progress.current_dir = "" return count + async def _hash_or_cached(self, root: Root, file_path: Path) -> IndexEntry | None: + """ + Cache-aware replacement for a bare _scan_file() call: skips the + content read entirely when this path's (size, mtime) still match + what was hashed last time — the difference between a redundant full + rehash of a 100+ GB library on every restart and a stat()-only pass. + The only place that decides to actually read a file's bytes. + """ + if not _is_indexable(file_path): + return None + try: + st = file_path.stat() + except OSError: + return None + + if self._cache is not None: + cached = await self._cache.lookup(str(file_path), st.st_size, st.st_mtime) + if cached is not None: + return IndexEntry( + id=cached.hash, + name=file_path.name, + path=_virtual_dir(root, file_path), + size=st.st_size, + type=cached.type, + added_at=cached.added_at, + ) + + loop = asyncio.get_event_loop() + entry = await loop.run_in_executor(self._executor, _scan_file, root, file_path) + if entry and self._cache is not None: + await self._cache.put(str(file_path), st.st_size, st.st_mtime, + entry.id, entry.type, entry.added_at) + return entry + def _report_collisions(self) -> None: """ Names that are the same file on a case-insensitive filesystem. @@ -324,20 +446,38 @@ class DirectoryIndexer: async def _reconcile_loop(self) -> None: while True: try: - await asyncio.sleep(self.RECONCILE_SECS) - await self.reconcile() + await asyncio.sleep(self._reconcile_delay) + changed = await self.reconcile() + if changed: + self._reconcile_delay = self.reconcile_secs + else: + self._reconcile_delay = min( + self._reconcile_delay * 2, self.RECONCILE_BACKOFF_CAP) except asyncio.CancelledError: raise except Exception: log.exception("Reconcile failed — continuing") - async def reconcile(self) -> None: + def note_activity(self) -> None: + """ + Called when something makes a prompt reconcile worth having again — + today, a peer completing the handshake for this group + (webrtc_server.py). Someone is looking, so the backstop should be at + its normal cadence rather than however far backoff had stretched it. + """ + self._reconcile_delay = self.reconcile_secs + + async def reconcile(self) -> bool: """ Re-check availability, and rescan roots that came back. The only place a root's entries are dropped: when the root is readable and the files are genuinely gone. A root that is not readable is left untouched, which is the whole point. + + Returns whether anything actually changed — _reconcile_loop uses this + to back off when a pass finds nothing to do, rather than running at + the same cadence forever regardless of how quiet the root is. """ changed = self.roots.refresh_availability() touched = False @@ -367,6 +507,8 @@ class DirectoryIndexer: if self.on_change: await self.on_change(self) + return touched + async def _sweep_available_roots(self) -> bool: """ Catch what the watcher missed: files gone, and files never announced. @@ -375,24 +517,18 @@ class DirectoryIndexer: absent has nothing to compare against, and comparing anyway is exactly the mistake this module exists to avoid. """ - loop = asyncio.get_event_loop() changed = False + loop = asyncio.get_event_loop() for root in self.roots: if not root.available: continue try: - on_disk = {p.resolve() for p in root.path.rglob("*") - if _is_indexable(p)} + on_disk, known = await loop.run_in_executor( + self._executor, self._sweep_scan_root, root) except OSError as e: log.warning("Cannot reconcile root %r: %s", root.name, e) continue - known: dict[Path, str] = {} - for entry in self._entries_under(root): - abs_path = self._entry_path(root, entry) - if abs_path: - known[abs_path] = entry.id - for missing in set(known) - on_disk: # Duplicate content is handled without a special case here: the # entry goes, and the add loop below re-indexes the surviving @@ -404,26 +540,46 @@ class DirectoryIndexer: log.info("Reconcile: %s is gone", missing) changed = True - for added in on_disk - set(known): - entry = await loop.run_in_executor( - self._executor, _scan_file, root, added) - if not entry: - continue - # The index is keyed by **content hash**, so two identical files - # at two paths are one entry and the path comparison above - # cannot see the second. Adding it anyway rewrites that entry's - # path every cycle, bumps the version, and pushes an index - # update to every connected peer once a minute — for ever. - # Measured on a live node: `clip.mp4` present at the root and in - # uploads/ with the same bytes. - if self._index.get_entry(entry.id) is not None: - log.debug("Reconcile: %s duplicates content already indexed " - "as %s — leaving the index alone", - added, entry.id[:8]) - continue - self._index.add_entry(entry) - log.info("Reconcile: %s appeared (missed event)", added) - changed = True + added_paths = on_disk - set(known) + if not added_paths: + continue + + added_sized: list[tuple[Path, int]] = [] + for p in added_paths: + try: + added_sized.append((p, p.stat().st_size)) + except OSError: + added_sized.append((p, 0)) + + self.progress.scanning = True + self.progress.scanned_bytes = 0 + self.progress.total_bytes = sum(size for _, size in added_sized) + self.progress.current_dir = "" + try: + for added, size in added_sized: + self.progress.current_dir = added.parent.name + entry = await self._hash_or_cached(root, added) + self.progress.scanned_bytes += size + if not entry: + continue + # The index is keyed by **content hash**, so two identical + # files at two paths are one entry and the path comparison + # above cannot see the second. Adding it anyway rewrites + # that entry's path every cycle, bumps the version, and + # pushes an index update to every connected peer once a + # minute — for ever. Measured on a live node: `clip.mp4` + # present at the root and in uploads/ with the same bytes. + if self._index.get_entry(entry.id) is not None: + log.debug("Reconcile: %s duplicates content already " + "indexed as %s — leaving the index alone", + added, entry.id[:8]) + continue + self._index.add_entry(entry) + log.info("Reconcile: %s appeared (missed event)", added) + changed = True + finally: + self.progress.scanning = False + self.progress.current_dir = "" return changed def _entries_under(self, root: Root) -> list[IndexEntry]: @@ -444,6 +600,21 @@ class DirectoryIndexer: except OSError: return None + def _sweep_scan_root(self, root: Root) -> tuple[set[Path], dict[Path, str]]: + """ + Blocking: the two disk-touching pieces of one reconcile pass for a + root, bundled so both run together in the executor rather than in + the asyncio loop — the tree walk, and resolving the real path of + every entry already known under this root (one syscall each). + """ + on_disk = {p.resolve() for p in root.path.rglob("*") if _is_indexable(p)} + known: dict[Path, str] = {} + for entry in self._entries_under(root): + abs_path = self._entry_path(root, entry) + if abs_path: + known[abs_path] = entry.id + return on_disk, known + def _restart_observer(self) -> None: """Re-schedule watches after roots appeared or disappeared.""" if self._observer: @@ -454,8 +625,6 @@ class DirectoryIndexer: # ── Internal update ─────────────────────────────────────────────────────── - _DEBOUNCE_SECS = 2.0 - def _schedule_update(self, file_path: Path, deleted: bool = False) -> None: """Called from watchdog thread — schedule debounced async update.""" if not self._loop: @@ -472,7 +641,7 @@ class DirectoryIndexer: self._pending_timers.pop(key, None) asyncio.ensure_future(self._update_entry(file_path, deleted)) - self._pending_timers[key] = self._loop.call_later(self._DEBOUNCE_SECS, fire) + self._pending_timers[key] = self._loop.call_later(self.debounce_secs, fire) def _remove_by_path(self, root: Root, file_path: Path) -> None: """Remove any existing entries that point at this file.""" @@ -511,9 +680,7 @@ class DirectoryIndexer: self._remove_by_path(root, file_path) if not deleted: - loop = asyncio.get_event_loop() - entry = await loop.run_in_executor( - self._executor, _scan_file, root, file_path) + entry = await self._hash_or_cached(root, file_path) if entry: self._index.add_entry(entry) log.debug("Indexed: %s (%s, %d bytes)", |