aboutsummaryrefslogtreecommitdiffstats
diff options
context:
space:
mode:
authorChristophe Besson <cbesson@gmail.com>2026-10-07 12:48:34 +0200
committerChristophe Besson <cbesson@gmail.com>2026-10-07 12:48:34 +0200
commit3d5a168ee592c28c696e445c623c9f8d2996b715 (patch)
treeadf258cabc4f86e7aa03b9feff7f4eec8c95d1ab
parent00caeb3e10bfbf86892e7f31ad81f91ce8e50da7 (diff)
downloadmeshbay-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.md1
-rw-r--r--packages/meshbay-node/src/meshbay_node/indexer/indexer.py62
-rw-r--r--packages/meshbay-node/tests/test_root_availability.py97
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()