diff options
Diffstat (limited to 'packages/meshbay-node/src')
| -rw-r--r-- | packages/meshbay-node/src/meshbay_node/transport/webrtc_server.py | 120 |
1 files changed, 120 insertions, 0 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 81db0e9..eab1cac 100644 --- a/packages/meshbay-node/src/meshbay_node/transport/webrtc_server.py +++ b/packages/meshbay-node/src/meshbay_node/transport/webrtc_server.py @@ -59,6 +59,7 @@ from meshbay_common.adminop import ( OP_DIR_DELETE, OP_FILE_DELETE, OP_INVITE_CREATE, + OP_MEMBER_REVOKE, admin_transcript, ) from meshbay_common.crypto import pk_to_b64, wrap_gek_aes @@ -171,6 +172,10 @@ def _extract_dtls_fingerprint(sdp: str) -> bytes: STREAM_SEGMENT_SIZE = 256 * 1024 +# What a client may ask for in one go, and how long the node waits for it to ask +# again before deciding nobody is watching any more. +STREAM_MAX_CREDIT = 256 +STREAM_CREDIT_TIMEOUT = 120 _H264_PROFILES = {"Baseline": "42", "Main": "4d", "High": "64", "High 10": "6e"} @@ -289,6 +294,9 @@ class WebRTCPeerSession: # Set from the roster: the key this node pinned for this account. Never # from the JWT — the hub picks what goes in there. self._pinned_pk: str = "" + # Flow control for video: how many segments the client says it can take. + self._stream_credit = 0 + self._stream_credit_evt = asyncio.Event() self._gek_challenge: bytes | None = None # Same value as the GEK challenge, but kept for the life of the connection: # a join_request is signed over it, and it must stay verifiable after the @@ -366,12 +374,16 @@ class WebRTCPeerSession: self._do_admin_response(msg) elif mtype == MNP.INVITE_CREATE: self._do_invite_create(msg) + elif mtype == MNP.MEMBER_REVOKE: + self._do_member_revoke(msg) elif mtype == MNP.KEYPAIR_BUNDLE_STORE: asyncio.ensure_future(self._do_keypair_bundle_store(msg)) elif mtype == MNP.KEYPAIR_BUNDLE_DELETE: asyncio.ensure_future(self._do_keypair_bundle_delete()) elif mtype == MNP.STREAM_REQUEST: asyncio.ensure_future(self._stream_video(msg)) + elif mtype == MNP.STREAM_MORE: + self._grant_stream_credit(msg) else: log.warning("Unknown MNP message type on DataChannel: %s", mtype) except Exception as e: @@ -999,6 +1011,69 @@ class WebRTCPeerSession: self._audit("dir_delete", rel) self._send({"type": MNP.DIR_DELETE_ACK, "v": MNP_VERSION, "dir": rel}) + def _do_member_revoke(self, msg: dict) -> None: + """ + Stop serving the group key to someone, at the operator's request. + + The same authority as an invite, and the same reason: the roster decides + who this node serves, so only a key the node pinned as an operator may + change it. Membership on the hub is not consulted — the hub can remove + someone from a group, and that stops them reaching the node at all, but + it cannot make the node forget them. + """ + user_id = str(msg.get("user_id", "")).strip() + if not user_id: + self._send({"type": "error", "detail": "Missing user_id"}) + return + if user_id == self._user_id: + # Removing yourself from your own node is not a member operation; + # it would leave the group with nobody able to invite. + self._send({"type": "error", "detail": "Cannot revoke yourself"}) + return + if not self._has_admin_authority(): + self._send({"type": "error", "detail": "No authorized key for this"}) + return + self._issue_admin_challenge(OP_MEMBER_REVOKE, user_id) + + async def _admin_exec_member_revoke( + self, pending: dict, transcript: bytes, sig: bytes, + ) -> None: + user_id = pending["subject"] + if not await self._verify_admin_sig(transcript, sig): + self._send({"type": "error", "detail": "Signature verification failed"}) + self._audit("admin_auth_failed", f"member_revoke:{user_id[:8]}") + return + + roster = self._ctx.get("roster") + if roster is None: + self._send({"type": "error", "detail": "Roster not available"}) + return + + group_id = self._group_id or "" + if not await roster.set_status(group_id, user_id, "revoked"): + self._send({"type": "error", "detail": "Not a member of this group"}) + return + + # Anyone connected right now keeps the key they already unwrapped; what + # they lose is the next one. Rotating it is the operator's call, and the + # ack says so rather than implying this undid anything already read. + peer = self._peer_registry().get(user_id) + if peer is not None: + try: + await peer.close() + except Exception: + pass + + log.info("Member revoked by %s: user=%s group=%s", + self._user_id[:8], user_id[:8], group_id[:8] or "-") + self._audit("member_revoke", user_id) + self._send({ + "type": MNP.MEMBER_REVOKE_ACK, "v": MNP_VERSION, + "user_id": user_id, + "reminder": "they still hold the current group key — rotate it with " + "meshbay-node gek-init", + }) + async def _do_keypair_bundle_delete(self) -> None: """ Withdraw our own key backup from this node. @@ -1542,6 +1617,9 @@ class WebRTCPeerSession: elif pending["op"] == OP_DIR_DELETE: asyncio.ensure_future( self._admin_exec_dir_delete(pending, transcript, sig_bytes)) + elif pending["op"] == OP_MEMBER_REVOKE: + asyncio.ensure_future( + self._admin_exec_member_revoke(pending, transcript, sig_bytes)) elif pending["op"] == OP_INVITE_CREATE: asyncio.ensure_future( self._admin_exec_invite_create(pending, transcript, sig_bytes)) @@ -1634,6 +1712,37 @@ class WebRTCPeerSession: "file_id": file_id, }) + def _grant_stream_credit(self, msg: dict) -> None: + """The client has room for more segments.""" + try: + n = int(msg.get("n", 1)) + except (TypeError, ValueError): + n = 1 + self._stream_credit += max(0, min(n, STREAM_MAX_CREDIT)) + self._stream_credit_evt.set() + + async def _await_stream_credit(self) -> bool: + """ + Block until the client has room. False if it stopped asking. + + Without this the node hands ffmpeg's entire output to the channel as + fast as it is produced, and the browser holds a four gigabyte film in a + JavaScript array while MediaSource consumes it a segment at a time. + """ + while self._stream_credit <= 0: + self._stream_credit_evt.clear() + try: + await asyncio.wait_for(self._stream_credit_evt.wait(), + timeout=STREAM_CREDIT_TIMEOUT) + except asyncio.TimeoutError: + log.info("Stream stalled: no credit from peer=%s", + (self._user_id or "?")[:8]) + return False + if self._channel is None or self._channel.readyState != "open": + return False + self._stream_credit -= 1 + return True + async def _stream_video(self, msg: dict) -> None: """Stream a video file as fMP4 segments via MSE-compatible output.""" # One ffmpeg per request with no cap lets any member exhaust the node's @@ -1692,9 +1801,20 @@ class WebRTCPeerSession: "duration": duration, }) + # A client that says nothing gets the old behaviour, which is why this + # defaults to unlimited rather than to zero: a stream that waits for + # credit from a peer that will never send any is a stream that hangs. + try: + self._stream_credit = int(msg.get("credits", 0) or 0) + except (TypeError, ValueError): + self._stream_credit = 0 + paced = self._stream_credit > 0 + index = 0 try: while True: + if paced and not await self._await_stream_credit(): + break data = await proc.stdout.read(STREAM_SEGMENT_SIZE) if not data: break |