diff options
| author | Christophe Besson <cbesson@gmail.com> | 2026-10-07 12:48:34 +0200 |
|---|---|---|
| committer | Christophe Besson <cbesson@gmail.com> | 2026-10-07 12:48:34 +0200 |
| commit | 3d5a168ee592c28c696e445c623c9f8d2996b715 (patch) | |
| tree | adf258cabc4f86e7aa03b9feff7f4eec8c95d1ab | |
| parent | 00caeb3e10bfbf86892e7f31ad81f91ce8e50da7 (diff) | |
| download | meshbay-3d5a168ee592c28c696e445c623c9f8d2996b715.tar.gz | |
fix(node): drop a directory moved out of the root from the index
Watchdog reports such a move as one "directory deleted" event and nothing
for the files, which the indexer ignored until the next reconcile. The
freeze rules still apply: root live, directory gone, parent present.
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
| -rw-r--r-- | docs/MESHBAY_DESIGN.md | 1 | ||||
| -rw-r--r-- | packages/meshbay-node/src/meshbay_node/indexer/indexer.py | 62 | ||||
| -rw-r--r-- | packages/meshbay-node/tests/test_root_availability.py | 97 |
3 files changed, 152 insertions, 8 deletions
diff --git a/docs/MESHBAY_DESIGN.md b/docs/MESHBAY_DESIGN.md index 1730641..cb2218e 100644 --- a/docs/MESHBAY_DESIGN.md +++ b/docs/MESHBAY_DESIGN.md @@ -4055,6 +4055,7 @@ process runs it — `systemctl --user` on Linux, Task Scheduler on Windows. | **Any request sent during a reconnect's own connect() fails at once** | `_sendAndWait` skips `waitForReconnect` while `_inReconnectAttempt` is set, and that flag is transport-wide, not the handshake's: a chat send or an index request made in that window throws "DataChannel not open (state: connecting)". File chunks now wait the reconnect out in `_fetchChunkResilient`; nothing else does | | **A failed download leaves its in-flight chunks rejecting unhandled** | When `pipelinedDownload` throws, the other promises of its window are abandoned and each prints "Uncaught (in promise)" in the console. Noise, but in the log a dropped connection is read from | | **A queued transfer logs "no answer" every minute** | The lease watchdog re-asks after 60 s of silence whatever the last answer was, so a lease the node has answered `queued` and that has not moved prints "no answer for transfer ... asking again" once a minute for as long as it waits, which reads as a fault | +| **"Prune stale entries" does not touch the index** | The Node page's maintenance button runs `prune_index_cache`, which drops rows of the hash cache and nothing else. An operator who sees a deleted file still listed reaches for it, and nothing changes: only the reconcile removes an index entry, and nothing lets the operator ask for one | --- diff --git a/packages/meshbay-node/src/meshbay_node/indexer/indexer.py b/packages/meshbay-node/src/meshbay_node/indexer/indexer.py index 5e139bb..732b499 100644 --- a/packages/meshbay-node/src/meshbay_node/indexer/indexer.py +++ b/packages/meshbay-node/src/meshbay_node/indexer/indexer.py @@ -1165,14 +1165,16 @@ class DirectoryIndexer: # ── Internal update ─────────────────────────────────────────────────────── - def _schedule_update(self, file_path: Path, deleted: bool = False) -> None: + def _schedule_update(self, file_path: Path, deleted: bool = False, + directory: bool = False) -> None: """Called from watchdog thread — schedule debounced async update.""" if not self._loop: return key = str(file_path) - self._loop.call_soon_threadsafe(self._debounce, key, file_path, deleted) + self._loop.call_soon_threadsafe(self._debounce, key, file_path, deleted, directory) - def _debounce(self, key: str, file_path: Path, deleted: bool) -> None: + def _debounce(self, key: str, file_path: Path, deleted: bool, + directory: bool = False) -> None: old = self._pending_timers.pop(key, None) if old: old.cancel() @@ -1204,7 +1206,7 @@ class DirectoryIndexer: def fire() -> None: self._pending_timers.pop(key, None) - spawn(self._update_entry(file_path, deleted)) + spawn(self._update_entry(file_path, deleted, directory)) self._pending_timers[key] = self._loop.call_later(self.debounce_secs, fire) @@ -1218,7 +1220,36 @@ class DirectoryIndexer: if self._entry_path(root, entry) == resolved: self._index.remove_entry(entry.id) - async def _update_entry(self, file_path: Path, deleted: bool) -> None: + def _entries_in_dir(self, root: Root, dir_path: Path) -> list[str]: + """Ids of the entries at or below `dir_path`, matched on the index path + alone: the directory is gone, so there is nothing left to resolve.""" + try: + rel = dir_path.relative_to(root.path).as_posix() + except ValueError: + return [] + if rel in ("", "."): + return [] + vdir = fold(f"{root.name}/{rel}") + return [e.id for e in self._entries_under(root) + if fold(e.path) == vdir or fold(e.path).startswith(vdir + "/")] + + @staticmethod + def _dir_left_the_root(dir_path: Path) -> bool: + """ + Blocking. The directory is gone and its parent is not. + + The parent is what tells a directory that was deleted, or moved out of + the root, from a volume that went away: a vanishing volume takes the + parent with it. A parent deleted as well reports its own event, which + covers this directory too. + """ + try: + return not dir_path.exists() and dir_path.parent.is_dir() + except OSError: + return False + + async def _update_entry(self, file_path: Path, deleted: bool, + directory: bool = False) -> None: try: root = self._root_for(file_path) if root is None: @@ -1243,7 +1274,19 @@ class DirectoryIndexer: if not root.available: return - self._remove_by_path(root, file_path) + if directory: + if file_path == root.path or not await off_disk( + self.roots, self._dir_left_the_root, file_path): + return + gone = self._entries_in_dir(root, file_path) + if not gone: + return + for entry_id in gone: + self._index.remove_entry(entry_id) + log.info("Watchdog: %s left the root, %d entries removed", + file_path, len(gone)) + else: + self._remove_by_path(root, file_path) if not deleted: entry = await self._hash_or_cached(root, file_path) @@ -1288,8 +1331,11 @@ class _WatchdogHandler(FileSystemEventHandler): 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) + # A directory moved out of the root arrives as this one event, with + # nothing for the files it held: without it they stayed indexed until + # the next reconcile, up to RECONCILE_BACKOFF_CAP later. + self._indexer._schedule_update(Path(event.src_path), deleted=True, + directory=event.is_directory) def on_moved(self, event: FileSystemEvent): if not event.is_directory: diff --git a/packages/meshbay-node/tests/test_root_availability.py b/packages/meshbay-node/tests/test_root_availability.py index 82bf6bd..e92ec95 100644 --- a/packages/meshbay-node/tests/test_root_availability.py +++ b/packages/meshbay-node/tests/test_root_availability.py @@ -13,6 +13,7 @@ an indexer that treats a vanished root as a set of deletions, which is what the straightforward implementation does. """ +import asyncio import os from pathlib import Path @@ -305,3 +306,99 @@ async def test_deleting_one_copy_keeps_the_other_listed(tmp_path): assert len(idx.index.entries) == 1, "the surviving copy was delisted" assert idx.index.entries[0].path == survivor + + +# ── A directory that leaves the root ───────────────────────────────────────── +# +# Moving a folder out of the root arrives from the watcher as one "directory +# deleted" event, with nothing for the files it held. Found live: a season moved +# out of a shared folder stayed listed, every episode "File not on disk", until +# the reconcile backstop came round. These are the same freeze rules as above, +# applied to that event. + +def _season(root: Path, name: str, *episodes: str) -> Path: + d = root / "Show" / name + d.mkdir(parents=True) + for ep in episodes: + (d / ep).write_bytes(ep.encode()) + return d + + +async def test_a_directory_moved_out_of_a_live_root_is_removed(tmp_path): + shows = tmp_path / "Shows" + low = _season(shows, "S3.LQ", "e1.mp4", "e2.mp4") + _season(shows, "S3", "e1.hd.mkv") + idx = await _indexer(_roots(shows)) + + low.rename(tmp_path / "S3.LQ") + await idx._update_entry(low, deleted=True, directory=True) + + # S3 shares a prefix with S3.LQ and must survive it. + assert _names(idx) == {"e1.hd.mkv"} + + +async def test_a_directory_event_from_a_vanished_root_freezes(tmp_path): + shows = tmp_path / "Shows" + season = _season(shows, "S1", "e1.mp4", "e2.mp4") + idx = await _indexer(_roots(shows)) + + os.rename(shows, tmp_path / "elsewhere") + await idx._update_entry(season, deleted=True, directory=True) + + assert _names(idx) == {"e1.mp4", "e2.mp4"} + assert idx.roots.roots[0].available is False + + +async def test_a_directory_whose_parent_is_gone_waits_for_the_parent(tmp_path): + """The parent going too is what a volume vanishing looks like below the + root. The parent's own event, when it comes, is the one that acts.""" + shows = tmp_path / "Shows" + season = _season(shows, "S1", "e1.mp4") + _season(shows, "S2", "e2.mp4") + idx = await _indexer(_roots(shows)) + + os.rename(shows / "Show", tmp_path / "Show") + await idx._update_entry(season, deleted=True, directory=True) + assert _names(idx) == {"e1.mp4", "e2.mp4"} + + await idx._update_entry(shows / "Show", deleted=True, directory=True) + assert _names(idx) == set() + + +async def test_a_directory_that_is_back_is_left_alone(tmp_path): + shows = tmp_path / "Shows" + season = _season(shows, "S1", "e1.mp4") + idx = await _indexer(_roots(shows)) + + # Moved out and back within the debounce: the event is stale. + await idx._update_entry(season, deleted=True, directory=True) + assert _names(idx) == {"e1.mp4"} + + +async def test_the_root_itself_is_never_removed_as_a_directory(tmp_path): + shows = tmp_path / "Shows" + _season(shows, "S1", "e1.mp4") + idx = await _indexer(_roots(shows)) + + await idx._update_entry(shows, deleted=True, directory=True) + assert _names(idx) == {"e1.mp4"} + + +async def test_the_watcher_reports_a_directory_moved_out(tmp_path): + """End to end, with the real observer: the event is the one watchdog emits.""" + shows = tmp_path / "Shows" + low = _season(shows, "S3.LQ", "e1.mp4") + _season(shows, "S3", "e1.hd.mkv") + idx = DirectoryIndexer(roots=_roots(shows), group_id="g" * 32, + sk_node=Ed25519PrivateKey.generate(), gek=None, + debounce_secs=0.05) + await idx.start() + try: + low.rename(tmp_path / "S3.LQ") + for _ in range(100): + if _names(idx) == {"e1.hd.mkv"}: + break + await asyncio.sleep(0.05) + assert _names(idx) == {"e1.hd.mkv"} + finally: + await idx.stop() |