aboutsummaryrefslogtreecommitdiffstats
path: root/packages/meshbay-node/src/meshbay_node/transport/webrtc
diff options
context:
space:
mode:
Diffstat (limited to 'packages/meshbay-node/src/meshbay_node/transport/webrtc')
-rw-r--r--packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/music.py2
-rw-r--r--packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/streaming.py76
-rw-r--r--packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/subtitles.py19
-rw-r--r--packages/meshbay-node/src/meshbay_node/transport/webrtc/channel.py4
-rw-r--r--packages/meshbay-node/src/meshbay_node/transport/webrtc/chat.py17
-rw-r--r--packages/meshbay-node/src/meshbay_node/transport/webrtc/core.py16
-rw-r--r--packages/meshbay-node/src/meshbay_node/transport/webrtc/limits.py16
-rw-r--r--packages/meshbay-node/src/meshbay_node/transport/webrtc/media_tools.py7
-rw-r--r--packages/meshbay-node/src/meshbay_node/transport/webrtc/upload_handlers.py40
9 files changed, 162 insertions, 35 deletions
diff --git a/packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/music.py b/packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/music.py
index 857db33..a2d9f47 100644
--- a/packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/music.py
+++ b/packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/music.py
@@ -120,7 +120,7 @@ class MusicMixin:
blob = await _transcode_audio_to_aac(file_path)
except Exception as e:
log.warning("Audio transcode failed for %s: %s", entry.id[:12], e)
- self._send({"type": "error", "detail": f"Transcode failed: {e}"})
+ self._send({"type": "error", "detail": "This track could not be converted"})
return
transcode_hash = blake3.blake3(blob).hexdigest()
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..d6a5248 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
@@ -246,7 +287,10 @@ class StreamingMixin:
try:
probe = await _probe_video(str(file_path))
except Exception as e:
- self._send({"type": "error", "detail": f"Probe failed: {e}"})
+ # The cause to the operator's log; to the member, that it failed.
+ # ffmpeg's own words carry the operator's paths and versions.
+ log.warning("stream: probe failed for %s: %s", entry.id[:12], e)
+ self._send({"type": "error", "detail": "This video could not be read"})
return
codec_str = probe.codec
duration = probe.duration
diff --git a/packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/subtitles.py b/packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/subtitles.py
index 70781eb..525f2a9 100644
--- a/packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/subtitles.py
+++ b/packages/meshbay-node/src/meshbay_node/transport/webrtc/apps/subtitles.py
@@ -10,6 +10,7 @@ from meshbay_common.protocol import MNP
from meshbay_node.media_probe import probe_video as _probe_video
from meshbay_node.roots import off_disk
+from meshbay_node.transport.webrtc.apps.streaming import MAX_SUBTITLE_JOBS_PER_ACCOUNT
from meshbay_node.transport.webrtc.disk import _locate
from meshbay_node.transport.webrtc.media_tools import (
_extract_subtitle_to_webvtt,
@@ -135,10 +136,16 @@ class SubtitlesMixin:
return
budget = _subtitle_timeout_for(entry.size)
- async with sem:
- log.info("subtitle: extracting file=%s track=%d (slot taken, up to %.0fs)",
- file_id[:12], ordinal, budget)
- blob = await _extract_subtitle_to_webvtt(file_path, ordinal, budget)
+ with self._account_share("subtitles", MAX_SUBTITLE_JOBS_PER_ACCOUNT) as ok:
+ if not ok:
+ log.info("subtitle: refused, account at its extraction share")
+ self._send({"type": "error",
+ "detail": "Subtitles are already being prepared, retry shortly"})
+ return
+ async with sem:
+ log.info("subtitle: extracting file=%s track=%d (slot taken, up to %.0fs)",
+ file_id[:12], ordinal, budget)
+ blob = await _extract_subtitle_to_webvtt(file_path, ordinal, budget)
subtitle_hash = blake3.blake3(blob).hexdigest()
await media_cache.put_thumb(subtitle_hash, synthetic_id, blob)
@@ -158,8 +165,10 @@ class SubtitlesMixin:
except BaseException as e:
log.warning("subtitle: extract failed file=%s track=%d after %.1fs: %r",
file_id[:12], ordinal, time.monotonic() - t0, e)
+ # The cause is in the log line above; ffmpeg's own words carry the
+ # operator's paths and versions.
self._send({"type": "error",
- "detail": f"Subtitle extraction failed: {e}"})
+ "detail": "These subtitles could not be extracted"})
if isinstance(e, asyncio.CancelledError):
raise
finally:
diff --git a/packages/meshbay-node/src/meshbay_node/transport/webrtc/channel.py b/packages/meshbay-node/src/meshbay_node/transport/webrtc/channel.py
index 107a43e..31cc83c 100644
--- a/packages/meshbay-node/src/meshbay_node/transport/webrtc/channel.py
+++ b/packages/meshbay-node/src/meshbay_node/transport/webrtc/channel.py
@@ -7,7 +7,7 @@ import struct
import msgpack
from aiortc import RTCPeerConnection
-from meshbay_node.transport.webrtc.limits import MAX_MSG
+from meshbay_node.transport.webrtc.limits import MAX_MSG, UNPACK_LIMITS
def _extract_dtls_fingerprint(sdp: str) -> bytes:
@@ -77,7 +77,7 @@ class _DataChannelBuffer:
break
msg_bytes = bytes(self._buf[4:4 + length])
del self._buf[:4 + length]
- yield msgpack.unpackb(msg_bytes, raw=False)
+ yield msgpack.unpackb(msg_bytes, raw=False, **UNPACK_LIMITS)
def _get_remote_ip(pc: RTCPeerConnection) -> str:
diff --git a/packages/meshbay-node/src/meshbay_node/transport/webrtc/chat.py b/packages/meshbay-node/src/meshbay_node/transport/webrtc/chat.py
index 26ec27c..493d7f0 100644
--- a/packages/meshbay-node/src/meshbay_node/transport/webrtc/chat.py
+++ b/packages/meshbay-node/src/meshbay_node/transport/webrtc/chat.py
@@ -283,7 +283,18 @@ class ChatMixin:
# anyone on the node.
gctx = self._group_ctx()
chat_store = gctx.get("chat_store")
+ # The two fields that travel in clear beside the ciphertext (the sealed
+ # envelope carries its own). Stored and relayed to every member, so they
+ # are what they claim to be and no larger: a name as long as a username,
+ # a thread id as long as a message id. Anything else is dropped.
sender_name = msg.get("sender_name", "")
+ if not isinstance(sender_name, str) or len(sender_name) > 64:
+ sender_name = ""
+ thread_id = msg.get("thread_id")
+ id_like = (isinstance(thread_id, int) and not isinstance(thread_id, bool)
+ or isinstance(thread_id, str) and len(thread_id) <= 64)
+ if thread_id is not None and not id_like:
+ thread_id = None
# Two shapes, and keeping them apart is what makes this deployable.
#
@@ -329,7 +340,7 @@ class ChatMixin:
self._spawn(self._store_chat_message(
chat_store,
iteration=msg.get("iteration", 0), payload=raw,
- thread_id=msg.get("thread_id"), sender_name=sender_name,
+ thread_id=thread_id, sender_name=sender_name,
format=fmt, epoch=epoch, device=device, nonce=nonce, sig=sig,
))
@@ -340,7 +351,7 @@ class ChatMixin:
"sender_id": self._user_id,
"sender_name": sender_name,
"payload": payload,
- "thread_id": msg.get("thread_id"),
+ "thread_id": thread_id,
"timestamp": time.time(),
"format": fmt,
"epoch": epoch,
@@ -598,7 +609,7 @@ class ChatMixin:
the client asks, the node produces on demand, the asking device
caches — nothing durable here).
- `linkpreview.safe_url` is the SSRF gate: the URL a *member* chose
+ `linkpreview.check_url` is the SSRF gate: the URL a *member* chose
decides an outbound request from the operator's machine, so http(s)
only and the resolved address must be globally routable. Failure of
any kind — blocked, unreachable, not HTML, nothing worth showing —
diff --git a/packages/meshbay-node/src/meshbay_node/transport/webrtc/core.py b/packages/meshbay-node/src/meshbay_node/transport/webrtc/core.py
index 882ddbc..ae3f0e1 100644
--- a/packages/meshbay-node/src/meshbay_node/transport/webrtc/core.py
+++ b/packages/meshbay-node/src/meshbay_node/transport/webrtc/core.py
@@ -139,7 +139,21 @@ class SessionCore:
log.info("WebRTC data received: %d bytes, msg #%d (peer=%s)",
len(message), self._msg_count, self._peer_id)
self._buffer.feed(message)
- for msg in self._buffer.messages():
+ decoded = self._buffer.messages()
+ while True:
+ try:
+ msg = next(decoded)
+ except StopIteration:
+ break
+ except ValueError as e:
+ # Over the size limit, or a container past its decode
+ # limit. The buffer still starts with that frame, so every
+ # later message would fail the same way: the session ends
+ # here. Only decoding is caught — a handler's own error is
+ # not a reason to drop the peer.
+ log.warning("Closing peer %s: %s", self._peer_id, e)
+ self._spawn(self.close())
+ break
self._handle_message(msg)
if _WEBRTC_TRACE:
diff --git a/packages/meshbay-node/src/meshbay_node/transport/webrtc/limits.py b/packages/meshbay-node/src/meshbay_node/transport/webrtc/limits.py
index 7d458f4..86412e0 100644
--- a/packages/meshbay-node/src/meshbay_node/transport/webrtc/limits.py
+++ b/packages/meshbay-node/src/meshbay_node/transport/webrtc/limits.py
@@ -2,7 +2,21 @@
CHUNK_SIZE = 1024 * 1024
-MAX_MSG = 64 * 1024 * 1024
+# The largest message a peer may send once it has proved the group key. The
+# largest a client really sends is a sealed playlist blob, 1 MiB (blobs.py);
+# chat is 64 KiB and an upload chunk 48 KiB. Eight times the largest, because a
+# message of many small objects decodes to several times its size in memory.
+MAX_MSG = 8 * 1024 * 1024
+
+# Per container, when a message is decoded: nothing a client sends comes near
+# them, and without them one message of tiny elements is one enormous list.
+UNPACK_LIMITS = {
+ "max_array_len": 100_000,
+ "max_map_len": 10_000,
+ "max_str_len": 1024 * 1024,
+ "max_bin_len": MAX_MSG,
+ "max_ext_len": 0,
+}
# What the `tr` on a chunk request turned out to be (see `_lease_of`).
diff --git a/packages/meshbay-node/src/meshbay_node/transport/webrtc/media_tools.py b/packages/meshbay-node/src/meshbay_node/transport/webrtc/media_tools.py
index 7ecb0a5..22c1690 100644
--- a/packages/meshbay-node/src/meshbay_node/transport/webrtc/media_tools.py
+++ b/packages/meshbay-node/src/meshbay_node/transport/webrtc/media_tools.py
@@ -165,6 +165,7 @@ async def _seek_lands_at(file_path: Path, t: float, map_args: list[str]) -> floa
fd, tmp_name = tempfile.mkstemp(suffix=".mp4")
os.close(fd)
tmp_path = Path(tmp_name)
+ proc = probe = None
try:
proc = await asyncio.create_subprocess_exec(
platform.ffmpeg_cmd(), "-hide_banner", "-loglevel", "error", "-y",
@@ -189,6 +190,12 @@ async def _seek_lands_at(file_path: Path, t: float, map_args: list[str]) -> floa
log.warning("stream: seek probe failed at %.1fs: %r", t, e)
return None
finally:
+ # A timed-out wait leaves its process running; it is stopped here, not
+ # left to finish a seek nobody is waiting for.
+ for p in (proc, probe):
+ if p is not None and p.returncode is None:
+ p.kill()
+ await p.wait()
await _discard_scratch(tmp_path)
text = stdout.decode(errors="replace").strip().rstrip(",")
try:
diff --git a/packages/meshbay-node/src/meshbay_node/transport/webrtc/upload_handlers.py b/packages/meshbay-node/src/meshbay_node/transport/webrtc/upload_handlers.py
index 1c6d1ce..02a50af 100644
--- a/packages/meshbay-node/src/meshbay_node/transport/webrtc/upload_handlers.py
+++ b/packages/meshbay-node/src/meshbay_node/transport/webrtc/upload_handlers.py
@@ -8,7 +8,14 @@ from pathlib import Path
from meshbay_common.protocol import UPLOAD_PROBE_INDEX, file_upload_ack_wire, file_upload_payload
from meshbay_node import uploads as uploads_mod
-from meshbay_node.roots import SAFE_UPLOAD_NAME, RootSet, _free_name, off_disk
+from meshbay_node.roots import (
+ SAFE_UPLOAD_NAME,
+ RootSet,
+ _free_name,
+ off_disk,
+ publish_upload,
+ shell_active,
+)
from meshbay_node.transport.webrtc.disk import _append_chunk
from meshbay_node.transport.webrtc.limits import LEASE_NONE, LEASE_QUEUED
@@ -212,6 +219,9 @@ class UploadMixin:
if not SAFE_UPLOAD_NAME.match(filename):
_refuse("Invalid filename", "invalid_filename")
return
+ if shell_active(filename):
+ _refuse("This type of file is not accepted", "file_type_refused")
+ return
roots: RootSet | None = ctx.get("roots")
if not roots:
@@ -306,9 +316,13 @@ class UploadMixin:
# A shared directory means two people can send the same name. Refusing the
# second is safe but silly — everyone's camera produces IMG_1234.jpg — so
# a free name is found instead. Never a replacement.
+ # Names other uploads into this directory will publish under are taken
+ # too: none of them is on disk yet.
+ reserved = uploads.reserved_names(rel_dir)
stored_name = (state.stored_name if state
- else await off_disk(roots, _free_name, target_dir, filename))
- tmp_path = target_dir / f"{stored_name}{uploads_mod.PART_SUFFIX}"
+ else await off_disk(roots, _free_name, target_dir, filename, reserved))
+ tmp_path = (state.part_path if state and state.part_path
+ else target_dir / uploads_mod.part_name(stored_name))
final_path = target_dir / stored_name
if chunk_index == UPLOAD_PROBE_INDEX:
@@ -361,6 +375,22 @@ class UploadMixin:
await off_disk(roots, _append_chunk, tmp_path, chunk_bytes, chunk_index == 0)
uploads.advance(user_id, rel_dir, filename, chunk_index, len(chunk_bytes))
+ last = chunk_index + 1 >= total_chunks
+ if last:
+ # Published before the last ack, so the ack names the file as it is
+ # on disk: publication never replaces a file, and may have had to
+ # take another free name for this one.
+ uploads.drop(user_id, rel_dir, filename)
+ try:
+ stored_name = await off_disk(roots, publish_upload, tmp_path, target_dir,
+ stored_name, filename,
+ uploads.reserved_names(rel_dir))
+ except OSError as e:
+ log.warning("Upload %s could not be published: %s", stored_name, e)
+ _refuse("The file could not be stored", "store_failed")
+ return
+ final_path = target_dir / stored_name
+
self._send(file_upload_ack_wire(
gek, self._group_id or "",
upload_id=upload_id,
@@ -372,9 +402,7 @@ class UploadMixin:
dir=rel_dir,
))
- if chunk_index + 1 >= total_chunks:
- uploads.drop(user_id, rel_dir, filename)
- await off_disk(roots, tmp_path.rename, final_path)
+ if last:
log.info("Upload complete: %s (%d chunks, %d bytes)",
stored_name, total_chunks, state.bytes)
self._audit("file_upload", f"{rel_dir}/{stored_name}")