diff options
Diffstat (limited to 'packages/meshbay-node/src')
16 files changed, 208 insertions, 543 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 4b619d7..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) @@ -719,6 +720,21 @@ class DirectoryIndexer: # The table now; the files when the scan ends. await self.on_change(self) + async def publish_roots(self) -> None: + """ + Tell every connected peer the table, after a root's flags were edited + in place (`ops.update_root`). + + A reload compares the edited set with itself and finds nothing to do, + so without this a directory made writable from the operator's own + machine stayed read-only on every open page until something else + happened to push the index. + """ + self._index.roots = self.roots.describe() + self._index.version = int(time.time()) + if self.on_change: + await self.on_change(self) + def _holds(self, root: Root) -> bool: return any(r.folded == root.folded and r.path == root.path for r in self.roots) @@ -810,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) @@ -1017,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 @@ -1024,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) @@ -1043,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/chat.py b/packages/meshbay-node/src/meshbay_node/ops/chat.py index 469e907..539284b 100644 --- a/packages/meshbay-node/src/meshbay_node/ops/chat.py +++ b/packages/meshbay-node/src/meshbay_node/ops/chat.py @@ -17,7 +17,7 @@ log = logging.getLogger("meshbay_node.ops") # # The key a group's chat archive is encrypted under. Generated here, by the # node, and never by a member — the C5b rule is about key material arriving from -# outside, and this is the same rule that lets `gek_rotate` be a signed +# outside, and this is the same rule that lets a group key rotation be an # instruction rather than a delivery. # # An *epoch* rather than a rotation, and the distinction is the whole design: diff --git a/packages/meshbay-node/src/meshbay_node/ops/roots.py b/packages/meshbay-node/src/meshbay_node/ops/roots.py index 480fe63..b976e8f 100644 --- a/packages/meshbay-node/src/meshbay_node/ops/roots.py +++ b/packages/meshbay-node/src/meshbay_node/ops/roots.py @@ -174,11 +174,14 @@ async def update_root(state: dict, group_id: str, root_name: str, *, writable=match.writable, removable=match.removable) # Update the live RootSet so GET /api/groups returns correct data - # immediately, without waiting for the async reload to finish. + # immediately, without waiting for the async reload to finish — and the + # indexer's, which is normally the same object but need not be, since it + # is the one the table pushed to every peer is read from. live_roots: RootSet | None = state.get("groups_ctx", {}).get( group_id, {}).get("roots") - if live_roots: - for lr in live_roots.roots: + indexer = state.get("indexers", {}).get(group_id) + for rootset in {id(x): x for x in (live_roots, indexer and indexer.roots) if x}.values(): + for lr in rootset.roots: lr_name = lr.name or str(Path(lr.path).name) if fold(lr_name) == target: if writable is not None: @@ -187,6 +190,9 @@ async def update_root(state: dict, group_id: str, root_name: str, *, lr.removable = removable break + if indexer: + await indexer.publish_roots() + # Built from config when there is no live set, never returned empty: an # empty list is a *valid answer* meaning "this group has no directories", # and the client cannot tell it from "the node could not say". It would @@ -234,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( @@ -282,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/ops/settings.py b/packages/meshbay-node/src/meshbay_node/ops/settings.py index eade6ca..906cd03 100644 --- a/packages/meshbay-node/src/meshbay_node/ops/settings.py +++ b/packages/meshbay-node/src/meshbay_node/ops/settings.py @@ -182,9 +182,9 @@ async def set_transfer_limits(state: dict, group_id: str, Same shape as every other operator setting: lives on the node (roster.db, not the hub and not node.toml, for the reason change 5 gives — a hub that - decided this would have authority over someone else's machine), signed - (webrtc_server checks the caller's admin authority before this runs), and - live, so the pools are updated in place rather than at the next restart. + decided this would have authority over someone else's machine), set on the + node's own machine (loopback API, CLI), and live, so the pools are updated + in place rather than at the next restart. """ roster = _roster(state) ctx = _group_ctx(state, group_id) 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/transport/webrtc/admin.py b/packages/meshbay-node/src/meshbay_node/transport/webrtc/admin.py index daaee62..9a6cbd1 100644 --- a/packages/meshbay-node/src/meshbay_node/transport/webrtc/admin.py +++ b/packages/meshbay-node/src/meshbay_node/transport/webrtc/admin.py @@ -17,27 +17,20 @@ from meshbay_common.adminop import ( OP_CHAT_LINK_PREVIEW, OP_DIR_DELETE, OP_FILE_DELETE, - OP_GEK_ROTATE, - OP_GROUP_ATTACH, - OP_GROUP_DETACH, OP_INVITE_CANCEL, OP_INVITE_CREATE, OP_INVITE_LINK_CREATE, OP_MEMBER_REVOKE, - OP_MEMBER_UNPIN, OP_MUSICBRAINZ_ENABLED, - OP_ROOT_ADD, OP_ROOT_EJECT, OP_ROOT_PLUG, OP_ROOT_REMOVE, - OP_ROOT_UPDATE, OP_SEARCH_LISTED, OP_SET_SCAN_SETTINGS, OP_TMDB_CONFIG, OP_TMDB_ENABLED, OP_TMDB_OVERRIDE, OP_TMDB_REMATCH, - OP_TRANSFER_LIMITS, admin_transcript, ) from meshbay_common.crypto import pk_to_b64 @@ -62,28 +55,21 @@ _ADMIN_EXECUTORS = { OP_INVITE_CREATE: "_admin_exec_invite_create", OP_INVITE_LINK_CREATE: "_admin_exec_invite_link_create", OP_INVITE_CANCEL: "_admin_exec_invite_cancel", - OP_GEK_ROTATE: "_admin_exec_gek_rotate", - OP_MEMBER_UNPIN: "_admin_exec_member_unpin", OP_APPS_ENABLED: "_admin_exec_apps_enabled", - OP_TRANSFER_LIMITS: "_admin_exec_transfer_limits", OP_SET_SCAN_SETTINGS: "_admin_exec_set_scan_settings", OP_TMDB_CONFIG: "_admin_exec_tmdb_config", OP_TMDB_ENABLED: "_admin_exec_tmdb_enabled", OP_TMDB_OVERRIDE: "_admin_exec_tmdb_override", OP_TMDB_REMATCH: "_admin_exec_tmdb_rematch", OP_MUSICBRAINZ_ENABLED: "_admin_exec_musicbrainz_enabled", - OP_ROOT_ADD: "_admin_exec_root_add", OP_ROOT_REMOVE: "_admin_exec_root_remove", OP_APP_DIRECTORIES: "_admin_exec_app_directories", OP_CHAT_DIRECTORY: "_admin_exec_chat_directory", OP_CHAT_LINK_PREVIEW: "_admin_exec_chat_link_preview", OP_SEARCH_LISTED: "_admin_exec_search_listed", OP_CHAT_EPOCH: "_admin_exec_chat_epoch", - OP_ROOT_UPDATE: "_admin_exec_root_update", OP_ROOT_EJECT: "_admin_exec_root_eject", OP_ROOT_PLUG: "_admin_exec_root_plug", - OP_GROUP_ATTACH: "_admin_exec_group_attach", - OP_GROUP_DETACH: "_admin_exec_group_detach", } @@ -180,49 +166,12 @@ class AdminMixin: says "the hub says you are the owner", which is the one thing NS4 and M3 rule out: a hub that can name the operator can install itself as node administrator. It rides the handshake ack so a client knows whether - to offer the Node page at all, and every operation is gated on - `_operator_device()` below. + to offer operator controls at all; every operation is gated on a + signature (`_verify_admin_sig`). """ node_user_id = self._ctx.get("node_user_id") return bool(node_user_id and self._user_id == node_user_id) - async def _operator_device(self) -> bool: - """ - Whether this connection may run the node's own controls. - - Two things, and the second is the one that cannot be forged: - - - the account is the one this node belongs to (`_is_node_admin`), which - is what keeps node-wide controls with the machine's owner rather than - with every paired operator of every group on it; and - - **the device on this connection proved a key the node pinned as an - operator**. `device_hello` is signed over a transcript naming this - node, this group and this connection's nonce, and `operator_pks()` is - rebuilt from the roster on each call, so an unpinned browser and a - revoked one are both refused at once. - - The second clause is the fix for the door this used to leave open. - `node_status`, `node_settings_set`, `roster_read`, `denylist_read`, - `denylist_clear` and `node_reload` were gated on the account id alone — - a value the hub chooses. An active hub that can also reach the group key - (which §3.5 concedes it can in an open-join group) could therefore mint - a token for the owner's account and read `node_status`, which lists - every group on the node with the operator's **absolute paths**, or clear - the denylist, which is the persisted revocation H4 exists to keep. - - It holds no user keys and cannot countersign anything, so it cannot - produce a `device_hello` — which is the same property device linking - rests on (§3.3), applied to the node's own surface. - """ - if not self._is_node_admin(): - return False - if not self._device_confirmed or not self._pinned_pk: - return False - roster = self._ctx.get("roster") - if roster is None: - return False - return self._pinned_pk in await roster.operator_pks() - def _has_admin_authority(self) -> bool: """ Cheap synchronous pre-check: is there anyone who could authorize this? diff --git a/packages/meshbay-node/src/meshbay_node/transport/webrtc/chat.py b/packages/meshbay-node/src/meshbay_node/transport/webrtc/chat.py index 493d7f0..dc43ae2 100644 --- a/packages/meshbay-node/src/meshbay_node/transport/webrtc/chat.py +++ b/packages/meshbay-node/src/meshbay_node/transport/webrtc/chat.py @@ -198,9 +198,9 @@ class ChatMixin: There is no switch to turn chat encryption on: MNP 2.0 has no plaintext chat to fall back to. What an operator may want to do deliberately is - move the key on — the same instruction as `gek_rotate`, and signed for - the same reason. The removals that matter (member revoke, member unpin, - device revoke, `gek_rotate`) already open one by themselves. + move the key on, which is signed like the rest. The removals that matter + (member revoke, member unpin, device revoke, a group key rotation) + already open one by themselves. """ group_id = str(msg.get("group_id", "")).strip() or self._group_id if not group_id: diff --git a/packages/meshbay-node/src/meshbay_node/transport/webrtc/dispatch.py b/packages/meshbay-node/src/meshbay_node/transport/webrtc/dispatch.py index ee2e3a1..4bf9dbb 100644 --- a/packages/meshbay-node/src/meshbay_node/transport/webrtc/dispatch.py +++ b/packages/meshbay-node/src/meshbay_node/transport/webrtc/dispatch.py @@ -48,7 +48,6 @@ _HANDLERS = { MNP.DEVICE_REVOKE: ("_do_device_revoke", SPAWNED), MNP.DEVICE_HELLO: ("_do_device_hello", SPAWNED), MNP.APPS_ENABLED: ("_do_apps_enabled", INLINE), - MNP.TRANSFER_LIMITS: ("_do_transfer_limits", INLINE), MNP.SET_SCAN_SETTINGS: ("_do_set_scan_settings", INLINE), MNP.TMDB_CONFIG: ("_do_tmdb_config", INLINE), MNP.TMDB_ENABLED: ("_do_tmdb_enabled", INLINE), @@ -68,21 +67,9 @@ _HANDLERS = { MNP.MUSIC_META_REQ: ("_do_music_meta_request", SPAWNED), MNP.AUDIO_TRANSCODE_REQ: ("_do_audio_transcode_request", SPAWNED), MNP.SUBTITLE_REQ: ("_do_subtitle_request", SPAWNED), - MNP.MEMBER_UNPIN: ("_do_member_unpin", INLINE), - MNP.GEK_ROTATE: ("_do_gek_rotate", INLINE), - MNP.NODE_STATUS: ("_do_node_status", SPAWNED), - MNP.ROOT_ADD: ("_do_root_add", INLINE), MNP.ROOT_REMOVE: ("_do_root_remove", INLINE), - MNP.ROOT_UPDATE: ("_do_root_update", INLINE), MNP.ROOT_EJECT: ("_do_root_eject", INLINE), MNP.ROOT_PLUG: ("_do_root_plug", INLINE), - MNP.ROSTER_READ: ("_do_roster_read", SPAWNED), - MNP.DENYLIST_READ: ("_do_denylist_read", SPAWNED), - MNP.DENYLIST_CLEAR: ("_do_denylist_clear", SPAWNED), - MNP.GROUP_ATTACH: ("_do_group_attach", INLINE), - MNP.GROUP_DETACH: ("_do_group_detach", INLINE), - MNP.NODE_SETTINGS_SET: ("_do_node_settings_set", SPAWNED), - MNP.NODE_RELOAD: ("_do_node_reload", SPAWNED), MNP.KEYPAIR_BUNDLE_STORE: ("_do_keypair_bundle_store", SPAWNED), MNP.KEYPAIR_BUNDLE_DELETE: ("_do_keypair_bundle_delete", SPAWNED), MNP.USER_BLOB_STORE: ("_do_user_blob_store", SPAWNED), diff --git a/packages/meshbay-node/src/meshbay_node/transport/webrtc/group_ops.py b/packages/meshbay-node/src/meshbay_node/transport/webrtc/group_ops.py index 3756eac..a94e2aa 100644 --- a/packages/meshbay-node/src/meshbay_node/transport/webrtc/group_ops.py +++ b/packages/meshbay-node/src/meshbay_node/transport/webrtc/group_ops.py @@ -5,9 +5,7 @@ from meshbay_common import MNP_VERSION from meshbay_common.adminop import ( OP_APP_DIRECTORIES, OP_APPS_ENABLED, - OP_GEK_ROTATE, OP_MEMBER_REVOKE, - OP_MEMBER_UNPIN, OP_SEARCH_LISTED, ) from meshbay_common.groupbox import PURPOSE_ROSTER, seal @@ -41,88 +39,6 @@ class GroupOpsMixin: return self._issue_admin_challenge(OP_MEMBER_REVOKE, user_id) - def _do_gek_rotate(self, msg: dict) -> None: - """ - Ask for a new group key. Operator only, and signed. - - This is what actually removes a revoked member's access: revocation - stops the node serving the *next* key, and they still hold the current - one. The node generates the replacement itself — nothing arriving here - contributes key material, which is what the C5b rule is about. - """ - group_id = str(msg.get("group_id", "")).strip() or self._group_id - if not group_id: - self._send({"type": "error", "detail": "No group on this connection"}) - return - if not self._has_admin_authority(): - self._send({"type": "error", "detail": "No authorized key for this"}) - return - self._issue_admin_challenge(OP_GEK_ROTATE, group_id, group_id=group_id) - - async def _admin_exec_gek_rotate( - self, pending: dict, transcript: bytes, sig: bytes, - ) -> None: - if not await self._verify_admin_sig(transcript, sig): - self._send({"type": "error", "detail": "Signature verification failed"}) - self._audit("admin_auth_failed", f"gek_rotate:{pending['subject'][:8]}") - return - try: - result = await self._run_op( - ops.set_gek, pending["subject"], rotate=True) - except ops.OpError as e: - self._send({"type": "error", "detail": e.message}) - return - # The operator is rotating because somebody left, and the chat archive - # key is not derived from the group key — so rotating that one does not - # move this one. Doing both here is what makes "rotate after a removal" - # mean the same thing for chat as it does for files. - await self._new_chat_epoch(pending["subject"], "gek_rotate") - self._audit("gek_rotate", pending["subject"]) - self._send({ - "type": MNP.GEK_ROTATE_ACK, "v": MNP_VERSION, - "group_id": pending["subject"], - "authorized_members": result.get("authorized_members", 0), - # Said plainly, because rotating is the step people skip: content - # already downloaded stays readable to whoever holds it. - "note": "members re-receive the key on their next connect; content " - "already downloaded is unaffected", - }) - - def _do_member_unpin(self, msg: dict) -> None: - """Forget a pinned identity, so someone can pair again with a new key.""" - user_id = str(msg.get("user_id", "")).strip() - if not user_id: - self._send({"type": "error", "detail": "Missing user_id"}) - return - if user_id == self._user_id: - # Unpinning yourself over the connection your pin authorizes would - # end that connection's authority mid-operation. - self._send({"type": "error", "detail": "Cannot unpin yourself"}) - return - if not self._has_admin_authority(): - self._send({"type": "error", "detail": "No authorized key for this"}) - return - self._issue_admin_challenge(OP_MEMBER_UNPIN, user_id) - - async def _admin_exec_member_unpin( - self, pending: dict, transcript: bytes, sig: bytes, - ) -> None: - user_id = pending["subject"] - if not await self._verify_admin_sig(transcript, sig): - self._send({"type": "error", "detail": "Signature verification failed"}) - self._audit("admin_auth_failed", f"member_unpin:{user_id[:8]}") - return - try: - # The new chat epochs and the closed sessions are the op's own - # (`ops.members._after_removal`), for every door alike. - await self._run_op(ops.unpin_member, user_id) - except ops.OpError as e: - self._send({"type": "error", "detail": e.message}) - return - self._audit("member_unpin", user_id) - self._send({"type": MNP.MEMBER_UNPIN_ACK, "v": MNP_VERSION, - "user_id": user_id}) - # Every "application" a group can show. Photos joins this set (and # apps.js's registry, client-side) when it lands; nothing else about # this handler changes. DEFAULT_APPS (roster.py) deliberately does not diff --git a/packages/meshbay-node/src/meshbay_node/transport/webrtc/node_ops.py b/packages/meshbay-node/src/meshbay_node/transport/webrtc/node_ops.py index f7bbbfa..237a359 100644 --- a/packages/meshbay-node/src/meshbay_node/transport/webrtc/node_ops.py +++ b/packages/meshbay-node/src/meshbay_node/transport/webrtc/node_ops.py @@ -1,21 +1,16 @@ -"""The operator's controls over the node itself: status and settings, roster -and denylist, roots, hosted groups, reload, scan pacing and transfer limits.""" +"""The operator's controls over a group's roots and scan pacing that MNP carries: +removing, ejecting and plugging a root. What widens the sharing, and the node's +own status, settings, roster and denylist, are the loopback API's and the CLI's +(MNP 6.0).""" import logging from meshbay_common import MNP_VERSION from meshbay_common.adminop import ( - OP_GROUP_ATTACH, - OP_GROUP_DETACH, - OP_ROOT_ADD, OP_ROOT_EJECT, OP_ROOT_PLUG, OP_ROOT_REMOVE, - OP_ROOT_UPDATE, OP_SET_SCAN_SETTINGS, - OP_TRANSFER_LIMITS, - group_attach_subject, - root_add_subject, ) from meshbay_common.protocol import MNP @@ -60,64 +55,6 @@ class NodeOpsMixin: self._issue_admin_challenge( OP_SET_SCAN_SETTINGS, f"{reconcile:g},{debounce:g}") - MIN_TRANSFER_LIMIT = 1 - MAX_TRANSFER_LIMIT = 32 - - def _do_transfer_limits(self, msg: dict) -> None: - """How many transfers one member may run at once in this group. - - Zero is not "unlimited" and is refused: a member who may not transfer at - all is a member the operator revokes, and reading 0 as no-limit would - make the most dangerous value the easiest to type by accident. - """ - try: - downloads = int(msg.get("downloads")) - uploads = int(msg.get("uploads")) - except (TypeError, ValueError): - self._send({"type": "error", "detail": "Invalid transfer limits"}) - return - for value in (downloads, uploads): - if not (self.MIN_TRANSFER_LIMIT <= value <= self.MAX_TRANSFER_LIMIT): - self._send({"type": "error", - "detail": f"transfer limits must be between " - f"{self.MIN_TRANSFER_LIMIT} and " - f"{self.MAX_TRANSFER_LIMIT}"}) - return - if not self._has_admin_authority(): - self._send({"type": "error", "detail": "No authorized key for this"}) - return - self._issue_admin_challenge(OP_TRANSFER_LIMITS, - f"d={downloads},u={uploads}") - - async def _admin_exec_transfer_limits( - self, pending: dict, transcript: bytes, sig: bytes, - ) -> None: - try: - parts = dict(p.split("=") for p in pending["subject"].split(",")) - downloads, uploads = int(parts["d"]), int(parts["u"]) - except (ValueError, KeyError): - self._send({"type": "error", "detail": "Invalid transfer limits"}) - return - if not await self._verify_admin_sig(transcript, sig): - self._send({"type": "error", "detail": "Signature verification failed"}) - self._audit("admin_auth_failed", f"transfer_limits:{pending['subject']}") - return - try: - result = await self._run_op( - ops.set_transfer_limits, self._group_id or "", downloads, uploads) - except ops.OpError as e: - self._send({"type": "error", "detail": e.message}) - return - self._audit("transfer_limits", pending["subject"]) - - notice = {"type": MNP.TRANSFER_LIMITS_ACK, "v": MNP_VERSION, - "limits": result["limits"]} - for session in list(self._peer_registry().values()): - try: - session._send(notice) - except Exception: - pass - async def _admin_exec_set_scan_settings( self, pending: dict, transcript: bytes, sig: bytes, ) -> None: @@ -146,245 +83,6 @@ class NodeOpsMixin: except Exception: pass - # ── Node management (D5) ───────────────────────────────────────────────── - - async def _do_node_status(self, msg: dict) -> None: - """All groups, roots, peers — the operator's overview. - - Including every root's absolute path, which is why this is gated on a - proved operator device and not on an account the hub named. - """ - node_uid = self._ctx.get("node_user_id") - log.info("node_status: user=%s node_user=%s owner=%s device=%s", - self._user_id, node_uid, self._is_node_admin(), - "confirmed" if self._device_confirmed else "unidentified") - if not await self._operator_device(): - self._send({"type": "error", "detail": "Not the node operator", - "code": "not_operator"}) - return - try: - result = await self._run_op(ops.list_groups) - self._send({"type": MNP.NODE_STATUS_ACK, "v": MNP_VERSION, **result}) - except ops.OpError as e: - self._send({"type": "error", "detail": e.message}) - except Exception as e: - log.error("node_status failed: %s", e, exc_info=True) - self._send({"type": "error", "detail": "Internal error"}) - - async def _do_node_settings_set(self, msg: dict) -> None: - if not await self._operator_device(): - self._send({"type": "error", "detail": "Not the node operator", - "code": "not_operator"}) - return - settings = msg.get("settings", {}) - if not settings: - self._send({"type": "error", "detail": "No settings provided"}) - return - try: - result = await self._run_op(ops.set_node_settings, settings) - self._send({"type": MNP.NODE_SETTINGS_SET_ACK, "v": MNP_VERSION, - **result}) - except ops.OpError as e: - self._send({"type": "error", "detail": e.message}) - except Exception as e: - log.error("node_settings_set failed: %s", e, exc_info=True) - self._send({"type": "error", "detail": "Internal error"}) - - async def _do_roster_read(self, msg: dict) -> None: - if not await self._operator_device(): - self._send({"type": "error", "detail": "Not the node operator", - "code": "not_operator"}) - return - group_id = str(msg.get("group_id", "")).strip() - try: - result = await self._run_op(ops.read_roster, group_id) - self._send({"type": MNP.ROSTER_READ_ACK, "v": MNP_VERSION, **result}) - except ops.OpError as e: - self._send({"type": "error", "detail": e.message}) - except Exception as e: - log.error("roster_read failed: %s", e, exc_info=True) - self._send({"type": "error", "detail": "Internal error"}) - - async def _do_denylist_read(self, msg: dict) -> None: - if not await self._operator_device(): - self._send({"type": "error", "detail": "Not the node operator", - "code": "not_operator"}) - return - try: - result = await self._run_op(ops.read_denylist) - self._send({"type": MNP.DENYLIST_READ_ACK, "v": MNP_VERSION, **result}) - except ops.OpError as e: - self._send({"type": "error", "detail": e.message}) - except Exception as e: - log.error("denylist_read failed: %s", e, exc_info=True) - self._send({"type": "error", "detail": "Internal error"}) - - async def _do_denylist_clear(self, msg: dict) -> None: - if not await self._operator_device(): - self._send({"type": "error", "detail": "Not the node operator", - "code": "not_operator"}) - return - subject = str(msg.get("subject", "")).strip() - try: - result = await self._run_op(ops.clear_denylist, subject=subject) - self._audit("denylist_clear", subject or "all") - self._send({"type": MNP.DENYLIST_CLEAR_ACK, "v": MNP_VERSION, **result}) - except ops.OpError as e: - self._send({"type": "error", "detail": e.message}) - except Exception as e: - log.error("denylist_clear failed: %s", e, exc_info=True) - self._send({"type": "error", "detail": "Internal error"}) - - def _do_group_attach(self, msg: dict) -> None: - name = str(msg.get("name", "")).strip() - shared_dir = str(msg.get("shared_dir", "")).strip() - if not name or not shared_dir: - self._send({"type": "error", "detail": "Missing name or shared_dir"}) - return - if not self._has_admin_authority(): - self._send({"type": "error", "detail": "No authorized key for this"}) - return - # `upload_dir` is not read here any more, and a client still sending it - # is ignored rather than obeyed: on load it forces every other root - # read-only, which is the model the RO/RW one replaced. A second - # writable directory is `root_add` with `writable`. - writable = bool(msg.get("writable", True)) - # The directory being exposed is signed, not only the group's name. - self._issue_admin_challenge( - OP_GROUP_ATTACH, group_attach_subject(name, shared_dir, writable), - payload={"name": name, "shared_dir": shared_dir, "writable": writable}, - group_id="") - - async def _admin_exec_group_attach( - self, pending: dict, transcript: bytes, sig: bytes, - ) -> None: - if not await self._verify_admin_sig(transcript, sig): - self._send({"type": "error", "detail": "Signature verification failed"}) - self._audit("admin_auth_failed", - f"group_attach:{pending['subject'][:16]}") - return - p = pending.get("payload") or {} - try: - result = await self._run_op( - ops.attach_group, p["name"], p["shared_dir"], - writable=bool(p.get("writable", True))) - except ops.OpError as e: - self._send({"type": "error", "detail": e.message}) - return - self._audit("group_attach", pending["subject"]) - self._send({"type": MNP.GROUP_ATTACH_ACK, "v": MNP_VERSION, **result}) - state = self._ctx.get("daemon_state") - reload_fn = state.get("reload_fn") if state else None - if reload_fn: - try: - await reload_fn() - except Exception as e: - log.error("Reload after group_attach failed: %s", e) - - def _do_group_detach(self, msg: dict) -> None: - name = str(msg.get("name", "")).strip() - if not name: - self._send({"type": "error", "detail": "Missing group name or id"}) - return - if not self._has_admin_authority(): - self._send({"type": "error", "detail": "No authorized key for this"}) - return - self._issue_admin_challenge( - OP_GROUP_DETACH, name, - payload={"name": name}, - group_id="") - - async def _admin_exec_group_detach( - self, pending: dict, transcript: bytes, sig: bytes, - ) -> None: - if not await self._verify_admin_sig(transcript, sig): - self._send({"type": "error", "detail": "Signature verification failed"}) - self._audit("admin_auth_failed", - f"group_detach:{pending['subject'][:16]}") - return - p = pending.get("payload") or {} - try: - result = await self._run_op(ops.detach_group, p["name"]) - except ops.OpError as e: - self._send({"type": "error", "detail": e.message}) - return - self._audit("group_detach", pending["subject"]) - self._send({"type": MNP.GROUP_DETACH_ACK, "v": MNP_VERSION, **result}) - state = self._ctx.get("daemon_state") - reload_fn = state.get("reload_fn") if state else None - if reload_fn: - try: - await reload_fn() - except Exception as e: - log.error("Reload after group_detach failed: %s", e) - - async def _do_node_reload(self, msg: dict) -> None: - if not await self._operator_device(): - self._send({"type": "error", "detail": "Not the node operator", - "code": "not_operator"}) - return - state = self._ctx.get("daemon_state") - reload_fn = state.get("reload_fn") if state else None - if not reload_fn: - self._send({"type": "error", "detail": "Reload not available"}) - return - try: - await reload_fn() - self._send({"type": MNP.NODE_RELOAD_ACK, "v": MNP_VERSION, - "status": "reloaded"}) - except Exception as e: - log.error("node_reload failed: %s", e, exc_info=True) - self._send({"type": "error", "detail": "Reload failed"}) - - def _do_root_add(self, msg: dict) -> None: - target_group = str(msg.get("group_id", "")).strip() - path = str(msg.get("path", "")).strip() - if not target_group or not path: - self._send({"type": "error", "detail": "Missing group_id or path"}) - return - if not self._has_admin_authority(): - self._send({"type": "error", "detail": "No authorized key for this"}) - return - payload = { - "group_id": target_group, "path": path, - "name": str(msg.get("name", ""))[:128], - "kind": str(msg.get("kind", "generic"))[:16], - "writable": bool(msg.get("writable", msg.get("upload", False))), - "removable": bool(msg.get("removable", False)), - } - # Everything the executor acts on is signed — `writable` decides whether - # every member may write there. The group is in the transcript itself. - self._issue_admin_challenge( - OP_ROOT_ADD, - root_add_subject(path, payload["name"], payload["kind"], - payload["writable"], payload["removable"]), - payload=payload, group_id=target_group) - - async def _admin_exec_root_add( - self, pending: dict, transcript: bytes, sig: bytes, - ) -> None: - if not await self._verify_admin_sig(transcript, sig): - self._send({"type": "error", "detail": "Signature verification failed"}) - self._audit("admin_auth_failed", f"root_add:{pending['subject'][:24]}") - return - p = pending["payload"] - try: - result = await self._run_op( - ops.add_root, p["group_id"], p["path"], - name=p.get("name", ""), kind=p.get("kind", "generic"), - writable=p.get("writable", False), - removable=p.get("removable", False)) - except ops.OpError as e: - self._send({"type": "error", "detail": e.message}) - return - except Exception as e: - log.error("root_add failed: %s", e, exc_info=True) - self._send({"type": "error", "detail": "Internal error"}) - return - self._audit("root_add", f"{p['path']}→{p['group_id'][:8]}") - await self._retarget_indexer(p["group_id"]) - self._send({"type": MNP.ROOT_ADD_ACK, "v": MNP_VERSION, **result}) - def _do_root_remove(self, msg: dict) -> None: target_group = str(msg.get("group_id", "")).strip() root_name = str(msg.get("root_name", "")).strip() @@ -421,59 +119,6 @@ class NodeOpsMixin: await self._retarget_indexer(p["group_id"]) self._send({"type": MNP.ROOT_REMOVE_ACK, "v": MNP_VERSION, **result}) - def _do_root_update(self, msg: dict) -> None: - target_group = str(msg.get("group_id", self._group_id or "")).strip() - root_name = str(msg.get("root_name", "")).strip() - if not target_group or not root_name: - self._send({"type": "error", "detail": "Missing group_id or root_name"}) - return - if not self._has_admin_authority(): - self._send({"type": "error", "detail": "No authorized key for this"}) - return - updates = [] - if "writable" in msg: - updates.append(f"rw={'on' if msg['writable'] else 'off'}") - if "removable" in msg: - updates.append(f"rem={'on' if msg['removable'] else 'off'}") - subject = f"{root_name}:{','.join(updates)}" if updates else root_name - self._issue_admin_challenge( - OP_ROOT_UPDATE, subject, - payload={ - "group_id": target_group, "root_name": root_name, - "writable": msg.get("writable"), - "removable": msg.get("removable"), - }, - group_id=target_group) - - async def _admin_exec_root_update( - self, pending: dict, transcript: bytes, sig: bytes, - ) -> None: - if not await self._verify_admin_sig(transcript, sig): - self._send({"type": "error", "detail": "Signature verification failed"}) - self._audit("admin_auth_failed", - f"root_update:{pending['subject'][:24]}") - return - p = pending["payload"] - try: - result = await self._run_op( - ops.update_root, p["group_id"], p["root_name"], - writable=p.get("writable"), removable=p.get("removable")) - except ops.OpError as e: - self._send({"type": "error", "detail": e.message}) - return - except Exception as e: - log.error("root_update failed: %s", e, exc_info=True) - self._send({"type": "error", "detail": "Internal error"}) - return - self._audit("root_update", pending["subject"]) - await self._retarget_indexer(p["group_id"]) - notice = {"type": MNP.ROOT_UPDATE_ACK, "v": MNP_VERSION, **result} - for uid, session in list(self._peer_registry().items()): - try: - session._send(notice) - except Exception: - pass - def _do_root_eject(self, msg: dict) -> None: target_group = str(msg.get("group_id", self._group_id or "")).strip() root_name = str(msg.get("root_name", "")).strip() @@ -558,24 +203,21 @@ class NodeOpsMixin: async def _retarget_indexer(self, group_id: str) -> None: """ - Pick up a root that was just added to or removed from node.toml. + Pick up a root that was just removed from node.toml. Through the daemon's own reload, which is what the loopback API has always done after the same operations (`ui/app.py`). This used to re-point the indexer at `groups_ctx[gid]["roots"]` instead — the very object the op had just edited — so `retarget` diffed a set against - itself, found no new names, scanned nothing, and dropped nothing. A - directory added over MNP reached node.toml and was invisible until a - restart; one removed kept serving its files. + itself, found no new names, scanned nothing, and dropped nothing: a + directory removed over MNP kept serving its files until a restart. Two front doors doing different things is the shape `ops.py` exists to prevent, and this was it: the loopback path worked and the MNP path did - not, which is why it survived until the operator added a directory from - a browser. + not. - Not awaited: a reload rescans, and a new library is minutes. The ack - the caller sends carries the set the node is moving to, and the - `index_sync` that follows the scan carries what it found. + Not awaited: a reload can be long. The ack the caller sends carries the + set the node is moving to. """ state = self._ctx.get("daemon_state") if not state: diff --git a/packages/meshbay-node/src/meshbay_node/ui/app.py b/packages/meshbay-node/src/meshbay_node/ui/app.py index f810903..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,24 +568,27 @@ 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): - # The same `ops.set_transfer_limits` the signed MNP handler calls. The - # op existed with only that one door, and nothing anywhere opened it — - # so the per-member cap sat at its default of 2 with no way to change - # it, which from outside is indistinguishable from a hardcoded 2. + """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. return await _op(lambda: ops.set_transfer_limits( state, group_id, int(payload.get("downloads", 0)), int(payload.get("uploads", 0)))) |