diff options
| author | Christophe Besson <cbesson@gmail.com> | 2026-10-01 10:34:26 +0200 |
|---|---|---|
| committer | Christophe Besson <cbesson@gmail.com> | 2026-10-01 10:34:26 +0200 |
| commit | 15e117673d2303bf476d4f78699e47913ce1aec0 (patch) | |
| tree | c70655684c1b0c41b5aeeb18ae4e4843ac72e375 | |
| parent | cf4fdda3523015248b3ff1f3e928ce9e69cc5a12 (diff) | |
| download | meshbay-15e117673d2303bf476d4f78699e47913ce1aec0.tar.gz | |
fix(node): one member holds a share of the node, sized past real use
128 peer sessions on the node, at most 64 per account (the hub names the
account with each offer; the node's own account is not counted). One account
plays at most half the stream slots, rounded up, and runs two subtitle
extractions at once. Frames after the handshake are 8 MiB (was 64), decoded
with per-container bounds, and a frame refused for either ends the session
instead of jamming its buffer (F-16).
Sized for the heaviest real member: twenty groups on one node, three devices
and a tab, up to 52 sessions. Measured: ~0.15 MiB and one fd per idle session.
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
11 files changed, 295 insertions, 35 deletions
diff --git a/docs/MESHBAY_DESIGN.md b/docs/MESHBAY_DESIGN.md index 609c751..1f2e43d 100644 --- a/docs/MESHBAY_DESIGN.md +++ b/docs/MESHBAY_DESIGN.md @@ -1898,7 +1898,7 @@ Node page: | `invite_ttl_hours` | 168 | how long a member invitation stays valid. Not an invitation link, which is fixed at seven days (§3.4) | | `pair_ttl_hours` | 24 | how long an operator pairing code stays valid | | `device_request_ttl_minutes` | 60 | how long a device request waits for approval. Comfort, not security: the code is bound to the keys by its hash | -| `max_concurrent_streams` | 8 | simultaneous video streams. One process per viewer, ~50 MB each; a slot is held for the length of a film, so this counts viewers | +| `max_concurrent_streams` | 8 | simultaneous video streams. One process per viewer, ~50 MB each; a slot is held for the length of a film, so this counts viewers. One account plays at most half of them, rounded up — three screens in one home fit, and no member alone takes every slot — and runs at most two subtitle extractions at once, which take the same slots | | `max_concurrent_downloads` | 8 | node-wide download leases (§5.5) | | `max_concurrent_uploads` | 8 | node-wide upload leases. A separate pool from downloads, because the two cost different things and one queue for both makes each cap meaningless | | `max_upload_gb` | 8 | the largest single upload, §6.4. Fractions are allowed, and it is read per chunk so a change reaches an upload already running | @@ -2020,8 +2020,13 @@ refused offer is indistinguishable, to the reader, from a node that is down: connected above all, fails at once, so a dead node costs no time. The node bounds its own total separately (`MAX_PEER_SESSIONS` in -`webrtc_server.py`), which is the limit that protects the machine whatever the -number of members. +`webrtc_server.py`, 128), which is the limit that protects the machine whatever +the number of members, and **one account holds at most half of it** (64, from +the account the hub names with each offer). The heaviest real member — twenty +groups on one node, three devices and a spare tab — holds up to 52: each device +keeps up to twelve connections for search and music, plus the open group page. +The node's own account is not counted; it is the operator's machine. An idle +connected session costs about 0.15 MiB and one file descriptor. **A node registered for no group shares one with nobody**, and is refused rather than exempted. Written as "check membership if the node claims any group", the diff --git a/docs/MESHBAY_NODE_PROTOCOL.md b/docs/MESHBAY_NODE_PROTOCOL.md index 9e51dc9..676b96f 100644 --- a/docs/MESHBAY_NODE_PROTOCOL.md +++ b/docs/MESHBAY_NODE_PROTOCOL.md @@ -102,7 +102,8 @@ Identical on every transport: | Bound | Value | Where | |---|---|---| | Max frame **before** the client's group-key proof | 64 KiB | `PRE_HANDSHAKE_MAX_MSG` | -| Max frame **after** the proof | 64 MiB | `MAX_MSG` | +| Max frame **after** the proof | 8 MiB | `MAX_MSG` — eight times the largest message a client sends (a sealed playlist blob, 1 MiB) | +| Decoding a frame | arrays 100 000, maps 10 000, strings 1 MiB, binaries 8 MiB, no extension types | `UNPACK_LIMITS` | | File chunk (plaintext) | 1 MiB | `CHUNK_SIZE` | | Video segment (plaintext, before encryption) | 256 KiB | `STREAM_SEGMENT_SIZE` | | Upload chunk sent by the browser | 48 KiB | fits the aiortc SCTP limit after msgpack overhead | @@ -115,6 +116,10 @@ into it, holding that much memory per connection for as long as it liked; a hund such connections is the node's memory, from peers that have proved nothing. Exceeding the limit is a hard protocol error and the buffer raises rather than truncating — truncating would hand a parser a valid-looking prefix of something it never received. +A frame refused for its size or its decoding ends the session: the buffer still starts +with it, so nothing after it could be read. The decoding bounds exist because a frame of +many tiny elements decodes to many times its size in memory; with them, an 8 MiB frame +costs tens of megabytes at worst. ### 3.3 Versioning field @@ -245,7 +250,7 @@ answer, and a discriminator of the message's own (`file_id`, `url`, `upload_id`, v +-------------------------+ | AUTHENTICATED | `_user_id` / `_group_id` set, - | full message set | frame limit raised to 64 MiB + | full message set | frame limit raised to 8 MiB +-----------+-------------+ | channel closes v @@ -416,7 +421,7 @@ fails if a transport skips a step. | compare_digest(proof); | | roster admits the user | | 9. session authenticated: | - | frame limit -> 64 MiB, | + | frame limit -> 8 MiB, | | peer registry, audit | | | | 10. handshake_ack | @@ -2422,7 +2427,7 @@ LP(x) = uint32be(len(x)) || x every field, no exceptions | Device attempts | 5 per connection | ” | | `MAX_DEVICES_PER_USER` | 5 | `roster.py` | | `MAX_LINK_INVITES_PER_GROUP` | 20 unredeemed invitation links | `roster.py` | -| `PRE_HANDSHAKE_MAX_MSG` / `MAX_MSG` | 64 KiB / 64 MiB | `webrtc/core.py` / `webrtc/limits.py` | +| `PRE_HANDSHAKE_MAX_MSG` / `MAX_MSG` | 64 KiB / 8 MiB | `webrtc/core.py` / `webrtc/limits.py` | | `CHUNK_SIZE` | 1 MiB | `webrtc/limits.py` | | `DOWNLOAD_BUFFER_HIGH` | 2 MiB | `webrtc/files.py` | | `MAX_UPLOAD_BYTES` | 8 GiB, default only — `max_upload_gb` overrides it per node | `webrtc/upload_handlers.py` | diff --git a/packages/meshbay-node/src/meshbay_node/daemon.py b/packages/meshbay-node/src/meshbay_node/daemon.py index 45a6b9a..095bcae 100644 --- a/packages/meshbay-node/src/meshbay_node/daemon.py +++ b/packages/meshbay-node/src/meshbay_node/daemon.py @@ -577,12 +577,12 @@ class NodeDaemon(EnrichmentMixin): log.info("QUIC server disabled ([node] quic_enabled = false)") # 8. Hub WebSocket (signaling + revocations + WebRTC offers) - async def on_webrtc_offer(sdp, peer_id, ice_candidates): + async def on_webrtc_offer(sdp, peer_id, ice_candidates, user_id=""): if not self._webrtc: return None try: answer_sdp, answer_ice = await self._webrtc.handle_offer( - sdp, peer_id) + sdp, peer_id, user_id) log.info("WebRTC answer for peer=%s (%d peers)", peer_id, self._webrtc.active_peers) return (answer_sdp, answer_ice) diff --git a/packages/meshbay-node/src/meshbay_node/hub_client.py b/packages/meshbay-node/src/meshbay_node/hub_client.py index 2ce8993..54f87f0 100644 --- a/packages/meshbay-node/src/meshbay_node/hub_client.py +++ b/packages/meshbay-node/src/meshbay_node/hub_client.py @@ -451,7 +451,8 @@ class HubClient: """Negotiate one WebRTC offer and return the answer, off the read loop.""" try: answer = await on_webrtc_offer( - msg["sdp"], msg["peer_id"], msg.get("ice_candidates", [])) + msg["sdp"], msg["peer_id"], msg.get("ice_candidates", []), + str(msg.get("user_id") or "")) except Exception as e: log.warning("WebRTC offer from %s failed: %s", str(msg.get("peer_id"))[:8], e) diff --git a/packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/streaming.py b/packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/streaming.py index 4337e24..7156e84 100644 --- a/packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/streaming.py +++ b/packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/streaming.py @@ -2,6 +2,7 @@ stream a session holds, and the ffmpeg pipeline behind it.""" import asyncio +import contextlib import logging import time @@ -34,6 +35,10 @@ log = logging.getLogger("meshbay_node.transport.webrtc_server") # so the operator sets `max_concurrent_streams` under [node] in node.toml. This # value applies when they have said nothing. MAX_CONCURRENT_TRANSCODES = 8 +# Subtitle extractions one account may run at once. The player asks for one +# track at a time; two covers a quick change of track. Each holds a transcode +# slot for up to fifteen minutes on a long film. +MAX_SUBTITLE_JOBS_PER_ACCOUNT = 2 STREAM_SEGMENT_SIZE = 256 * 1024 @@ -179,6 +184,36 @@ class StreamingMixin: self._stream_task = asyncio.current_task() await self._stream_video(msg) + @contextlib.contextmanager + def _account_share(self, kind: str, limit: int): + """ + Hold one of this account's `limit` places for `kind`, or yield False. + + Counted on the node, across every session of the account: a member's + devices and tabs share one allowance. The node's own account is not + counted — it is the operator's machine. + """ + user = getattr(self, "_user_id", "") or "" + if not user or user == self._ctx.get("node_user_id"): + yield True + return + held = self._ctx.setdefault(f"_{kind}_by_account", {}) + if held.get(user, 0) >= limit: + yield False + return + held[user] = held.get(user, 0) + 1 + try: + yield True + finally: + held[user] -= 1 + if held[user] <= 0: + held.pop(user, None) + + def _streams_per_account(self) -> int: + """Half the node's viewers, rounded up: three screens in one home fit, + and no member alone takes every slot the operator set.""" + return max(1, -(-self._stream_capacity() // 2)) + def _transcode_semaphore(self) -> asyncio.Semaphore: """The node's stream budget, shared across every peer. @@ -205,22 +240,28 @@ class StreamingMixin: ctx = self._ctx log.info("stream: waiting for a slot (%d of %d in use)", ctx.get("_streams_in_flight", 0), self._stream_capacity()) - async with sem: - # Counted here rather than read back out of the semaphore's private - # `_value`: `set_capacity` needs to know how many slots are held in - # order to resize without letting the pool overshoot, and a number - # this code maintains itself is one that survives the semaphore - # object being replaced underneath it. - ctx["_streams_in_flight"] = ctx.get("_streams_in_flight", 0) + 1 - log.info("stream: slot acquired (%d of %d in use)", - ctx["_streams_in_flight"], self._stream_capacity()) - try: - await self._stream_video_inner(msg) - finally: - ctx["_streams_in_flight"] = max( - 0, ctx.get("_streams_in_flight", 1) - 1) - log.info("stream: slot released (%d of %d in use)", + with self._account_share("streams", self._streams_per_account()) as ok: + if not ok: + self._send({"type": "error", + "detail": "Too many videos playing from this account, " + "stop one and retry"}) + return + async with sem: + # Counted here rather than read back out of the semaphore's private + # `_value`: `set_capacity` needs to know how many slots are held in + # order to resize without letting the pool overshoot, and a number + # this code maintains itself is one that survives the semaphore + # object being replaced underneath it. + ctx["_streams_in_flight"] = ctx.get("_streams_in_flight", 0) + 1 + log.info("stream: slot acquired (%d of %d in use)", ctx["_streams_in_flight"], self._stream_capacity()) + try: + await self._stream_video_inner(msg) + finally: + ctx["_streams_in_flight"] = max( + 0, ctx.get("_streams_in_flight", 1) - 1) + log.info("stream: slot released (%d of %d in use)", + ctx["_streams_in_flight"], self._stream_capacity()) def _stream_capacity(self) -> int: return self._ctx.get("max_concurrent_streams") or MAX_CONCURRENT_TRANSCODES diff --git a/packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/subtitles.py b/packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/subtitles.py index 70781eb..35211af 100644 --- a/packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/subtitles.py +++ b/packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/subtitles.py @@ -10,6 +10,7 @@ from meshbay_common.protocol import MNP from meshbay_node.media_probe import probe_video as _probe_video from meshbay_node.roots import off_disk +from meshbay_node.transport.webrtc.apps.streaming import MAX_SUBTITLE_JOBS_PER_ACCOUNT from meshbay_node.transport.webrtc.disk import _locate from meshbay_node.transport.webrtc.media_tools import ( _extract_subtitle_to_webvtt, @@ -135,10 +136,16 @@ class SubtitlesMixin: return budget = _subtitle_timeout_for(entry.size) - async with sem: - log.info("subtitle: extracting file=%s track=%d (slot taken, up to %.0fs)", - file_id[:12], ordinal, budget) - blob = await _extract_subtitle_to_webvtt(file_path, ordinal, budget) + with self._account_share("subtitles", MAX_SUBTITLE_JOBS_PER_ACCOUNT) as ok: + if not ok: + log.info("subtitle: refused, account at its extraction share") + self._send({"type": "error", + "detail": "Subtitles are already being prepared, retry shortly"}) + return + async with sem: + log.info("subtitle: extracting file=%s track=%d (slot taken, up to %.0fs)", + file_id[:12], ordinal, budget) + blob = await _extract_subtitle_to_webvtt(file_path, ordinal, budget) subtitle_hash = blake3.blake3(blob).hexdigest() await media_cache.put_thumb(subtitle_hash, synthetic_id, blob) diff --git a/packages/meshbay-node/src/meshbay_node/transport/webrtc/channel.py b/packages/meshbay-node/src/meshbay_node/transport/webrtc/channel.py index 107a43e..31cc83c 100644 --- a/packages/meshbay-node/src/meshbay_node/transport/webrtc/channel.py +++ b/packages/meshbay-node/src/meshbay_node/transport/webrtc/channel.py @@ -7,7 +7,7 @@ import struct import msgpack from aiortc import RTCPeerConnection -from meshbay_node.transport.webrtc.limits import MAX_MSG +from meshbay_node.transport.webrtc.limits import MAX_MSG, UNPACK_LIMITS def _extract_dtls_fingerprint(sdp: str) -> bytes: @@ -77,7 +77,7 @@ class _DataChannelBuffer: break msg_bytes = bytes(self._buf[4:4 + length]) del self._buf[:4 + length] - yield msgpack.unpackb(msg_bytes, raw=False) + yield msgpack.unpackb(msg_bytes, raw=False, **UNPACK_LIMITS) def _get_remote_ip(pc: RTCPeerConnection) -> str: diff --git a/packages/meshbay-node/src/meshbay_node/transport/webrtc/core.py b/packages/meshbay-node/src/meshbay_node/transport/webrtc/core.py index 882ddbc..ae3f0e1 100644 --- a/packages/meshbay-node/src/meshbay_node/transport/webrtc/core.py +++ b/packages/meshbay-node/src/meshbay_node/transport/webrtc/core.py @@ -139,7 +139,21 @@ class SessionCore: log.info("WebRTC data received: %d bytes, msg #%d (peer=%s)", len(message), self._msg_count, self._peer_id) self._buffer.feed(message) - for msg in self._buffer.messages(): + decoded = self._buffer.messages() + while True: + try: + msg = next(decoded) + except StopIteration: + break + except ValueError as e: + # Over the size limit, or a container past its decode + # limit. The buffer still starts with that frame, so every + # later message would fail the same way: the session ends + # here. Only decoding is caught — a handler's own error is + # not a reason to drop the peer. + log.warning("Closing peer %s: %s", self._peer_id, e) + self._spawn(self.close()) + break self._handle_message(msg) if _WEBRTC_TRACE: diff --git a/packages/meshbay-node/src/meshbay_node/transport/webrtc/limits.py b/packages/meshbay-node/src/meshbay_node/transport/webrtc/limits.py index 7d458f4..86412e0 100644 --- a/packages/meshbay-node/src/meshbay_node/transport/webrtc/limits.py +++ b/packages/meshbay-node/src/meshbay_node/transport/webrtc/limits.py @@ -2,7 +2,21 @@ CHUNK_SIZE = 1024 * 1024 -MAX_MSG = 64 * 1024 * 1024 +# The largest message a peer may send once it has proved the group key. The +# largest a client really sends is a sealed playlist blob, 1 MiB (blobs.py); +# chat is 64 KiB and an upload chunk 48 KiB. Eight times the largest, because a +# message of many small objects decodes to several times its size in memory. +MAX_MSG = 8 * 1024 * 1024 + +# Per container, when a message is decoded: nothing a client sends comes near +# them, and without them one message of tiny elements is one enormous list. +UNPACK_LIMITS = { + "max_array_len": 100_000, + "max_map_len": 10_000, + "max_str_len": 1024 * 1024, + "max_bin_len": MAX_MSG, + "max_ext_len": 0, +} # What the `tr` on a chunk request turned out to be (see `_lease_of`). diff --git a/packages/meshbay-node/src/meshbay_node/transport/webrtc_server.py b/packages/meshbay-node/src/meshbay_node/transport/webrtc_server.py index 20636fa..0ca4b3d 100644 --- a/packages/meshbay-node/src/meshbay_node/transport/webrtc_server.py +++ b/packages/meshbay-node/src/meshbay_node/transport/webrtc_server.py @@ -70,7 +70,15 @@ log = logging.getLogger(__name__) # so the cost to an operator grew with the number of people in their groups. # Sized to be unreachable in ordinary use: a browser holds one connection per # open group, and a handshake unfinished after a minute is not going to finish. -MAX_PEER_SESSIONS = 64 +MAX_PEER_SESSIONS = 128 +# One account's share of them: half. A member of twenty groups hosted here, +# with three devices and a spare tab, holds up to 52 (each device keeps up to +# twelve connections for search and music, plus the open group page), so the +# share never refuses real use — and no single member can hold more than half +# of what the node will take. The node's own account is not counted against it: +# it is the operator's machine. Measured, an idle connected session costs about +# 0.15 MiB and one file descriptor. +MAX_PEER_SESSIONS_PER_ACCOUNT = MAX_PEER_SESSIONS // 2 UNAUTHENTICATED_SESSION_TIMEOUT = 60 # seconds @@ -238,8 +246,13 @@ class WebRTCTransport: pass return + def _sessions_of(self, user_id: str) -> int: + return sum(1 for s in self._sessions.values() + if (getattr(s, "_offer_user", "") or getattr(s, "_user_id", "") or "") + == user_id) + async def handle_offer( - self, offer_sdp: str, peer_id: str, + self, offer_sdp: str, peer_id: str, user_id: str = "", ) -> tuple[str, list[dict]]: """ Process a WebRTC SDP offer from a browser client. @@ -267,9 +280,18 @@ class WebRTCTransport: log.warning("Refusing WebRTC offer: %d peer sessions already open", len(self._sessions)) raise RuntimeError("Node is at its peer-connection limit") + # The account the hub authenticated for this offer. A hub that lied + # could only move the count between accounts; it can refuse offers + # outright already. + if (user_id and user_id != self._ctx.get("node_user_id") + and self._sessions_of(user_id) >= MAX_PEER_SESSIONS_PER_ACCOUNT): + log.warning("Refusing WebRTC offer: account %s already holds %d sessions", + user_id[:8], MAX_PEER_SESSIONS_PER_ACCOUNT) + raise RuntimeError("Account is at its peer-connection share") pc = RTCPeerConnection(configuration=config) session = WebRTCPeerSession(pc, self._ctx, peer_id=peer_id) + session._offer_user = user_id self._sessions[peer_id] = session self._reap_if_unauthenticated(peer_id) diff --git a/packages/meshbay-node/tests/test_member_capacity.py b/packages/meshbay-node/tests/test_member_capacity.py new file mode 100644 index 0000000..16d112d --- /dev/null +++ b/packages/meshbay-node/tests/test_member_capacity.py @@ -0,0 +1,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()) |