diff options
Diffstat (limited to 'packages/meshbay-hub/tests/test_signaling_limits.py')
| -rw-r--r-- | packages/meshbay-hub/tests/test_signaling_limits.py | 172 |
1 files changed, 172 insertions, 0 deletions
diff --git a/packages/meshbay-hub/tests/test_signaling_limits.py b/packages/meshbay-hub/tests/test_signaling_limits.py new file mode 100644 index 0000000..2326011 --- /dev/null +++ b/packages/meshbay-hub/tests/test_signaling_limits.py @@ -0,0 +1,172 @@ +""" +The hub's ceilings on offers must bound a member without failing an ordinary one. + +They did fail one. `MAX_PENDING_PER_USER` was 3 and Search dialled three groups +at once, so a single tile or reconnection beside the sweep was refused with 429, +and the page reported the group's node as unreachable while that node was +answering every other offer in under a second. The per-address limit, 30 a +minute per node, was spent by four or five reloads of a page on a phone whose +groups share one node. Measured on meshbay.org on 2026-09-23: 93 offers refused +in half an hour, all from one phone, all to nodes that were up. + +What these tests hold is the shape the limits must admit — twenty groups on one +node, dialled together — and that each refusal says when to come back, because +the browser now does (transport.js, `postOffer`). +""" + +import asyncio +import json +import time + +import pytest +from test_availability_between_members import ( + _add_member, + _announce_node, + _make_group, + _make_user, +) + + +class _SlowNode: + """A node that answers each offer once `release` is set.""" + + def __init__(self, node_id: str): + self.node_id = node_id + self.release = asyncio.Event() + self.offers = 0 + + async def send_text(self, text): + from meshbay_hub.api.signaling import handle_webrtc_answer + + msg = json.loads(text) + self.offers += 1 + + async def answer(): + await self.release.wait() + handle_webrtc_answer({"peer_id": msg["peer_id"], "sdp": "v=0\r\nanswer", + "ice_candidates": []}, self.node_id) + + asyncio.get_running_loop().create_task(answer()) + + +async def _member_with_node(client, prefix: str): + owner = await _make_user(client, f"{prefix}_owner") + member = await _make_user(client, f"{prefix}_member") + group_id = await _make_group(client, owner, f"{prefix}-group") + await _add_member(client, owner, group_id, member) + node_id = await _announce_node(client, owner) + return member, group_id, node_id + + +def _offer(client, node_id, user): + return client.post(f"/v1/nodes/{node_id}/webrtc/offer", + json={"sdp": "v=0\r\noffer", "ice_candidates": []}, + headers={"Authorization": f"Bearer {user['token']}"}) + + +@pytest.mark.asyncio +async def test_twenty_groups_dialled_together_are_all_brokered(client): + """An account with twenty groups on one node, all dialled at once — more + than one Search page ever does — is refused nothing.""" + from meshbay_hub.api import revocation as rev + + member, group_id, node_id = await _member_with_node(client, "sig_twenty") + node = _SlowNode(node_id) + rev._connected_nodes[node_id] = node + rev._node_groups[node_id] = [group_id] + try: + pending = [asyncio.ensure_future(_offer(client, node_id, member)) for _ in range(20)] + while node.offers < 20: + await asyncio.sleep(0.01) + node.release.set() + statuses = [r.status_code for r in await asyncio.gather(*pending)] + assert statuses == [200] * 20, statuses + finally: + rev._connected_nodes.pop(node_id, None) + rev._node_groups.pop(node_id, None) + + +@pytest.mark.asyncio +async def test_the_pending_ceiling_refuses_with_a_time_to_come_back(client): + from meshbay_hub.api import revocation as rev + from meshbay_hub.api.signaling import MAX_PENDING_PER_USER + + assert MAX_PENDING_PER_USER >= 32, "sized for three devices of a twenty-group account" + member, group_id, node_id = await _member_with_node(client, "sig_pending") + node = _SlowNode(node_id) + rev._connected_nodes[node_id] = node + rev._node_groups[node_id] = [group_id] + try: + held = [asyncio.ensure_future(_offer(client, node_id, member)) + for _ in range(MAX_PENDING_PER_USER)] + while node.offers < MAX_PENDING_PER_USER: + await asyncio.sleep(0.01) + refused = await _offer(client, node_id, member) + assert refused.status_code == 429, refused.text + assert refused.headers.get("Retry-After") == "1" + node.release.set() + assert all(r.status_code == 200 for r in await asyncio.gather(*held)) + # Once they are answered the account has its places back. + node.release = asyncio.Event() + node.release.set() + again = await _offer(client, node_id, member) + assert again.status_code == 200, again.text + finally: + rev._connected_nodes.pop(node_id, None) + rev._node_groups.pop(node_id, None) + + +@pytest.mark.asyncio +async def test_an_exhausted_node_budget_says_when_and_spares_other_nodes(client): + """The per-node budget bounds what one member costs one node, and only that + node: an account that has spent it on one machine still reaches the others.""" + from meshbay_hub.api import revocation as rev + from meshbay_hub.api import signaling + + member, group_id, node_id = await _member_with_node(client, "sig_budget") + # A second machine hosting the same group, as a group with two nodes has. + _, _, other_node = await _member_with_node(client, "sig_budget2") + node, other = _SlowNode(node_id), _SlowNode(other_node) + node.release.set() + other.release.set() + rev._connected_nodes[node_id] = node + rev._node_groups[node_id] = [group_id] + rev._connected_nodes[other_node] = other + rev._node_groups[other_node] = [group_id] + try: + signaling._offer_buckets[(member["user_id"], node_id)] = (0.0, time.monotonic()) + refused = await _offer(client, node_id, member) + assert refused.status_code == 429, refused.text + assert int(refused.headers["Retry-After"]) >= 1 + assert node.offers == 0, "a refused offer still reached the node" + + ok = await _offer(client, other_node, member) + assert ok.status_code == 200, ok.text + finally: + signaling._offer_buckets.pop((member["user_id"], node_id), None) + for n in (node_id, other_node): + rev._connected_nodes.pop(n, None) + rev._node_groups.pop(n, None) + + +def test_the_node_budget_admits_a_burst_and_refills(): + """The arithmetic, without a clock: a twenty-group page reloaded six times + in a row fits, the offer after the burst waits half a second, and a minute + of quiet gives the whole burst back.""" + from meshbay_hub.api import signaling + + key = ("budget-user", "budget-node") + signaling._offer_buckets.pop(key, None) + try: + t0 = 1000.0 + assert signaling.OFFER_BURST >= 120 + for _ in range(signaling.OFFER_BURST): + assert signaling._take_offer(*key, t0) is None + wait = signaling._take_offer(*key, t0) + assert wait == pytest.approx(1 / signaling.OFFER_REFILL_PER_S) + assert signaling._take_offer(*key, t0 + wait) is None + later = t0 + signaling.OFFER_BURST / signaling.OFFER_REFILL_PER_S + wait + for _ in range(signaling.OFFER_BURST): + assert signaling._take_offer(*key, later) is None + finally: + signaling._offer_buckets.pop(key, None) |