""" Integration test: QuicChunkServer ↔ QuicChunkClient over QUIC/UDP loopback. Same structure as test_transport.py but uses QUIC instead of TCP+TLS. """ import asyncio import os import time import jwt import pytest from pathlib import Path from cryptography.hazmat.primitives.asymmetric.ed25519 import Ed25519PrivateKey from cryptography.hazmat.primitives import serialization from meshbay_common.crypto import generate_gek, pk_to_b64 from meshbay_node.indexer import DirectoryIndexer, GroupIndex from meshbay_node.transport.quic_server import QuicChunkServer, Denylist from meshbay_node.transport.quic_client import QuicChunkClient @pytest.fixture def sk_node(): return Ed25519PrivateKey.generate() @pytest.fixture def sk_hub(): 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.mp4").write_bytes(os.urandom(2 * 1024 * 1024)) # 2 MB (d / "small.txt").write_bytes(b"hello quic " * 100) return d def make_jwt(sk_hub, pk_node_b64, ttl=3600, groups=None): 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_node_b64, "hub_id": "test-hub", "jti": "test-jti", "iat": now, "exp": now + ttl, # group_id is mandatory (M1), so default tokens are members of "g". "groups": groups if groups is not None else ["g"], }, sk_pem, algorithm="EdDSA") @pytest.mark.asyncio async def test_quic_chunk_roundtrip(sk_node, sk_hub, gek, shared_dir, tmp_path): """Full QUIC roundtrip: server serves chunk, client verifies and decrypts.""" hub_pk_pem = sk_hub.public_key().public_bytes( serialization.Encoding.PEM, serialization.PublicFormat.SubjectPublicKeyInfo) indexer = DirectoryIndexer(root=shared_dir, group_id="g", sk_node=sk_node, gek=gek) await indexer.initial_scan() cert_path = tmp_path / "node.crt" key_path = tmp_path / "node.key" server = QuicChunkServer( sk_node=sk_node, hub_pk_pem=hub_pk_pem, gek=gek, shared_root=shared_dir, index=indexer.index, host="127.0.0.1", port=19100, cert_path=cert_path, key_path=key_path, ) await server.start() token = make_jwt(sk_hub, pk_to_b64(sk_node.public_key())) entry = next(e for e in indexer.index.entries if e.name == "test.mp4") async with QuicChunkClient( host="127.0.0.1", port=19100, jwt_token=token, gek=gek, pk_node_b64=pk_to_b64(sk_node.public_key()), group_id="g", ) as client: chunk0 = await client.fetch_chunk(entry.id, chunk_index=0) chunk1 = await client.fetch_chunk(entry.id, chunk_index=1) original = (shared_dir / "test.mp4").read_bytes() assert chunk0 + chunk1 == original await server.stop() @pytest.mark.asyncio async def test_quic_fetch_index(sk_node, sk_hub, gek, shared_dir, tmp_path): """QUIC index sync returns deserializable GroupIndex.""" hub_pk_pem = sk_hub.public_key().public_bytes( serialization.Encoding.PEM, serialization.PublicFormat.SubjectPublicKeyInfo) indexer = DirectoryIndexer(root=shared_dir, group_id="g", sk_node=sk_node, gek=gek) await indexer.initial_scan() cert_path = tmp_path / "node.crt" key_path = tmp_path / "node.key" server = QuicChunkServer( sk_node=sk_node, hub_pk_pem=hub_pk_pem, gek=gek, shared_root=shared_dir, index=indexer.index, host="127.0.0.1", port=19101, cert_path=cert_path, key_path=key_path, ) await server.start() token = make_jwt(sk_hub, pk_to_b64(sk_node.public_key())) async with QuicChunkClient( host="127.0.0.1", port=19101, jwt_token=token, gek=gek, pk_node_b64=pk_to_b64(sk_node.public_key()), group_id="g", ) as client: wire = await client.fetch_index() recovered = GroupIndex.deserialize(wire, sk_node=sk_node, gek=gek) assert recovered.count == 2 await server.stop() @pytest.mark.asyncio async def test_quic_invalid_jwt_rejected(sk_node, sk_hub, gek, shared_dir, tmp_path): """QUIC server rejects connections with tokens signed by wrong hub key.""" hub_pk_pem = sk_hub.public_key().public_bytes( serialization.Encoding.PEM, serialization.PublicFormat.SubjectPublicKeyInfo) indexer = DirectoryIndexer(root=shared_dir, group_id="g", sk_node=sk_node, gek=gek) await indexer.initial_scan() cert_path = tmp_path / "node.crt" key_path = tmp_path / "node.key" server = QuicChunkServer( sk_node=sk_node, hub_pk_pem=hub_pk_pem, gek=gek, shared_root=shared_dir, index=indexer.index, host="127.0.0.1", port=19102, cert_path=cert_path, key_path=key_path, ) await server.start() sk_other = Ed25519PrivateKey.generate() bad_token = make_jwt(sk_other, pk_to_b64(sk_node.public_key())) with pytest.raises(Exception): async with QuicChunkClient( host="127.0.0.1", port=19102, jwt_token=bad_token, gek=gek, pk_node_b64=pk_to_b64(sk_node.public_key()), ) as client: await client.fetch_index() await server.stop() @pytest.mark.asyncio async def test_quic_wrong_group_rejected(sk_node, sk_hub, gek, shared_dir, tmp_path): """QUIC server rejects a client whose JWT groups don't include the requested group_id.""" hub_pk_pem = sk_hub.public_key().public_bytes( serialization.Encoding.PEM, serialization.PublicFormat.SubjectPublicKeyInfo) indexer = DirectoryIndexer(root=shared_dir, group_id="g", sk_node=sk_node, gek=gek) await indexer.initial_scan() cert_path = tmp_path / "node.crt" key_path = tmp_path / "node.key" server = QuicChunkServer( sk_node=sk_node, hub_pk_pem=hub_pk_pem, gek=gek, shared_root=shared_dir, index=indexer.index, host="127.0.0.1", port=19103, cert_path=cert_path, key_path=key_path, ) await server.start() token = make_jwt(sk_hub, pk_to_b64(sk_node.public_key()), groups=["group-a"]) with pytest.raises(ConnectionError, match="rejected"): async with QuicChunkClient( host="127.0.0.1", port=19103, jwt_token=token, gek=gek, pk_node_b64=pk_to_b64(sk_node.public_key()), group_id="group-b", ) as client: await client.fetch_index() await server.stop() @pytest.mark.asyncio async def test_quic_session_resumption(sk_node, sk_hub, gek, shared_dir, tmp_path): """QUIC 0-RTT: connect, save session ticket, reconnect with ticket.""" hub_pk_pem = sk_hub.public_key().public_bytes( serialization.Encoding.PEM, serialization.PublicFormat.SubjectPublicKeyInfo) indexer = DirectoryIndexer(root=shared_dir, group_id="g", sk_node=sk_node, gek=gek) await indexer.initial_scan() cert_path = tmp_path / "node.crt" key_path = tmp_path / "node.key" server = QuicChunkServer( sk_node=sk_node, hub_pk_pem=hub_pk_pem, gek=gek, shared_root=shared_dir, index=indexer.index, host="127.0.0.1", port=19104, cert_path=cert_path, key_path=key_path, ) await server.start() token = make_jwt(sk_hub, pk_to_b64(sk_node.public_key())) pk_b64 = pk_to_b64(sk_node.public_key()) # First connection — captures session ticket saved_ticket = None async with QuicChunkClient( host="127.0.0.1", port=19104, jwt_token=token, gek=gek, pk_node_b64=pk_b64, group_id="g", ) as client: wire = await client.fetch_index() assert GroupIndex.deserialize(wire, sk_node=sk_node, gek=gek).count == 2 saved_ticket = client.session_ticket saved_cert = client.peer_cert_der # Allow server to process the close await asyncio.sleep(0.1) # Second connection — reuses session ticket (0-RTT) async with QuicChunkClient( host="127.0.0.1", port=19104, jwt_token=token, gek=gek, pk_node_b64=pk_b64, session_ticket=saved_ticket, # A resumed session carries no certificate, so the binding anchor from the # original handshake travels with the ticket (11.5.6). peer_cert_der=saved_cert, group_id="g", ) as client: wire = await client.fetch_index() assert GroupIndex.deserialize(wire, sk_node=sk_node, gek=gek).count == 2 await server.stop() @pytest.mark.asyncio async def test_quic_denylist_blocks_user(sk_node, sk_hub, gek, shared_dir, tmp_path): """QUIC server rejects a connection when the user is on the denylist.""" hub_pk_pem = sk_hub.public_key().public_bytes( serialization.Encoding.PEM, serialization.PublicFormat.SubjectPublicKeyInfo) indexer = DirectoryIndexer(root=shared_dir, group_id="g", sk_node=sk_node, gek=gek) await indexer.initial_scan() cert_path = tmp_path / "node.crt" key_path = tmp_path / "node.key" denylist = Denylist() server = QuicChunkServer( sk_node=sk_node, hub_pk_pem=hub_pk_pem, gek=gek, shared_root=shared_dir, index=indexer.index, host="127.0.0.1", port=19105, cert_path=cert_path, key_path=key_path, denylist=denylist, ) await server.start() token = make_jwt(sk_hub, pk_to_b64(sk_node.public_key())) pk_b64 = pk_to_b64(sk_node.public_key()) # Connection works before denylisting async with QuicChunkClient( host="127.0.0.1", port=19105, jwt_token=token, gek=gek, pk_node_b64=pk_b64, group_id="g", ) as client: wire = await client.fetch_index() assert GroupIndex.deserialize(wire, sk_node=sk_node, gek=gek).count == 2 # Add user to denylist denylist.deny_user("user-001") # Connection now rejected with pytest.raises(Exception): async with QuicChunkClient( host="127.0.0.1", port=19105, jwt_token=token, gek=gek, pk_node_b64=pk_b64, group_id="g", ) as client: await client.fetch_index() await server.stop()