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)
|