aboutsummaryrefslogtreecommitdiffstats
path: root/packages/meshbay-node/tests/test_member_capacity.py
blob: 16d112df6af8f66c7e124594048d8df926de6174 (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
"""
What one member may hold of a node: a share, sized so that real use never
meets it.

The heaviest real member — twenty groups on this node, three devices and a
spare tab — holds up to 52 peer sessions (twelve per device for search and
music, plus the open group page). The node holds 128, and one account at most
half. One account may play half the node's video slots and run two subtitle
extractions. The node's own account is the operator's machine and is not
counted. And a message is at most 8 MiB, decoded with a bound on every
container, because one message of tiny elements decodes to many times its size.
"""

import struct
from unittest.mock import MagicMock

import msgpack
import pytest
from meshbay_node.transport.webrtc.channel import _DataChannelBuffer
from meshbay_node.transport.webrtc.limits import MAX_MSG
from meshbay_node.transport.webrtc_server import (
    MAX_PEER_SESSIONS,
    MAX_PEER_SESSIONS_PER_ACCOUNT,
    WebRTCPeerSession,
    WebRTCTransport,
)


class _Held:
    def __init__(self, user):
        self._offer_user = user
        self._user_id = user

    async def close(self):
        pass


def _transport(**ctx) -> WebRTCTransport:
    tp = WebRTCTransport(sk_node=MagicMock(), hub_pk_pem=b"", gek=None,
                         roots=None, index=None)
    tp._ctx.update(ctx)
    return tp


def test_the_shares_are_what_was_agreed():
    assert MAX_PEER_SESSIONS == 128
    assert MAX_PEER_SESSIONS_PER_ACCOUNT == 64
    assert MAX_MSG == 8 * 1024 * 1024
    # The heaviest real member fits with room to spare.
    assert 4 * (12 + 1) < MAX_PEER_SESSIONS_PER_ACCOUNT


@pytest.mark.asyncio
async def test_one_account_cannot_hold_more_than_its_share():
    tp = _transport()
    for i in range(MAX_PEER_SESSIONS_PER_ACCOUNT):
        tp._sessions[f"a-{i}"] = _Held("alice")

    with pytest.raises(RuntimeError, match="share"):
        await tp.handle_offer("v=0", "alice-one-more", "alice")
    assert "alice-one-more" not in tp._sessions

    # Somebody else still gets in: what stops this offer is not the share.
    try:
        await tp.handle_offer("v=0", "bob-first", "bob")
    except RuntimeError as e:
        assert "share" not in str(e) and "limit" not in str(e)
    except Exception:
        pass


@pytest.mark.asyncio
async def test_the_operators_own_account_is_not_counted():
    tp = _transport(node_user_id="operator")
    for i in range(MAX_PEER_SESSIONS_PER_ACCOUNT):
        tp._sessions[f"o-{i}"] = _Held("operator")
    try:
        await tp.handle_offer("v=0", "operator-more", "operator")
    except RuntimeError as e:
        assert "share" not in str(e)
    except Exception:
        pass


def _session(ctx, user):
    s = WebRTCPeerSession.__new__(WebRTCPeerSession)
    s._ctx = ctx
    s._user_id = user
    s.sent = []
    s._send = s.sent.append
    return s


def test_an_accounts_devices_share_one_allowance():
    ctx = {}
    phone, laptop = _session(ctx, "alice"), _session(ctx, "alice")
    with phone._account_share("streams", 2) as a, laptop._account_share("streams", 2) as b:
        assert a and b
        with phone._account_share("streams", 2) as c:
            assert c is False
        with _session(ctx, "bob")._account_share("streams", 2) as d:
            assert d, "another member is not counted against alice"
    with phone._account_share("streams", 2) as e:
        assert e, "a place is given back when its stream ends"


def test_the_operator_is_not_counted_for_streams_either():
    ctx = {"node_user_id": "operator"}
    s = _session(ctx, "operator")
    with s._account_share("streams", 1) as a, s._account_share("streams", 1) as b:
        assert a and b


@pytest.mark.asyncio
async def test_a_stream_past_the_accounts_share_is_refused_by_name():
    ctx = {"max_concurrent_streams": 8, "_streams_by_account": {"alice": 4}}
    s = _session(ctx, "alice")
    assert s._streams_per_account() == 4
    await s._stream_video({"file_id": "x"})
    assert s.sent[-1]["type"] == "error"
    assert "Too many videos" in s.sent[-1]["detail"]


@pytest.mark.parametrize("cap,share", [(1, 1), (2, 1), (3, 2), (8, 4), (9, 5)])
def test_the_stream_share_is_half_rounded_up(cap, share):
    assert _session({"max_concurrent_streams": cap}, "a")._streams_per_account() == share


def _frame(obj) -> bytes:
    body = msgpack.packb(obj, use_bin_type=True)
    return struct.pack(">I", len(body)) + body


def test_the_largest_real_message_passes():
    buf = _DataChannelBuffer()
    buf.feed(_frame({"type": "user_blob_store", "blob": b"x" * (1024 * 1024 + 64)}))
    assert len(list(buf.messages())) == 1


def test_a_message_of_tiny_elements_is_refused():
    buf = _DataChannelBuffer()
    buf.feed(_frame({"type": "x", "items": [None] * 200_000}))
    with pytest.raises(ValueError):
        list(buf.messages())


def test_a_message_over_the_limit_is_refused():
    buf = _DataChannelBuffer()
    buf.feed(struct.pack(">I", MAX_MSG + 1) + b"x")
    with pytest.raises(ValueError):
        list(buf.messages())