diff options
Diffstat (limited to 'packages/meshbay-node')
5 files changed, 325 insertions, 25 deletions
diff --git a/packages/meshbay-node/src/meshbay_node/config.py b/packages/meshbay-node/src/meshbay_node/config.py index e4444d3..3de2473 100644 --- a/packages/meshbay-node/src/meshbay_node/config.py +++ b/packages/meshbay-node/src/meshbay_node/config.py @@ -37,6 +37,13 @@ ui_port = 18000 # local admin UI (127.0.0.1 only) invite_ttl_hours = 168 # 7 days pair_ttl_hours = 24 +# How many people may watch a video at once. One ffmpeg runs per viewer for as +# long as they watch — it remuxes rather than re-encodes, so it costs little CPU +# and around 50 MB of memory, and spends most of the film idle. Past this, a +# viewer is told the server is busy. Raise it on a machine with memory to spare; +# lower it on a Pi. +max_concurrent_streams = 8 + # Browser and native clients reach this node over WebRTC DataChannel via hub # signaling — no inbound port to open. QUIC is the optional direct path. @@ -82,6 +89,11 @@ class NodeConfig: # the SSH session that printed it. invite_ttl_hours: int = 168 # 7 days pair_ttl_hours: int = 24 + # How many people may watch a video at the same time. One ffmpeg runs per + # viewer for as long as they watch, so this is the knob that decides when + # the node answers "server busy" — see MAX_CONCURRENT_TRANSCODES in + # transport/webrtc_server.py for what one costs. + max_concurrent_streams: int = 8 @dataclass @@ -121,6 +133,30 @@ class Config: return self.groups[0] if self.groups else GroupConfig() +def _positive(value: object, default: int, name: str) -> int: + """A count that must be at least one, or the default with a word about it. + + Zero is the dangerous one: `asyncio.Semaphore(0)` is not "no limit", it is + a node where no video ever plays and nothing in the log says why. + """ + # bool before int: TOML `true` is a bool, and `int(True)` is 1 — a node + # where exactly one person may watch, arrived at by a typo and announced + # nowhere. + if isinstance(value, bool) or not isinstance(value, (int, str)): + log.warning("%s = %r is not a count — using %d", name, value, default) + return default + try: + n = int(value) + except (TypeError, ValueError): + log.warning("%s = %r is not a number — using %d", name, value, default) + return default + if n < 1: + log.warning("%s = %d would stop the feature entirely — using %d", + name, n, default) + return default + return n + + def load_config(path: Path = DEFAULT_CONFIG_PATH) -> Config: """ Load config from TOML file. Supports both single [group] and @@ -144,6 +180,9 @@ def load_config(path: Path = DEFAULT_CONFIG_PATH) -> Config: nd.get("invite_ttl_hours", cfg.node.invite_ttl_hours)) cfg.node.pair_ttl_hours = int( nd.get("pair_ttl_hours", cfg.node.pair_ttl_hours)) + cfg.node.max_concurrent_streams = _positive( + nd.get("max_concurrent_streams", cfg.node.max_concurrent_streams), + cfg.node.max_concurrent_streams, "max_concurrent_streams") # Multi-group: [[groups]] array if "groups" in raw: @@ -189,6 +228,10 @@ def load_config(path: Path = DEFAULT_CONFIG_PATH) -> Config: cfg.hub.username = user if port := os.environ.get("MESHBAY_QUIC_PORT"): cfg.node.quic_port = int(port) + if streams := os.environ.get("MESHBAY_MAX_CONCURRENT_STREAMS"): + cfg.node.max_concurrent_streams = _positive( + streams, cfg.node.max_concurrent_streams, + "MESHBAY_MAX_CONCURRENT_STREAMS") return cfg diff --git a/packages/meshbay-node/src/meshbay_node/daemon.py b/packages/meshbay-node/src/meshbay_node/daemon.py index 6fafc74..8b4a1d7 100644 --- a/packages/meshbay-node/src/meshbay_node/daemon.py +++ b/packages/meshbay-node/src/meshbay_node/daemon.py @@ -274,6 +274,7 @@ class NodeDaemon: index=first["index"], groups=groups_ctx, denylist=denylist, + max_concurrent_streams=self._config.node.max_concurrent_streams, ) # No global chat_store here: each group's store lives in # groups_ctx[gid]["chat_store"] and is resolved per session via 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 91d801c..5f5c9ff 100644 --- a/packages/meshbay-node/src/meshbay_node/transport/webrtc_server.py +++ b/packages/meshbay-node/src/meshbay_node/transport/webrtc_server.py @@ -90,7 +90,20 @@ MAX_UPLOAD_BYTES = 4 * 1024 * 1024 * 1024 # 4 GB per file PRE_HANDSHAKE_MAX_MSG = 64 * 1024 # ffmpeg is spawned per stream request; without a cap any member can fork-bomb # the node by requesting many streams at once (H6). -MAX_CONCURRENT_TRANSCODES = 2 +# +# Two was sized when a stream was a burst: the client took segments as fast as +# it could append them, so a slot was held for the minute it took to push the +# file and then came back. Now that the client only pulls ninety seconds ahead +# of the playhead, a slot is held for as long as the film runs — so two slots +# means two people can watch anything at all, and the third is refused for the +# next hour and a half. The work behind a slot has not changed and is small: +# ffmpeg runs `-c copy`, a remux with no encoding in it, and spends most of the +# film blocked on a pipe nobody is reading. +# +# This is the default, not the policy: the right number depends on the machine, +# so the operator sets `max_concurrent_streams` under [node] in node.toml. This +# value applies when they have said nothing. +MAX_CONCURRENT_TRANSCODES = 8 # Bundle fetches are served in the pre-proof window (C4). Bounded and audited # until the native client removes remote keypair bundles entirely. MAX_PRE_PROOF_FETCHES = 4 @@ -326,6 +339,12 @@ class WebRTCPeerSession: self._stream_credit = 0 self._stream_credit_evt = asyncio.Event() self._stream_stopped = False + # When the peer last said anything about this stream. See + # _await_stream_credit: silence is what ends a stream, not stinginess. + self._stream_heard_at = 0.0 + # Diagnostics: how many `stream_more n=0` the peer sent. See + # _grant_stream_credit — it tells a paced client from an unpaced one. + self._stream_keepalives = 0 # The stream this session currently owns. One viewer plays one film at # a time, so a second request means the first is over — see # _replace_stream for why waiting for it to time out is not an option. @@ -434,6 +453,28 @@ class WebRTCPeerSession: self._spawn(self._replace_stream(msg)) elif mtype == MNP.STREAM_MORE: self._grant_stream_credit(msg) + elif mtype == "client_diag": + # Diagnostics only. The node acts on none of it — it writes it + # next to its own view of the same stream, which is the only + # place the two halves can be compared when the client is a + # phone with no console. + # Every field is peer-controlled, so each is stringified and + # cut short: this is a log line, not a channel for writing + # whatever one likes into the operator's file. + def _f(key: str, n: int = 24) -> str: + return str(msg.get(key))[:n].replace("\n", " ") + # Debug: one line every five seconds per viewer. Run the daemon + # with --log-level debug to see inside a player that is + # misbehaving — it is the only view of the browser there is + # when the browser is a phone. + log.debug( + "stream: client t=%ss ahead=%ss ready=%s paused=%s " + "stalled=%s q=%s inflight=%s appending=%s updating=%s " + "quota=%s ms=%s err=%s ranges=[%s] (sent=%d)", + _f("t"), _f("ahead"), _f("ready"), _f("paused"), + _f("stalled"), _f("q"), _f("inflight"), _f("appending"), + _f("updating"), _f("quota"), _f("ms"), _f("err", 80), + _f("ranges", 120), self._stream_segments) elif mtype == MNP.STREAM_STOP: age = (time.monotonic() - self._stream_started_at if self._stream_started_at else -1) @@ -1325,10 +1366,7 @@ class WebRTCPeerSession: self._send({"type": "error", "detail": "File not on disk"}) return - sem = self._ctx.get("_transcode_sem") - if sem is None: - sem = asyncio.Semaphore(MAX_CONCURRENT_TRANSCODES) - self._ctx["_transcode_sem"] = sem + sem = self._transcode_semaphore() try: async with sem: @@ -1830,14 +1868,30 @@ class WebRTCPeerSession: }) def _grant_stream_credit(self, msg: dict) -> None: + """ + The client has room for more segments. + + `n` of zero is a keepalive, not a no-op: a viewer whose buffer is + already a minute and a half ahead of the playhead deliberately grants + nothing, and must still be able to say it is there. Without that, the + stall timeout below cannot tell a paused film from a closed tab. + """ log.debug("stream credit +%s (had %d, sent %d)", msg.get("n"), self._stream_credit, self._stream_segments) - """The client has room for more segments.""" try: n = int(msg.get("n", 1)) except (TypeError, ValueError): n = 1 + if n == 0: + # The fingerprint of a client that bounds its read-ahead. A client + # that never sends one is granting credit per append — which is + # what fills the browser's buffer ceiling and wedges the player. + self._stream_keepalives += 1 + if self._stream_keepalives == 1: + log.info("stream: peer is pacing itself (first keepalive at " + "%d segments)", self._stream_segments) self._stream_credit += max(0, min(n, STREAM_MAX_CREDIT)) + self._stream_heard_at = time.monotonic() self._stream_credit_evt.set() def _stop_stream(self) -> None: @@ -1860,8 +1914,22 @@ class WebRTCPeerSession: fast as it is produced, and the browser holds a four gigabyte film in a JavaScript array while MediaSource consumes it a segment at a time. """ - waited = 0.0 + # Measured from the last thing the peer said, not from the start of the + # wait: a viewer that is buffered well ahead sends keepalives and grants + # nothing for minutes at a time, and that is a watched film, not a + # stalled one. + self._stream_heard_at = time.monotonic() + waiting_since = 0.0 while self._stream_credit <= 0: + if waiting_since == 0.0: + waiting_since = time.monotonic() + # Debug: a paced viewer runs out of credit between every + # window, so this is one line per eight segments — hundreds + # per film. It is worth having, but not by default. + log.debug("stream: out of credit at %d segments (%.0f MB) — " + "waiting for the peer", + self._stream_segments, + self._stream_segments * STREAM_SEGMENT_SIZE / 1048576) if self._stream_stopped: return False # Checked before the wait as well as after it: a peer that vanishes @@ -1877,16 +1945,23 @@ class WebRTCPeerSession: await asyncio.wait_for(self._stream_credit_evt.wait(), timeout=STREAM_CREDIT_POLL) except asyncio.TimeoutError: - waited += STREAM_CREDIT_POLL - if waited >= STREAM_CREDIT_TIMEOUT: - log.info("Stream stalled: no credit from peer=%s", - (self._user_id or "?")[:8]) + silent = time.monotonic() - self._stream_heard_at + if silent >= STREAM_CREDIT_TIMEOUT: + log.info("Stream stalled: nothing from peer=%s for %.0fs", + (self._user_id or "?")[:8], silent) return False continue if self._stream_stopped: return False if self._channel is None or self._channel.readyState != "open": return False + if waiting_since: + waited_for = time.monotonic() - waiting_since + # Only a wait long enough to be a symptom. Normal pacing puts a + # gap of a few seconds between windows; a minute means the viewer + # is buffered right up and playing, or has stopped watching. + level = log.info if waited_for >= 10 else log.debug + level("stream: credit arrived after %.1fs", waited_for) self._stream_credit -= 1 return True @@ -1924,15 +1999,26 @@ class WebRTCPeerSession: self._stream_task = asyncio.current_task() await self._stream_video(msg) - async def _stream_video(self, msg: dict) -> None: - """Stream a video file as fMP4 segments via MSE-compatible output.""" - # One ffmpeg per request with no cap lets any member exhaust the node's - # CPU and process table (H6). The semaphore lives on the transport context - # so it is shared across all peers, not per-session. + def _transcode_semaphore(self) -> asyncio.Semaphore: + """The node's stream budget, shared across every peer. + + One ffmpeg per request with no cap lets any member exhaust the node's + CPU and process table (H6). The semaphore lives on the transport + context rather than the session so that it counts the node's viewers + and not one browser's, and it is created once: rebuilding it per call + would hand every caller its own budget and cap nothing at all. + """ sem = self._ctx.get("_transcode_sem") if sem is None: - sem = asyncio.Semaphore(MAX_CONCURRENT_TRANSCODES) + n = self._ctx.get("max_concurrent_streams") or MAX_CONCURRENT_TRANSCODES + sem = asyncio.Semaphore(n) self._ctx["_transcode_sem"] = sem + log.info("stream: %d concurrent viewers allowed", n) + return sem + + async def _stream_video(self, msg: dict) -> None: + """Stream a video file as fMP4 segments via MSE-compatible output.""" + sem = self._transcode_semaphore() if sem.locked() and sem._value <= 0: self._send({"type": "error", "detail": "Server busy, retry shortly"}) return @@ -2029,6 +2115,14 @@ class WebRTCPeerSession: }) index += 1 self._stream_segments = index + if index % 100 == 0: + # A stream that stops shows up here as a last line, and the + # numbers on it say which side stopped it. + log.info("stream: %d segments (%.0f MB), credit=%d, " + "keepalives=%d, %.0fs in", + index, index * STREAM_SEGMENT_SIZE / 1048576, + self._stream_credit, self._stream_keepalives, + time.monotonic() - self._stream_started_at) await asyncio.sleep(0) except Exception as e: log.error("Stream error: %s", e) @@ -2150,6 +2244,7 @@ class WebRTCTransport: groups: dict[str, dict] | None = None, denylist: Any | None = None, stun_servers: list[str] | None = None, + max_concurrent_streams: int | None = None, ): self._ctx: dict[str, Any] = { "sk_node": sk_node, @@ -2158,6 +2253,9 @@ class WebRTCTransport: "shared_root": shared_root, "index": index, "_peers": {}, + # None means "the operator said nothing" — the default applies. It + # is read once, when the first stream builds the semaphore. + "max_concurrent_streams": max_concurrent_streams, } if groups: self._ctx["groups"] = groups diff --git a/packages/meshbay-node/tests/test_stream_capacity_config.py b/packages/meshbay-node/tests/test_stream_capacity_config.py new file mode 100644 index 0000000..7c33a2f --- /dev/null +++ b/packages/meshbay-node/tests/test_stream_capacity_config.py @@ -0,0 +1,148 @@ +""" +How many people may watch at once, and who decides. + +Bounding the client's read-ahead to ninety seconds of film changed what a +transcode slot is. It used to be a burst — the browser took segments as fast as +it could append them, so a slot came back within the minute whatever the length +of the film. Now a slot is held for as long as someone is watching, so the cap +is a cap on simultaneous viewers, and the right number stopped being a property +of the code: it depends on the machine the node runs on. + +So it belongs to the operator. `max_concurrent_streams` under [node] in +node.toml, `MESHBAY_MAX_CONCURRENT_STREAMS` in the environment, and the +constant in the source as the default when neither says anything. + +The tests below follow the value along that whole path rather than checking +that the field parses, because every join in it has been wrong at least once: +the daemon reads `self._config`, not `self.cfg`, and nothing about the +attribute that does not exist fails until a video is played. +""" + +import asyncio +import textwrap +from pathlib import Path + +import pytest + +from meshbay_node.config import load_config +from meshbay_node.transport.webrtc_server import ( + MAX_CONCURRENT_TRANSCODES, + WebRTCPeerSession, + WebRTCTransport, +) + + +def _cfg(tmp_path: Path, body: str): + p = tmp_path / "node.toml" + p.write_text(textwrap.dedent(body)) + return load_config(p) + + +# ── What the operator writes ────────────────────────────────────────────────── + +def test_the_operator_sets_it(tmp_path): + cfg = _cfg(tmp_path, """ + [node] + max_concurrent_streams = 3 + """) + assert cfg.node.max_concurrent_streams == 3 + + +def test_saying_nothing_gets_the_default(tmp_path): + cfg = _cfg(tmp_path, """ + [node] + quic_port = 19010 + """) + assert cfg.node.max_concurrent_streams == MAX_CONCURRENT_TRANSCODES, ( + "the config default and the source default disagree, so the number " + "depends on whether a node.toml happens to mention it") + + +def test_the_environment_wins_over_the_file(tmp_path, monkeypatch): + monkeypatch.setenv("MESHBAY_MAX_CONCURRENT_STREAMS", "5") + cfg = _cfg(tmp_path, """ + [node] + max_concurrent_streams = 3 + """) + assert cfg.node.max_concurrent_streams == 5 + + +@pytest.mark.parametrize("value", ["0", "-4", '"lots"', "true"]) +def test_a_value_that_would_break_streaming_is_refused(tmp_path, value, caplog): + """Zero is the dangerous one. + + `asyncio.Semaphore(0)` is not "no limit". It is a node where every video + waits forever, with nothing in the log to say why — so the operator gets + the default and a warning naming the setting instead. + """ + cfg = _cfg(tmp_path, f""" + [node] + max_concurrent_streams = {value} + """) + assert cfg.node.max_concurrent_streams == MAX_CONCURRENT_TRANSCODES + assert "max_concurrent_streams" in caplog.text, ( + "the value was silently discarded — the operator has no way to learn " + "their setting is not in effect") + + +# ── That the number reaches the thing it limits ─────────────────────────────── + +class _FakePC: + def on(self, *a, **k): + return lambda f: f + + +def _semaphore_size(n): + t = WebRTCTransport( + sk_node=None, hub_pk_pem=b"", gek=b"\0" * 32, + shared_root=Path("/tmp"), index=None, max_concurrent_streams=n) + s = WebRTCPeerSession(_FakePC(), t._ctx, peer_id="p") + return s._transcode_semaphore()._value + + +@pytest.mark.parametrize("n,expect", [(None, MAX_CONCURRENT_TRANSCODES), (3, 3), (20, 20)]) +def test_the_configured_number_is_the_semaphore(n, expect): + assert asyncio.run(_run(n)) == expect + + +async def _run(n): + return _semaphore_size(n) + + +def test_the_budget_is_shared_between_peers(): + """One budget for the node, not one per browser. + + Building it per call would cap nothing: every viewer would arrive with a + full allowance and the node would spawn ffmpeg without limit. + """ + t = WebRTCTransport( + sk_node=None, hub_pk_pem=b"", gek=b"\0" * 32, + shared_root=Path("/tmp"), index=None, max_concurrent_streams=2) + + async def go(): + a = WebRTCPeerSession(_FakePC(), t._ctx, peer_id="a") + b = WebRTCPeerSession(_FakePC(), t._ctx, peer_id="b") + sem_a, sem_b = a._transcode_semaphore(), b._transcode_semaphore() + assert sem_a is sem_b, "each peer got its own budget, so there is no cap" + await sem_a.acquire() + assert sem_b._value == 1, "one peer's stream did not spend the node's budget" + + asyncio.run(go()) + + +def test_the_daemon_passes_it( ): + """The join that syntax checking cannot see. + + `self.cfg` parses and imports perfectly well; it raises AttributeError the + first time somebody plays a video, which is not where anyone would look. + """ + daemon = (Path(__file__).resolve().parents[1] / "src" / "meshbay_node" + / "daemon.py").read_text() + i = daemon.index("WebRTCTransport(") + call = daemon[i:daemon.index(")", daemon.index("denylist=denylist", i))] + assert "max_concurrent_streams=" in call, ( + "the daemon builds the transport without the operator's setting, so " + "node.toml is read and then ignored") + assert "self._config.node.max_concurrent_streams" in call, ( + "the daemon holds its config in _config; any other attribute is an " + "AttributeError deferred until someone plays a video") diff --git a/packages/meshbay-node/tests/test_task_lifetime.py b/packages/meshbay-node/tests/test_task_lifetime.py index 8897e71..e1ebfed 100644 --- a/packages/meshbay-node/tests/test_task_lifetime.py +++ b/packages/meshbay-node/tests/test_task_lifetime.py @@ -8,9 +8,10 @@ pending!" and nothing else happens. For `_stream_video` that was expensive. It holds a transcode slot for its whole life with `async with sem`, and a destroyed task never reaches `__aexit__`. The -node allows two, so two abandoned streams left it answering "Server busy" to -every request from then on: videos stopped playing entirely, first try included, -until the daemon was restarted. +node allows only a handful, so a few abandoned streams left it answering +"Server busy" to every request from then on: videos stopped playing entirely, +first try included, until the daemon was restarted. It was two slots at the +time, which is how few it took. Seen in the wild on 2026-08-16 after a viewer switched films mid-stream. """ @@ -86,12 +87,21 @@ def test_closing_a_session_releases_its_tasks(session): assert "gather" in fn, "cancelling without awaiting does not run the exits" -def test_two_slots_is_the_whole_margin(source): - """States the number the failure hinged on, so a change is deliberate.""" +def test_the_slot_count_matches_what_a_slot_now_costs(source): + """States the number, so a change is deliberate rather than drifted into. + + It was two, and two was right while a stream was a burst: the client took + segments as fast as it could append them and the slot came back within the + minute. Bounding the read-ahead to ninety seconds of film (the buffer + ceiling fix) changed what a slot is — it is now held for as long as someone + is watching, so the count is a count of simultaneous viewers. + """ n = int(re.search(r"MAX_CONCURRENT_TRANSCODES\s*=\s*(\d+)", source).group(1)) - assert n == 2, ( - f"the cap is now {n}; the leak above emptied it in {n} abandoned " - "streams, so if this moves the comments explaining it should too") + assert n >= 8, ( + f"the cap is {n}; with a slot held for the length of a film, that is " + f"{n} people watching before the node refuses everyone else") + assert "for as long as the film runs" in source, ( + "the number moved but the comment explaining what a slot costs did not") # ── One viewer, one stream ──────────────────────────────────────────────────── |