From 88cfc139333ac3fe5789f39f3970065181df9043 Mon Sep 17 00:00:00 2001 From: Christophe Besson Date: Sun, 9 Aug 2026 05:13:51 +0200 Subject: feat(node): add QUIC transport (MNP v2) — quic_server + quic_client MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit QuicChunkServer/QuicChunkClient: same MNP protocol over QUIC/UDP. Enables hole-punching (Spike 4 Cone NAT validated). Uses aioquic 1.3.0. Bug found+fixed: asyncio.Event race condition in client recv loop (quic_event_received overwrote _stream_events[0] after _recv created it). Fixed with asyncio.Queue (no shared mutable state). Server uses synchronous handlers in quic_event_received (avoids ensure_future transmit timing issue). 3/3 tests. Full suite: 50/50. Co-Authored-By: Claude Sonnet 4.6 (1M context) --- packages/meshbay-node/tests/test_quic_transport.py | 157 +++++++++++++++++++++ 1 file changed, 157 insertions(+) create mode 100644 packages/meshbay-node/tests/test_quic_transport.py (limited to 'packages/meshbay-node/tests') diff --git a/packages/meshbay-node/tests/test_quic_transport.py b/packages/meshbay-node/tests/test_quic_transport.py new file mode 100644 index 0000000..2abd465 --- /dev/null +++ b/packages/meshbay-node/tests/test_quic_transport.py @@ -0,0 +1,157 @@ +""" +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 +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): + 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, + }, 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()), + ) 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()), + ) 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() -- cgit v1.2.3