summaryrefslogtreecommitdiffstats
path: root/packages/meshbay-node/tests/test_stream_capacity.py
diff options
context:
space:
mode:
authorChristophe Besson <cbesson@gmail.com>2026-09-09 14:28:40 +0200
committerChristophe Besson <cbesson@gmail.com>2026-09-09 14:28:40 +0200
commit7e2d078fe1870d256ae47781bee6ac4f454edf24 (patch)
tree05689049fd48f01ebbdb9995d5189cba15cee052 /packages/meshbay-node/tests/test_stream_capacity.py
parent813d18424ec57963bb56e6f40824a2db0ccce50d (diff)
parente6f895c473a0b19e7b186889c1836d3945bc880b (diff)
downloadmeshbay-7e2d078fe1870d256ae47781bee6ac4f454edf24.tar.gz
Merge branch 'fix/large-download-paths'
Concurrent-transfer limits, with the queue, the pause and the flag day. A node now caps how many transfers it runs at once (8 downloads, 8 uploads, node-wide) and how many one member may run in one group (2 by default, operator-signed). Beyond that the node answers "queued" and the client waits its turn, visibly, in the transfers panel — and a slot that frees starts whatever is next, skipping past a member who is at their own cap rather than letting them stall everyone behind them. Browsing is never subject to a slot: not the poster grid, not the covers, not opening a photo to look at it. That is structural — a transfer is what the transfers widget shows — and the exemption is bounded rather than open, at two files in flight per session, because an exemption with no bound is a leaseless branch under another name. Transfers can be cancelled, and now paused and resumed. A paused one holds nothing: its slot goes back at once and resuming rejoins the queue at the tail. Uploads survive the connection that started them and resume where the node stopped, asked for inside the seal rather than on a clear message. What they leave behind when they are abandoned is reaped, which closes a disk leak that predates this work. MNP 3.0 makes the lease compulsory and refuses 2.x at the handshake, with the desktop client checking `client.minimum` before connecting so an un-updated one says "update" instead of failing every connection in a protocol vocabulary. Fourteen defects were found on the way, eight of them by a person clicking Download and pasting a console — none of which 2075 tests could reach. Section 12 of ~/next/improve-downloads.md is that report, including the three this work introduced itself and the one that turned out to be caused by an instruction to hard-reload after each deployment. Node suite 1209 passed, hub suite 866 passed. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01HCGdheDLxGReuKHga3BtST
Diffstat (limited to 'packages/meshbay-node/tests/test_stream_capacity.py')
-rw-r--r--packages/meshbay-node/tests/test_stream_capacity.py155
1 files changed, 155 insertions, 0 deletions
diff --git a/packages/meshbay-node/tests/test_stream_capacity.py b/packages/meshbay-node/tests/test_stream_capacity.py
new file mode 100644
index 0000000..a35ece8
--- /dev/null
+++ b/packages/meshbay-node/tests/test_stream_capacity.py
@@ -0,0 +1,155 @@
+"""
+`max_concurrent_streams` must take effect without a restart.
+
+`ops.set_node_settings` did this by assigning `webrtc._stream_sem` — an
+attribute that has never existed. The pool is `ctx["_transcode_sem"]`, so
+`hasattr(webrtc, "_stream_sem")` was always False, the branch never ran, and the
+setting only ever applied on a restart. Draft-v6 §2.11 says it applies live, the
+Node page offers it as a live setting, and it did nothing: an operator lowering
+the cap on a struggling machine, or raising it after "Server busy", saw no
+change and had no way to know why.
+
+Nothing here mocks the pool. `set_capacity` is called on a real
+`WebRTCTransport` and the assertions read what a stream request would actually
+find.
+"""
+
+import asyncio
+
+import pytest
+
+from meshbay_node.transport.webrtc_server import (
+ MAX_CONCURRENT_TRANSCODES, WebRTCPeerSession, WebRTCTransport,
+)
+
+
+def _pool(transport) -> asyncio.Semaphore:
+ """The pool a stream request would acquire, built the way one builds it."""
+ session = WebRTCPeerSession.__new__(WebRTCPeerSession)
+ session._ctx = transport._ctx
+ return session._transcode_semaphore()
+
+
+@pytest.fixture
+def transport(tmp_path):
+ """A real WebRTCTransport. Its keys and index are genuine but incidental —
+ nothing below the capacity code reads them."""
+ from cryptography.hazmat.primitives.asymmetric.ed25519 import Ed25519PrivateKey
+
+ from conftest import one_root
+ from meshbay_common.crypto import generate_gek
+ from meshbay_node.indexer.group_index import GroupIndex
+
+ sk_node = Ed25519PrivateKey.generate()
+ gek = generate_gek()
+ shared = tmp_path / "shared"
+ shared.mkdir()
+ return WebRTCTransport(
+ sk_node=sk_node, hub_pk_pem=b"", gek=gek,
+ roots=one_root(shared),
+ index=GroupIndex(group_id="g", sk_node=sk_node, gek=gek),
+ stun_servers=[])
+
+
+def test_raising_the_cap_is_visible_to_the_next_stream(transport):
+ """The bug, at its simplest: the number changes and nothing happens."""
+ pool = _pool(transport)
+ assert pool._value == MAX_CONCURRENT_TRANSCODES
+ transport.set_capacity(max_concurrent_streams=16)
+ assert _pool(transport)._value == 16, (
+ "the setting was accepted and the pool never changed — this is the "
+ "no-op that shipped")
+
+
+def test_lowering_the_cap_does_not_interrupt_what_is_running(transport):
+ """
+ A slot is held for the length of a film, so lowering the cap cannot take a
+ viewer's film away. It stops the next one starting, and the replacement pool
+ carries only the permits that remain.
+ """
+ _pool(transport)
+ transport._ctx["_streams_in_flight"] = 3
+ transport.set_capacity(max_concurrent_streams=4)
+ assert _pool(transport)._value == 1, (
+ "a full set of permits would let more viewers in than either the old "
+ "cap or the new one, on top of the three still watching")
+
+
+def test_lowering_below_what_is_running_refuses_the_next_one(transport):
+ _pool(transport)
+ transport._ctx["_streams_in_flight"] = 6
+ transport.set_capacity(max_concurrent_streams=2)
+ assert _pool(transport)._value == 0, "the pool must not go negative"
+
+
+def test_the_value_is_kept_for_a_pool_not_yet_built(transport):
+ """Nothing has streamed, so there is nothing to resize — but the number has
+ to be there when the first request builds the pool."""
+ transport.set_capacity(max_concurrent_streams=3)
+ assert transport._ctx.get("_transcode_sem") is None
+ assert _pool(transport)._value == 3
+
+
+def test_a_cap_below_one_is_refused(transport):
+ for bad in (0, -1):
+ with pytest.raises(ValueError):
+ transport.set_capacity(max_concurrent_streams=bad)
+
+
+def test_nothing_changes_when_nothing_is_passed(transport):
+ _pool(transport)
+ before = transport._ctx["_transcode_sem"]
+ assert transport.set_capacity() == {}
+ assert transport._ctx["_transcode_sem"] is before
+
+
+@pytest.mark.asyncio
+async def test_in_flight_is_counted_by_the_streaming_path_itself(transport):
+ """
+ `set_capacity` resizes against `_streams_in_flight`, so that counter has to
+ be maintained where slots are actually taken — not set by a test. Drives the
+ real `_stream_video`, with the work under it stubbed: what is being checked
+ is the accounting around the slot, which is where flow control in this repo
+ has gone wrong before.
+ """
+ session = WebRTCPeerSession.__new__(WebRTCPeerSession)
+ session._ctx = transport._ctx
+ session._send = lambda msg: None
+
+ seen = []
+ release = asyncio.Event()
+
+ async def _inner(_msg):
+ seen.append(transport._ctx.get("_streams_in_flight"))
+ await release.wait()
+
+ session._stream_video_inner = _inner
+ task = asyncio.create_task(session._stream_video({"file_id": "x"}))
+ await asyncio.sleep(0)
+ await asyncio.sleep(0)
+ assert seen == [1], "the slot was taken without being counted"
+
+ release.set()
+ await task
+ assert transport._ctx["_streams_in_flight"] == 0, (
+ "a slot that is not given back is a viewer nobody can replace — the "
+ "class of bug _replace_stream and shutdown_tasks exist for")
+
+
+def test_ops_calls_the_real_mechanism():
+ """
+ The dead branch, pinned. `hasattr(webrtc, '_stream_sem')` is False for every
+ WebRTCTransport that has ever existed, so a test that only checked
+ "set_node_settings does not raise" passed throughout.
+ """
+ import inspect
+
+ from meshbay_node import ops
+
+ src = inspect.getsource(ops.set_node_settings)
+ # Comments stripped: this function now *explains* the dead attribute, and a
+ # test that matched the prose would fail on its own documentation.
+ code = "\n".join(line.split("#", 1)[0] for line in src.splitlines())
+ assert "_stream_sem" not in code, "the attribute that never existed is back"
+ assert "set_capacity" in code, "the setting must reach the pool that exists"
+ assert not hasattr(WebRTCTransport, "_stream_sem")