aboutsummaryrefslogtreecommitdiffstats
path: root/packages/meshbay-node/src
diff options
context:
space:
mode:
Diffstat (limited to 'packages/meshbay-node/src')
-rw-r--r--packages/meshbay-node/src/meshbay_node/__init__.py2
-rw-r--r--packages/meshbay-node/src/meshbay_node/cli/status.py3
-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.py80
-rw-r--r--packages/meshbay-node/src/meshbay_node/ops/chat.py2
-rw-r--r--packages/meshbay-node/src/meshbay_node/ops/roots.py15
-rw-r--r--packages/meshbay-node/src/meshbay_node/ops/settings.py6
-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/src/meshbay_node/transport/webrtc/admin.py55
-rw-r--r--packages/meshbay-node/src/meshbay_node/transport/webrtc/chat.py6
-rw-r--r--packages/meshbay-node/src/meshbay_node/transport/webrtc/dispatch.py13
-rw-r--r--packages/meshbay-node/src/meshbay_node/transport/webrtc/group_ops.py84
-rw-r--r--packages/meshbay-node/src/meshbay_node/transport/webrtc/node_ops.py378
-rw-r--r--packages/meshbay-node/src/meshbay_node/ui/app.py62
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))))