diff options
| -rw-r--r-- | docs/MESHBAY_DESIGN.md | 16 | ||||
| -rw-r--r-- | packages/meshbay-node/src/meshbay_node/daemon.py | 5 | ||||
| -rw-r--r-- | packages/meshbay-node/src/meshbay_node/indexer/cache.py | 14 | ||||
| -rw-r--r-- | packages/meshbay-node/src/meshbay_node/indexer/indexer.py | 65 | ||||
| -rw-r--r-- | packages/meshbay-node/src/meshbay_node/ops/roots.py | 3 | ||||
| -rw-r--r-- | packages/meshbay-node/src/meshbay_node/roots.py | 6 | ||||
| -rw-r--r-- | packages/meshbay-node/src/meshbay_node/roster.py | 20 | ||||
| -rw-r--r-- | packages/meshbay-node/tests/test_root_eject.py | 128 |
8 files changed, 247 insertions, 10 deletions
diff --git a/docs/MESHBAY_DESIGN.md b/docs/MESHBAY_DESIGN.md index 2cfe7a4..9c87372 100644 --- a/docs/MESHBAY_DESIGN.md +++ b/docs/MESHBAY_DESIGN.md @@ -1565,7 +1565,8 @@ the whole tree or presents an empty directory to the next scan. Both propagate a though the owner erased their library. So a root has two independent runtime states: -- **`ejected`** — operator-controlled, persisted in `roster.db`. +- **`ejected`** — set by the operator, or by the safety net below; persisted in + `roster.db`, with which of the two set it. - **`available`** — computed as `not ejected and is_live()`. This is what clients and the indexer see. @@ -1586,12 +1587,21 @@ let the following scan read the empty mount point as an erased library. It lives hand-written config must not be rewritten because a USB drive was unplugged. **Auto-eject is the safety net.** If a `removable` root's path disappears, the -availability sweep sets `ejected` as though the operator had clicked it, and -reports it so the daemon persists it. Nothing is deleted: index entries, cached +availability sweep sets `ejected` and reports it so the daemon persists it, marked +as the safety net's. Nothing is deleted: index entries, cached metadata, thumbnails, chat history referencing those files and app directory configurations all survive, the last flagged as temporarily invalid rather than wrong. +**The safety net's eject undoes itself; the operator's never does.** At startup +and at every reconcile, an auto-ejected root whose path is readable again is +checked against what the hash cache knows was under it: a few of those files, +at the same path with the same size and mtime. One found, and the root is +plugged back and rescanned, like a plug. None found, and it stays ejected: an +empty mount point or another drive mounted in its place is exactly what the eject +protects the index from. The case this serves is ordinary: a node started with +the session, before the desktop has mounted its USB drives. + ### 6.3 Indexing The index is **content-addressed**: `GroupIndex` is keyed by blake3, so the same 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() |