diff options
Diffstat (limited to 'packages/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 | ||||
| -rw-r--r-- | packages/meshbay-hub/tests/test_hub_work_is_bounded.py | 51 |
3 files changed, 113 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 diff --git a/packages/meshbay-hub/tests/test_hub_work_is_bounded.py b/packages/meshbay-hub/tests/test_hub_work_is_bounded.py new file mode 100644 index 0000000..40ea81d --- /dev/null +++ b/packages/meshbay-hub/tests/test_hub_work_is_bounded.py @@ -0,0 +1,51 @@ +""" +What an authenticated caller can make the hub do, bounded where it was not. + +An offer wrote an IP-log row — kept a year — before any check, for whatever +string the caller named as a node, and carried an ICE list of any length to the +node. A node could send `update_groups` as fast as it liked, each one a +database read, with a list of any length. +""" + +import pytest +from meshbay_hub.db.models import IPLog +from sqlalchemy import func, select +from test_availability_between_members import _make_user + + +def _offer(client, user, node_id, candidates): + return client.post(f"/v1/nodes/{node_id}/webrtc/offer", + json={"sdp": "v=0\r\n", "ice_candidates": candidates}, + headers={"Authorization": f"Bearer {user['token']}"}) + + +@pytest.mark.asyncio +async def test_an_offer_that_goes_nowhere_writes_no_log_row(client, db_session): + user = await _make_user(client, "offer_nowhere") + for i in range(5): + r = await _offer(client, user, f"not-a-node-{i}", []) + assert r.status_code == 404 + rows = await db_session.scalar( + select(func.count()).select_from(IPLog).where(IPLog.event == "webrtc_offer")) + assert rows == 0 + + +@pytest.mark.asyncio +async def test_the_ice_list_is_bounded(client): + from meshbay_hub.api.signaling import MAX_ICE_CANDIDATES + user = await _make_user(client, "offer_ice") + one = {"candidate": "candidate:1 1 udp 2122260223 192.0.2.1 50000 typ host", + "sdpMid": "0", "sdpMLineIndex": 0} + assert (await _offer(client, user, "nowhere", [one] * MAX_ICE_CANDIDATES)).status_code == 404 + too_many = [one] * (MAX_ICE_CANDIDATES + 1) + assert (await _offer(client, user, "nowhere", too_many)).status_code == 422 + big = {"candidate": "x" * 40_000} + assert (await _offer(client, user, "nowhere", [big])).status_code == 422 + + +def test_a_node_reloading_is_budgeted(): + from meshbay_hub.api.revocation import UPDATE_GROUPS_BURST, _update_budget, _update_window + _update_window.clear() + assert all(_update_budget("node-a") for _ in range(UPDATE_GROUPS_BURST)) + assert not _update_budget("node-a") + assert _update_budget("node-b"), "one node's budget is not another's" |