diff options
| author | Christophe Besson <cbesson@gmail.com> | 2026-09-08 14:04:47 +0200 |
|---|---|---|
| committer | Christophe Besson <cbesson@gmail.com> | 2026-09-08 14:04:47 +0200 |
| commit | 038066c43caa8ee76dd1e04761234271e4e67ecd (patch) | |
| tree | 502fb9b090197dc75135b968efcd715b386a8713 /packages/meshbay-node/src/meshbay_node/transport/webrtc_server.py | |
| parent | 28f1b5686c7ab200aeda6782f5f6e829c24759dd (diff) | |
| download | meshbay-038066c43caa8ee76dd1e04761234271e4e67ecd.tar.gz | |
fix(node): make max_concurrent_streams take effect without a restart
`ops.set_node_settings` hot-swapped the stream pool by assigning
`webrtc._stream_sem`. That attribute has never existed on WebRTCTransport — the
pool is `ctx["_transcode_sem"]` — so `hasattr(webrtc, '_stream_sem')` was always
False and the branch never ran. The setting was accepted, written to roster.db
and node.toml, and applied only on the next restart, which is exactly what
draft-v6 §2.11 says it does not need. An operator lowering the cap on a
struggling machine, or raising it after "Server busy", saw nothing happen and
had no way to find out why.
`WebRTCTransport.set_capacity()` is the one implementation, on the object that
owns the state, so the download and upload caps the transfer-slots plan adds
next do not each grow their own copy of the mistake.
Resizing has semantics worth stating: the new cap governs new streams and never
interrupts one that is running, because a slot is held for the length of a film
and lowering a number must not take somebody's film away. The replacement pool
is built with the permits that remain (`new - in_flight`, floored at zero) — a
full set would briefly allow more concurrent viewers than either the old cap or
the new one.
That needs a count of slots in use, so `_stream_video` now maintains one instead
of the code reading the semaphore's private `_value`: a number this code keeps
itself survives the semaphore object being replaced underneath it, and the same
counter makes the "N of M in use" log lines mean something.
test_stream_capacity.py drives the real transport and the real `_stream_video`;
`test_ops_calls_the_real_mechanism` fails if the dead attribute comes back.
Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01HCGdheDLxGReuKHga3BtST
Diffstat (limited to 'packages/meshbay-node/src/meshbay_node/transport/webrtc_server.py')
| -rw-r--r-- | packages/meshbay-node/src/meshbay_node/transport/webrtc_server.py | 63 |
1 files changed, 60 insertions, 3 deletions
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 dfabe9b..33f5474 100644 --- a/packages/meshbay-node/src/meshbay_node/transport/webrtc_server.py +++ b/packages/meshbay-node/src/meshbay_node/transport/webrtc_server.py @@ -5231,13 +5231,28 @@ class WebRTCPeerSession: if sem.locked() and sem._value <= 0: self._send({"type": "error", "detail": "Server busy, retry shortly"}) return - log.info("stream: waiting for a slot (free=%s)", sem._value) + 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: - log.info("stream: slot acquired (free=%s)", sem._value) + # 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: - log.info("stream: slot released (free=%s)", sem._value + 1) + 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 async def _stream_video_inner(self, msg: dict) -> None: ctx = self._group_ctx() @@ -5605,6 +5620,48 @@ class WebRTCTransport: self._stun = stun_servers or list(DEFAULT_STUN_SERVERS) self._sessions: dict[str, WebRTCPeerSession] = {} + def set_capacity(self, *, max_concurrent_streams: int | None = None) -> dict: + """Resize a live pool without restarting the daemon. + + `ops.set_node_settings` used to do this by assigning + `webrtc._stream_sem`, an attribute that has never existed — the pool is + `ctx["_transcode_sem"]`, and `hasattr(webrtc, "_stream_sem")` is always + False. So the hot-swap was a no-op and **`max_concurrent_streams` has + never taken effect from the Node page without a restart**, contrary to + draft-v6 §2.11. This is the one implementation, on the object that owns + the state, so the next two caps do not each grow their own copy of the + mistake. + + What resizing means, stated because it is a decision and not a + detail: **the new cap governs new streams; the ones already running are + never interrupted.** A slot is held for the length of a film, so + lowering the cap below what is in flight cannot take a viewer's film + away — it stops the next one starting. The replacement pool is therefore + created with the permits that remain (`new - in_flight`, floored at + zero), not with a full set, or lowering the cap would briefly allow more + viewers than either the old value or the new one. + """ + changed: dict = {} + if max_concurrent_streams is not None: + n = int(max_concurrent_streams) + if n < 1: + raise ValueError("max_concurrent_streams must be positive") + before = self._ctx.get("max_concurrent_streams") + self._ctx["max_concurrent_streams"] = n + if self._ctx.get("_transcode_sem") is not None: + in_flight = self._ctx.get("_streams_in_flight", 0) + self._ctx["_transcode_sem"] = asyncio.Semaphore( + max(0, n - in_flight)) + log.info("stream: capacity %s -> %d (%d in flight, %d free now)", + before, n, in_flight, max(0, n - in_flight)) + else: + # Nothing has streamed yet; the pool is built from this value on + # first use, so there is nothing to resize. + log.info("stream: capacity %s -> %d (no pool built yet)", + before, n) + changed["max_concurrent_streams"] = n + return changed + async def handle_offer( self, offer_sdp: str, peer_id: str, ) -> tuple[str, list[dict]]: |