""" 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 sign_chunk, pk_to_b64 from meshbay_common.webcrypto import chunk_key_aes as derive_chunk_key, encrypt_chunk_aes as encrypt_chunk 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 = bytes.fromhex(entry.id) 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