diff options
Diffstat (limited to 'packages/meshbay-node/src')
6 files changed, 106 insertions, 7 deletions
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: |