summaryrefslogtreecommitdiffstats
path: root/packages/meshbay-hub/src/meshbay_hub/api/signaling.py
diff options
context:
space:
mode:
Diffstat (limited to 'packages/meshbay-hub/src/meshbay_hub/api/signaling.py')
-rw-r--r--packages/meshbay-hub/src/meshbay_hub/api/signaling.py56
1 files changed, 53 insertions, 3 deletions
diff --git a/packages/meshbay-hub/src/meshbay_hub/api/signaling.py b/packages/meshbay-hub/src/meshbay_hub/api/signaling.py
index bd343c9..8f84163 100644
--- a/packages/meshbay-hub/src/meshbay_hub/api/signaling.py
+++ b/packages/meshbay-hub/src/meshbay_hub/api/signaling.py
@@ -17,11 +17,15 @@ import json
import logging
import uuid
-from fastapi import APIRouter, Depends, HTTPException
+from fastapi import APIRouter, Depends, HTTPException, Request
from pydantic import BaseModel
+from sqlalchemy import select
+from sqlalchemy.ext.asyncio import AsyncSession
from meshbay_hub.api.deps import get_current_user
-from meshbay_hub.db.models import User
+from meshbay_hub.api.middleware import limiter
+from meshbay_hub.db.engine import get_db
+from meshbay_hub.db.models import Group, GroupMember, User
log = logging.getLogger(__name__)
@@ -41,25 +45,66 @@ class WebRTCOfferResponse(BaseModel):
peer_id: str
+MAX_SDP_BYTES = 16 * 1024 # an SDP offer is ~2 KB
+MAX_PENDING_PER_USER = 3 # concurrent in-flight offers per account
+
+_pending_per_user: dict[str, int] = {}
+
+
@router.post("/{node_id}/webrtc/offer", response_model=WebRTCOfferResponse)
+@limiter.limit("30/minute")
async def webrtc_offer(
node_id: str,
body: WebRTCOfferRequest,
+ request: Request,
current_user: User = Depends(get_current_user),
+ db: AsyncSession = Depends(get_db),
):
"""
Browser sends WebRTC SDP offer for a node. Hub relays via WebSocket.
Returns the node's SDP answer once received.
+
+ Finding H6: this was reachable by any authenticated user, for any node, with no
+ rate limit and no membership check. Each call makes the node allocate an
+ aiortc RTCPeerConnection and gather ICE, so it was a remote resource-exhaustion
+ primitive against an arbitrary third party's machine.
+
+ Finding H4: it also ignored group status, so "suspend a group" did not stop new
+ connections from being brokered to nodes hosting it.
"""
- from meshbay_hub.api.revocation import _connected_nodes
+ from meshbay_hub.api.revocation import _connected_nodes, _node_groups
+
+ if len(body.sdp) > MAX_SDP_BYTES:
+ raise HTTPException(status_code=413, detail="SDP too large")
ws = _connected_nodes.get(node_id)
if not ws:
raise HTTPException(status_code=404, detail="Node not connected")
+ # The caller must share at least one active group with the target node.
+ node_group_ids = set(_node_groups.get(node_id, []))
+ if node_group_ids:
+ result = await db.execute(
+ select(GroupMember.group_id).where(
+ GroupMember.user_id == current_user.id,
+ GroupMember.group_id.in_(node_group_ids),
+ ))
+ shared = [gid for (gid,) in result.all()]
+ if not shared:
+ raise HTTPException(status_code=403, detail="Not a member of any group on this node")
+
+ active = await db.execute(
+ select(Group.id).where(Group.id.in_(shared), Group.status == "active"))
+ if not active.first():
+ raise HTTPException(status_code=403, detail="Group is not active")
+
+ if _pending_per_user.get(current_user.id, 0) >= MAX_PENDING_PER_USER:
+ raise HTTPException(status_code=429, detail="Too many pending connections")
+
peer_id = str(uuid.uuid4())
answer_future: asyncio.Future = asyncio.get_event_loop().create_future()
_webrtc_answers[peer_id] = answer_future
+ _pending_per_user[current_user.id] = _pending_per_user.get(current_user.id, 0) + 1
try:
await ws.send_text(json.dumps({
@@ -83,6 +128,11 @@ async def webrtc_offer(
)
finally:
_webrtc_answers.pop(peer_id, None)
+ remaining = _pending_per_user.get(current_user.id, 1) - 1
+ if remaining > 0:
+ _pending_per_user[current_user.id] = remaining
+ else:
+ _pending_per_user.pop(current_user.id, None)
def handle_webrtc_answer(msg: dict) -> None: