aboutsummaryrefslogtreecommitdiffstats
path: root/packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/streaming.py
diff options
context:
space:
mode:
Diffstat (limited to 'packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/streaming.py')
-rw-r--r--packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/streaming.py71
1 files changed, 56 insertions, 15 deletions
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