""" 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())