aboutsummaryrefslogtreecommitdiffstats
path: root/packages
diff options
context:
space:
mode:
Diffstat (limited to 'packages')
-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
-rw-r--r--packages/meshbay-node/tests/test_root_eject.py128
7 files changed, 234 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:
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()