aboutsummaryrefslogtreecommitdiffstats
path: root/packages/meshbay-node/src/meshbay_node/transport/webrtc_server.py
diff options
context:
space:
mode:
authorChristophe Besson <cbesson@gmail.com>2026-09-24 10:36:54 +0200
committerChristophe Besson <cbesson@gmail.com>2026-09-24 16:45:38 +0200
commit41b814483cf990d855188899332b90fc525135b7 (patch)
tree7049a38f8907f14f76addc3b7faf9922016ef36e /packages/meshbay-node/src/meshbay_node/transport/webrtc_server.py
parent600a6dfe3779b9cd82d4e0cb94c6f6e7df21ab20 (diff)
downloadmeshbay-41b814483cf990d855188899332b90fc525135b7.tar.gz
refactor(node): move the MNP handshake out of webrtc_server
HandshakeMixin in transport/webrtc/handshake.py: challenge, proof, channel binding, and the sealed configuration a peer receives once admitted. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
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.py357
1 files changed, 3 insertions, 354 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 0657b57..d67426d 100644
--- a/packages/meshbay-node/src/meshbay_node/transport/webrtc_server.py
+++ b/packages/meshbay-node/src/meshbay_node/transport/webrtc_server.py
@@ -70,24 +70,9 @@ from meshbay_common.adminop import (
)
from meshbay_common.crypto import pk_to_b64
from meshbay_common.groupbox import (
- PURPOSE_ACK,
PURPOSE_ROSTER,
seal,
)
-from meshbay_common.handshake import (
- MNP_MIN_SUPPORTED,
- NONCE_LEN,
- ROLE_CLIENT,
- ROLE_NODE,
- HandshakeError,
- authorize_token,
- challenge_transcript,
- check_version,
- handshake_transcript,
- make_proof,
- verify_proof,
- webrtc_binding,
-)
from meshbay_common.protocol import (
MNP,
)
@@ -95,7 +80,6 @@ from meshbay_common.protocol import (
from meshbay_node import ops
from meshbay_node import transfers as transfers_mod
from meshbay_node.indexer import GroupIndex
-from meshbay_node.indexer.indexer import DirectoryIndexer
# Re-imported under its original name: every call site and existing test in
# this module still refers to it as `_probe_video`. The implementation lives
@@ -113,13 +97,12 @@ from meshbay_node.transport.webrtc.blobs import BlobsMixin
from meshbay_node.transport.webrtc.channel import (
_REPLY_TO,
_DataChannelBuffer,
- _extract_dtls_fingerprint,
_get_remote_ip,
_pack,
)
from meshbay_node.transport.webrtc.chat import ChatMixin
from meshbay_node.transport.webrtc.files import FilesMixin
-from meshbay_node.transport.webrtc.limits import MAX_MSG
+from meshbay_node.transport.webrtc.handshake import HandshakeMixin
from meshbay_node.transport.webrtc.transfer_handlers import TransferMixin
from meshbay_node.transport.webrtc.upload_handlers import UploadMixin
@@ -161,7 +144,8 @@ _WEBRTC_TRACE_INTERVAL_S = 30.0
class WebRTCPeerSession(
- AdmissionMixin, BlobsMixin, ChatMixin, FilesMixin, TransferMixin, UploadMixin,
+ AdmissionMixin, BlobsMixin, ChatMixin, FilesMixin, HandshakeMixin,
+ TransferMixin, UploadMixin,
StreamingMixin, VideoMetaMixin, MusicMixin, SubtitlesMixin,
):
"""One WebRTC peer connection, handling MNP over a DataChannel."""
@@ -542,291 +526,6 @@ class WebRTCPeerSession(
detail=detail,
))
- def _channel_binding(self) -> bytes:
- """Both DTLS fingerprints, so a proof is valid on this connection only."""
- offer_fp = b""
- answer_fp = b""
- if self._pc.remoteDescription:
- offer_fp = _extract_dtls_fingerprint(self._pc.remoteDescription.sdp)
- if self._pc.localDescription:
- answer_fp = _extract_dtls_fingerprint(self._pc.localDescription.sdp)
- if not offer_fp or not answer_fp:
- return b""
- return webrtc_binding(offer_fp, answer_fp)
-
- 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)
- # Before the token, and before anything is decided from it: a peer we
- # cannot speak to is refused with a code it can act on, rather than
- # served messages it will misread as missing fields (L2).
- try:
- check_version(msg.get("v", ""), msg.get("v_min", ""))
- except HandshakeError as refusal:
- self._send({"type": "error", "detail": str(refusal),
- "code": refusal.code})
- return
- try:
- peer = authorize_token(
- msg.get("token", ""),
- self._ctx["hub_pk_pem"],
- group_id=group_id,
- hosted_groups=self._ctx.get("groups"),
- denylist=self._ctx.get("denylist"),
- )
- except HandshakeError as refusal:
- # HandshakeError messages are authored to be peer-safe, unlike arbitrary
- # exception text (L3) — the client needs to know *why* it was refused.
- self._send({"type": "error", "detail": str(refusal),
- "code": getattr(refusal, "code", "")})
- self._audit_auth_failed(group_id, str(refusal))
- return
-
- try:
- self._nonce_client = base64.b64decode(msg.get("nonce", ""))
- except Exception:
- self._nonce_client = b""
- if len(self._nonce_client) < NONCE_LEN:
- # The client nonce is what makes the NODE's proof fresh (C3). Without
- # it a recorded ack could be replayed by an impersonating peer.
- self._send({"type": "error", "detail": "Client nonce required"})
- return
-
- # Decoded, but NOT authenticated: that happens on the GEK proof.
- self._pending_sub = peer.user_id
- self._pending_group = peer.group_id
- self._pending_username = peer.username
-
- 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",
- })
- return
-
- 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,
- # Our half of the range. The client refuses us on this rather than
- # discovering the mismatch when a field it expected is not there.
- "v_min": MNP_MIN_SUPPORTED,
- "nonce": base64.b64encode(self._gek_challenge).decode(),
- # Announced here because a first-time joiner needs it *before* the
- # ack: join_request signs a transcript naming this node, and someone
- # who has never held the GEK cannot complete the handshake to learn
- # it. Unverified at this point — the ack proves it, the client checks
- # the two match, and a wrong value only makes our own verification
- # fail. It is never a substitute for the ack's proof and signature.
- "node_pk": self._node_pk_b64(),
- # ...except that since 3.4 it is signed, so a client that already
- # knows which key to expect can check it before it sends a code.
- **self._challenge_sig(peer.group_id, self._channel_binding()),
- })
-
- def _challenge_sig(self, group_id: str, binding: bytes) -> dict:
- """
- `{"sig": ...}` over the challenge transcript, or nothing (MNP 3.4).
-
- What makes `node_pk` above more than an announcement: a client about to
- send an invitation code can check this node holds the key it was told
- to expect, before the code leaves. No binding means no signature rather
- than an unbound one — a signature that is not tied to the channel is one
- somebody can relay, and the handshake proof refuses that case anyway.
- """
- if not binding:
- return {}
- transcript = challenge_transcript(
- group_id, self._nonce_client, self._gek_challenge, binding)
- return {"sig": base64.b64encode(self._ctx["sk_node"].sign(transcript)).decode()}
-
- def _do_handshake_response(self, msg: dict) -> None:
- if not self._gek_challenge or not hasattr(self, "_pending_sub"):
- self._send({"type": "error", "detail": "No pending handshake challenge"})
- return
-
- group_id = self._pending_group
- gctx = self._ctx["groups"][group_id] if "groups" in self._ctx else self._ctx
- gek = gctx.get("gek")
- if not gek:
- self._send({"type": "error", "detail": "Group encryption not initialized"})
- self._gek_challenge = None
- return
-
- try:
- proof_bytes = base64.b64decode(msg.get("proof", ""))
- except Exception:
- self._send({"type": "error", "detail": "Invalid proof encoding"})
- return
-
- binding = self._channel_binding()
- if not binding:
- # Refuse rather than fall back to an unbound proof (L4).
- self._send({"type": "error", "detail": "Channel binding unavailable"})
- self._gek_challenge = None
- self._audit_auth_failed(group_id, "no channel binding")
- return
-
- if not verify_proof(gek, proof_bytes, ROLE_CLIENT, group_id,
- self._nonce_client, self._gek_challenge, binding):
- self._send({"type": "error", "detail": "GEK proof failed"})
- self._gek_challenge = None
- self._audit_auth_failed(group_id, "GEK HMAC mismatch")
- return
-
- self._complete_handshake(gek, binding)
- self._gek_challenge = None
-
- def _complete_handshake(self, gek: bytes, binding: bytes) -> None:
- # Authenticated peers may send large frames (file uploads); unauthenticated
- # ones may not (H6).
- self._buffer.max_message = MAX_MSG
- self._user_id = self._pending_sub
- self._group_id = self._pending_group
- self._username = self._pending_username
- self._spawn(self._load_pinned_pk())
-
- self._register_peer()
-
- node_user_id = self._ctx.get("node_user_id")
- log.info("WebRTC handshake OK — user=%s group=%s",
- self._user_id[:8],
- self._group_id[:8] if self._group_id else "none")
- # The node proves itself too (C3): possession of the GEK over the client's
- # nonce, plus a signature over the same transcript with its long-term key.
- # Previously the client received an unverifiable node_pk and trusted
- # is_node_admin from whoever answered — so a peer that had hijacked
- # signaling could serve a forged index, chat history and permissions.
- node_transcript = handshake_transcript(
- ROLE_NODE, self._group_id or "", self._nonce_client,
- self._gek_challenge or b"", binding)
- node_proof = make_proof(
- gek, ROLE_NODE, self._group_id or "", self._nonce_client,
- self._gek_challenge or b"", binding)
-
- # Everything the client needs in order to *authenticate* us stays in clear —
- # node_pk, proof and sig are what it checks before it would trust a
- # decryption, so they cannot themselves be behind one. The configuration
- # below is sealed under a GEK-derived subkey, which gives it an
- # authentication tag from a key the hub does not hold. Until MNP 1.0 the
- # signed transcript named no ack field at all, so is_node_admin,
- # enabled_apps, the app directories and the rest were authenticated by DTLS
- # channel and nothing else.
- config = {
- "is_node_admin": self._is_node_admin(),
- # Which group "applications" to show. Absent/empty falls back to
- # every registered one client-side, so a node that predates this
- # setting (or one whose context has not loaded it yet) hides
- # nothing.
- "enabled_apps": list(self._group_ctx().get("enabled_apps") or []),
- # Read once and kept current in place by the signed op, and
- # surfaced here rather than only via tmdb_enabled_ack, so a client
- # that connects after the operator configured it does not have to
- # wait for a live change to find out.
- "tmdb_enabled": bool(self._group_ctx().get("tmdb_enabled", True)),
- # Token/language stay node-wide (one shared credential/cache) —
- # via daemon_state, kept current by tmdb_config_ack.
- "tmdb_token_customized": bool(
- self._ctx.get("daemon_state", {}).get("tmdb_token_customized", False)),
- "tmdb_language": str(
- self._ctx.get("daemon_state", {}).get("tmdb_language") or ""),
- # Music app (docs/MESHBAY_DESIGN.md §9.8) — same shape as the TMDB
- # fields above. No language field: MusicBrainz search doesn't
- # take one the way TMDB does.
- "musicbrainz_enabled": bool(self._group_ctx().get("musicbrainz_enabled", True)),
- # Every app's configured folders — `<app>_directories`, keyed by
- # the app's own registry name, always a list. Built by daemon.py's
- # `_app_directories_ctx`, and the only form on the wire: the
- # `video_root` / `audio_root` / `photo_roots` scalars that used to
- # ride here are gone. One folder was never the general case, and two
- # spellings of one answer meant whichever the reader consulted first
- # decided it.
- #
- # `chat_directory` below is the one surviving second name, and it is
- # safe for the reason those were not: it is *derived* from this list
- # on every build rather than stored beside it, so the two cannot
- # drift apart.
- **self._app_directories_ack(),
- # Where chat attachments are written — the singular form, because
- # Chat genuinely has one destination. "" means the operator has not
- # chosen, and the paperclip says so.
- "chat_directory": self._group_ctx().get("chat_directory") or "",
- # Whether the node unfurls links members post here. Absent means
- # on, which is what it did before this existed.
- "chat_link_preview": bool(
- self._group_ctx().get("chat_link_preview", True)),
- # Whether the reader's cross-group Search should list this group.
- # Presentation only: the index below is served to Search and to the
- # group page alike, and this cannot tell them apart. Sealed like the
- # rest, so the hub cannot flip it. Absent means listed.
- "search_listed": bool(self._group_ctx().get("search_listed", True)),
- # Which chat epoch key a client should be sealing under. Inside
- # the sealed part of the ack like every other configuration field,
- # so it carries an authentication tag from a key the hub does not
- # hold — a forged epoch would have a client sealing under a key the
- # group has retired.
- #
- # No `chat_encrypted` beside it: there is no switch. A peer that
- # reached this point speaks MNP 2.0, and 2.0 has no plaintext chat.
- "chat_epoch": int(self._group_ctx().get("chat_epoch", 0) or 0),
- # This member's own transfer caps in this group, so the interface
- # can say "2 of 2 of your slots are busy" rather than draw a bare
- # spinner. Absent reads as "no limit known" and the hint is simply
- # not drawn — never as "unlimited", which would have the interface
- # contradicting the node.
- "transfer_limits": {
- "download": self._slots().member_cap(
- transfers_mod.DOWNLOAD,
- (self._group_id or "", self._user_id or "")),
- "upload": self._slots().member_cap(
- transfers_mod.UPLOAD,
- (self._group_id or "", self._user_id or "")),
- },
- # So a client that connects mid-scan shows the indexing state
- # immediately, instead of waiting for the next periodic
- # INDEX_PROGRESS push. Never a path or filename — see
- # IndexProgress in indexer.py.
- "indexing": self._indexing_status(),
- # Current values only — not enforced from here, just shown to
- # the operator in Settings so the number on screen matches what
- # the indexer is actually doing (set_scan_settings, ops.py).
- "scan_settings": {
- "reconcile_interval_secs": self._group_ctx().get(
- "reconcile_interval_secs", DirectoryIndexer.DEFAULT_RECONCILE_SECS),
- "debounce_secs": self._group_ctx().get(
- "debounce_secs", DirectoryIndexer.DEFAULT_DEBOUNCE_SECS),
- },
- }
- if node_user_id:
- config["node_user_id"] = node_user_id
- pk_x_b64 = self._ctx.get("pk_x25519_b64")
- if pk_x_b64:
- config["node_pk_x25519"] = pk_x_b64
-
- ack = {
- "type": MNP.HANDSHAKE_ACK,
- "v": MNP_VERSION,
- "node_pk": pk_to_b64(self._ctx["sk_node"].public_key()),
- "proof": base64.b64encode(node_proof).decode(),
- "sig": base64.b64encode(
- self._ctx["sk_node"].sign(node_transcript)).decode(),
- **seal(gek, PURPOSE_ACK, MNP.HANDSHAKE_ACK, self._group_id or "", config),
- }
- self._send(ack)
- self._audit("handshake")
-
- # Someone is here now — reconcile's backstop should be prompt again
- # rather than however far its backoff had stretched while nobody
- # was connected (indexer.py DirectoryIndexer.note_activity).
- note_activity = self._group_ctx().get("note_activity")
- if note_activity:
- note_activity()
-
# ── Per-account blobs (playlists, docs/playlists.md §8.2) ────────────────
#
# The node stores bytes it cannot read and hands them back. `self._user_id`
@@ -1778,33 +1477,6 @@ class WebRTCPeerSession(
"reminder": result.get("reminder", ""),
})
- def _audit_pre_proof_fetch(self, mtype: str) -> None:
- """Record bundle access made before the GEK proof (C4)."""
- audit = self._ctx.get("audit_store")
- if not audit:
- return
- self._remote_ip = self._remote_ip or _get_remote_ip(self._pc)
- self._spawn(audit.log_event(
- user_id=getattr(self, "_pending_sub", "unknown"),
- event="pre_proof_fetch",
- ip=self._remote_ip,
- username=self._username or getattr(self, "_pending_username", ""),
- group_id=getattr(self, "_pending_group", "") or "",
- detail=mtype,
- ))
-
- def _audit_auth_failed(self, group_id: str, reason: str) -> None:
- audit = self._ctx.get("audit_store")
- if audit:
- self._remote_ip = _get_remote_ip(self._pc)
- self._spawn(audit.log_event(
- user_id="unknown",
- event="auth_failed",
- ip=self._remote_ip,
- group_id=group_id,
- detail=reason,
- ))
-
def _spawn(self, coro) -> asyncio.Task:
"""Run a coroutine in the background and hold on to it.
@@ -1854,29 +1526,6 @@ class WebRTCPeerSession(
return self._ctx["groups"].get(self._group_id) or {}
return self._ctx
- def _indexing_status(self) -> dict:
- """
- The counters of `_push_index_progress` (daemon.py), for the handshake
- ack — never a path, a filename or a root name, which stay local to the
- operator's own admin UI. Absent "progress" (context not loaded, or a
- group with no indexer at all) reads as idle rather than erroring.
- """
- progress = self._group_ctx().get("progress")
- if progress is None:
- return {"scanning": False, "scanned_bytes": 0, "total_bytes": 0,
- "files_done": 0, "files_total": 0, "kind": "",
- "root_pos": -1, "queued": 0}
- return {
- "scanning": progress.scanning,
- "scanned_bytes": progress.scanned_bytes,
- "total_bytes": progress.total_bytes,
- "files_done": progress.files_done,
- "files_total": progress.files_total,
- "kind": progress.kind,
- "root_pos": progress.root_pos,
- "queued": len(progress.queued),
- }
-
# ── Upload ceiling and transfer slots ────────────────────────────────────
def _register_peer(self) -> None: