aboutsummaryrefslogtreecommitdiffstats
path: root/packages/meshbay-hub/src/meshbay_hub/api/signaling.py
blob: bd343c991c8e84bc4ce40e51060895fbabb2091a (plain) (blame)
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
"""
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)