aboutsummaryrefslogtreecommitdiffstats
path: root/packages/meshbay-node/src/meshbay_node/transport/webrtc_server.py
diff options
context:
space:
mode:
Diffstat (limited to 'packages/meshbay-node/src/meshbay_node/transport/webrtc_server.py')
-rw-r--r--packages/meshbay-node/src/meshbay_node/transport/webrtc_server.py310
1 files changed, 304 insertions, 6 deletions
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"})