summaryrefslogtreecommitdiffstats
path: root/packages/meshbay-node/src/meshbay_node
diff options
context:
space:
mode:
Diffstat (limited to 'packages/meshbay-node/src/meshbay_node')
-rw-r--r--packages/meshbay-node/src/meshbay_node/ops.py9
-rw-r--r--packages/meshbay-node/src/meshbay_node/transport/webrtc_server.py63
2 files changed, 67 insertions, 5 deletions
diff --git a/packages/meshbay-node/src/meshbay_node/ops.py b/packages/meshbay-node/src/meshbay_node/ops.py
index 1bad487..7557302 100644
--- a/packages/meshbay-node/src/meshbay_node/ops.py
+++ b/packages/meshbay-node/src/meshbay_node/ops.py
@@ -1361,8 +1361,13 @@ async def set_node_settings(state: dict, settings: dict) -> dict:
_update_node_toml(conf_path, updated)
if "max_concurrent_streams" in updated:
webrtc = state.get("webrtc")
- if webrtc and hasattr(webrtc, '_stream_sem'):
- webrtc._stream_sem = asyncio.Semaphore(updated["max_concurrent_streams"])
+ # `webrtc._stream_sem` was assigned here for months. That attribute
+ # has never existed -- the pool is `ctx["_transcode_sem"]` -- so the
+ # `hasattr` guard was always False and the setting only ever took
+ # effect on a restart, which draft-v6 §2.11 says it does not need.
+ if webrtc is not None:
+ webrtc.set_capacity(
+ max_concurrent_streams=updated["max_concurrent_streams"])
if "stun_servers" in updated:
webrtc = state.get("webrtc")
if webrtc and hasattr(webrtc, '_stun'):
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]]: