diff options
Diffstat (limited to 'packages/meshbay-hub/src/meshbay_hub')
| -rw-r--r-- | packages/meshbay-hub/src/meshbay_hub/api/revocation.py | 53 | ||||
| -rw-r--r-- | packages/meshbay-hub/src/meshbay_hub/api/signaling.py | 34 |
2 files changed, 62 insertions, 25 deletions
diff --git a/packages/meshbay-hub/src/meshbay_hub/api/revocation.py b/packages/meshbay-hub/src/meshbay_hub/api/revocation.py index 6a6baa7..5a33d77 100644 --- a/packages/meshbay-hub/src/meshbay_hub/api/revocation.py +++ b/packages/meshbay-hub/src/meshbay_hub/api/revocation.py @@ -87,6 +87,16 @@ NOTIFY_WINDOW_SECONDS = 60 _notify_window: dict[str, tuple[float, int]] = {} # node_id → (window start, count) +# `update_groups` re-reads the node's groups from the database. A node sends one +# when its configuration is reloaded — an operator attaching a group, an owner's +# approval arriving — so ten a minute is far past real use, and the same budget +# rule as chat_notify keeps one node from spending the hub's database for others. +UPDATE_GROUPS_BURST = 10 +_update_window: dict[str, tuple[float, int]] = {} +# The groups one node may claim in one message. An operator with fifty groups +# is a large one. +MAX_CLAIMED_GROUPS = 1000 + def forget_node(node_id: str) -> None: """Drop everything a disconnected node's socket owned. @@ -106,25 +116,35 @@ def forget_node(node_id: str) -> None: _node_users.pop(node_id, None) -def _notify_budget(node_id: str) -> bool: - """True if this node may send one more chat_notify now.""" +def _spend(window: dict[str, tuple[float, int]], node_id: str, burst: int) -> bool: + """True if this node may send one more message of a budgeted kind now.""" now = time.monotonic() - if len(_notify_window) > 1000: + if len(window) > 1000: # Swept here rather than on disconnect, which would let a node refill # its budget by reconnecting — the same token stays valid for an hour. - for nid, (started, _) in list(_notify_window.items()): + for nid, (started, _) in list(window.items()): if now - started >= NOTIFY_WINDOW_SECONDS: - _notify_window.pop(nid, None) - start, count = _notify_window.get(node_id, (now, 0)) + window.pop(nid, None) + start, count = window.get(node_id, (now, 0)) if now - start >= NOTIFY_WINDOW_SECONDS: start, count = now, 0 - if count >= NOTIFY_BURST: - _notify_window[node_id] = (start, count) + if count >= burst: + window[node_id] = (start, count) return False - _notify_window[node_id] = (start, count + 1) + window[node_id] = (start, count + 1) return True +def _notify_budget(node_id: str) -> bool: + """True if this node may send one more chat_notify now.""" + return _spend(_notify_window, node_id, NOTIFY_BURST) + + +def _update_budget(node_id: str) -> bool: + """True if this node may send one more update_groups now.""" + return _spend(_update_window, node_id, UPDATE_GROUPS_BURST) + + async def _mark_hosted(group_ids: list[str]) -> None: """Stamp the first time a node announced it hosts each of these groups. @@ -461,8 +481,8 @@ async def node_websocket(ws: WebSocket): await _reject(ws, "Node already connected", 4009) return - resolved_id, result = await _authorize_node_ws( - msg["token"], claimed_id, msg.get("group_ids")) + claimed = [str(g) for g in (msg.get("group_ids") or [])][:MAX_CLAIMED_GROUPS] + resolved_id, result = await _authorize_node_ws(msg["token"], claimed_id, claimed) if resolved_id is None: await _reject(ws, result, 4003) return @@ -472,7 +492,7 @@ async def node_websocket(ws: WebSocket): node_id = resolved_id _connected_nodes[node_id] = ws _node_groups[node_id] = group_ids - _node_claims[node_id] = list(msg.get("group_ids") or []) + _node_claims[node_id] = claimed _node_users[node_id] = user_id await _mark_hosted(group_ids) log.info("Node WS connected: %s (user=%s, groups=%d)", @@ -493,13 +513,16 @@ async def node_websocket(ws: WebSocket): from meshbay_hub.api.signaling import handle_webrtc_answer handle_webrtc_answer(msg, node_id) elif msg.get("type") == "update_groups": + if not _update_budget(node_id): + log.warning("Node %s exceeded its update_groups rate", node_id[:8]) + continue # Through the same gate as the registration above. This used to # assign the message's list verbatim, so the ceiling that makes # C2 hold at authentication could be stepped over one message # later: a node had only to reload to claim any group on the hub. - _node_claims[node_id] = list(msg.get("group_ids") or []) - new_gids = await resolve_node_groups( - node_id, user_id, msg.get("group_ids")) + claimed = [str(g) for g in (msg.get("group_ids") or [])][:MAX_CLAIMED_GROUPS] + _node_claims[node_id] = claimed + new_gids = await resolve_node_groups(node_id, user_id, claimed) _node_groups[node_id] = new_gids await _mark_hosted(new_gids) log.info("Node %s updated groups: %d", node_id[:8], len(new_gids)) diff --git a/packages/meshbay-hub/src/meshbay_hub/api/signaling.py b/packages/meshbay-hub/src/meshbay_hub/api/signaling.py index 6c9699b..56e0e4b 100644 --- a/packages/meshbay-hub/src/meshbay_hub/api/signaling.py +++ b/packages/meshbay-hub/src/meshbay_hub/api/signaling.py @@ -20,7 +20,7 @@ import time import uuid from fastapi import APIRouter, Depends, HTTPException, Request -from pydantic import BaseModel +from pydantic import BaseModel, Field, field_validator from sqlalchemy import select from sqlalchemy.ext.asyncio import AsyncSession @@ -41,9 +41,23 @@ _webrtc_answers: dict[str, asyncio.Future] = {} _answer_owner: dict[str, str] = {} +# A browser offers a handful of candidates — a host and a reflexive one per +# interface — and embeds them in the SDP anyway. The list is relayed to the node +# as it came, so it is bounded like the SDP beside it. +MAX_ICE_CANDIDATES = 64 +MAX_ICE_BYTES = 32 * 1024 + + class WebRTCOfferRequest(BaseModel): sdp: str - ice_candidates: list[dict] = [] + ice_candidates: list[dict] = Field(default_factory=list, max_length=MAX_ICE_CANDIDATES) + + @field_validator("ice_candidates") + @classmethod + def _bounded(cls, v: list[dict]) -> list[dict]: + if len(json.dumps(v)) > MAX_ICE_BYTES: + raise ValueError("ICE candidates too large") + return v class WebRTCOfferResponse(BaseModel): @@ -176,14 +190,6 @@ async def webrtc_offer( if len(body.sdp) > MAX_SDP_BYTES: raise HTTPException(status_code=413, detail="SDP too large") - # Logged here because this is the moment a browser starts a peer connection, - # and the address it starts it from is this one — the hub's own view of the - # TCP connection. Whatever address the peers then discover through STUN is - # theirs to negotiate and is not what a log should record. - db.add(IPLog(user_id=current_user.id, event="webrtc_offer", - ip_address=client_ip(request), detail=node_id[:8])) - await db.commit() - ws = _connected_nodes.get(node_id) if not ws: raise HTTPException(status_code=404, detail="Node not connected") @@ -205,6 +211,14 @@ async def webrtc_offer( raise HTTPException(status_code=429, detail="Too many connections to this node", headers={"Retry-After": str(max(1, math.ceil(wait)))}) + # Logged once the offer is going to a node, not before: the address a peer + # connection starts from is the hub's own view of this TCP connection, and an + # IP log row is kept a year — written before the checks above, any account + # could add rows for any string it named as a node. + db.add(IPLog(user_id=current_user.id, event="webrtc_offer", + ip_address=client_ip(request), detail=node_id[:8])) + await db.commit() + peer_id = str(uuid.uuid4()) answer_future: asyncio.Future = asyncio.get_event_loop().create_future() _webrtc_answers[peer_id] = answer_future |