summaryrefslogtreecommitdiffstats
path: root/packages/meshbay-hub/src/meshbay_hub/api/revocation.py
diff options
context:
space:
mode:
Diffstat (limited to 'packages/meshbay-hub/src/meshbay_hub/api/revocation.py')
-rw-r--r--packages/meshbay-hub/src/meshbay_hub/api/revocation.py53
1 files changed, 38 insertions, 15 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))