aboutsummaryrefslogtreecommitdiffstats
path: root/packages
diff options
context:
space:
mode:
Diffstat (limited to 'packages')
-rw-r--r--packages/meshbay-hub/src/meshbay_hub/api/revocation.py53
-rw-r--r--packages/meshbay-hub/src/meshbay_hub/api/signaling.py34
-rw-r--r--packages/meshbay-hub/tests/test_hub_work_is_bounded.py51
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"