""" WebRTC signaling — relay SDP/ICE between browser and node. The hub NEVER touches content. This is pure signaling: < 1 KB per message, stateless relay. After the SDP exchange completes, the browser and node communicate P2P via WebRTC DataChannel — hub is out of the loop. Flow: Browser → Hub : POST /v1/nodes/{node_id}/webrtc/offer {sdp, ice_candidates} Hub → Node : WS push {type: "webrtc_offer", sdp, ice_candidates, peer_id} Node → Hub : WS reply {type: "webrtc_answer", sdp, ice_candidates, peer_id} Hub → Browser : HTTP response {sdp, ice_candidates} """ import asyncio import json import logging import uuid from fastapi import APIRouter, Depends, HTTPException from pydantic import BaseModel from meshbay_hub.api.deps import get_current_user from meshbay_hub.db.models import User log = logging.getLogger(__name__) router = APIRouter(prefix="/v1/nodes", tags=["signaling"]) _webrtc_answers: dict[str, asyncio.Future] = {} class WebRTCOfferRequest(BaseModel): sdp: str ice_candidates: list[dict] = [] class WebRTCOfferResponse(BaseModel): sdp: str ice_candidates: list[dict] = [] peer_id: str @router.post("/{node_id}/webrtc/offer", response_model=WebRTCOfferResponse) async def webrtc_offer( node_id: str, body: WebRTCOfferRequest, current_user: User = Depends(get_current_user), ): """ Browser sends WebRTC SDP offer for a node. Hub relays via WebSocket. Returns the node's SDP answer once received. """ from meshbay_hub.api.revocation import _connected_nodes ws = _connected_nodes.get(node_id) if not ws: raise HTTPException(status_code=404, detail="Node not connected") peer_id = str(uuid.uuid4()) answer_future: asyncio.Future = asyncio.get_event_loop().create_future() _webrtc_answers[peer_id] = answer_future try: await ws.send_text(json.dumps({ "type": "webrtc_offer", "peer_id": peer_id, "user_id": current_user.id, "sdp": body.sdp, "ice_candidates": body.ice_candidates, })) try: answer = await asyncio.wait_for(answer_future, timeout=15.0) except asyncio.TimeoutError: raise HTTPException( status_code=504, detail="Node did not respond with WebRTC answer") return WebRTCOfferResponse( sdp=answer["sdp"], ice_candidates=answer.get("ice_candidates", []), peer_id=peer_id, ) finally: _webrtc_answers.pop(peer_id, None) def handle_webrtc_answer(msg: dict) -> None: """Called from the node WebSocket message loop when a webrtc_answer arrives.""" peer_id = msg.get("peer_id") if not peer_id: log.warning("webrtc_answer without peer_id") return future = _webrtc_answers.get(peer_id) if future and not future.done(): future.set_result({ "sdp": msg.get("sdp", ""), "ice_candidates": msg.get("ice_candidates", []), }) else: log.warning("webrtc_answer for unknown peer_id: %s", peer_id)