diff options
3 files changed, 564 insertions, 1 deletions
diff --git a/packages/meshbay-node/src/meshbay_node/transport/__init__.py b/packages/meshbay-node/src/meshbay_node/transport/__init__.py index e7b9573..ddd11f5 100644 --- a/packages/meshbay-node/src/meshbay_node/transport/__init__.py +++ b/packages/meshbay-node/src/meshbay_node/transport/__init__.py @@ -1,5 +1,6 @@ """TCP+TLS transport layer (MNP v1). QUIC added in v2.""" from .server import ChunkServer from .client import ChunkClient +from .http_server import create_http_app -__all__ = ["ChunkServer", "ChunkClient"] +__all__ = ["ChunkServer", "ChunkClient", "create_http_app"] diff --git a/packages/meshbay-node/src/meshbay_node/transport/http_server.py b/packages/meshbay-node/src/meshbay_node/transport/http_server.py new file mode 100644 index 0000000..7554086 --- /dev/null +++ b/packages/meshbay-node/src/meshbay_node/transport/http_server.py @@ -0,0 +1,335 @@ +""" +MeshBay Node — HTTP file API (port 19001, public content). + +Serves public group content over standard HTTP so browsers can +access files without any special protocol. + +Endpoints: + GET / node info (JSON) + GET /index public Mesh Group Index (JSON) + GET /file/{file_id} full file download (streaming) + GET /file/{file_id}/{chunk} single encrypted chunk (JSON) + GET /hls/{file_id}/playlist.m3u8 HLS playlist + GET /hls/{file_id}/{segment}.ts HLS segment (binary TS) + +Auth: Bearer JWT in Authorization header (or ?token= query param). +For public groups: auth optional (anonymous browse allowed). +For chunk download: auth required (JWT verified offline with hub PK). + +Note: this server handles PUBLIC content only (no GEK decryption). +Private group content requires a client that can do ChaCha20 (Phase 5). +""" + +import asyncio +import base64 +import json +import logging +import os +import struct +import subprocess +import tempfile +from pathlib import Path + +import blake3 +import jwt +from fastapi import FastAPI, Header, HTTPException, Query, Request +from fastapi.responses import FileResponse, JSONResponse, StreamingResponse +from cryptography.hazmat.primitives.asymmetric.ed25519 import Ed25519PrivateKey + +from meshbay_common import MNP_VERSION +from meshbay_common.crypto import chunk_key as derive_chunk_key, encrypt_chunk, sign_chunk, pk_to_b64 +from meshbay_node import __version__ +from meshbay_node.indexer import GroupIndex +from meshbay_node.indexer.group_index import GroupIndex + +log = logging.getLogger(__name__) + +CHUNK_SIZE = 1024 * 1024 # 1 MB +HLS_SEGMENT_DURATION = 4 # seconds per HLS segment + + +def create_http_app( + sk_node: Ed25519PrivateKey, + hub_pk_pem: bytes, + shared_root: Path, + index: GroupIndex, + group_id: str, + group_name: str, + gek: bytes | None = None, # None for public groups +) -> FastAPI: + """ + Create the node's public HTTP API FastAPI app. + Bind to 0.0.0.0:19001 (or configured port) for external access. + """ + app = FastAPI( + title="MeshBay Node HTTP API", + version=__version__, + docs_url=None, + redoc_url=None, + ) + + # ── Auth helper ─────────────────────────────────────────────────────────── + + def _verify_token_optional( + authorization: str | None, + token_param: str | None, + ) -> dict | None: + """Verify JWT if provided. Returns decoded payload or None.""" + raw = None + if authorization and authorization.lower().startswith("bearer "): + raw = authorization[7:] + elif token_param: + raw = token_param + if not raw: + return None + try: + return jwt.decode(raw, hub_pk_pem, algorithms=["EdDSA"]) + except Exception: + return None + + def _require_token( + authorization: str | None, + token_param: str | None, + ) -> dict: + decoded = _verify_token_optional(authorization, token_param) + if decoded is None: + raise HTTPException(status_code=401, detail="Authentication required") + return decoded + + # ── Node info ───────────────────────────────────────────────────────────── + + @app.get("/") + async def node_info(): + return { + "node_version": __version__, + "mnp_version": MNP_VERSION, + "group_id": group_id, + "group_name": group_name, + "file_count": index.count, + "pk_node": pk_to_b64(sk_node.public_key()), + } + + # ── Public index ────────────────────────────────────────────────────────── + + @app.get("/index") + async def get_index( + authorization: str | None = Header(default=None), + token: str | None = Query(default=None), + ): + """Public Mesh Group Index as JSON. No auth required for public groups.""" + entries = [ + { + "id": e.id, + "name": e.name, + "path": e.path, + "size": e.size, + "type": e.type, + "duration": e.duration, + } + for e in index.entries + ] + return { + "group_id": group_id, + "group_name": group_name, + "version": index.version, + "entries": entries, + } + + # ── Full file download (streaming) ──────────────────────────────────────── + + @app.get("/file/{file_id}") + async def download_file( + file_id: str, + authorization: str | None = Header(default=None), + token: str | None = Query(default=None), + ): + """Stream an entire file. Public groups: no auth needed.""" + entry = index.get_entry(file_id) + if not entry: + raise HTTPException(status_code=404, detail="File not found in index") + + file_path = shared_root / entry.path / entry.name + if not file_path.exists(): + raise HTTPException(status_code=404, detail="File not on disk") + + return FileResponse( + path=str(file_path), + filename=entry.name, + media_type=_media_type(entry.name), + ) + + # ── Chunk endpoint (encrypted, for MNP-aware clients) ──────────────────── + + @app.get("/file/{file_id}/{chunk_index}") + async def get_chunk( + file_id: str, + chunk_index: int, + authorization: str | None = Header(default=None), + token: str | None = Query(default=None), + ): + """ + Serve one encrypted chunk (JSON). Auth required. + Clients that understand MNP can decrypt with the GEK they got from the hub. + """ + _require_token(authorization, token) + + entry = index.get_entry(file_id) + if not entry: + raise HTTPException(status_code=404, detail="File not found") + + file_path = shared_root / entry.path / entry.name + if not file_path.exists(): + raise HTTPException(status_code=404, detail="File not on disk") + + # Read chunk + with open(file_path, "rb") as f: + f.seek(chunk_index * CHUNK_SIZE) + plaintext = f.read(CHUNK_SIZE) + + if not plaintext: + raise HTTPException(status_code=416, detail="Chunk out of range") + + file_hash = blake3.blake3(file_path.read_bytes()).digest() + pt_hash = blake3.blake3(plaintext).digest() + + if gek: + # Private group: encrypt chunk + ckey = derive_chunk_key(gek, file_hash, chunk_index) + nonce, ct = encrypt_chunk(ckey, plaintext) + ct_hash = blake3.blake3(ct).digest() + sig = sign_chunk(sk_node, chunk_index, nonce, ct_hash) + return { + "chunk_index": chunk_index, + "plaintext_size": len(plaintext), + "encrypted": True, + "nonce_b64": base64.b64encode(nonce).decode(), + "ct_b64": base64.b64encode(ct).decode(), + "ct_hash_b64": base64.b64encode(ct_hash).decode(), + "pt_hash_b64": base64.b64encode(pt_hash).decode(), + "sig_b64": base64.b64encode(sig).decode(), + "pk_node_b64": pk_to_b64(sk_node.public_key()), + "file_hash_b64": base64.b64encode(file_hash).decode(), + } + else: + # Public group: serve plaintext chunk (TLS provides transport encryption) + pt_hash_b = blake3.blake3(plaintext).digest() + sig_payload = chunk_index.to_bytes(4, "big") + bytes(12) + pt_hash_b + sig = sk_node.sign(sig_payload) + return { + "chunk_index": chunk_index, + "plaintext_size": len(plaintext), + "encrypted": False, + "data_b64": base64.b64encode(plaintext).decode(), + "pt_hash_b64": base64.b64encode(pt_hash).decode(), + "sig_b64": base64.b64encode(sig).decode(), + "pk_node_b64": pk_to_b64(sk_node.public_key()), + } + + # ── HLS streaming ───────────────────────────────────────────────────────── + + @app.get("/hls/{file_id}/playlist.m3u8") + async def hls_playlist( + file_id: str, + authorization: str | None = Header(default=None), + token: str | None = Query(default=None), + ): + """Generate HLS playlist for a video file.""" + entry = index.get_entry(file_id) + if not entry or entry.type != "video": + raise HTTPException(status_code=404, detail="Video file not found") + + file_path = shared_root / entry.path / entry.name + if not file_path.exists(): + raise HTTPException(status_code=404, detail="File not on disk") + + duration = entry.duration or _probe_duration(file_path) + if not duration: + raise HTTPException(status_code=422, detail="Cannot determine video duration") + + n_segments = max(1, int(duration / HLS_SEGMENT_DURATION) + 1) + token_param = f"?token={token}" if token else "" + + lines = [ + "#EXTM3U", + "#EXT-X-VERSION:3", + f"#EXT-X-TARGETDURATION:{HLS_SEGMENT_DURATION}", + "#EXT-X-MEDIA-SEQUENCE:0", + ] + for i in range(n_segments): + seg_dur = min(HLS_SEGMENT_DURATION, duration - i * HLS_SEGMENT_DURATION) + if seg_dur <= 0: + break + lines.append(f"#EXTINF:{seg_dur:.3f},") + lines.append(f"/hls/{file_id}/{i}.ts{token_param}") + lines.append("#EXT-X-ENDLIST") + + return StreamingResponse( + iter(["\n".join(lines)]), + media_type="application/vnd.apple.mpegurl", + ) + + @app.get("/hls/{file_id}/{segment_index}.ts") + async def hls_segment( + file_id: str, + segment_index: int, + authorization: str | None = Header(default=None), + token: str | None = Query(default=None), + ): + """Serve one HLS segment as MPEG-TS via ffmpeg transcoding.""" + entry = index.get_entry(file_id) + if not entry or entry.type != "video": + raise HTTPException(status_code=404, detail="Video not found") + + file_path = shared_root / entry.path / entry.name + if not file_path.exists(): + raise HTTPException(status_code=404, detail="File not on disk") + + start_time = segment_index * HLS_SEGMENT_DURATION + + async def generate(): + proc = await asyncio.create_subprocess_exec( + "ffmpeg", "-hide_banner", "-loglevel", "error", + "-ss", str(start_time), + "-i", str(file_path), + "-t", str(HLS_SEGMENT_DURATION), + "-c:v", "copy", "-c:a", "copy", + "-f", "mpegts", "pipe:1", + stdout=asyncio.subprocess.PIPE, + stderr=asyncio.subprocess.DEVNULL, + ) + assert proc.stdout + while chunk := await proc.stdout.read(65536): + yield chunk + await proc.wait() + + return StreamingResponse(generate(), media_type="video/mp2t") + + return app + + +# ── Helpers ─────────────────────────────────────────────────────────────────── + +def _media_type(filename: str) -> str: + ext = Path(filename).suffix.lower() + return { + ".mp4": "video/mp4", ".mkv": "video/x-matroska", + ".webm": "video/webm", ".avi": "video/x-msvideo", + ".mp3": "audio/mpeg", ".flac": "audio/flac", + ".ogg": "audio/ogg", ".opus": "audio/opus", + ".jpg": "image/jpeg", ".png": "image/png", + ".pdf": "application/pdf", + }.get(ext, "application/octet-stream") + + +def _probe_duration(path: Path) -> float | None: + """Use ffprobe to get video duration in seconds.""" + try: + result = subprocess.run( + ["ffprobe", "-v", "quiet", "-print_format", "json", + "-show_format", str(path)], + capture_output=True, text=True, timeout=10, + ) + data = json.loads(result.stdout) + return float(data["format"]["duration"]) + except Exception: + return None diff --git a/packages/meshbay-node/tests/test_http_server.py b/packages/meshbay-node/tests/test_http_server.py new file mode 100644 index 0000000..1483a1d --- /dev/null +++ b/packages/meshbay-node/tests/test_http_server.py @@ -0,0 +1,227 @@ +"""Tests for the node HTTP file API.""" + +import asyncio +import base64 +import json +import os +import time +import pytest +import jwt +import httpx +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 +from meshbay_node.transport.http_server import create_http_app + + +@pytest.fixture +def sk_node(): + return Ed25519PrivateKey.generate() + +@pytest.fixture +def sk_hub(): + return Ed25519PrivateKey.generate() + +@pytest.fixture +def hub_pk_pem(sk_hub): + return sk_hub.public_key().public_bytes( + serialization.Encoding.PEM, serialization.PublicFormat.SubjectPublicKeyInfo) + +@pytest.fixture +def gek(): + return generate_gek() + +@pytest.fixture +def shared_dir(tmp_path): + d = tmp_path / "shared" + d.mkdir() + (d / "video.mp4").write_bytes(os.urandom(3 * 1024 * 1024)) # 3MB + (d / "doc.pdf").write_bytes(os.urandom(512 * 1024)) + (d / "song.mp3").write_bytes(os.urandom(256 * 1024)) + return d + +def make_token(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_node_info(sk_node, sk_hub, hub_pk_pem, gek, shared_dir): + indexer = DirectoryIndexer(root=shared_dir, group_id="g", sk_node=sk_node, gek=gek) + await indexer.initial_scan() + + app = create_http_app( + sk_node=sk_node, hub_pk_pem=hub_pk_pem, + shared_root=shared_dir, index=indexer.index, + group_id="test-group", group_name="Test Group", + ) + async with httpx.AsyncClient( + transport=httpx.ASGITransport(app=app), base_url="http://test" + ) as c: + r = await c.get("/") + assert r.status_code == 200 + data = r.json() + assert data["group_id"] == "test-group" + assert data["file_count"] == 3 + assert "pk_node" in data + + +@pytest.mark.asyncio +async def test_public_index(sk_node, sk_hub, hub_pk_pem, shared_dir): + """Public group: index accessible without auth.""" + indexer = DirectoryIndexer(root=shared_dir, group_id="g", sk_node=sk_node, gek=None) + await indexer.initial_scan() + + app = create_http_app( + sk_node=sk_node, hub_pk_pem=hub_pk_pem, + shared_root=shared_dir, index=indexer.index, + group_id="pub-group", group_name="Public Group", + gek=None, + ) + async with httpx.AsyncClient( + transport=httpx.ASGITransport(app=app), base_url="http://test" + ) as c: + r = await c.get("/index") + assert r.status_code == 200 + data = r.json() + assert len(data["entries"]) == 3 + names = {e["name"] for e in data["entries"]} + assert "video.mp4" in names + assert "doc.pdf" in names + + +@pytest.mark.asyncio +async def test_file_download(sk_node, sk_hub, hub_pk_pem, shared_dir): + """Full file download via HTTP.""" + indexer = DirectoryIndexer(root=shared_dir, group_id="g", sk_node=sk_node, gek=None) + await indexer.initial_scan() + + app = create_http_app( + sk_node=sk_node, hub_pk_pem=hub_pk_pem, + shared_root=shared_dir, index=indexer.index, + group_id="g", group_name="G", + ) + entry = next(e for e in indexer.index.entries if e.name == "doc.pdf") + original = (shared_dir / "doc.pdf").read_bytes() + + async with httpx.AsyncClient( + transport=httpx.ASGITransport(app=app), base_url="http://test" + ) as c: + r = await c.get(f"/file/{entry.id}") + assert r.status_code == 200 + assert r.content == original + + +@pytest.mark.asyncio +async def test_chunk_public_group(sk_node, sk_hub, hub_pk_pem, shared_dir): + """Public group chunk: plaintext, signed, auth required.""" + indexer = DirectoryIndexer(root=shared_dir, group_id="g", sk_node=sk_node, gek=None) + await indexer.initial_scan() + + app = create_http_app( + sk_node=sk_node, hub_pk_pem=hub_pk_pem, + shared_root=shared_dir, index=indexer.index, + group_id="g", group_name="G", gek=None, + ) + entry = next(e for e in indexer.index.entries if e.name == "video.mp4") + token = make_token(sk_hub, pk_to_b64(sk_node.public_key())) + + async with httpx.AsyncClient( + transport=httpx.ASGITransport(app=app), base_url="http://test" + ) as c: + r = await c.get(f"/file/{entry.id}/0", + headers={"Authorization": f"Bearer {token}"}) + assert r.status_code == 200 + chunk = r.json() + assert chunk["encrypted"] is False + assert chunk["chunk_index"] == 0 + assert "data_b64" in chunk + + # Verify the chunk data matches original + original = (shared_dir / "video.mp4").read_bytes() + data = base64.b64decode(chunk["data_b64"]) + assert data == original[:len(data)] + + +@pytest.mark.asyncio +async def test_chunk_private_group(sk_node, sk_hub, hub_pk_pem, gek, shared_dir): + """Private group chunk: encrypted with GEK.""" + from meshbay_common.crypto import chunk_key as derive_chunk_key, decrypt_chunk + import blake3 + + indexer = DirectoryIndexer(root=shared_dir, group_id="g", sk_node=sk_node, gek=gek) + await indexer.initial_scan() + + app = create_http_app( + sk_node=sk_node, hub_pk_pem=hub_pk_pem, + shared_root=shared_dir, index=indexer.index, + group_id="g", group_name="G", gek=gek, + ) + entry = next(e for e in indexer.index.entries if e.name == "doc.pdf") + token = make_token(sk_hub, pk_to_b64(sk_node.public_key())) + + async with httpx.AsyncClient( + transport=httpx.ASGITransport(app=app), base_url="http://test" + ) as c: + r = await c.get(f"/file/{entry.id}/0", + headers={"Authorization": f"Bearer {token}"}) + assert r.status_code == 200 + chunk = r.json() + assert chunk["encrypted"] is True + + # Decrypt and verify + file_hash = base64.b64decode(chunk["file_hash_b64"]) + nonce = base64.b64decode(chunk["nonce_b64"]) + ct = base64.b64decode(chunk["ct_b64"]) + ckey = derive_chunk_key(gek, file_hash, 0) + plaintext = decrypt_chunk(ckey, nonce, ct) + original = (shared_dir / "doc.pdf").read_bytes() + assert plaintext == original[:len(plaintext)] + + +@pytest.mark.asyncio +async def test_chunk_requires_auth(sk_node, hub_pk_pem, shared_dir): + """Chunk endpoint rejects unauthenticated requests.""" + indexer = DirectoryIndexer(root=shared_dir, group_id="g", sk_node=sk_node, gek=None) + await indexer.initial_scan() + app = create_http_app( + sk_node=sk_node, hub_pk_pem=hub_pk_pem, + shared_root=shared_dir, index=indexer.index, + group_id="g", group_name="G", + ) + entry = indexer.index.entries[0] + + async with httpx.AsyncClient( + transport=httpx.ASGITransport(app=app), base_url="http://test" + ) as c: + r = await c.get(f"/file/{entry.id}/0") # no token + assert r.status_code == 401 + + +@pytest.mark.asyncio +async def test_unknown_file_404(sk_node, hub_pk_pem, sk_hub, shared_dir): + indexer = DirectoryIndexer(root=shared_dir, group_id="g", sk_node=sk_node, gek=None) + await indexer.initial_scan() + app = create_http_app( + sk_node=sk_node, hub_pk_pem=hub_pk_pem, + shared_root=shared_dir, index=indexer.index, + group_id="g", group_name="G", + ) + token = make_token(sk_hub, pk_to_b64(sk_node.public_key())) + + async with httpx.AsyncClient( + transport=httpx.ASGITransport(app=app), base_url="http://test" + ) as c: + r = await c.get("/file/nonexistent-hash/0", + headers={"Authorization": f"Bearer {token}"}) + assert r.status_code == 404 |