""" Integration test: ChunkServer ↔ ChunkClient over TLS. Starts a real TLS server on localhost, connects a client, fetches index and a chunk, verifies signature+hash+decryption. """ import asyncio import base64 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.server import ChunkServer from meshbay_node.transport.client import ChunkClient @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 meshbay " * 100) return d def make_jwt(sk_hub, pk_node_b64, user_id="user-001", 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_id, "pk_user": pk_node_b64, "hub_id": "test-hub", "jti": "test-jti", "iat": now, "exp": now + ttl, "groups": groups or [], }, sk_pem, algorithm="EdDSA") @pytest.mark.asyncio async def test_chunk_server_client_roundtrip( sk_node, sk_hub, gek, shared_dir, tmp_path): """Full integration: server serves a chunk, client verifies and decrypts.""" # Build index indexer = DirectoryIndexer( root=shared_dir, group_id="g", sk_node=sk_node, gek=gek) await indexer.initial_scan() assert indexer.index.count == 2 # Hub PK for JWT verification hub_pk_pem = sk_hub.public_key().public_bytes( serialization.Encoding.PEM, serialization.PublicFormat.SubjectPublicKeyInfo) # TLS cert in tmp dir cert_path = tmp_path / "node.crt" key_path = tmp_path / "node.key" server = ChunkServer( 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=0, # OS picks a free port cert_path=cert_path, key_path=key_path, ) await server.start() port = server._server.sockets[0].getsockname()[1] token = make_jwt(sk_hub, pk_to_b64(sk_node.public_key())) # Find the large test file in the index entry = next(e for e in indexer.index.entries if e.name == "test.mp4") async with ChunkClient( host="127.0.0.1", port=port, jwt_token=token, gek=gek, pk_node_b64=pk_to_b64(sk_node.public_key()), ) as client: # Fetch first chunk chunk0 = await client.fetch_chunk(entry.id, chunk_index=0) assert len(chunk0) == 1024 * 1024 # first 1MB of 2MB file # Fetch second chunk chunk1 = await client.fetch_chunk(entry.id, chunk_index=1) assert len(chunk1) == 1024 * 1024 # second 1MB # Reassembled file matches original original = (shared_dir / "test.mp4").read_bytes() assert chunk0 + chunk1 == original await server.stop() @pytest.mark.asyncio async def test_invalid_jwt_rejected(sk_node, sk_hub, gek, shared_dir, tmp_path): 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 = ChunkServer( 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=0, cert_path=cert_path, key_path=key_path, ) await server.start() port = server._server.sockets[0].getsockname()[1] # Use a different hub key to sign the token sk_other_hub = Ed25519PrivateKey.generate() bad_token = make_jwt(sk_other_hub, pk_to_b64(sk_node.public_key())) with pytest.raises(Exception): async with ChunkClient( host="127.0.0.1", port=port, jwt_token=bad_token, gek=gek, pk_node_b64=pk_to_b64(sk_node.public_key()), ) as client: pass await server.stop() @pytest.mark.asyncio async def test_wrong_group_rejected(sk_node, sk_hub, gek, shared_dir, tmp_path): """TCP+TLS 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 = ChunkServer( 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=0, cert_path=cert_path, key_path=key_path, ) await server.start() port = server._server.sockets[0].getsockname()[1] token = make_jwt(sk_hub, pk_to_b64(sk_node.public_key()), groups=["group-a"]) with pytest.raises(ConnectionError, match="rejected"): async with ChunkClient( host="127.0.0.1", port=port, jwt_token=token, gek=gek, pk_node_b64=pk_to_b64(sk_node.public_key()), group_id="group-b", ) as client: pass await server.stop() @pytest.mark.asyncio async def test_fetch_index(sk_node, sk_hub, gek, shared_dir, tmp_path): 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 = ChunkServer( 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=0, cert_path=cert_path, key_path=key_path, ) await server.start() port = server._server.sockets[0].getsockname()[1] token = make_jwt(sk_hub, pk_to_b64(sk_node.public_key())) async with ChunkClient( host="127.0.0.1", port=port, 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()