From 6cdc6016d72dcfb7530ac38a8fa92232418ac305 Mon Sep 17 00:00:00 2001 From: Christophe Besson Date: Mon, 5 Oct 2026 11:10:50 +0200 Subject: fix(node): plug an auto-ejected removable root back once its files return A node started with the desktop session runs before the session has mounted its USB drives. The safety net then auto-ejected every removable root and persisted it exactly like an operator's eject, so after each reboot those roots stayed ejected until someone plugged them by hand (seen on a node whose /media drives were mounted a minute after it started). An auto-eject is now stored as such ("auto" in roster.db). At startup and at every reconcile, an auto-ejected root whose path is readable again is checked against a few files the hash cache knows under it, at the same path with the same size and mtime; one found and the root is plugged back and rescanned. An empty mount point or another drive in its place is not recognised and stays ejected. An operator's eject is never undone automatically. Co-Authored-By: Claude Opus 5.5 --- packages/meshbay-node/src/meshbay_node/daemon.py | 5 +- .../meshbay-node/src/meshbay_node/indexer/cache.py | 14 +++ .../src/meshbay_node/indexer/indexer.py | 65 ++++++++++- .../meshbay-node/src/meshbay_node/ops/roots.py | 3 +- packages/meshbay-node/src/meshbay_node/roots.py | 6 + packages/meshbay-node/src/meshbay_node/roster.py | 20 +++- packages/meshbay-node/tests/test_root_eject.py | 128 +++++++++++++++++++++ 7 files changed, 234 insertions(+), 7 deletions(-) (limited to 'packages/meshbay-node') diff --git a/packages/meshbay-node/src/meshbay_node/daemon.py b/packages/meshbay-node/src/meshbay_node/daemon.py index 1777bc3..d6293ee 100644 --- a/packages/meshbay-node/src/meshbay_node/daemon.py +++ b/packages/meshbay-node/src/meshbay_node/daemon.py @@ -1259,11 +1259,13 @@ class NodeDaemon(EnrichmentMixin): specs = [asdict(r) for r in group_cfg.roots] if self._roster: ejected = await self._roster.ejected_roots(group_cfg.id) + auto = await self._roster.auto_ejected_roots(group_cfg.id) if ejected: for spec in specs: name = spec.get("name") or Path(spec.get("path", "")).name if fold(name) in ejected: spec["ejected"] = True + spec["ejected_auto"] = fold(name) in auto return RootSet.build(specs) # Every application that keeps directories. This is the one list, and it @@ -1305,9 +1307,10 @@ class NodeDaemon(EnrichmentMixin): """`on_root_ejected` bound to one group, for that group's indexer.""" async def persist(root_name: str, ejected: bool) -> None: if self._roster: + # The indexer only reports the safety net's own changes. await self._roster.set_root_ejected( group_id, root_name, ejected, - set_by=self._state.get("node_user_id", "")) + set_by=self._state.get("node_user_id", ""), auto=ejected) return persist async def _on_index_change(self, indexer: DirectoryIndexer) -> None: diff --git a/packages/meshbay-node/src/meshbay_node/indexer/cache.py b/packages/meshbay-node/src/meshbay_node/indexer/cache.py index 6c87f40..5e328ef 100644 --- a/packages/meshbay-node/src/meshbay_node/indexer/cache.py +++ b/packages/meshbay-node/src/meshbay_node/indexer/cache.py @@ -23,6 +23,7 @@ only daemon.py's wiring changed. """ import logging +import os import time from dataclasses import dataclass from pathlib import Path @@ -191,6 +192,19 @@ class IndexCache: row = await cur.fetchone() return row[0] if row else 0 + async def sample_under(self, directory: str, limit: int) -> list[tuple[str, int, float]]: + """Up to `limit` cached files under `directory`, as (path, size, mtime). + + A prefix compared with `substr`, not `LIKE`: `_` and `%` are wildcards + there, and both are ordinary in a folder name. + """ + prefix = directory.rstrip("/\\") + os.sep + async with self._db.execute( + "SELECT path, size, mtime FROM files WHERE substr(path, 1, ?) = ? LIMIT ?", + (len(prefix), prefix, limit)) as cur: + rows = await cur.fetchall() + return [(row[0], row[1], row[2]) for row in rows] + async def all_paths(self) -> list[str]: """Every cached path, for a caller that decides staleness itself — this cache has no notion of which paths are still claimed by a diff --git a/packages/meshbay-node/src/meshbay_node/indexer/indexer.py b/packages/meshbay-node/src/meshbay_node/indexer/indexer.py index f0505f7..5e139bb 100644 --- a/packages/meshbay-node/src/meshbay_node/indexer/indexer.py +++ b/packages/meshbay-node/src/meshbay_node/indexer/indexer.py @@ -389,6 +389,7 @@ class DirectoryIndexer: async def _initial_scan(self) -> None: await off_disk(self.roots, self.roots.refresh_availability) + await self._plug_back_recognised() total = 0 waiting = [r.name for r in self.roots if r.available] self._queue(waiting) @@ -825,6 +826,8 @@ class DirectoryIndexer: except Exception: log.exception("Could not persist the auto-eject of root %r", name) + changed += [(root, True) for root in await self._plug_back_recognised()] + for root, available in changed: if available: log.info("Root %r is back — rescanning", root.name) @@ -1032,6 +1035,65 @@ class DirectoryIndexer: self._observer = None self._start_observer() + # How many of a root's known files are looked for before it is plugged back + # automatically. One found is enough: another drive, or an empty mount + # point, holds none of them at the same path with the same size and mtime. + RECOGNISE_SAMPLE = 5 + + async def _plug_back_recognised(self) -> list[Root]: + """ + Un-eject the roots the safety net ejected, once they hold their own files + again: a drive not mounted yet when the node started, or unplugged and + plugged back. An operator's eject is never undone here. + """ + back = [] + for root in [r for r in self.roots if r.ejected and r.auto]: + if not await off_disk(self.roots, self._recognises, root, + await self._known_files(root)): + continue + root.ejected = root.auto = False + root.available = True + log.info("Root %r is back with its files — plugged automatically", root.name) + if self.on_root_ejected: + try: + await self.on_root_ejected(root.name, False) + except Exception: + log.exception("Could not persist the return of root %r", root.name) + back.append(root) + return back + + async def _known_files(self, root: Root) -> list[tuple[str, int, float | None]]: + """Files this root is known to hold, as (path, size, mtime or None). + + The hash cache first: it survives a restart, when an ejected root has no + entry in the index at all. The index is the fallback for a node without + a cache. + """ + if self._cache is not None: + known = await self._cache.sample_under(str(root.path), self.RECOGNISE_SAMPLE) + if known: + return known + return [(str(path), e.size, None) + for e in self._entries_under(root)[:self.RECOGNISE_SAMPLE] + if (path := self._entry_path(root, e)) is not None] + + @staticmethod + def _recognises(root: Root, known: list[tuple[str, int, float | None]]) -> bool: + """Blocking: is this the root's own content, readable again?""" + if not root.is_live(): + return False + if not known: + # Nothing known under it, so nothing a wrong disk could pass for. + return True + for path, size, mtime in known: + try: + st = Path(path).stat() + except OSError: + continue + if st.st_size == size and (mtime is None or st.st_mtime == mtime): + return True + return False + def eject_root(self, root_name: str) -> None: """Stop watching a root without touching its entries.""" from meshbay_common.paths import fold @@ -1039,6 +1101,7 @@ class DirectoryIndexer: for root in self.roots: if fold(root.name) == target: root.ejected = True + root.auto = False root.available = False frozen = len(self._entries_under(root)) log.info("Root %r ejected — %d entries frozen", root.name, frozen) @@ -1058,7 +1121,7 @@ class DirectoryIndexer: break if root is None: return - root.ejected = False + root.ejected = root.auto = False root.available = await off_disk(self.roots, root.is_live) if not root.available: await self._finish_plug(None) diff --git a/packages/meshbay-node/src/meshbay_node/ops/roots.py b/packages/meshbay-node/src/meshbay_node/ops/roots.py index 9dbade4..b976e8f 100644 --- a/packages/meshbay-node/src/meshbay_node/ops/roots.py +++ b/packages/meshbay-node/src/meshbay_node/ops/roots.py @@ -240,6 +240,7 @@ async def eject_root(state: dict, group_id: str, root_name: str) -> dict: if indexer: indexer.eject_root(root_name) root.ejected = True + root.auto = False root.available = False await _roster(state).set_root_ejected( @@ -288,7 +289,7 @@ async def plug_root(state: dict, group_id: str, root_name: str) -> dict: indexer = state.get("indexers", {}).get(group_id) if indexer: await indexer.plug_root(root_name) - root.ejected = False + root.ejected = root.auto = False root.available = await off_disk(roots, root.is_live) log.info("Root plugged: %s in group %s", root_name, group_id[:8]) diff --git a/packages/meshbay-node/src/meshbay_node/roots.py b/packages/meshbay-node/src/meshbay_node/roots.py index 2de0708..2b48292 100644 --- a/packages/meshbay-node/src/meshbay_node/roots.py +++ b/packages/meshbay-node/src/meshbay_node/roots.py @@ -178,6 +178,10 @@ class Root: writable: bool = False removable: bool = False ejected: bool = False + # Ejected by the safety net in `refresh_availability`, not by the operator: + # such a root comes back on its own once its files are there again + # (`DirectoryIndexer._reconcile`). An operator's eject never does. + auto: bool = False available: bool = True @property @@ -323,6 +327,7 @@ class RootSet: writable=writable, removable=bool(spec.get("removable", False)), ejected=bool(spec.get("ejected", False)), + auto=bool(spec.get("ejected_auto", False)), available=not bool(spec.get("ejected", False))) _refuse_nesting(root, roots) roots.append(root) @@ -441,6 +446,7 @@ class RootSet: live = root.is_live() if not live and root.removable: root.ejected = True + root.auto = True # Recorded for the caller to persist. A flag that only lives # in memory would be forgotten on the next restart, and the # rescan that followed would read an empty mount point as an diff --git a/packages/meshbay-node/src/meshbay_node/roster.py b/packages/meshbay-node/src/meshbay_node/roster.py index 0116aaa..eb8c031 100644 --- a/packages/meshbay-node/src/meshbay_node/roster.py +++ b/packages/meshbay-node/src/meshbay_node/roster.py @@ -686,9 +686,11 @@ class Roster: return cls.SETTING_ROOT_EJECTED_PREFIX + fold(root_name) async def set_root_ejected(self, group_id: str, root_name: str, - ejected: bool, set_by: str = "") -> None: - await self.set_setting(group_id, self.root_ejected_key(root_name), - "1" if ejected else "0", set_by) + ejected: bool, set_by: str = "", + auto: bool = False) -> None: + """`auto`: ejected by the safety net, not by the operator.""" + value = ("auto" if auto else "1") if ejected else "0" + await self.set_setting(group_id, self.root_ejected_key(root_name), value, set_by) async def ejected_roots(self, group_id: str) -> set[str]: """ @@ -705,7 +707,17 @@ class Roster: (group_id,)) as cur: rows = await cur.fetchall() return {r["key"][len(prefix):] for r in rows - if r["key"].startswith(prefix) and r["value"] == "1"} + if r["key"].startswith(prefix) and r["value"] in ("1", "auto")} + + async def auto_ejected_roots(self, group_id: str) -> set[str]: + """The folded names of the roots the safety net ejected, a subset of the above.""" + prefix = self.SETTING_ROOT_EJECTED_PREFIX + async with self._db.execute( + "SELECT key, value FROM group_settings WHERE group_id = ?", + (group_id,)) as cur: + rows = await cur.fetchall() + return {r["key"][len(prefix):] for r in rows + if r["key"].startswith(prefix) and r["value"] == "auto"} async def get_setting(self, group_id: str, key: str, default: str | None = None) -> str | None: diff --git a/packages/meshbay-node/tests/test_root_eject.py b/packages/meshbay-node/tests/test_root_eject.py index d73e71c..b6b50aa 100644 --- a/packages/meshbay-node/tests/test_root_eject.py +++ b/packages/meshbay-node/tests/test_root_eject.py @@ -264,3 +264,131 @@ async def test_the_ejected_key_is_case_folded(tmp_path): assert Roster.root_ejected_key("Films") == Roster.root_ejected_key("FILMS") finally: await roster.close() + + +# ── The safety net's eject undoes itself; the operator's does not ─────────── + +def _vanish(films: Path) -> None: + for f in films.iterdir(): + f.unlink() + films.rmdir() + + +async def test_an_auto_ejected_root_comes_back_with_its_own_files(tmp_path): + """ + The drive that was not mounted yet when the node started (found on a node + started with the session, its USB drives mounted a minute later): the + safety net ejected it, and once the same files are readable at the same + place it is plugged back without anyone having to. + """ + films = tmp_path / "Films" + films.mkdir() + (films / "a.mkv").write_bytes(b"a") + seen: list[tuple[str, bool]] = [] + + async def record(name: str, ejected: bool) -> None: + seen.append((name, ejected)) + + roots = _roots(films) + idx = await _indexer(roots, on_root_ejected=record) + hidden = tmp_path / "unmounted" + films.rename(hidden) + await idx.reconcile() + assert roots.roots[0].ejected is True + + hidden.rename(films) + await idx.reconcile() + assert roots.roots[0].ejected is False + assert roots.roots[0].available is True + assert _names(idx) == {"a.mkv"} + assert seen == [("Films", True), ("Films", False)] + + +@pytest.mark.parametrize("what_came_back", ["empty", "another drive"]) +async def test_an_auto_ejected_root_stays_out_when_its_files_are_not_there( + tmp_path, what_came_back): + """ + What the safety net exists for: an empty mount point, or another drive + mounted at the same place, is not the library. Plugging it back would + rescan it, and the rescan would read the library as erased. + """ + films = tmp_path / "Films" + films.mkdir() + (films / "a.mkv").write_bytes(b"a") + roots = _roots(films) + idx = await _indexer(roots) + _vanish(films) + await idx.reconcile() + + films.mkdir() + if what_came_back == "another drive": + (films / "other.mkv").write_bytes(b"something else") + await idx.reconcile() + assert roots.roots[0].ejected is True + assert _names(idx) == {"a.mkv"}, "the library was treated as erased" + + +async def test_an_operator_eject_is_never_undone_automatically(tmp_path): + films = tmp_path / "Films" + films.mkdir() + (films / "a.mkv").write_bytes(b"a") + roots = _roots(films) + idx = await _indexer(roots) + + idx.eject_root("Films") + await idx.reconcile() + assert roots.roots[0].ejected is True + + +async def test_after_a_restart_the_cache_recognises_the_drive(tmp_path): + """ + A restarted node has no index entry for an ejected root: the index is + rebuilt by scanning, and an ejected root is not scanned. What it does have is + the hash cache, with the size and mtime of every file it read there. + """ + import os + + from meshbay_node.indexer.cache import IndexCache + + films = tmp_path / "Films" + films.mkdir() + movie = films / "a.mkv" + movie.write_bytes(b"a") + st = movie.stat() + + def restarted_ejected() -> RootSet: + return RootSet.build([{"path": str(films), "removable": True, + "ejected": True, "ejected_auto": True}]) + + async with IndexCache(tmp_path / "cache.db") as cache: + await _indexer(_roots(films), cache=cache) + + # Another drive at the same place, with a file of the same name. + movie.write_bytes(b"another drive") + roots = restarted_ejected() + await _indexer(roots, cache=cache) + assert roots.roots[0].ejected is True + + # The drive itself. + movie.write_bytes(b"a") + os.utime(movie, ns=(st.st_atime_ns, st.st_mtime_ns)) + roots = restarted_ejected() + idx = await _indexer(roots, cache=cache) + assert roots.roots[0].ejected is False + assert _names(idx) == {"a.mkv"} + + +async def test_the_roster_keeps_an_auto_eject_apart(tmp_path): + roster = Roster(db_path=tmp_path / "roster.db") + await roster.open() + try: + await roster.set_root_ejected("g1", "Films", True, set_by="op", auto=True) + await roster.set_root_ejected("g1", "Music", True, set_by="op") + assert await roster.ejected_roots("g1") == {"films", "music"} + assert await roster.auto_ejected_roots("g1") == {"films"} + + # An operator eject of the same root replaces the safety net's. + await roster.set_root_ejected("g1", "Films", True, set_by="op") + assert await roster.auto_ejected_roots("g1") == set() + finally: + await roster.close() -- cgit v1.2.3