diff options
Diffstat (limited to 'packages/meshbay-hub/src/meshbay_hub/api/revocation.py')
| -rw-r--r-- | packages/meshbay-hub/src/meshbay_hub/api/revocation.py | 53 |
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)) |