aboutsummaryrefslogtreecommitdiffstats
path: root/packages/meshbay-node/src
diff options
context:
space:
mode:
authorChristophe Besson <cbesson@gmail.com>2026-08-20 08:56:23 +0200
committerChristophe Besson <cbesson@gmail.com>2026-08-20 08:56:23 +0200
commitc8af746c846b5dbc792f7e4f0d806647d513cc5c (patch)
tree56c54ec5c703ec042589697b22b53e42451a61e4 /packages/meshbay-node/src
parent90c3c6a01f46c9b29c5d92702d7383f32e8951d7 (diff)
downloadmeshbay-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')
-rw-r--r--packages/meshbay-node/src/meshbay_node/daemon.py207
-rw-r--r--packages/meshbay-node/src/meshbay_node/ops.py217
-rw-r--r--packages/meshbay-node/src/meshbay_node/roster.py2
-rw-r--r--packages/meshbay-node/src/meshbay_node/transport/webrtc_server.py310
-rw-r--r--packages/meshbay-node/src/meshbay_node/ui/app.py18
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):