""" Integration test: WebRTC DataChannel transport for browser clients. Phase 9 milestone 9.1 — spike: validate aiortc WebRTC DataChannel works for MNP protocol exchange (handshake, index_sync, file_request, file_chunk). Uses local loopback (no STUN/ICE needed for localhost). """ import asyncio import base64 import hashlib import hmac import os import struct import time import jwt import msgpack import pytest from cryptography.hazmat.primitives import serialization from cryptography.hazmat.primitives.asymmetric.ed25519 import Ed25519PrivateKey from aiortc import RTCPeerConnection, RTCSessionDescription from meshbay_common import MNP_VERSION from meshbay_common.crypto import ( generate_gek, pk_to_b64, wrap_gek, wrap_gek_aes, unwrap_gek, ) from meshbay_common.webcrypto import chunk_key_aes, decrypt_chunk_aes from meshbay_common.protocol import MNP from meshbay_node.bundle_store import BundleStore from meshbay_node.indexer import DirectoryIndexer from meshbay_node.transport.webrtc_server import WebRTCTransport @pytest.fixture def sk_node(): return Ed25519PrivateKey.generate() @pytest.fixture def sk_hub(): return Ed25519PrivateKey.generate() @pytest.fixture def sk_user(): return Ed25519PrivateKey.generate() @pytest.fixture def gek(): return generate_gek() @pytest.fixture def shared_dir(tmp_path): d = tmp_path / "shared" d.mkdir() (d / "test.bin").write_bytes(os.urandom(2048)) (d / "hello.txt").write_bytes(b"hello webrtc " * 50) return d def _hub_pk_pem(sk_hub): return sk_hub.public_key().public_bytes( serialization.Encoding.PEM, serialization.PublicFormat.SubjectPublicKeyInfo) def _make_jwt(sk_hub, groups=None, pk_user="test"): sk_pem = sk_hub.private_bytes( serialization.Encoding.PEM, serialization.PrivateFormat.PKCS8, serialization.NoEncryption(), ) now = int(time.time()) return jwt.encode({ "iss": "test-hub", "sub": "user-001", "pk_user": pk_user, "hub_id": "test-hub", "jti": "test-jti-webrtc", "iat": now, "exp": now + 3600, "groups": groups or [], }, sk_pem, algorithm="EdDSA") def _pack(obj: dict) -> bytes: data = msgpack.packb(obj, use_bin_type=True) return struct.pack(">I", len(data)) + data def _unpack(raw: bytes) -> dict: length = struct.unpack(">I", raw[:4])[0] return msgpack.unpackb(raw[4:4 + length], raw=False) def _extract_dtls_fp(sdp: str) -> bytes: for line in sdp.splitlines(): if line.startswith("a=fingerprint:sha-256 "): return bytes.fromhex(line.split(" ", 1)[1].replace(":", "")) return b"" async def _handshake_with_gek_proof(channel, received, sk_hub, gek, groups=None, browser_pc=None): """Send handshake, handle GEK challenge, return handshake_ack.""" token = _make_jwt(sk_hub, groups=groups) channel.send(_pack({ "type": MNP.HANDSHAKE, "v": MNP_VERSION, "token": token, })) msg = await asyncio.wait_for(received.get(), timeout=5.0) if msg["type"] == MNP.HANDSHAKE_CHALLENGE: nonce = base64.b64decode(msg["nonce"]) offer_fp = b"" answer_fp = b"" if browser_pc: offer_fp = _extract_dtls_fp(browser_pc.localDescription.sdp) answer_fp = _extract_dtls_fp(browser_pc.remoteDescription.sdp) proof = hmac.new(gek, nonce + offer_fp + answer_fp, hashlib.sha256).digest() channel.send(_pack({ "type": MNP.HANDSHAKE_RESPONSE, "v": MNP_VERSION, "proof": base64.b64encode(proof).decode(), })) msg = await asyncio.wait_for(received.get(), timeout=5.0) assert msg["type"] == MNP.HANDSHAKE_ACK return msg async def _setup_peer(transport, sk_hub, gek, peer_id, jwt_sub="user-001", sk_user=None): """Create a peer connection, perform handshake with GEK proof, return (pc, channel, queue).""" pc = RTCPeerConnection() q = asyncio.Queue() buf = bytearray() ch = pc.createDataChannel("mnp") ready = asyncio.Event() @ch.on("open") def on_open(): ready.set() @ch.on("message") def on_msg(message): if isinstance(message, str): message = message.encode() buf.extend(message) while len(buf) >= 4: length = struct.unpack(">I", buf[:4])[0] if len(buf) < 4 + length: break msg_bytes = bytes(buf[4:4 + length]) del buf[:4 + length] q.put_nowait(msgpack.unpackb(msg_bytes, raw=False)) offer = await pc.createOffer() await pc.setLocalDescription(offer) answer_sdp, _ = await transport.handle_offer(pc.localDescription.sdp, peer_id) await pc.setRemoteDescription(RTCSessionDescription(sdp=answer_sdp, type="answer")) await asyncio.wait_for(ready.wait(), timeout=5.0) pk_user = "test" if sk_user: pk_user = base64.b64encode( sk_user.public_key().public_bytes( serialization.Encoding.Raw, serialization.PublicFormat.Raw) ).decode() sk_h_pem = sk_hub.private_bytes( serialization.Encoding.PEM, serialization.PrivateFormat.PKCS8, serialization.NoEncryption(), ) now = int(time.time()) token = jwt.encode({ "iss": "test-hub", "sub": jwt_sub, "pk_user": pk_user, "hub_id": "test-hub", "jti": f"jti-{peer_id}", "iat": now, "exp": now + 3600, "groups": [], }, sk_h_pem, algorithm="EdDSA") ch.send(_pack({"type": MNP.HANDSHAKE, "v": MNP_VERSION, "token": token})) msg = await asyncio.wait_for(q.get(), timeout=5.0) if msg["type"] == MNP.HANDSHAKE_CHALLENGE: nonce = base64.b64decode(msg["nonce"]) offer_fp = _extract_dtls_fp(pc.localDescription.sdp) answer_fp = _extract_dtls_fp(pc.remoteDescription.sdp) proof = hmac.new(gek, nonce + offer_fp + answer_fp, hashlib.sha256).digest() ch.send(_pack({ "type": MNP.HANDSHAKE_RESPONSE, "v": MNP_VERSION, "proof": base64.b64encode(proof).decode(), })) msg = await asyncio.wait_for(q.get(), timeout=5.0) assert msg["type"] == MNP.HANDSHAKE_ACK return pc, ch, q @pytest.mark.asyncio async def test_webrtc_datachannel_handshake(sk_node, sk_hub, gek, shared_dir): """WebRTC DataChannel: browser sends MNP handshake, node responds with handshake_ack.""" hub_pk_pem = _hub_pk_pem(sk_hub) indexer = DirectoryIndexer(root=shared_dir, group_id="g", sk_node=sk_node, gek=gek) await indexer.initial_scan() transport = WebRTCTransport( sk_node=sk_node, hub_pk_pem=hub_pk_pem, gek=gek, shared_root=shared_dir, index=indexer.index, stun_servers=[], ) browser_pc = RTCPeerConnection() received = asyncio.Queue() channel = browser_pc.createDataChannel("mnp") @channel.on("message") def on_msg(message): if isinstance(message, str): message = message.encode() received.put_nowait(_unpack(message)) offer = await browser_pc.createOffer() await browser_pc.setLocalDescription(offer) answer_sdp, ice_candidates = await transport.handle_offer( browser_pc.localDescription.sdp, "peer-001") answer = RTCSessionDescription(sdp=answer_sdp, type="answer") await browser_pc.setRemoteDescription(answer) await asyncio.sleep(0.5) msg = await _handshake_with_gek_proof(channel, received, sk_hub, gek, browser_pc=browser_pc) assert msg["v"] == MNP_VERSION assert "node_pk" in msg await browser_pc.close() await transport.close_all() @pytest.mark.asyncio async def test_webrtc_datachannel_file_transfer(sk_node, sk_hub, gek, shared_dir): """WebRTC DataChannel: full file transfer — handshake, index, fetch chunk, decrypt.""" hub_pk_pem = _hub_pk_pem(sk_hub) indexer = DirectoryIndexer(root=shared_dir, group_id="g", sk_node=sk_node, gek=gek) await indexer.initial_scan() transport = WebRTCTransport( sk_node=sk_node, hub_pk_pem=hub_pk_pem, gek=gek, shared_root=shared_dir, index=indexer.index, stun_servers=[], ) browser_pc = RTCPeerConnection() received = asyncio.Queue() channel = browser_pc.createDataChannel("mnp") channel_ready = asyncio.Event() @channel.on("open") def on_open(): channel_ready.set() buf = bytearray() @channel.on("message") def on_msg(message): if isinstance(message, str): message = message.encode() buf.extend(message) while len(buf) >= 4: length = struct.unpack(">I", buf[:4])[0] if len(buf) < 4 + length: break msg_bytes = bytes(buf[4:4 + length]) del buf[:4 + length] received.put_nowait(msgpack.unpackb(msg_bytes, raw=False)) offer = await browser_pc.createOffer() await browser_pc.setLocalDescription(offer) answer_sdp, _ = await transport.handle_offer( browser_pc.localDescription.sdp, "peer-002") await browser_pc.setRemoteDescription( RTCSessionDescription(sdp=answer_sdp, type="answer")) await asyncio.wait_for(channel_ready.wait(), timeout=5.0) # 1) Handshake with GEK proof ack = await _handshake_with_gek_proof(channel, received, sk_hub, gek, browser_pc=browser_pc) assert ack["type"] == MNP.HANDSHAKE_ACK # 2) Request index channel.send(_pack({"type": MNP.INDEX_SYNC, "v": MNP_VERSION})) idx_msg = await asyncio.wait_for(received.get(), timeout=5.0) assert idx_msg["type"] == MNP.INDEX_SYNC assert "entries" in idx_msg assert len(idx_msg["entries"]) > 0 # 3) Request file chunk entry = next(e for e in indexer.index.entries if e.name == "test.bin") channel.send(_pack({ "type": MNP.FILE_REQUEST, "v": MNP_VERSION, "file_id": entry.id, "chunk_index": 0, })) chunk_msg = await asyncio.wait_for(received.get(), timeout=5.0) assert chunk_msg["type"] == MNP.FILE_CHUNK # 4) Verify and decrypt (binary fields — no base64, minimal envelope) ct = chunk_msg["ct"] nonce = chunk_msg["nonce"] file_hash = bytes.fromhex(entry.id) ckey = chunk_key_aes(gek, file_hash, 0) plaintext = decrypt_chunk_aes(ckey, nonce, ct) original = (shared_dir / "test.bin").read_bytes() assert plaintext == original await browser_pc.close() await transport.close_all() @pytest.mark.asyncio async def test_webrtc_invalid_jwt_rejected(sk_node, sk_hub, gek, shared_dir): """WebRTC DataChannel: invalid JWT is rejected with error.""" hub_pk_pem = _hub_pk_pem(sk_hub) indexer = DirectoryIndexer(root=shared_dir, group_id="g", sk_node=sk_node, gek=gek) await indexer.initial_scan() transport = WebRTCTransport( sk_node=sk_node, hub_pk_pem=hub_pk_pem, gek=gek, shared_root=shared_dir, index=indexer.index, stun_servers=[], ) browser_pc = RTCPeerConnection() received = asyncio.Queue() channel = browser_pc.createDataChannel("mnp") @channel.on("message") def on_msg(message): if isinstance(message, str): message = message.encode() received.put_nowait(_unpack(message)) offer = await browser_pc.createOffer() await browser_pc.setLocalDescription(offer) answer_sdp, _ = await transport.handle_offer( browser_pc.localDescription.sdp, "peer-003") await browser_pc.setRemoteDescription( RTCSessionDescription(sdp=answer_sdp, type="answer")) await asyncio.sleep(0.5) channel.send(_pack({ "type": MNP.HANDSHAKE, "v": MNP_VERSION, "token": "invalid.jwt.token", })) msg = await asyncio.wait_for(received.get(), timeout=5.0) assert msg["type"] == "error" assert "JWT" in msg["detail"] or "Invalid" in msg["detail"] await browser_pc.close() await transport.close_all() @pytest.mark.asyncio async def test_webrtc_request_before_handshake_rejected(sk_node, sk_hub, gek, shared_dir): """WebRTC DataChannel: request without handshake is rejected.""" hub_pk_pem = _hub_pk_pem(sk_hub) indexer = DirectoryIndexer(root=shared_dir, group_id="g", sk_node=sk_node, gek=gek) await indexer.initial_scan() transport = WebRTCTransport( sk_node=sk_node, hub_pk_pem=hub_pk_pem, gek=gek, shared_root=shared_dir, index=indexer.index, stun_servers=[], ) browser_pc = RTCPeerConnection() received = asyncio.Queue() channel = browser_pc.createDataChannel("mnp") @channel.on("message") def on_msg(message): if isinstance(message, str): message = message.encode() received.put_nowait(_unpack(message)) offer = await browser_pc.createOffer() await browser_pc.setLocalDescription(offer) answer_sdp, _ = await transport.handle_offer( browser_pc.localDescription.sdp, "peer-004") await browser_pc.setRemoteDescription( RTCSessionDescription(sdp=answer_sdp, type="answer")) await asyncio.sleep(0.5) channel.send(_pack({"type": MNP.INDEX_SYNC, "v": MNP_VERSION})) msg = await asyncio.wait_for(received.get(), timeout=5.0) assert msg["type"] == "error" assert "Handshake required" in msg["detail"] await browser_pc.close() await transport.close_all() @pytest.mark.asyncio async def test_webrtc_chat_send_and_history(sk_node, sk_hub, gek, shared_dir, tmp_path): """WebRTC DataChannel: send chat message, then retrieve history.""" from meshbay_node.chat.store import ChatStore hub_pk_pem = _hub_pk_pem(sk_hub) indexer = DirectoryIndexer(root=shared_dir, group_id="g", sk_node=sk_node, gek=gek) await indexer.initial_scan() chat_store = ChatStore(db_path=tmp_path / "chat_test.db") await chat_store.open() transport = WebRTCTransport( sk_node=sk_node, hub_pk_pem=hub_pk_pem, gek=gek, shared_root=shared_dir, index=indexer.index, stun_servers=[], ) transport._ctx["chat_store"] = chat_store browser_pc, channel, received = await _setup_peer( transport, sk_hub, gek, "peer-chat") channel.send(_pack({ "type": MNP.CHAT_MESSAGE, "v": MNP_VERSION, "payload": "hello from browser", })) chat_ack = await asyncio.wait_for(received.get(), timeout=5.0) assert chat_ack["type"] == "ack" await asyncio.sleep(0.2) channel.send(_pack({ "type": MNP.CHAT_HISTORY, "v": MNP_VERSION, "since": 0, "limit": 50, })) hist = await asyncio.wait_for(received.get(), timeout=5.0) assert hist["type"] == MNP.CHAT_HISTORY_RESPONSE assert len(hist["messages"]) == 1 assert hist["messages"][0]["payload"] == "hello from browser" assert hist["messages"][0]["sender_id"] == "user-001" await chat_store.close() await browser_pc.close() await transport.close_all() @pytest.mark.asyncio async def test_webrtc_chat_history_no_store(sk_node, sk_hub, gek, shared_dir): """WebRTC DataChannel: chat history without chat_store returns empty list.""" hub_pk_pem = _hub_pk_pem(sk_hub) indexer = DirectoryIndexer(root=shared_dir, group_id="g", sk_node=sk_node, gek=gek) await indexer.initial_scan() transport = WebRTCTransport( sk_node=sk_node, hub_pk_pem=hub_pk_pem, gek=gek, shared_root=shared_dir, index=indexer.index, stun_servers=[], ) browser_pc, channel, received = await _setup_peer( transport, sk_hub, gek, "peer-no-store") channel.send(_pack({ "type": MNP.CHAT_HISTORY, "v": MNP_VERSION, "since": 0, "limit": 50, })) hist = await asyncio.wait_for(received.get(), timeout=5.0) assert hist["type"] == MNP.CHAT_HISTORY_RESPONSE assert hist["messages"] == [] await browser_pc.close() await transport.close_all() @pytest.mark.asyncio async def test_webrtc_chat_broadcast(sk_node, sk_hub, gek, shared_dir, tmp_path): """WebRTC DataChannel: chat message from peer A is broadcast to peer B.""" from meshbay_node.chat.store import ChatStore hub_pk_pem = _hub_pk_pem(sk_hub) indexer = DirectoryIndexer(root=shared_dir, group_id="g", sk_node=sk_node, gek=gek) await indexer.initial_scan() chat_store = ChatStore(db_path=tmp_path / "chat_bc.db") await chat_store.open() transport = WebRTCTransport( sk_node=sk_node, hub_pk_pem=hub_pk_pem, gek=gek, shared_root=shared_dir, index=indexer.index, stun_servers=[], ) transport._ctx["chat_store"] = chat_store pc_a, ch_a, q_a = await _setup_peer(transport, sk_hub, gek, "peer-A", "user-A") pc_b, ch_b, q_b = await _setup_peer(transport, sk_hub, gek, "peer-B", "user-B") ch_a.send(_pack({ "type": MNP.CHAT_MESSAGE, "v": MNP_VERSION, "payload": "hi from A", })) ack_a = await asyncio.wait_for(q_a.get(), timeout=5.0) assert ack_a["type"] == "ack" broadcast = await asyncio.wait_for(q_b.get(), timeout=5.0) assert broadcast["type"] == MNP.CHAT_MESSAGE assert broadcast["sender_id"] == "user-A" assert broadcast["payload"] == "hi from A" await chat_store.close() await pc_a.close() await pc_b.close() await transport.close_all() @pytest.mark.asyncio async def test_webrtc_group_membership_enforced(sk_node, sk_hub, gek, shared_dir): """WebRTC DataChannel: JWT without matching group claim is rejected.""" hub_pk_pem = _hub_pk_pem(sk_hub) indexer = DirectoryIndexer(root=shared_dir, group_id="g", sk_node=sk_node, gek=gek) await indexer.initial_scan() transport = WebRTCTransport( sk_node=sk_node, hub_pk_pem=hub_pk_pem, gek=gek, shared_root=shared_dir, index=indexer.index, stun_servers=[], ) browser_pc = RTCPeerConnection() received = asyncio.Queue() channel = browser_pc.createDataChannel("mnp") @channel.on("message") def on_msg(message): if isinstance(message, str): message = message.encode() received.put_nowait(_unpack(message)) offer = await browser_pc.createOffer() await browser_pc.setLocalDescription(offer) answer_sdp, _ = await transport.handle_offer( browser_pc.localDescription.sdp, "peer-group-test") await browser_pc.setRemoteDescription( RTCSessionDescription(sdp=answer_sdp, type="answer")) await asyncio.sleep(0.5) token = _make_jwt(sk_hub, groups=["other-group"]) channel.send(_pack({ "type": MNP.HANDSHAKE, "v": MNP_VERSION, "token": token, "group_id": "my-group", })) msg = await asyncio.wait_for(received.get(), timeout=5.0) assert msg["type"] == "error" assert "Not a member" in msg["detail"] await browser_pc.close() await transport.close_all() @pytest.mark.asyncio async def test_webrtc_peer_cleanup_on_close(sk_node, sk_hub, gek, shared_dir): """WebRTC DataChannel: peer removed from _peers dict on session close.""" hub_pk_pem = _hub_pk_pem(sk_hub) indexer = DirectoryIndexer(root=shared_dir, group_id="g", sk_node=sk_node, gek=gek) await indexer.initial_scan() transport = WebRTCTransport( sk_node=sk_node, hub_pk_pem=hub_pk_pem, gek=gek, shared_root=shared_dir, index=indexer.index, stun_servers=[], ) browser_pc, channel, received = await _setup_peer( transport, sk_hub, gek, "peer-cleanup") assert "user-001" in transport._ctx["_peers"] assert transport.active_peers == 1 await transport.close_peer("peer-cleanup") assert "user-001" not in transport._ctx["_peers"] assert transport.active_peers == 0 await browser_pc.close() @pytest.mark.asyncio async def test_webrtc_stream_segment_missing_file(sk_node, sk_hub, gek, shared_dir): """WebRTC DataChannel: stream_segment for non-existent file returns error.""" hub_pk_pem = _hub_pk_pem(sk_hub) indexer = DirectoryIndexer(root=shared_dir, group_id="g", sk_node=sk_node, gek=gek) await indexer.initial_scan() transport = WebRTCTransport( sk_node=sk_node, hub_pk_pem=hub_pk_pem, gek=gek, shared_root=shared_dir, index=indexer.index, stun_servers=[], ) browser_pc, channel, received = await _setup_peer( transport, sk_hub, gek, "peer-stream") channel.send(_pack({ "type": MNP.STREAM_SEGMENT, "v": MNP_VERSION, "file_id": "nonexistent-file-id", "segment_index": 0, "segment_duration": 4, })) msg = await asyncio.wait_for(received.get(), timeout=5.0) assert msg["type"] == "error" assert "not found" in msg["detail"].lower() await browser_pc.close() await transport.close_all() @pytest.mark.asyncio async def test_webrtc_wrong_gek_proof_rejected(sk_node, sk_hub, gek, shared_dir): """WebRTC DataChannel: wrong GEK proof is rejected — hub admin can't fake membership.""" hub_pk_pem = _hub_pk_pem(sk_hub) indexer = DirectoryIndexer(root=shared_dir, group_id="g", sk_node=sk_node, gek=gek) await indexer.initial_scan() transport = WebRTCTransport( sk_node=sk_node, hub_pk_pem=hub_pk_pem, gek=gek, shared_root=shared_dir, index=indexer.index, stun_servers=[], ) browser_pc = RTCPeerConnection() received = asyncio.Queue() channel = browser_pc.createDataChannel("mnp") @channel.on("message") def on_msg(message): if isinstance(message, str): message = message.encode() received.put_nowait(_unpack(message)) offer = await browser_pc.createOffer() await browser_pc.setLocalDescription(offer) answer_sdp, _ = await transport.handle_offer( browser_pc.localDescription.sdp, "peer-fake") await browser_pc.setRemoteDescription( RTCSessionDescription(sdp=answer_sdp, type="answer")) await asyncio.sleep(0.5) token = _make_jwt(sk_hub) channel.send(_pack({ "type": MNP.HANDSHAKE, "v": MNP_VERSION, "token": token, })) challenge = await asyncio.wait_for(received.get(), timeout=5.0) assert challenge["type"] == MNP.HANDSHAKE_CHALLENGE fake_gek = os.urandom(32) nonce = base64.b64decode(challenge["nonce"]) offer_fp = _extract_dtls_fp(browser_pc.localDescription.sdp) answer_fp = _extract_dtls_fp(browser_pc.remoteDescription.sdp) bad_proof = hmac.new(fake_gek, nonce + offer_fp + answer_fp, hashlib.sha256).digest() channel.send(_pack({ "type": MNP.HANDSHAKE_RESPONSE, "v": MNP_VERSION, "proof": base64.b64encode(bad_proof).decode(), })) msg = await asyncio.wait_for(received.get(), timeout=5.0) assert msg["type"] == "error" assert "GEK proof failed" in msg["detail"] await browser_pc.close() await transport.close_all() @pytest.mark.asyncio async def test_webrtc_dtls_channel_binding_detects_mitm(sk_node, sk_hub, gek, shared_dir): """WebRTC: DTLS channel binding detects fingerprint substitution (simulated MitM).""" hub_pk_pem = _hub_pk_pem(sk_hub) indexer = DirectoryIndexer(root=shared_dir, group_id="g", sk_node=sk_node, gek=gek) await indexer.initial_scan() transport = WebRTCTransport( sk_node=sk_node, hub_pk_pem=hub_pk_pem, gek=gek, shared_root=shared_dir, index=indexer.index, stun_servers=[], ) browser_pc = RTCPeerConnection() received = asyncio.Queue() channel = browser_pc.createDataChannel("mnp") @channel.on("message") def on_msg(message): if isinstance(message, str): message = message.encode() received.put_nowait(_unpack(message)) offer = await browser_pc.createOffer() await browser_pc.setLocalDescription(offer) answer_sdp, _ = await transport.handle_offer( browser_pc.localDescription.sdp, "peer-mitm") await browser_pc.setRemoteDescription( RTCSessionDescription(sdp=answer_sdp, type="answer")) await asyncio.sleep(0.5) token = _make_jwt(sk_hub) channel.send(_pack({ "type": MNP.HANDSHAKE, "v": MNP_VERSION, "token": token, })) challenge = await asyncio.wait_for(received.get(), timeout=5.0) assert challenge["type"] == MNP.HANDSHAKE_CHALLENGE nonce = base64.b64decode(challenge["nonce"]) # Correct GEK but fake fingerprints — simulates MitM substituting DTLS certs fake_fp = os.urandom(32) proof = hmac.new(gek, nonce + fake_fp + fake_fp, hashlib.sha256).digest() channel.send(_pack({ "type": MNP.HANDSHAKE_RESPONSE, "v": MNP_VERSION, "proof": base64.b64encode(proof).decode(), })) msg = await asyncio.wait_for(received.get(), timeout=5.0) assert msg["type"] == "error" assert "GEK proof failed" in msg["detail"] await browser_pc.close() await transport.close_all() @pytest.mark.asyncio async def test_webrtc_admin_challenge_response(sk_node, sk_hub, gek, shared_dir): """WebRTC DataChannel: admin file delete requires Ed25519 challenge-response.""" hub_pk_pem = _hub_pk_pem(sk_hub) indexer = DirectoryIndexer(root=shared_dir, group_id="g", sk_node=sk_node, gek=gek) await indexer.initial_scan() sk_admin = Ed25519PrivateKey.generate() transport = WebRTCTransport( sk_node=sk_node, hub_pk_pem=hub_pk_pem, gek=gek, shared_root=shared_dir, index=indexer.index, stun_servers=[], ) transport._ctx["admin_pk_ed25519"] = sk_admin.public_key() transport._ctx["node_user_id"] = "user-001" browser_pc, channel, received = await _setup_peer( transport, sk_hub, gek, "peer-admin") entry = indexer.index.entries[0] channel.send(_pack({ "type": MNP.FILE_DELETE, "v": MNP_VERSION, "file_id": entry.id, })) challenge_msg = await asyncio.wait_for(received.get(), timeout=5.0) assert challenge_msg["type"] == MNP.ADMIN_CHALLENGE assert challenge_msg["file_id"] == entry.id challenge = base64.b64decode(challenge_msg["challenge"]) signature = sk_admin.sign(challenge) channel.send(_pack({ "type": MNP.ADMIN_RESPONSE, "v": MNP_VERSION, "file_id": entry.id, "signature": base64.b64encode(signature).decode(), })) ack = await asyncio.wait_for(received.get(), timeout=5.0) assert ack["type"] == MNP.FILE_DELETE_ACK assert ack["file_id"] == entry.id assert indexer.index.get_entry(entry.id) is None await browser_pc.close() await transport.close_all() @pytest.mark.asyncio async def test_webrtc_admin_bad_signature_rejected(sk_node, sk_hub, gek, shared_dir): """WebRTC DataChannel: wrong Ed25519 signature is rejected — hub can't fake admin.""" hub_pk_pem = _hub_pk_pem(sk_hub) indexer = DirectoryIndexer(root=shared_dir, group_id="g", sk_node=sk_node, gek=gek) await indexer.initial_scan() sk_admin = Ed25519PrivateKey.generate() sk_attacker = Ed25519PrivateKey.generate() transport = WebRTCTransport( sk_node=sk_node, hub_pk_pem=hub_pk_pem, gek=gek, shared_root=shared_dir, index=indexer.index, stun_servers=[], ) transport._ctx["admin_pk_ed25519"] = sk_admin.public_key() transport._ctx["node_user_id"] = "user-001" browser_pc, channel, received = await _setup_peer( transport, sk_hub, gek, "peer-attacker") entry = indexer.index.entries[0] channel.send(_pack({ "type": MNP.FILE_DELETE, "v": MNP_VERSION, "file_id": entry.id, })) challenge_msg = await asyncio.wait_for(received.get(), timeout=5.0) assert challenge_msg["type"] == MNP.ADMIN_CHALLENGE challenge = base64.b64decode(challenge_msg["challenge"]) bad_sig = sk_attacker.sign(challenge) channel.send(_pack({ "type": MNP.ADMIN_RESPONSE, "v": MNP_VERSION, "file_id": entry.id, "signature": base64.b64encode(bad_sig).decode(), })) msg = await asyncio.wait_for(received.get(), timeout=5.0) assert msg["type"] == "error" assert "signature" in msg["detail"].lower() or "verification" in msg["detail"].lower() assert indexer.index.get_entry(entry.id) is not None await browser_pc.close() await transport.close_all() @pytest.mark.asyncio async def test_webrtc_stream_request_missing_file(sk_node, sk_hub, gek, shared_dir): """WebRTC DataChannel: stream_request for non-existent file returns error.""" hub_pk_pem = _hub_pk_pem(sk_hub) indexer = DirectoryIndexer(root=shared_dir, group_id="g", sk_node=sk_node, gek=gek) await indexer.initial_scan() transport = WebRTCTransport( sk_node=sk_node, hub_pk_pem=hub_pk_pem, gek=gek, shared_root=shared_dir, index=indexer.index, stun_servers=[], ) browser_pc, channel, received = await _setup_peer( transport, sk_hub, gek, "peer-mse") channel.send(_pack({ "type": MNP.STREAM_REQUEST, "v": MNP_VERSION, "file_id": "nonexistent-file-id", })) msg = await asyncio.wait_for(received.get(), timeout=5.0) assert msg["type"] == "error" assert "not found" in msg["detail"].lower() await browser_pc.close() await transport.close_all() @pytest.mark.asyncio async def test_webrtc_uploader_delete_requires_challenge(sk_node, sk_hub, gek, shared_dir): """Uploader must prove Ed25519 key ownership to delete — no uploader shortcut.""" hub_pk_pem = _hub_pk_pem(sk_hub) indexer = DirectoryIndexer(root=shared_dir, group_id="g", sk_node=sk_node, gek=gek) await indexer.initial_scan() sk_uploader = Ed25519PrivateKey.generate() pk_uploader_b64 = base64.b64encode( sk_uploader.public_key().public_bytes( serialization.Encoding.Raw, serialization.PublicFormat.Raw) ).decode() transport = WebRTCTransport( sk_node=sk_node, hub_pk_pem=hub_pk_pem, gek=gek, shared_root=shared_dir, index=indexer.index, stun_servers=[], ) # No admin_pk configured — only uploader_pk should authorize deletion browser_pc, channel, received = await _setup_peer( transport, sk_hub, gek, "peer-uploader-del", sk_user=sk_uploader) # Tag an existing entry with the uploader's public key entry = indexer.index.entries[0] entry.uploader_id = "user-001" entry.uploader_pk = pk_uploader_b64 # Request deletion — should get a challenge (no shortcut) channel.send(_pack({ "type": MNP.FILE_DELETE, "v": MNP_VERSION, "file_id": entry.id, })) challenge_msg = await asyncio.wait_for(received.get(), timeout=5.0) assert challenge_msg["type"] == MNP.ADMIN_CHALLENGE assert challenge_msg["file_id"] == entry.id # Sign with uploader's Ed25519 key challenge = base64.b64decode(challenge_msg["challenge"]) signature = sk_uploader.sign(challenge) channel.send(_pack({ "type": MNP.ADMIN_RESPONSE, "v": MNP_VERSION, "file_id": entry.id, "signature": base64.b64encode(signature).decode(), })) ack = await asyncio.wait_for(received.get(), timeout=5.0) assert ack["type"] == MNP.FILE_DELETE_ACK assert ack["file_id"] == entry.id # Verify file was removed from index assert indexer.index.get_entry(entry.id) is None await browser_pc.close() await transport.close_all() @pytest.mark.asyncio async def test_webrtc_uploader_impersonation_blocked(sk_node, sk_hub, gek, shared_dir): """Hub-forged JWT with same sub cannot delete — wrong Ed25519 key is rejected.""" hub_pk_pem = _hub_pk_pem(sk_hub) indexer = DirectoryIndexer(root=shared_dir, group_id="g", sk_node=sk_node, gek=gek) await indexer.initial_scan() # User A uploaded the file sk_user_a = Ed25519PrivateKey.generate() pk_a_b64 = base64.b64encode( sk_user_a.public_key().public_bytes( serialization.Encoding.Raw, serialization.PublicFormat.Raw) ).decode() # User B is the attacker (different Ed25519 key, but hub forges JWT with same sub) sk_user_b = Ed25519PrivateKey.generate() transport = WebRTCTransport( sk_node=sk_node, hub_pk_pem=hub_pk_pem, gek=gek, shared_root=shared_dir, index=indexer.index, stun_servers=[], ) # No admin_pk — only uploader_pk matters # Tag entry with user A's public key entry = indexer.index.entries[0] entry.uploader_id = "user-001" entry.uploader_pk = pk_a_b64 # Connect as user B (same jwt_sub "user-001" via hub forgery, but B's Ed25519 key) browser_pc, channel, received = await _setup_peer( transport, sk_hub, gek, "peer-impersonator", jwt_sub="user-001", sk_user=sk_user_b) # Request deletion — should get a challenge channel.send(_pack({ "type": MNP.FILE_DELETE, "v": MNP_VERSION, "file_id": entry.id, })) challenge_msg = await asyncio.wait_for(received.get(), timeout=5.0) assert challenge_msg["type"] == MNP.ADMIN_CHALLENGE # Sign with user B's key (wrong key) challenge = base64.b64decode(challenge_msg["challenge"]) bad_sig = sk_user_b.sign(challenge) channel.send(_pack({ "type": MNP.ADMIN_RESPONSE, "v": MNP_VERSION, "file_id": entry.id, "signature": base64.b64encode(bad_sig).decode(), })) msg = await asyncio.wait_for(received.get(), timeout=5.0) assert msg["type"] == "error" assert "verification" in msg["detail"].lower() or "signature" in msg["detail"].lower() # File must still exist in the index assert indexer.index.get_entry(entry.id) is not None await browser_pc.close() await transport.close_all() # ── GEK bundle P2P exchange tests ────────────────────────────────────────── @pytest.fixture def x25519_keypair(): from cryptography.hazmat.primitives.asymmetric.x25519 import X25519PrivateKey sk = X25519PrivateKey.generate() sk_raw = sk.private_bytes( serialization.Encoding.Raw, serialization.PrivateFormat.Raw, serialization.NoEncryption()) pk_raw = sk.public_key().public_bytes( serialization.Encoding.Raw, serialization.PublicFormat.Raw) return sk_raw, pk_raw @pytest.mark.asyncio async def test_gek_bundle_store_and_fetch(sk_node, sk_hub, gek, shared_dir, tmp_path, x25519_keypair): """GEK bundle stored on node via DataChannel, then fetched during handshake.""" hub_pk_pem = _hub_pk_pem(sk_hub) indexer = DirectoryIndexer(root=shared_dir, group_id="g", sk_node=sk_node, gek=gek) await indexer.initial_scan() bundle_store = BundleStore(db_path=tmp_path / "bundles.db") await bundle_store.open() transport = WebRTCTransport( sk_node=sk_node, hub_pk_pem=hub_pk_pem, gek=gek, shared_root=shared_dir, index=indexer.index, stun_servers=[], ) transport._ctx["bundle_store"] = bundle_store # Connect as admin and store a GEK bundle for user-002 pc_admin, ch_admin, q_admin = await _setup_peer( transport, sk_hub, gek, "peer-admin") sk_x_raw, pk_x_raw = x25519_keypair bundle = wrap_gek(gek, pk_x_raw) ch_admin.send(_pack({ "type": MNP.GEK_BUNDLE_STORE, "v": MNP_VERSION, "user_id": "user-002", "group_id": "g", "pk_eph_b64": bundle["pk_eph_b64"], "nonce_b64": bundle["nonce_b64"], "wrapped_b64": bundle["wrapped_b64"], })) ack = await asyncio.wait_for(q_admin.get(), timeout=5.0) assert ack["type"] == "ack" assert ack["detail"] == "gek_bundle_stored" # Verify bundle was persisted stored = await bundle_store.fetch("g", "user-002") assert stored is not None assert stored["pk_eph_b64"] == bundle["pk_eph_b64"] # Unwrap to verify it's correct recovered = unwrap_gek(stored, sk_x_raw, pk_x_raw) assert recovered == gek await bundle_store.close() await pc_admin.close() await transport.close_all() @pytest.mark.asyncio async def test_gek_bundle_fetch_during_handshake(sk_node, sk_hub, gek, shared_dir, tmp_path, x25519_keypair): """Browser fetches GEK bundle from node during the handshake challenge window.""" hub_pk_pem = _hub_pk_pem(sk_hub) indexer = DirectoryIndexer(root=shared_dir, group_id="g", sk_node=sk_node, gek=gek) await indexer.initial_scan() sk_x_raw, pk_x_raw = x25519_keypair bundle_store = BundleStore(db_path=tmp_path / "bundles.db") await bundle_store.open() # Pre-populate a bundle for user-001 in group "g" bundle = wrap_gek(gek, pk_x_raw) await bundle_store.store("g", "user-001", bundle["pk_eph_b64"], bundle["nonce_b64"], bundle["wrapped_b64"]) transport = WebRTCTransport( sk_node=sk_node, hub_pk_pem=hub_pk_pem, gek=gek, shared_root=shared_dir, index=indexer.index, stun_servers=[], ) transport._ctx["bundle_store"] = bundle_store # Connect manually: handshake → challenge → gek_bundle_fetch → response browser_pc = RTCPeerConnection() received = asyncio.Queue() buf = bytearray() channel = browser_pc.createDataChannel("mnp") ready = asyncio.Event() @channel.on("open") def on_open(): ready.set() @channel.on("message") def on_msg(message): if isinstance(message, str): message = message.encode() buf.extend(message) while len(buf) >= 4: length = struct.unpack(">I", buf[:4])[0] if len(buf) < 4 + length: break msg_bytes = bytes(buf[4:4 + length]) del buf[:4 + length] received.put_nowait(msgpack.unpackb(msg_bytes, raw=False)) offer = await browser_pc.createOffer() await browser_pc.setLocalDescription(offer) answer_sdp, _ = await transport.handle_offer( browser_pc.localDescription.sdp, "peer-fetch") await browser_pc.setRemoteDescription( RTCSessionDescription(sdp=answer_sdp, type="answer")) await asyncio.wait_for(ready.wait(), timeout=5.0) # Step 1: Send handshake with group_id so _pending_group is set token = _make_jwt(sk_hub, groups=["g"]) channel.send(_pack({ "type": MNP.HANDSHAKE, "v": MNP_VERSION, "token": token, "group_id": "g", })) msg = await asyncio.wait_for(received.get(), timeout=5.0) assert msg["type"] == MNP.HANDSHAKE_CHALLENGE # Step 2: Fetch GEK bundle from node (during challenge window) channel.send(_pack({"type": MNP.GEK_BUNDLE_FETCH, "v": MNP_VERSION})) bundle_resp = await asyncio.wait_for(received.get(), timeout=5.0) assert bundle_resp["type"] == MNP.GEK_BUNDLE_RESP assert bundle_resp["found"] is True # Step 3: Unwrap GEK and compute HMAC proof recovered_gek = unwrap_gek(bundle_resp, sk_x_raw, pk_x_raw) assert recovered_gek == gek nonce = base64.b64decode(msg["nonce"]) offer_fp = _extract_dtls_fp(browser_pc.localDescription.sdp) answer_fp = _extract_dtls_fp(browser_pc.remoteDescription.sdp) proof = hmac.new(recovered_gek, nonce + offer_fp + answer_fp, hashlib.sha256).digest() # Step 4: Complete handshake channel.send(_pack({ "type": MNP.HANDSHAKE_RESPONSE, "v": MNP_VERSION, "proof": base64.b64encode(proof).decode(), })) ack = await asyncio.wait_for(received.get(), timeout=5.0) assert ack["type"] == MNP.HANDSHAKE_ACK await bundle_store.close() await browser_pc.close() await transport.close_all() # ── Keypair bundle P2P tests ───────────────────────────────────────────────── @pytest.mark.asyncio async def test_keypair_bundle_store_and_fetch(sk_node, sk_hub, gek, shared_dir, tmp_path): """Keypair bundle stored on node, then fetched during handshake window.""" hub_pk_pem = _hub_pk_pem(sk_hub) indexer = DirectoryIndexer(root=shared_dir, group_id="g", sk_node=sk_node, gek=gek) await indexer.initial_scan() bundle_store = BundleStore(db_path=tmp_path / "bundles.db") await bundle_store.open() transport = WebRTCTransport( sk_node=sk_node, hub_pk_pem=hub_pk_pem, gek=gek, shared_root=shared_dir, index=indexer.index, stun_servers=[], ) transport._ctx["bundle_store"] = bundle_store # Connect and store a keypair bundle pc1, ch1, q1 = await _setup_peer(transport, sk_hub, gek, "peer-kp-store") ch1.send(_pack({ "type": MNP.KEYPAIR_BUNDLE_STORE, "v": MNP_VERSION, "bundle_enc": "encrypted-keypair-data-base64", })) ack = await asyncio.wait_for(q1.get(), timeout=5.0) assert ack["type"] == "ack" assert ack["detail"] == "keypair_bundle_stored" # Verify in DB stored = await bundle_store.fetch_keypair("user-001") assert stored == "encrypted-keypair-data-base64" await pc1.close() # New connection: fetch during handshake window browser_pc = RTCPeerConnection() received = asyncio.Queue() buf = bytearray() channel = browser_pc.createDataChannel("mnp") ready = asyncio.Event() @channel.on("open") def on_open(): ready.set() @channel.on("message") def on_msg(message): if isinstance(message, str): message = message.encode() buf.extend(message) while len(buf) >= 4: length = struct.unpack(">I", buf[:4])[0] if len(buf) < 4 + length: break msg_bytes = bytes(buf[4:4 + length]) del buf[:4 + length] received.put_nowait(msgpack.unpackb(msg_bytes, raw=False)) offer = await browser_pc.createOffer() await browser_pc.setLocalDescription(offer) answer_sdp, _ = await transport.handle_offer( browser_pc.localDescription.sdp, "peer-kp-fetch") await browser_pc.setRemoteDescription( RTCSessionDescription(sdp=answer_sdp, type="answer")) await asyncio.wait_for(ready.wait(), timeout=5.0) token = _make_jwt(sk_hub, groups=["g"]) channel.send(_pack({ "type": MNP.HANDSHAKE, "v": MNP_VERSION, "token": token, "group_id": "g", })) msg = await asyncio.wait_for(received.get(), timeout=5.0) assert msg["type"] == MNP.HANDSHAKE_CHALLENGE # Fetch keypair bundle during challenge window channel.send(_pack({"type": MNP.KEYPAIR_BUNDLE_FETCH, "v": MNP_VERSION})) kp_resp = await asyncio.wait_for(received.get(), timeout=5.0) assert kp_resp["type"] == MNP.KEYPAIR_BUNDLE_RESP assert kp_resp["found"] is True assert kp_resp["bundle_enc"] == "encrypted-keypair-data-base64" await bundle_store.close() await browser_pc.close() await transport.close_all() @pytest.mark.asyncio async def test_keypair_bundle_fetch_not_found(sk_node, sk_hub, gek, shared_dir, tmp_path): """Keypair bundle fetch returns found=false when no bundle exists.""" hub_pk_pem = _hub_pk_pem(sk_hub) indexer = DirectoryIndexer(root=shared_dir, group_id="g", sk_node=sk_node, gek=gek) await indexer.initial_scan() bundle_store = BundleStore(db_path=tmp_path / "bundles.db") await bundle_store.open() transport = WebRTCTransport( sk_node=sk_node, hub_pk_pem=hub_pk_pem, gek=gek, shared_root=shared_dir, index=indexer.index, stun_servers=[], ) transport._ctx["bundle_store"] = bundle_store browser_pc = RTCPeerConnection() received = asyncio.Queue() buf = bytearray() channel = browser_pc.createDataChannel("mnp") ready = asyncio.Event() @channel.on("open") def on_open(): ready.set() @channel.on("message") def on_msg(message): if isinstance(message, str): message = message.encode() buf.extend(message) while len(buf) >= 4: length = struct.unpack(">I", buf[:4])[0] if len(buf) < 4 + length: break msg_bytes = bytes(buf[4:4 + length]) del buf[:4 + length] received.put_nowait(msgpack.unpackb(msg_bytes, raw=False)) offer = await browser_pc.createOffer() await browser_pc.setLocalDescription(offer) answer_sdp, _ = await transport.handle_offer( browser_pc.localDescription.sdp, "peer-kp-none") await browser_pc.setRemoteDescription( RTCSessionDescription(sdp=answer_sdp, type="answer")) await asyncio.wait_for(ready.wait(), timeout=5.0) token = _make_jwt(sk_hub, groups=["g"]) channel.send(_pack({ "type": MNP.HANDSHAKE, "v": MNP_VERSION, "token": token, "group_id": "g", })) msg = await asyncio.wait_for(received.get(), timeout=5.0) assert msg["type"] == MNP.HANDSHAKE_CHALLENGE channel.send(_pack({"type": MNP.KEYPAIR_BUNDLE_FETCH, "v": MNP_VERSION})) resp = await asyncio.wait_for(received.get(), timeout=5.0) assert resp["type"] == MNP.KEYPAIR_BUNDLE_RESP assert resp["found"] is False await bundle_store.close() await browser_pc.close() await transport.close_all() @pytest.mark.asyncio async def test_gek_auto_activate_on_node_bundle_store(sk_node, sk_hub, gek, shared_dir, tmp_path, x25519_keypair): """Storing the node operator's GEK bundle auto-activates GEK (AES variant).""" hub_pk_pem = _hub_pk_pem(sk_hub) indexer = DirectoryIndexer(root=shared_dir, group_id="g", sk_node=sk_node, gek=gek) await indexer.initial_scan() sk_x_raw, pk_x_raw = x25519_keypair bundle_store = BundleStore(db_path=tmp_path / "bundles.db") await bundle_store.open() new_gek = generate_gek() transport = WebRTCTransport( sk_node=sk_node, hub_pk_pem=hub_pk_pem, gek=gek, shared_root=shared_dir, index=indexer.index, stun_servers=[], ) transport._ctx["bundle_store"] = bundle_store transport._ctx["node_user_id"] = "node-operator" transport._ctx["sk_x25519_raw"] = sk_x_raw transport._ctx["pk_x25519_raw"] = pk_x_raw transport._ctx["pk_x25519_b64"] = base64.b64encode(pk_x_raw).decode() pc_admin, ch_admin, q_admin = await _setup_peer( transport, sk_hub, gek, "peer-setup-admin") # Store GEK bundle wrapped with AES-GCM (browser-compatible) node_bundle = wrap_gek_aes(new_gek, pk_x_raw) ch_admin.send(_pack({ "type": MNP.GEK_BUNDLE_STORE, "v": MNP_VERSION, "user_id": "node-operator", "group_id": "g", "pk_eph_b64": node_bundle["pk_eph_b64"], "nonce_b64": node_bundle["nonce_b64"], "wrapped_b64": node_bundle["wrapped_b64"], })) ack = await asyncio.wait_for(q_admin.get(), timeout=5.0) assert ack["type"] == "ack" await asyncio.sleep(0.2) assert transport._ctx.get("gek") == new_gek await bundle_store.close() await pc_admin.close() await transport.close_all() @pytest.mark.asyncio async def test_webrtc_no_gek_connection_refused(sk_node, sk_hub, shared_dir): """WebRTC DataChannel: connection refused when GEK is not initialized.""" hub_pk_pem = _hub_pk_pem(sk_hub) indexer = DirectoryIndexer(root=shared_dir, group_id="g", sk_node=sk_node, gek=None) await indexer.initial_scan() transport = WebRTCTransport( sk_node=sk_node, hub_pk_pem=hub_pk_pem, gek=None, shared_root=shared_dir, index=indexer.index, stun_servers=[], ) browser_pc = RTCPeerConnection() received = asyncio.Queue() channel = browser_pc.createDataChannel("mnp") @channel.on("message") def on_msg(message): if isinstance(message, str): message = message.encode() received.put_nowait(_unpack(message)) offer = await browser_pc.createOffer() await browser_pc.setLocalDescription(offer) answer_sdp, _ = await transport.handle_offer( browser_pc.localDescription.sdp, "peer-no-gek") await browser_pc.setRemoteDescription( RTCSessionDescription(sdp=answer_sdp, type="answer")) await asyncio.sleep(0.5) token = _make_jwt(sk_hub) channel.send(_pack({ "type": MNP.HANDSHAKE, "v": MNP_VERSION, "token": token, })) msg = await asyncio.wait_for(received.get(), timeout=5.0) assert msg["type"] == "error" assert "not initialized" in msg["detail"].lower() await browser_pc.close() await transport.close_all() @pytest.mark.asyncio async def test_gek_bundle_fetch_not_found(sk_node, sk_hub, gek, shared_dir, tmp_path): """GEK bundle fetch returns found=false when no bundle exists.""" hub_pk_pem = _hub_pk_pem(sk_hub) indexer = DirectoryIndexer(root=shared_dir, group_id="g", sk_node=sk_node, gek=gek) await indexer.initial_scan() bundle_store = BundleStore(db_path=tmp_path / "bundles.db") await bundle_store.open() transport = WebRTCTransport( sk_node=sk_node, hub_pk_pem=hub_pk_pem, gek=gek, shared_root=shared_dir, index=indexer.index, stun_servers=[], ) transport._ctx["bundle_store"] = bundle_store browser_pc = RTCPeerConnection() received = asyncio.Queue() buf = bytearray() channel = browser_pc.createDataChannel("mnp") ready = asyncio.Event() @channel.on("open") def on_open(): ready.set() @channel.on("message") def on_msg(message): if isinstance(message, str): message = message.encode() buf.extend(message) while len(buf) >= 4: length = struct.unpack(">I", buf[:4])[0] if len(buf) < 4 + length: break msg_bytes = bytes(buf[4:4 + length]) del buf[:4 + length] received.put_nowait(msgpack.unpackb(msg_bytes, raw=False)) offer = await browser_pc.createOffer() await browser_pc.setLocalDescription(offer) answer_sdp, _ = await transport.handle_offer( browser_pc.localDescription.sdp, "peer-nofound") await browser_pc.setRemoteDescription( RTCSessionDescription(sdp=answer_sdp, type="answer")) await asyncio.wait_for(ready.wait(), timeout=5.0) token = _make_jwt(sk_hub, groups=["g"]) channel.send(_pack({ "type": MNP.HANDSHAKE, "v": MNP_VERSION, "token": token, "group_id": "g", })) msg = await asyncio.wait_for(received.get(), timeout=5.0) assert msg["type"] == MNP.HANDSHAKE_CHALLENGE channel.send(_pack({"type": MNP.GEK_BUNDLE_FETCH, "v": MNP_VERSION})) resp = await asyncio.wait_for(received.get(), timeout=5.0) assert resp["type"] == MNP.GEK_BUNDLE_RESP assert resp["found"] is False await bundle_store.close() await browser_pc.close() await transport.close_all()