aboutsummaryrefslogtreecommitdiffstats
path: root/packages/meshbay-node
diff options
context:
space:
mode:
Diffstat (limited to 'packages/meshbay-node')
-rw-r--r--packages/meshbay-node/src/meshbay_node/daemon.py4
-rw-r--r--packages/meshbay-node/src/meshbay_node/hub_client.py3
-rw-r--r--packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/streaming.py71
-rw-r--r--packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/subtitles.py15
-rw-r--r--packages/meshbay-node/src/meshbay_node/transport/webrtc/channel.py4
-rw-r--r--packages/meshbay-node/src/meshbay_node/transport/webrtc/core.py16
-rw-r--r--packages/meshbay-node/src/meshbay_node/transport/webrtc/limits.py16
-rw-r--r--packages/meshbay-node/src/meshbay_node/transport/webrtc_server.py26
-rw-r--r--packages/meshbay-node/tests/test_member_capacity.py151
9 files changed, 278 insertions, 28 deletions
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())