diff options
| author | Christophe Besson <cbesson@gmail.com> | 2026-08-09 04:56:07 +0200 |
|---|---|---|
| committer | Christophe Besson <cbesson@gmail.com> | 2026-08-09 04:56:07 +0200 |
| commit | 9b503785afe10849f40a26735aa08e0af24a6efc (patch) | |
| tree | f760ccaa98813fec57c92efe23fc6967de90f6b0 /packages/meshbay-node/src | |
| parent | a4522fe8253fcf6702cf6af7c25284b9519f3969 (diff) | |
| download | meshbay-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>
Diffstat (limited to 'packages/meshbay-node/src')
| -rw-r--r-- | packages/meshbay-node/src/meshbay_node/transport/__init__.py | 3 | ||||
| -rw-r--r-- | packages/meshbay-node/src/meshbay_node/transport/http_server.py | 335 |
2 files changed, 337 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 |