diff options
| author | Christophe Besson <cbesson@gmail.com> | 2026-08-20 08:56:23 +0200 |
|---|---|---|
| committer | Christophe Besson <cbesson@gmail.com> | 2026-08-20 08:56:23 +0200 |
| commit | c8af746c846b5dbc792f7e4f0d806647d513cc5c (patch) | |
| tree | 56c54ec5c703ec042589697b22b53e42451a61e4 /packages/meshbay-node/src/meshbay_node | |
| parent | 90c3c6a01f46c9b29c5d92702d7383f32e8951d7 (diff) | |
| download | meshbay-c8af746c846b5dbc792f7e4f0d806647d513cc5c.tar.gz | |
feat(node): full Node admin panel — CLI parity, hot-reload, group lifecycle
Node admin panel (NodePage) now covers every CLI operation over MNP:
group attach/detach, roster, member unpin, GEK rotate, denylist, reload.
Daemon hot-loads new groups and tears down removed ones on config reload
instead of requiring a full restart. Group attach/detach via MNP or
local API triggers an automatic reload so the group is live immediately.
Fixed GroupPage hang on first visit to a newly created group: the JWT
issued at login didn't include the new group, the node rejected with
not_a_member, and the token-refresh path returned without re-triggering
the connect effect (Boolean(token) didn't change). Now bumps retryKey
after a successful refresh so the effect re-runs with the fresh token.
NodePage marks groups hosted by the node but absent from the hub with a
"not on hub" badge so stale groups are visible and easy to remove.
Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
Diffstat (limited to 'packages/meshbay-node/src/meshbay_node')
| -rw-r--r-- | packages/meshbay-node/src/meshbay_node/daemon.py | 207 | ||||
| -rw-r--r-- | packages/meshbay-node/src/meshbay_node/ops.py | 217 | ||||
| -rw-r--r-- | packages/meshbay-node/src/meshbay_node/roster.py | 2 | ||||
| -rw-r--r-- | packages/meshbay-node/src/meshbay_node/transport/webrtc_server.py | 310 | ||||
| -rw-r--r-- | packages/meshbay-node/src/meshbay_node/ui/app.py | 18 |
5 files changed, 708 insertions, 46 deletions
diff --git a/packages/meshbay-node/src/meshbay_node/daemon.py b/packages/meshbay-node/src/meshbay_node/daemon.py index 6d4f172..7e3ebc1 100644 --- a/packages/meshbay-node/src/meshbay_node/daemon.py +++ b/packages/meshbay-node/src/meshbay_node/daemon.py @@ -111,6 +111,7 @@ class NodeDaemon: "quic_port": config.node.quic_port, "endpoint_hint": None, "indexes": {}, + "indexers": {}, } self._quic_server = None self._webrtc = None @@ -243,6 +244,7 @@ class NodeDaemon: await indexer.start() self._indexers.append(indexer) self._state["indexes"][group_cfg.id] = indexer.index + self._state["indexers"][group_cfg.id] = indexer log.info("Indexing group %s: %s (%d files)", group_cfg.name, ", ".join(f"{r.name}={r.path}" for r in roots), @@ -416,12 +418,15 @@ class NodeDaemon: self._state["webrtc"] = self._webrtc self._state["quic_server"] = self._quic_server self._state["hub"] = hub + self._state["reload_fn"] = self._reload_config # Rotating a key has to reach every transport holding a copy of it, # and clearing the denylist has to reach the one the handshake # consults — so both are published rather than reachable only # through the object that happens to own them. self._state["denylist"] = self._denylist self._state["pk_x25519_raw"] = pk_x_raw + self._state["sk_x25519_raw"] = sk_x_raw + self._state["sk_ed25519"] = keys.sk_ed25519 self._state["status"] = "running" log.info("Node ready — %d groups, WebRTC=%s, QUIC=%s", @@ -460,19 +465,13 @@ class NodeDaemon: async def _reload_config(self) -> None: """ - Re-read node.toml on SIGHUP. + Re-read node.toml and reconcile groups. - Deliberately narrow: it picks up **root changes on groups already - hosted**, which is what an operator adjusts day to day, and reports - anything else as needing a restart. Adding or removing a whole group - means new indexers, chat stores, GEK loads and transport contexts, and - doing that under a live daemon is how a half-built group ends up serving - content. Saying "restart for that" is honest and costs one restart. - - Nothing here touches connections: a member watching a film keeps - watching it. + Handles root changes on existing groups, hot-loads new groups, and + tears down removed groups. Existing connections are untouched: a member + watching a film keeps watching it. """ - log.info("SIGHUP — re-reading %s", self._config_path) + log.info("Reloading config from %s", self._config_path) try: fresh = load_config(self._config_path) except Exception as e: @@ -482,13 +481,8 @@ class NodeDaemon: groups_ctx = self._state.get("groups_ctx") or {} hosted = set(groups_ctx) incoming = {g.id for g in fresh.groups if g.id} - if incoming != hosted: - added = ", ".join(sorted(incoming - hosted)) or "none" - removed = ", ".join(sorted(hosted - incoming)) or "none" - log.warning("Group set changed (added: %s, removed: %s) — restart the " - "daemon for that; roots of existing groups reloaded anyway", - added, removed) + # ── Root changes on existing groups ────────────────────────────── changed = 0 for group_cfg in fresh.groups: ctx = groups_ctx.get(group_cfg.id) @@ -515,9 +509,112 @@ class NodeDaemon: ctx["roots"] = roots changed += 1 + # ── Hot-load new groups ────────────────────────────────────────── + added_names = [] + sk_ed = self._state.get("sk_ed25519") + sk_x_raw = self._state.get("sk_x25519_raw") + pk_x_raw = self._state.get("pk_x25519_raw") + node_user_id = self._state.get("node_user_id") + data_dir = fresh.data_dir + + for group_cfg in fresh.groups: + if group_cfg.id in hosted: + continue + if not group_cfg.id or not group_cfg.roots: + log.warning("New group %r has no id or roots — skipping", + group_cfg.name) + continue + if not sk_ed: + log.warning("Cannot hot-load %r — signing key not available", + group_cfg.name) + continue + + try: + roots = RootSet.build([asdict(r) for r in group_cfg.roots]) + except RootError as e: + log.error("New group %r: %s — skipping", group_cfg.name, e) + continue + roots.refresh_availability() + + gek = None + if group_cfg.visibility == "private" and sk_x_raw and pk_x_raw: + gek = await self._load_gek( + group_cfg.id, node_user_id, sk_x_raw, pk_x_raw) + if gek: + log.info("GEK loaded for new group %s", group_cfg.id[:8]) + + indexer = DirectoryIndexer( + roots=roots, + group_id=group_cfg.id, + sk_node=sk_ed, + gek=gek, + on_change=self._on_index_change, + ) + await indexer.start() + self._indexers.append(indexer) + self._state["indexes"][group_cfg.id] = indexer.index + self._state["indexers"][group_cfg.id] = indexer + + data_dir.mkdir(parents=True, exist_ok=True) + chat_db = data_dir / group_cfg.id[:16] / "chat.db" + store = ChatStore(db_path=chat_db) + await store.open() + self._chat_stores[group_cfg.id] = store + + new_ctx = { + "gek": gek, + "roots": roots, + "index": indexer.index, + "visibility": group_cfg.visibility, + "join_policy": group_cfg.join_policy, + "member_upload": ( + await self._roster.member_upload_allowed(group_cfg.id) + if self._roster else True), + "chat_store": store, + } + groups_ctx[group_cfg.id] = new_ctx + + if self._webrtc: + self._webrtc._ctx["groups"][group_cfg.id] = new_ctx + log.info("Hot-loaded group %s (%s, %d roots)", + group_cfg.name, group_cfg.id[:8], len(roots)) + added_names.append(group_cfg.name) + + # ── Tear down removed groups ───────────────────────────────────── + removed_names = [] + for gid in hosted - incoming: + indexer = next((i for i in self._indexers + if i.group_id == gid), None) + if indexer: + try: + await indexer.stop() + except Exception: + pass + self._indexers.remove(indexer) + store = self._chat_stores.pop(gid, None) + if store: + try: + await store.close() + except Exception: + pass + self._state["indexes"].pop(gid, None) + self._state["indexers"].pop(gid, None) + old_name = gid[:8] + for g_cfg in self._config.groups: + if g_cfg.id == gid: + old_name = g_cfg.name + break + groups_ctx.pop(gid, None) + if self._webrtc and self._webrtc._ctx.get("groups") is not groups_ctx: + self._webrtc._ctx["groups"].pop(gid, None) + log.info("Unloaded group %s (%s)", old_name, gid[:8]) + removed_names.append(old_name) + self._config = fresh self._state["config"] = fresh - log.info("Reload complete — %d group(s) re-rooted", changed) + self._state["groups"] = [g.name for g in fresh.groups] + log.info("Reload complete — %d re-rooted, %d added, %d removed", + changed, len(added_names), len(removed_names)) async def _login_with_retry(self, hub: HubClient): """Login to hub, retrying if the node key hasn't been linked yet.""" @@ -782,16 +879,18 @@ def main() -> None: parser.add_argument("command", nargs="?", choices=["init", "status", "ui", "gek-init", "gek", "operator", "member", "group", "file", - "denylist", "reload", "calibrate-argon2"], + "denylist", "reload", "restart-daemon", + "calibrate-argon2"], help="init: write example config | status: node state and keys " "| ui: print the admin UI URL | operator pair: pair a " "browser with this node | member list|invite|revoke|unpin " - "| group list|add | gek init|rotate | file list|rm " + "| group list|add|remove | gek init|rotate | file list|rm " "| denylist show|clear | reload: re-read node.toml " + "| restart-daemon: full stop + start " "| calibrate-argon2: benchmark") parser.add_argument("subcommand", nargs="?", help="'pair' for operator; list|invite|revoke|unpin for " - "member; list|add for group; init|rotate for gek; " + "member; list|add|remove for group; init|rotate for gek; " "list|rm for file; show|clear for denylist") parser.add_argument("target", nargs="?", help="username for member invite|revoke|unpin; group name " @@ -813,7 +912,8 @@ def main() -> None: # Query commands print a report; library logging would interleave with it. quiet = args.command in ("status", "ui", "gek-init", "gek", "operator", - "member", "group", "file", "denylist", "reload") + "member", "group", "file", "denylist", "reload", + "restart-daemon") logging.basicConfig( level=logging.ERROR if quiet else getattr(logging, args.log_level), format="%(asctime)s %(levelname)-8s %(name)s: %(message)s", @@ -1051,6 +1151,49 @@ def main() -> None: print("watch the result: tail -f /tmp/meshbay-node.log") return + if args.command == "restart-daemon": + import os as _os + import signal as _signal + import subprocess as _subprocess + cfg = load_config(args.config or DEFAULT_CONFIG_PATH) + pid_out = _subprocess.run( + ["pgrep", "-f", "--", r"-m meshbay_node\.daemon$"], + capture_output=True, text=True) + pids = [int(x) for x in pid_out.stdout.split()] + if pids: + for pid in pids: + _os.kill(pid, _signal.SIGTERM) + print(f"stopped {len(pids)} daemon process(es)") + for pid in pids: + try: + _os.waitpid(pid, 0) + except ChildProcessError: + import time as _time + _time.sleep(2) + else: + print("no running daemon found — starting fresh") + log_path = "/tmp/meshbay-node.log" + config_flag = ["--config", str(args.config)] if args.config else [] + _subprocess.Popen( + [sys.executable, "-m", "meshbay_node.daemon"] + config_flag, + stdout=open(log_path, "a"), + stderr=_subprocess.STDOUT, + start_new_session=True, + ) + import time as _time + _time.sleep(3) + pid_out2 = _subprocess.run( + ["pgrep", "-f", "--", r"-m meshbay_node\.daemon$"], + capture_output=True, text=True) + new_pids = [int(x) for x in pid_out2.stdout.split()] + if new_pids: + print(f"daemon started (PID {new_pids[0]})") + print(f"log: tail -f {log_path}") + else: + print(f"daemon may have failed to start — check {log_path}") + sys.exit(1) + return + if args.command == "denylist": cfg = load_config(args.config or DEFAULT_CONFIG_PATH) sub = args.subcommand or "show" @@ -1159,8 +1302,26 @@ def main() -> None: f"--group {g['name']}") return + if args.subcommand == "remove": + if not args.target: + print("usage: meshbay-node group remove <name>") + sys.exit(1) + cfg = load_config(args.config or DEFAULT_CONFIG_PATH) + if not args.yes: + answer = input(f"Remove group '{args.target}' from this node? [y/N] ") + if answer.lower() not in ("y", "yes"): + print("cancelled") + return + out = _daemon_api(cfg, "/api/groups/detach", method="POST", + body={"name": args.target}) + print(f"{out['name']} ({out['group_id'][:8]}) removed from {out['config']}") + print() + print("Restart the daemon to stop hosting it:") + print(" meshbay-node restart-daemon") + return + if args.subcommand != "add": - print("usage: meshbay-node group list|add <name> --dir <path> [--upload-dir <path>]") + print("usage: meshbay-node group list|add|remove <name>") sys.exit(1) if not args.target or not args.dir: print("usage: meshbay-node group add <name> --dir <path> [--upload-dir <path>]") diff --git a/packages/meshbay-node/src/meshbay_node/ops.py b/packages/meshbay-node/src/meshbay_node/ops.py index 20345bb..2a76290 100644 --- a/packages/meshbay-node/src/meshbay_node/ops.py +++ b/packages/meshbay-node/src/meshbay_node/ops.py @@ -23,6 +23,7 @@ in the adapter. from __future__ import annotations import logging +import re from dataclasses import asdict from pathlib import Path from typing import Any @@ -93,9 +94,11 @@ async def read_roster(state: dict, group_id: str = "") -> dict: roster = state.get("roster") if not roster: return {"identities": [], "members": [], "invites": []} + members = await roster.list_members(group_id or None) + members = [m for m in members if m.get("pk_ed25519") is not None] return { "identities": await roster.list_identities(), - "members": await roster.list_members(group_id or None), + "members": members, "invites": await roster.list_invites(), } @@ -387,6 +390,125 @@ async def attach_group(state: dict, name: str, shared_dir: str, return result +async def detach_group(state: dict, name: str) -> dict: + """ + Remove a [[groups]] block from node.toml. + + Does not touch the hub — only stops this node from hosting the group + after the next reload or restart. + """ + if not name: + raise OpError("group name or id is required") + config = _config(state) + + match = [g for g in config.groups if g.id == name or g.name == name] + if not match: + raise OpError(f"No hosted group matches {name!r}", status=404, + extra={"available": [{"name": g.name, "id": g.id} + for g in config.groups]}) + group = match[0] + + conf_path = Path(state.get("config_path") or DEFAULT_CONFIG_PATH) + text = conf_path.read_text() + lines = text.split("\n") + + rng = _find_group_range(lines, group.id) + if rng is None: + raise OpError(f"Group {group.id[:8]} not found in {conf_path}") + + start, end = rng + while end < len(lines) and lines[end].strip() == "": + end += 1 + + new_lines = lines[:start] + lines[end:] + conf_path.write_text("\n".join(new_lines)) + log.info("Group detached: %s (%s) removed from %s", group.name, group.id[:8], conf_path) + + return {"group_id": group.id, "name": group.name, "config": str(conf_path), + "note": "restart the node to stop hosting it"} + + +def _find_group_range(lines: list[str], group_id: str) -> tuple[int, int] | None: + """Line range of a [[groups]] block by id: (start, end_exclusive).""" + id_re = re.compile(r'^\s*id\s*=\s*"([^"]*)"') + block_starts: list[int] = [] + for i, line in enumerate(lines): + if line.strip() == "[[groups]]": + block_starts.append(i) + + for j, start in enumerate(block_starts): + boundary = block_starts[j + 1] if j + 1 < len(block_starts) else len(lines) + for k in range(start + 1, boundary): + s = lines[k].strip() + if s.startswith("[") and s != "[[groups.roots]]": + boundary = k + break + for k in range(start + 1, boundary): + m = id_re.match(lines[k]) + if m and m.group(1) == group_id: + return (start, boundary) + return None + + +def _insert_roots_block(conf_path: Path, group_id: str, + root_block: str) -> None: + """Append a [[groups.roots]] block inside the matching [[groups]] section.""" + text = conf_path.read_text() + lines = text.split("\n") + + rng = _find_group_range(lines, group_id) + if rng is None: + raise OpError(f"Group {group_id[:8]} not found in {conf_path}") + + _start, end = rng + insert_at = end + while insert_at > _start + 1 and lines[insert_at - 1].strip() == "": + insert_at -= 1 + + new_lines = (lines[:insert_at] + + [""] + + root_block.rstrip("\n").split("\n") + + lines[insert_at:]) + conf_path.write_text("\n".join(new_lines)) + + +def _remove_roots_block(conf_path: Path, group_id: str, + resolved_path: str) -> None: + """Remove a [[groups.roots]] block whose resolved path matches.""" + text = conf_path.read_text() + lines = text.split("\n") + + rng = _find_group_range(lines, group_id) + if rng is None: + raise OpError(f"Group {group_id[:8]} not found in {conf_path}") + + start, end = rng + path_re = re.compile(r'^\s*path\s*=\s*"([^"]*)"') + roots_starts: list[int] = [] + for i in range(start + 1, end): + if lines[i].strip() == "[[groups.roots]]": + roots_starts.append(i) + + for j, rs in enumerate(roots_starts): + rs_end = roots_starts[j + 1] if j + 1 < len(roots_starts) else end + for k in range(rs, rs_end): + m = path_re.match(lines[k]) + if m: + try: + p = str(Path(m.group(1)).expanduser().resolve()) + except OSError: + continue + if p == resolved_path: + rm_start = rs + if rm_start > 0 and lines[rm_start - 1].strip() == "": + rm_start -= 1 + new_lines = lines[:rm_start] + lines[rs_end:] + conf_path.write_text("\n".join(new_lines)) + return + + raise OpError(f"Root path not found in config", status=404) + + async def add_root(state: dict, group_id: str, path: str, *, name: str = "", kind: str = "generic", upload: bool = False) -> dict: @@ -410,22 +532,85 @@ async def add_root(state: dict, group_id: str, path: str, *, raise OpError(str(e)) from e added = built.roots[-1] + + try: + added.path.mkdir(parents=True, exist_ok=True) + except OSError as e: + raise OpError(f"Cannot create {added.path}: {e}") from e + conf_path = Path(state.get("config_path") or DEFAULT_CONFIG_PATH) - raise OpError( - # Writing into the middle of a hand-written TOML file means finding the - # right [[groups]] block and appending inside it, which a text append - # cannot do. Until that is written, say so plainly rather than appending - # to the wrong group. - f"Add this to {conf_path} under the [[groups]] block for " - f"{cfg.name!r}, then restart the node:\n\n" - f' [[groups.roots]]\n' - f' path = "{added.path}"\n' - + (f' name = "{added.name}"\n' if name else "") - + (f' kind = "{added.kind}"\n' if kind != "generic" else "") - + (f' upload = true\n' if upload else ""), - status=501, - extra={"validated": True, "name": added.name, "path": str(added.path)}, - ) + root_block = f' [[groups.roots]]\n path = "{added.path}"' + if name: + root_block += f'\n name = "{added.name}"' + if kind != "generic": + root_block += f'\n kind = "{added.kind}"' + if upload: + root_block += f'\n upload = true' + _insert_roots_block(conf_path, group_id, root_block) + + from meshbay_node.config import RootSpec + cfg.roots.append(RootSpec( + path=str(added.path), name=added.name, kind=added.kind, + upload=added.upload, direct=added.direct)) + groups_ctx = state.get("groups_ctx", {}) + if group_id in groups_ctx: + groups_ctx[group_id]["roots"] = built + + log.info("Root added: %s → group %s", added.name, group_id[:8]) + return {"status": "added", "name": added.name, "path": str(added.path), + "group_id": group_id, "roots": built.describe()} + + +async def remove_root(state: dict, group_id: str, root_name: str) -> dict: + """Remove a named root from a group. At least one root must remain.""" + config = _config(state) + cfg = next((g for g in config.groups if g.id == group_id), None) + if cfg is None: + raise OpError("Group not configured on this node", status=404) + + from meshbay_common.paths import fold + from meshbay_node.roots import derive_name + target = fold(root_name) + match_idx = None + for i, r in enumerate(cfg.roots): + try: + rname = r.name or derive_name(Path(r.path).expanduser().resolve()) + except Exception: + continue + if fold(rname) == target: + match_idx = i + break + + if match_idx is None: + raise OpError(f"No root named {root_name!r} in this group", status=404) + if len(cfg.roots) < 2: + raise OpError("Cannot remove the only root", status=400) + + removed = cfg.roots[match_idx] + if removed.upload: + raise OpError( + "Cannot remove the upload root — file uploads and chat " + "attachments are stored there", status=400) + resolved = str(Path(removed.path).expanduser().resolve()) + + conf_path = Path(state.get("config_path") or DEFAULT_CONFIG_PATH) + _remove_roots_block(conf_path, group_id, resolved) + + cfg.roots.pop(match_idx) + remaining = [asdict(r) for r in cfg.roots] + try: + built = RootSet.build(remaining) + except RootError: + built = None + if built is not None: + groups_ctx = state.get("groups_ctx", {}) + if group_id in groups_ctx: + groups_ctx[group_id]["roots"] = built + + log.info("Root removed: %s from group %s", root_name, group_id[:8]) + return {"status": "removed", "name": root_name, "group_id": group_id, + "roots": built.describe() if built else [], + "note": "restart recommended to update the file index"} # ── Files ──────────────────────────────────────────────────────────────────── diff --git a/packages/meshbay-node/src/meshbay_node/roster.py b/packages/meshbay-node/src/meshbay_node/roster.py index c811016..9a52b53 100644 --- a/packages/meshbay-node/src/meshbay_node/roster.py +++ b/packages/meshbay-node/src/meshbay_node/roster.py @@ -359,6 +359,8 @@ class Roster: assert self._db cur = await self._db.execute( "DELETE FROM identities WHERE user_id = ?", (user_id,)) + await self._db.execute( + "DELETE FROM members WHERE user_id = ?", (user_id,)) await self._db.commit() return cur.rowcount > 0 diff --git a/packages/meshbay-node/src/meshbay_node/transport/webrtc_server.py b/packages/meshbay-node/src/meshbay_node/transport/webrtc_server.py index f182ab2..1f1f2d2 100644 --- a/packages/meshbay-node/src/meshbay_node/transport/webrtc_server.py +++ b/packages/meshbay-node/src/meshbay_node/transport/webrtc_server.py @@ -63,6 +63,10 @@ from meshbay_common.adminop import ( OP_GEK_ROTATE, OP_MEMBER_UNPIN, OP_MEMBER_UPLOAD, + OP_ROOT_ADD, + OP_ROOT_REMOVE, + OP_GROUP_ATTACH, + OP_GROUP_DETACH, admin_transcript, ) from meshbay_common.crypto import pk_to_b64, wrap_gek_aes @@ -385,11 +389,16 @@ class WebRTCPeerSession: def _setup_channel(self, channel: RTCDataChannel) -> None: self._channel = channel + self._msg_count = 0 @channel.on("message") def on_message(message): if isinstance(message, str): message = message.encode() + self._msg_count += 1 + if self._msg_count <= 3: + log.info("WebRTC data received: %d bytes, msg #%d (peer=%s)", + len(message), self._msg_count, self._peer_id) self._buffer.feed(message) for msg in self._buffer.messages(): self._handle_message(msg) @@ -476,6 +485,24 @@ class WebRTCPeerSession: self._do_member_unpin(msg) elif mtype == MNP.GEK_ROTATE: self._do_gek_rotate(msg) + elif mtype == MNP.NODE_STATUS: + self._spawn(self._do_node_status(msg)) + elif mtype == MNP.ROOT_ADD: + self._do_root_add(msg) + elif mtype == MNP.ROOT_REMOVE: + self._do_root_remove(msg) + elif mtype == MNP.ROSTER_READ: + self._spawn(self._do_roster_read(msg)) + elif mtype == MNP.DENYLIST_READ: + self._spawn(self._do_denylist_read(msg)) + elif mtype == MNP.DENYLIST_CLEAR: + self._spawn(self._do_denylist_clear(msg)) + elif mtype == MNP.GROUP_ATTACH: + self._do_group_attach(msg) + elif mtype == MNP.GROUP_DETACH: + self._do_group_detach(msg) + elif mtype == MNP.NODE_RELOAD: + self._spawn(self._do_node_reload(msg)) elif mtype == MNP.KEYPAIR_BUNDLE_STORE: self._spawn(self._do_keypair_bundle_store(msg)) elif mtype == MNP.KEYPAIR_BUNDLE_DELETE: @@ -565,6 +592,8 @@ class WebRTCPeerSession: def _do_handshake(self, msg: dict) -> None: group_id = msg.get("group_id", "") + log.info("WebRTC handshake request: group=%s (peer=%s)", + group_id[:8] if group_id else "none", self._peer_id) try: peer = authorize_token( msg.get("token", ""), @@ -598,6 +627,7 @@ class WebRTCPeerSession: gctx = self._ctx["groups"][peer.group_id] if "groups" in self._ctx else self._ctx if not gctx.get("gek"): + log.warning("Handshake refused — no GEK for group=%s", peer.group_id[:8]) self._send({ "type": "error", "detail": "Group encryption not initialized — contact node operator", @@ -606,6 +636,7 @@ class WebRTCPeerSession: self._gek_challenge = os.urandom(NONCE_LEN) self._nonce_node = self._gek_challenge + log.info("WebRTC handshake challenge sent (peer=%s)", self._peer_id) self._send({ "type": MNP.HANDSHAKE_CHALLENGE, "v": MNP_VERSION, @@ -1510,13 +1541,14 @@ class WebRTCPeerSession: one. The node generates the replacement itself — nothing arriving here contributes key material, which is what the C5b rule is about. """ - if not self._group_id: + 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, self._group_id) + 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, @@ -1627,6 +1659,241 @@ class WebRTCPeerSession: except Exception: pass + # ── Node management (D5) ───────────────────────────────────────────────── + + async def _do_node_status(self, msg: dict) -> None: + """All groups, roots, peers — the operator's overview.""" + node_uid = self._ctx.get("node_user_id") + log.info("node_status: user=%s node_user=%s admin=%s", + self._user_id, node_uid, self._is_node_admin()) + if not self._is_node_admin(): + self._send({"type": "error", "detail": "Not the node 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_roster_read(self, msg: dict) -> None: + if not self._is_node_admin(): + self._send({"type": "error", "detail": "Not the node 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 self._is_node_admin(): + self._send({"type": "error", "detail": "Not the node 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 self._is_node_admin(): + self._send({"type": "error", "detail": "Not the node 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 = str(msg.get("upload_dir", "")).strip() + self._issue_admin_challenge( + OP_GROUP_ATTACH, name, + payload={"name": name, "shared_dir": shared_dir, + "upload_dir": upload_dir}, + 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"], p.get("upload_dir", "")) + 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 self._is_node_admin(): + self._send({"type": "error", "detail": "Not the node 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 + self._issue_admin_challenge( + OP_ROOT_ADD, path, + payload={ + "group_id": target_group, "path": path, + "name": str(msg.get("name", ""))[:128], + "kind": str(msg.get("kind", "generic"))[:16], + "upload": bool(msg.get("upload", False)), + }, + 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"), + upload=p.get("upload", 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() + 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 + self._issue_admin_challenge( + OP_ROOT_REMOVE, root_name, + payload={"group_id": target_group, "root_name": root_name}, + group_id=target_group) + + async def _admin_exec_root_remove( + 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_remove:{pending['subject'][:24]}") + return + p = pending["payload"] + try: + result = await self._run_op( + ops.remove_root, p["group_id"], p["root_name"]) + except ops.OpError as e: + self._send({"type": "error", "detail": e.message}) + return + except Exception as e: + log.error("root_remove failed: %s", e, exc_info=True) + self._send({"type": "error", "detail": "Internal error"}) + return + self._audit("root_remove", f"{p['root_name']}←{p['group_id'][:8]}") + await self._retarget_indexer(p["group_id"]) + self._send({"type": MNP.ROOT_REMOVE_ACK, "v": MNP_VERSION, **result}) + async def _run_op(self, fn, *args, **kwargs): """ Call an operation from `meshbay_node.ops` with the daemon's own view. @@ -1642,6 +1909,16 @@ class WebRTCPeerSession: raise ops.OpError("Node state not available", status=503) return await fn(state, *args, **kwargs) + async def _retarget_indexer(self, group_id: str) -> None: + """Tell the indexer to rescan after roots changed.""" + state = self._ctx.get("daemon_state") + if not state: + return + indexer = state.get("indexers", {}).get(group_id) + roots = state.get("groups_ctx", {}).get(group_id, {}).get("roots") + if indexer and roots: + await indexer.retarget(roots) + async def _admin_exec_member_revoke( self, pending: dict, transcript: bytes, sig: bytes, ) -> None: @@ -1740,7 +2017,11 @@ class WebRTCPeerSession: """ task = asyncio.ensure_future(coro) self._tasks.add(task) - task.add_done_callback(self._tasks.discard) + def _on_done(t): + self._tasks.discard(t) + if not t.cancelled() and t.exception(): + log.error("Spawned task failed: %s", t.exception(), exc_info=t.exception()) + task.add_done_callback(_on_done) return task def _group_ctx(self) -> dict: @@ -2213,6 +2494,7 @@ class WebRTCPeerSession: def _issue_admin_challenge( self, op: str, subject: str, payload: dict | None = None, + group_id: str | None = None, ) -> None: """ Ask the client to authorize `op` on `subject` with its Ed25519 identity key. @@ -2221,13 +2503,17 @@ class WebRTCPeerSession: rebuild and inspect what it signs. The node keeps the authoritative copy and rebuilds the transcript itself at verification time — nothing signed is ever taken from the response message. + + `group_id` overrides the connection's group for cross-group operations + (e.g. root management from a NodePage connection). """ + gid = group_id if group_id is not None else (self._group_id or "") nonce = os.urandom(32) ts = int(time.time()) op_id = base64.b64encode(os.urandom(16)).decode() self._admin_ops[op_id] = { "op": op, "subject": subject, "nonce": nonce, "ts": ts, - "payload": payload or {}, + "payload": payload or {}, "group_id": gid, } self._send({ "type": MNP.ADMIN_CHALLENGE, @@ -2238,7 +2524,7 @@ class WebRTCPeerSession: "nonce": base64.b64encode(nonce).decode(), "ts": ts, "node_pk": self._node_pk_b64(), - "group_id": self._group_id or "", + "group_id": gid, }) @staticmethod @@ -2327,7 +2613,7 @@ class WebRTCPeerSession: transcript = admin_transcript( op=pending["op"], node_pk_b64=self._node_pk_b64(), - group_id=self._group_id or "", + group_id=pending["group_id"] if pending.get("group_id") is not None else (self._group_id or ""), subject=pending["subject"], nonce=pending["nonce"], ts=pending["ts"], @@ -2354,6 +2640,18 @@ class WebRTCPeerSession: elif pending["op"] == OP_MEMBER_UPLOAD: self._spawn( self._admin_exec_member_upload(pending, transcript, sig_bytes)) + elif pending["op"] == OP_ROOT_ADD: + self._spawn( + self._admin_exec_root_add(pending, transcript, sig_bytes)) + elif pending["op"] == OP_ROOT_REMOVE: + self._spawn( + self._admin_exec_root_remove(pending, transcript, sig_bytes)) + elif pending["op"] == OP_GROUP_ATTACH: + self._spawn( + self._admin_exec_group_attach(pending, transcript, sig_bytes)) + elif pending["op"] == OP_GROUP_DETACH: + self._spawn( + self._admin_exec_group_detach(pending, transcript, sig_bytes)) else: self._send({"type": "error", "detail": "Unknown admin operation"}) diff --git a/packages/meshbay-node/src/meshbay_node/ui/app.py b/packages/meshbay-node/src/meshbay_node/ui/app.py index 74a7c8a..a068a8d 100644 --- a/packages/meshbay-node/src/meshbay_node/ui/app.py +++ b/packages/meshbay-node/src/meshbay_node/ui/app.py @@ -12,6 +12,7 @@ Served only on 127.0.0.1 — not exposed to the network. Gated by a per-run session token (11.5.3) — printed at daemon startup. """ +import asyncio import base64 import json import logging @@ -132,12 +133,27 @@ def create_ui_app(state: dict) -> FastAPI: return await _op(lambda: ops.list_groups(state)) @app.post("/api/groups/attach") async def attach_group(payload: dict): - return await _op(lambda: ops.attach_group( + result = await _op(lambda: ops.attach_group( state, (payload.get("name") or "").strip(), (payload.get("shared_dir") or "").strip(), upload_dir=(payload.get("upload_dir") or "").strip(), )) + reload_fn = state.get("reload_fn") + if reload_fn: + asyncio.ensure_future(reload_fn()) + return result + + @app.post("/api/groups/detach") + async def detach_group(payload: dict): + result = await _op(lambda: ops.detach_group( + state, + (payload.get("name") or payload.get("group_id") or "").strip(), + )) + reload_fn = state.get("reload_fn") + if reload_fn: + asyncio.ensure_future(reload_fn()) + return result @app.delete("/api/groups/{group_id}/files/{file_id}") async def delete_file(group_id: str, file_id: str): |