aboutsummaryrefslogtreecommitdiffstats
path: root/packages/meshbay-node
diff options
context:
space:
mode:
Diffstat (limited to 'packages/meshbay-node')
-rw-r--r--packages/meshbay-node/src/meshbay_node/indexer/indexer.py62
-rw-r--r--packages/meshbay-node/tests/test_root_availability.py97
2 files changed, 151 insertions, 8 deletions
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()