diff options
Diffstat (limited to 'packages/meshbay-node/src/meshbay_node/transport/webrtc')
5 files changed, 99 insertions, 23 deletions
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`). |