""" Directory indexer — watches a directory and maintains a GroupIndex. Uses watchdog for filesystem events. On any change (create/modify/delete/move), the affected file is re-scanned and the GroupIndex is updated. File metadata (blake3 hash, size, type, duration) is computed on first scan. Heavy operations (hashing large files) run in a thread pool to avoid blocking. """ import asyncio import logging import mimetypes import time from concurrent.futures import ThreadPoolExecutor from pathlib import Path from typing import Callable, Awaitable import blake3 from cryptography.hazmat.primitives.asymmetric.ed25519 import Ed25519PrivateKey from watchdog.events import FileSystemEvent, FileSystemEventHandler from watchdog.observers import Observer from meshbay_common.protocol import IndexEntry from meshbay_node.indexer.group_index import GroupIndex log = logging.getLogger(__name__) # File types we include in the index (skip hidden files, temp files, etc.) EXCLUDED_PREFIXES = (".", "~", "#") EXCLUDED_SUFFIXES = (".tmp", ".part", ".crdownload", ".download") MEDIA_EXTENSIONS = { "video": {".mp4", ".mkv", ".avi", ".mov", ".wmv", ".flv", ".webm", ".m4v"}, "audio": {".mp3", ".flac", ".ogg", ".wav", ".aac", ".m4a", ".opus"}, "image": {".jpg", ".jpeg", ".png", ".gif", ".webp", ".svg", ".bmp", ".tiff"}, "document": {".pdf", ".epub", ".mobi", ".txt", ".md", ".docx", ".odt"}, "archive": {".zip", ".tar", ".gz", ".bz2", ".xz", ".7z", ".rar"}, } def _detect_type(path: Path) -> str: suffix = path.suffix.lower() for ftype, exts in MEDIA_EXTENSIONS.items(): if suffix in exts: return ftype return "other" def _is_indexable(path: Path) -> bool: if not path.is_file(): return False name = path.name return ( not any(name.startswith(p) for p in EXCLUDED_PREFIXES) and not any(name.endswith(s) for s in EXCLUDED_SUFFIXES) ) _HASH_CHUNK = 8 * 1024 * 1024 # 8 MB streaming hash chunks def _scan_file(root: Path, 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.) don't require loading the whole file into memory.""" if not _is_indexable(file_path): return None try: stat = file_path.stat() hasher = blake3.blake3() with open(file_path, "rb") as f: while chunk := f.read(_HASH_CHUNK): hasher.update(chunk) file_id = hasher.hexdigest() rel_path = str(file_path.parent.relative_to(root)) if rel_path == ".": rel_path = "" return IndexEntry( id=file_id, name=file_path.name, path=rel_path, size=stat.st_size, type=_detect_type(file_path), added_at=int(stat.st_mtime), ) except (OSError, PermissionError) as e: log.warning("Cannot index %s: %s", file_path, e) return None class DirectoryIndexer: """ Watches a directory and keeps a GroupIndex up to date. Usage: indexer = DirectoryIndexer( root=Path("/home/user/shared"), group_id="my-group", sk_node=sk, gek=gek_bytes, on_change=async_callback, ) await indexer.start() # ... later await indexer.stop() """ def __init__( self, root: Path, group_id: str, sk_node: Ed25519PrivateKey, gek: bytes | None, on_change: Callable[["DirectoryIndexer"], Awaitable[None]] | None = None, ): self.root = root.resolve() self.group_id = group_id self.sk_node = sk_node self.gek = gek self.on_change = on_change self._index = GroupIndex(group_id=group_id, sk_node=sk_node, gek=gek) self._executor = ThreadPoolExecutor(max_workers=2, thread_name_prefix="indexer") self._observer: Observer | None = None self._loop: asyncio.AbstractEventLoop | None = None @property def index(self) -> GroupIndex: return self._index # ── Initial scan ────────────────────────────────────────────────────────── async def initial_scan(self) -> None: """Scan the entire directory tree. Run once at startup.""" log.info("Scanning %s ...", self.root) loop = asyncio.get_event_loop() files = [p for p in self.root.rglob("*") if p.is_file()] count = 0 for file_path in files: entry = await loop.run_in_executor( self._executor, _scan_file, self.root, file_path) if entry: self._index.add_entry(entry) count += 1 self._index.version = int(time.time()) log.info("Initial scan complete: %d files indexed", count) # ── Watchdog integration ────────────────────────────────────────────────── async def start(self) -> None: """Start initial scan + filesystem watcher.""" self._loop = asyncio.get_event_loop() await self.initial_scan() handler = _WatchdogHandler(self) self._observer = Observer() self._observer.schedule(handler, str(self.root), recursive=True) self._observer.start() log.info("Watching %s for changes", self.root) async def stop(self) -> None: """Stop the filesystem watcher.""" if self._observer: self._observer.stop() self._observer.join() self._observer = None self._executor.shutdown(wait=False) log.info("Indexer stopped") # ── 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: return key = str(file_path.resolve()) self._loop.call_soon_threadsafe( self._debounce, key, file_path, deleted) def _debounce(self, key: str, file_path: Path, deleted: bool) -> None: if not hasattr(self, "_pending_timers"): self._pending_timers: dict[str, asyncio.TimerHandle] = {} old = self._pending_timers.pop(key, None) if old: old.cancel() handle = self._loop.call_later( self._DEBOUNCE_SECS, lambda: asyncio.ensure_future(self._update_entry(file_path, deleted)), ) self._pending_timers[key] = handle def _remove_by_path(self, file_path: Path) -> None: """Remove any existing entries that match this file's path + name.""" resolved = file_path.resolve() to_remove = [ e.id for e in self._index.entries if (self.root / e.path / e.name).resolve() == resolved ] for fid in to_remove: self._index.remove_entry(fid) async def _update_entry(self, file_path: Path, deleted: bool) -> None: self._remove_by_path(file_path) if not deleted: loop = asyncio.get_event_loop() entry = await loop.run_in_executor( self._executor, _scan_file, self.root, file_path) if entry: self._index.add_entry(entry) log.debug("Indexed: %s (%s, %d bytes)", file_path.name, entry.id[:8], entry.size) self._index.version = int(time.time()) if self.on_change: await self.on_change(self) class _WatchdogHandler(FileSystemEventHandler): def __init__(self, indexer: DirectoryIndexer): self._indexer = indexer def on_created(self, event: FileSystemEvent): if not event.is_directory: self._indexer._schedule_update(Path(event.src_path)) def on_modified(self, event: FileSystemEvent): if not event.is_directory: self._indexer._schedule_update(Path(event.src_path)) def on_deleted(self, event: FileSystemEvent): if not event.is_directory: self._indexer._schedule_update(Path(event.src_path), deleted=True) def on_moved(self, event: FileSystemEvent): if not event.is_directory: self._indexer._schedule_update(Path(event.src_path), deleted=True) self._indexer._schedule_update(Path(event.dest_path))