aboutsummaryrefslogtreecommitdiffstats
path: root/packages/meshbay-node/src
diff options
context:
space:
mode:
Diffstat (limited to 'packages/meshbay-node/src')
-rw-r--r--packages/meshbay-node/src/meshbay_node/daemon.py5
-rw-r--r--packages/meshbay-node/src/meshbay_node/indexer/cache.py14
-rw-r--r--packages/meshbay-node/src/meshbay_node/indexer/indexer.py65
-rw-r--r--packages/meshbay-node/src/meshbay_node/ops/roots.py3
-rw-r--r--packages/meshbay-node/src/meshbay_node/roots.py6
-rw-r--r--packages/meshbay-node/src/meshbay_node/roster.py20
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: