summaryrefslogtreecommitdiffstats
path: root/packages/meshbay-hub/tests/test_signaling_limits.py
blob: 23260110cf41d82ce7441e991f738eee43ba608f (plain) (blame)
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
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)