aboutsummaryrefslogtreecommitdiffstats
diff options
context:
space:
mode:
authorChristophe Besson <cbesson@gmail.com>2026-08-09 04:56:07 +0200
committerChristophe Besson <cbesson@gmail.com>2026-08-09 04:56:07 +0200
commit9b503785afe10849f40a26735aa08e0af24a6efc (patch)
treef760ccaa98813fec57c92efe23fc6967de90f6b0
parenta4522fe8253fcf6702cf6af7c25284b9519f3969 (diff)
downloadmeshbay-9b503785afe10849f40a26735aa08e0af24a6efc.tar.gz
feat(node): add HTTP file API for browser access (Phase 4)
FastAPI app on port 19001: GET / (node info), GET /index (public group index JSON), GET /file/{id} (full download), GET /file/{id}/{n} (encrypted or plaintext chunk), GET /hls/{id}/playlist.m3u8 + GET /hls/{id}/{n}.ts (HLS streaming via ffmpeg). Public groups: index browsable without auth, files downloadable. Private groups: chunks encrypted with GEK, auth required. 7/7 tests passing. Full suite: 47/47. Co-Authored-By: Claude Sonnet 4.6 (1M context) <noreply@anthropic.com>
-rw-r--r--packages/meshbay-node/src/meshbay_node/transport/__init__.py3
-rw-r--r--packages/meshbay-node/src/meshbay_node/transport/http_server.py335
-rw-r--r--packages/meshbay-node/tests/test_http_server.py227
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