diff options
| author | Christophe Besson <cbesson@gmail.com> | 2026-09-24 10:36:54 +0200 |
|---|---|---|
| committer | Christophe Besson <cbesson@gmail.com> | 2026-09-24 16:45:38 +0200 |
| commit | 41b814483cf990d855188899332b90fc525135b7 (patch) | |
| tree | 7049a38f8907f14f76addc3b7faf9922016ef36e /packages/meshbay-node/src/meshbay_node/transport/webrtc_server.py | |
| parent | 600a6dfe3779b9cd82d4e0cb94c6f6e7df21ab20 (diff) | |
| download | meshbay-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.py | 357 |
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: |