diff options
Diffstat (limited to 'packages/meshbay-node/src/meshbay_node')
9 files changed, 162 insertions, 11 deletions
diff --git a/packages/meshbay-node/src/meshbay_node/__init__.py b/packages/meshbay-node/src/meshbay_node/__init__.py index c3c18ff..c58bc8a 100644 --- a/packages/meshbay-node/src/meshbay_node/__init__.py +++ b/packages/meshbay-node/src/meshbay_node/__init__.py @@ -1,3 +1,3 @@ """MeshBay Node — local file host, streaming server, and group daemon.""" -__version__ = "0.17.0" +__version__ = "0.18.0" diff --git a/packages/meshbay-node/src/meshbay_node/cli/status.py b/packages/meshbay-node/src/meshbay_node/cli/status.py index 162bd46..ebc6c77 100644 --- a/packages/meshbay-node/src/meshbay_node/cli/status.py +++ b/packages/meshbay-node/src/meshbay_node/cli/status.py @@ -34,7 +34,8 @@ def status(args) -> None: live = None if live: - print(f"daemon running — {live.get('status')}") + state = live.get("status") + print("daemon running" + (f" — {state}" if state and state != "running" else "")) print(f"node_id {live.get('endpoint_hint') or '—'}") print(f"groups {live.get('group_count', 0)}" f" files {live.get('total_files', 0)}" 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/src/meshbay_node/ui/app.py b/packages/meshbay-node/src/meshbay_node/ui/app.py index db11a09..03db8b6 100644 --- a/packages/meshbay-node/src/meshbay_node/ui/app.py +++ b/packages/meshbay-node/src/meshbay_node/ui/app.py @@ -120,6 +120,10 @@ def create_ui_app(state: dict) -> FastAPI: @app.get("/api/status") async def api_status(): + """ + The daemon's state, and what it still needs: a linked key, a group, an operator, a group + key. + """ indexes = state.get("indexes", {}) total_files = sum(idx.count for idx in indexes.values()) groups_ctx = state.get("groups_ctx", {}) @@ -163,6 +167,7 @@ def create_ui_app(state: dict) -> FastAPI: @app.delete("/api/unlink") async def api_unlink(): + """Unlink the node's key from its hub account.""" hub = state.get("hub") if not hub: raise HTTPException(status_code=503, detail="Hub not connected") @@ -171,9 +176,13 @@ def create_ui_app(state: dict) -> FastAPI: @app.get("/api/groups") async def api_groups(): + """The groups this node hosts, with live status, and whether an operator is paired.""" return await _op(lambda: ops.list_groups(state)) @app.post("/api/groups/attach") async def attach_group(payload: dict): + """ + Host a group that exists on the hub: add it to node.toml with its first folder, then reload. + """ result = await _op(lambda: ops.attach_group( state, (payload.get("name") or "").strip(), @@ -188,6 +197,7 @@ def create_ui_app(state: dict) -> FastAPI: @app.post("/api/groups/detach") async def detach_group(payload: dict): + """Stop hosting a group: remove it from node.toml, then reload.""" result = await _op(lambda: ops.detach_group( state, (payload.get("name") or payload.get("group_id") or "").strip(), @@ -199,31 +209,41 @@ def create_ui_app(state: dict) -> FastAPI: @app.delete("/api/groups/{group_id}/files/{file_id}") async def delete_file(group_id: str, file_id: str): - """Milestone 14.11 — the last operator action that needed a browser.""" + """ + Delete a file from the group's folder on disk. + + Milestone 14.11 — the last operator action that needed a browser. + """ return await _op(lambda: ops.delete_file(state, group_id, file_id)) @app.get("/api/denylist") async def api_denylist(): + """What the node currently refuses.""" return await _op(lambda: ops.read_denylist(state)) @app.post("/api/denylist/clear") async def api_denylist_clear(subject: str = ""): + """Drop denylist entries: all of them, or one identifier.""" return await _op(lambda: ops.clear_denylist(state, subject=subject)) @app.get("/api/index-cache") async def api_index_cache_stats(): + """Size of the index cache.""" return await _op(lambda: ops.index_cache_stats(state)) @app.post("/api/index-cache/prune") async def api_index_cache_prune(): + """Drop index cache rows that no longer match a file on disk.""" return await _op(lambda: ops.prune_index_cache(state)) @app.post("/api/groups/{group_id}/video/rematch") async def api_video_rematch(group_id: str): + """Forget the automatic matches of the group's videos, so they are looked up again.""" return await _op(lambda: ops.rematch_video(state, group_id)) @app.get("/api/groups/{group_id}/files") async def api_group_files(group_id: str): + """The group's files, from its index.""" groups_ctx = state.get("groups_ctx", {}) ctx = groups_ctx.get(group_id) if not ctx: @@ -247,6 +267,7 @@ def create_ui_app(state: dict) -> FastAPI: @app.get("/api/peers") async def api_peers(): + """The connected peers.""" webrtc = state.get("webrtc") if not webrtc: return {"peers": []} @@ -275,6 +296,7 @@ def create_ui_app(state: dict) -> FastAPI: user_id: str | None = Query(default=None), event: str | None = Query(default=None), ): + """The audit log, filtered by time, account and event.""" audit = state.get("audit_store") if not audit: return {"entries": [], "offset": 0, "limit": limit, "has_more": False} @@ -316,66 +338,80 @@ def create_ui_app(state: dict) -> FastAPI: @app.post("/api/operator/pair") async def operator_pair(): + """A one-time code that pairs an application as this node's operator.""" return await _op(lambda: ops.pair_operator(state)) @app.get("/api/roster") async def api_roster(group_id: str = ""): + """The pinned identities, for one group or all.""" return await _op(lambda: ops.read_roster(state, group_id)) @app.post("/api/groups/{group_id}/invites") async def create_invite(group_id: str, username: str): + """An invitation code for one account, for this group.""" return await _op(lambda: ops.create_invite(state, group_id, username)) # Both halves, node and hub, for the CLI: an operator at the machine gets a # whole link, not a code without a ticket. @app.post("/api/groups/{group_id}/invite-links") async def create_link_invite(group_id: str, email: str = ""): + """A whole invitation link: the node's code, then the hub's ticket.""" return await _op(lambda: ops.create_link_invitation(state, group_id, email)) @app.delete("/api/groups/{group_id}/invite-links/{invite_id}") async def cancel_invite(group_id: str, invite_id: str): + """Take an invitation link back, on the node and on the hub.""" return await _op(lambda: ops.cancel_link_invitation(state, group_id, invite_id)) @app.get("/api/resolve") async def resolve_user(username: str): + """Map a username to an account id, through the hub.""" return await _op(lambda: ops.resolve_user(state, username)) @app.post("/api/members/{user_id}/revoke") async def revoke_member(user_id: str, group_id: str): + """Stop serving the group key to a member.""" return await _op(lambda: ops.revoke_member(state, user_id, group_id)) @app.post("/api/members/{user_id}/unpin") async def unpin_member(user_id: str): + """Forget a pinned identity, so the person can pair again with a new key.""" return await _op(lambda: ops.unpin_member(state, user_id)) # ── Chat encryption (operator only, localhost) ───────────────────────── @app.get("/api/groups/{group_id}/chat") async def chat_status(group_id: str): + """What the operator needs to decide anything about the group's chat.""" return await _op(lambda: ops.chat_status(state, group_id)) @app.post("/api/groups/{group_id}/chat/epoch") async def rotate_chat_epoch(group_id: str): + """Open a new chat epoch.""" return await _op(lambda: ops.open_chat_epoch(state, group_id)) @app.post("/api/groups/{group_id}/chat/encrypt-history") async def encrypt_chat_history(group_id: str): + """Re-encrypt the messages written before the group's chat was encrypted.""" return await _op(lambda: ops.encrypt_chat_history(state, group_id)) @app.post("/api/groups/{group_id}/chat/prune") async def prune_chat(group_id: str, max_age_days: int): + """Delete chat messages older than a number of days.""" return await _op(lambda: ops.prune_chat(state, group_id, max_age_days)) # ── GEK initialization (operator only, localhost) ────────────────────── @app.post("/api/groups/{group_id}/gek") async def init_gek(group_id: str, rotate: bool = False): + """Generate the group key, or rotate it with ?rotate=true.""" return await _op(lambda: ops.set_gek(state, group_id, rotate=rotate)) # ── Roots management (operator only, localhost) ──────────────────────── @app.post("/api/groups/{group_id}/roots") async def add_root(group_id: str, payload: dict): + """Add a folder to a group.""" result = await _op(lambda: ops.add_root( state, group_id, (payload.get("path") or "").strip(), @@ -392,6 +428,7 @@ def create_ui_app(state: dict) -> FastAPI: @app.patch("/api/groups/{group_id}/roots/{root_name}") async def update_root(group_id: str, root_name: str, payload: dict): + """Make a folder writable or removable, or not.""" result = await _op(lambda: ops.update_root( state, group_id, root_name, writable=payload.get("writable"), @@ -404,14 +441,17 @@ def create_ui_app(state: dict) -> FastAPI: @app.put("/api/groups/{group_id}/roots/{root_name}/eject") async def eject_root(group_id: str, root_name: str): + """Eject a removable folder so its disk can be unplugged.""" return await _op(lambda: ops.eject_root(state, group_id, root_name)) @app.put("/api/groups/{group_id}/roots/{root_name}/plug") async def plug_root(group_id: str, root_name: str): + """Bring an ejected folder back.""" return await _op(lambda: ops.plug_root(state, group_id, root_name)) @app.delete("/api/groups/{group_id}/roots/{root_name}") async def remove_root(group_id: str, root_name: str): + """Remove a folder from a group. At least one must remain.""" result = await _op(lambda: ops.remove_root(state, group_id, root_name)) reload_fn = state.get("reload_fn") if reload_fn: @@ -427,6 +467,8 @@ def create_ui_app(state: dict) -> FastAPI: @app.get("/api/groups/{group_id}/index-status") async def index_status(group_id: str): """ + One group's indexing progress. + Polled by the Create Group wizard and by "add a directory" in Settings — the same source either way, since both just start a scan on this group's indexer. `current_dir` is a basename only, and is @@ -454,7 +496,9 @@ def create_ui_app(state: dict) -> FastAPI: @app.get("/api/index-status") async def index_status_all(): """ - Every group's indexing at once, for the client's progress band — which + Every group's indexing progress. + + All groups at once, for the client's progress band — which is on screen whatever page the operator is on, so it cannot ask per group. Names roots, like `current_dir` above: loopback only, the operator's own screen. Reads state["indexers"] for the same reason. @@ -488,6 +532,7 @@ def create_ui_app(state: dict) -> FastAPI: @app.put("/api/groups/{group_id}/apps") async def set_enabled_apps(group_id: str, payload: dict): + """Which applications members see for the group.""" apps = payload.get("apps") if not isinstance(apps, list) or not apps: raise HTTPException(400, "apps must be a non-empty list") @@ -497,6 +542,7 @@ def create_ui_app(state: dict) -> FastAPI: @app.post("/api/reload") async def reload_config(): + """Reload node.toml. Returns before the reload finishes.""" # start_reload, not reload_config: this must return before a # brand-new group's synchronous initial scan finishes (minutes, not # seconds, on a real library) — see ops.start_reload for why. @@ -506,6 +552,7 @@ def create_ui_app(state: dict) -> FastAPI: @app.post("/api/shutdown") async def shutdown(): + """Stop the daemon.""" # The graceful stop every front door tries first (the desktop app, the # CLI, the installer): it reaches a node in any session -- a service # node runs in session 0, where taskkill and CTRL_BREAK from the user's @@ -521,20 +568,24 @@ def create_ui_app(state: dict) -> FastAPI: @app.get("/api/node-settings") async def get_node_settings(): + """The node's effective settings.""" return await _op(lambda: ops.get_node_settings(state)) @app.put("/api/node-settings") async def update_node_settings(payload: dict): + """Change node settings, written to roster.db and node.toml.""" return await _op(lambda: ops.set_node_settings(state, payload)) # ── Transfers (operator only, localhost) ─────────────────────────────── @app.get("/api/transfers") async def get_transfers(): + """Live transfer leases and queue depth.""" return await _op(lambda: ops.list_transfers(state)) @app.put("/api/groups/{group_id}/transfer-limits") async def set_transfer_limits(group_id: str, payload: dict): + """How many transfers one member may run at once in this group.""" # The only door to `ops.set_transfer_limits` (the CLI uses it). It once # had only a signed MNP message, which nothing anywhere sent — so the # per-member cap sat at its default of 2 with no way to change it. |