aboutsummaryrefslogtreecommitdiffstats
path: root/packages/meshbay-node/src/meshbay_node/transport
diff options
context:
space:
mode:
Diffstat (limited to 'packages/meshbay-node/src/meshbay_node/transport')
-rw-r--r--packages/meshbay-node/src/meshbay_node/transport/quic_client.py19
-rw-r--r--packages/meshbay-node/src/meshbay_node/transport/quic_server.py76
-rw-r--r--packages/meshbay-node/src/meshbay_node/transport/webrtc_server.py309
3 files changed, 171 insertions, 233 deletions
diff --git a/packages/meshbay-node/src/meshbay_node/transport/quic_client.py b/packages/meshbay-node/src/meshbay_node/transport/quic_client.py
index b22b8df..af87b70 100644
--- a/packages/meshbay-node/src/meshbay_node/transport/quic_client.py
+++ b/packages/meshbay-node/src/meshbay_node/transport/quic_client.py
@@ -317,22 +317,3 @@ class QuicChunkClient:
# substitutes it would otherwise choose which key we decrypt with.
return file_chunk_plaintext(
self._gek, msg, file_hash=bytes.fromhex(file_id))
-
- async def fetch_stream_segment(
- self, file_id: str, segment_index: int, segment_duration: int = 4,
- ) -> bytes:
- """Fetch one HLS segment (MPEG-TS bytes) over QUIC."""
- sid = self._new_stream()
- self._proto._send(sid, {
- "type": MNP.STREAM_SEGMENT,
- "v": MNP_VERSION,
- "file_id": file_id,
- "segment_index": segment_index,
- "segment_duration": segment_duration,
- })
- msg = await self._proto._recv(sid, timeout=30.0)
-
- if msg.get("type") == "error":
- raise LookupError(msg.get("detail", "Unknown error"))
-
- return base64.b64decode(msg["data_b64"])
diff --git a/packages/meshbay-node/src/meshbay_node/transport/quic_server.py b/packages/meshbay-node/src/meshbay_node/transport/quic_server.py
index 284b488..e34153c 100644
--- a/packages/meshbay-node/src/meshbay_node/transport/quic_server.py
+++ b/packages/meshbay-node/src/meshbay_node/transport/quic_server.py
@@ -22,7 +22,6 @@ import base64
import logging
import os
import struct
-import subprocess
import uuid
from pathlib import Path
from typing import Any, Callable
@@ -54,7 +53,6 @@ from meshbay_common.groupbox import PURPOSE_ACK, seal
from meshbay_common.protocol import MNP, file_chunk_wire
from meshbay_node.indexer import GroupIndex
from meshbay_node.transport.wire import index_sync_message
-from meshbay_node import platform
log = logging.getLogger(__name__)
@@ -62,16 +60,6 @@ CHUNK_SIZE = 1024 * 1024
MAX_MSG = 64 * 1024 * 1024
ALPN = ["meshbay-mnp"]
-# ffmpeg is spawned per STREAM_SEGMENT request, and `_extract_segment` runs
-# `subprocess.run` synchronously — so without a bound, an authenticated peer can
-# both fork-bomb the node and block its event loop for up to 30 s per request
-# (finding M2c). Extraction now runs in a thread and passes through this
-# semaphore. Small on purpose: the QUIC path has no shipping client yet, this is
-# parity work with the WebRTC transcode cap.
-_MAX_CONCURRENT_SEGMENTS = 4
-_segment_sem = asyncio.Semaphore(_MAX_CONCURRENT_SEGMENTS)
-
-
class Denylist:
"""
Denylist for revoked users, groups and invalidated JWTs.
@@ -247,8 +235,6 @@ class _MNPServerProtocol(QuicConnectionProtocol):
self._do_index_sync_sync(stream_id)
elif mtype == MNP.FILE_REQUEST:
self._do_file_request_sync(stream_id, msg)
- elif mtype == MNP.STREAM_SEGMENT:
- self._spawn(self._do_stream_segment(stream_id, msg))
elif mtype == MNP.CHAT_MESSAGE:
self._do_chat_message_sync(stream_id, msg)
elif mtype == MNP.PING:
@@ -435,49 +421,6 @@ class _MNPServerProtocol(QuicConnectionProtocol):
ctx["gek"], file_path, chunk_index, file_hash, entry.id)
self._send(stream_id, chunk_data)
- async def _do_stream_segment(self, stream_id: int, msg: dict) -> None:
- """
- Extract and serve one segment via ffmpeg — off the event loop and behind
- a concurrency bound, so one request can neither stall the whole node nor
- fork-bomb it (finding M2c). The WebRTC path has had both since Phase 11.5.
- """
- try:
- ctx = self._group_ctx()
- file_id = msg["file_id"]
- segment_index = msg["segment_index"]
- segment_duration = msg.get("segment_duration", 4)
-
- entry = ctx["index"].get_entry(file_id)
- if not entry:
- self._send(stream_id, {"type": "error", "detail": "File not found"})
- return
-
- file_path = entry_abs_path(ctx["roots"], entry)
- if not file_path.exists():
- self._send(stream_id, {"type": "error", "detail": "File not on disk"})
- return
-
- start_time = segment_index * segment_duration
- loop = asyncio.get_event_loop()
- async with _segment_sem:
- segment_data = await loop.run_in_executor(
- None, _extract_segment, file_path, start_time, segment_duration)
- if segment_data is None:
- self._send(stream_id, {"type": "error", "detail": "Segment extraction failed"})
- return
-
- self._send(stream_id, {
- "type": MNP.STREAM_SEGMENT,
- "v": MNP_VERSION,
- "file_id": file_id,
- "segment_index": segment_index,
- "data_b64": base64.b64encode(segment_data).decode(),
- "size": len(segment_data),
- })
- except Exception as e:
- log.error("stream_segment: %s", e)
- self._send(stream_id, {"type": "error", "detail": "Segment extraction failed"})
-
def _do_chat_message_sync(self, stream_id: int, msg: dict) -> None:
"""
Store a chat message and broadcast it to the rest of THIS group.
@@ -548,25 +491,6 @@ def _read_and_encrypt(
return file_chunk_wire(gek, plaintext, chunk_index, file_hash, file_id)
-def _extract_segment(file_path: Path, start_time: float, duration: float) -> bytes | None:
- """Extract one HLS segment via ffmpeg. Returns MPEG-TS bytes or None on failure."""
- try:
- result = subprocess.run(
- [platform.ffmpeg_cmd(), "-hide_banner", "-loglevel", "error",
- "-ss", str(start_time),
- "-i", str(file_path),
- "-t", str(duration),
- "-c:v", "copy", "-c:a", "copy",
- "-f", "mpegts", "pipe:1"],
- capture_output=True, timeout=30,
- )
- if result.returncode == 0 and result.stdout:
- return result.stdout
- return None
- except Exception:
- return None
-
-
# ── QuicChunkServer ────────────────────────────────────────────────────────────
class QuicChunkServer:
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 9774831..4e4a23f 100644
--- a/packages/meshbay-node/src/meshbay_node/transport/webrtc_server.py
+++ b/packages/meshbay-node/src/meshbay_node/transport/webrtc_server.py
@@ -24,6 +24,7 @@ Signaling flow (handled externally by the hub):
import asyncio
import base64
+import contextvars
import hashlib
import hmac
import logging
@@ -105,9 +106,18 @@ from meshbay_common.join import (
ROLE_OPERATOR,
join_transcript,
)
-from meshbay_common.protocol import MNP, chunk_ciphertext, file_chunk_wire
-from meshbay_common.chatbox import NONCE_LEN as CHAT_NONCE_LEN, SIG_LEN as CHAT_SIG_LEN
-from meshbay_node.chat import FORMAT_PLAIN, FORMAT_SEALED_V1, ReplayedMessage
+from meshbay_common.chatbox import (
+ NONCE_LEN as CHAT_NONCE_LEN,
+ SIG_LEN as CHAT_SIG_LEN,
+)
+from meshbay_common.protocol import (
+ MNP,
+ chunk_ciphertext,
+ file_chunk_wire,
+ file_upload_ack_wire,
+ file_upload_payload,
+)
+from meshbay_node.chat import FORMAT_SEALED_V1, ReplayedMessage
from meshbay_node.transport.wire import index_sync_message
from meshbay_node.indexer import GroupIndex
from meshbay_node.indexer.indexer import DirectoryIndexer
@@ -263,6 +273,32 @@ def _pack(obj: dict) -> bytes:
_WEBRTC_TRACE = os.environ.get("MESHBAY_WEBRTC_TRACE") == "1"
_WEBRTC_TRACE_INTERVAL_S = 30.0
+# The request this session is currently answering, as (session, req_id).
+#
+# MNP has never carried a correlation id: a reply named its own type and
+# nothing else, so a client with more than one request outstanding had to guess
+# which one a message answered — by arrival order, for every reply the client
+# could not key off a field of its own. The guess is wrong whenever two replies
+# reorder, and catastrophically wrong for the replies that name *nothing*: this
+# module sends `{"type": "error"}` from 240 places and two of them name what
+# they are about. A refusal therefore reached no caller at all, and the request
+# it belonged to waited out the client's 30s timeout while some unrelated
+# request was resolved with the refusal instead. Live symptom, found 2026-09-06:
+# the Chat composer is disabled while a send is in flight, so a chat message
+# whose reply went astray froze the tab for 30 seconds.
+#
+# `req_id` closes it: whatever the caller put on the request is stamped on the
+# reply. A ContextVar rather than a parameter because the alternative is
+# threading an argument through all 240 send sites — and asyncio copies the
+# current context into a task, so a handler that `_spawn`s its real work still
+# answers under the id of the request that started it.
+#
+# The session is held alongside the id because a handler may send to *other*
+# sessions as well as its own (a chat broadcast, an index push): those are not
+# replies to anything and must not be stamped. _send checks the owner.
+_REPLY_TO: contextvars.ContextVar[tuple] = contextvars.ContextVar(
+ "meshbay_reply_to", default=(None, None))
+
class _DataChannelBuffer:
"""
@@ -416,6 +452,22 @@ class WebRTCPeerSession:
)
def _handle_message(self, msg: dict) -> None:
+ """Answer one MNP message, under the correlation id it carries.
+
+ The id is published for the whole handler — see _REPLY_TO — so that
+ every reply _send puts on the wire, including the ones a spawned task
+ sends much later and the generic refusal below, names the request it
+ answers. Resetting on the way out only clears it for *this* call: a
+ task spawned in between captured its own copy of the context when it
+ was created and keeps answering under the right id.
+ """
+ token = _REPLY_TO.set((self, msg.get("req_id")))
+ try:
+ self._dispatch_message(msg)
+ finally:
+ _REPLY_TO.reset(token)
+
+ def _dispatch_message(self, msg: dict) -> None:
mtype = msg.get("type")
log.debug("WebRTC recv: %s", mtype)
try:
@@ -459,8 +511,6 @@ class WebRTCPeerSession:
# Chunks are matched by file and index on the client, so
# answering out of order is safe.
self._spawn(self._do_file_request(msg))
- elif mtype == MNP.STREAM_SEGMENT:
- self._do_stream_segment(msg)
elif mtype == MNP.CHAT_MESSAGE:
self._do_chat_message(msg)
elif mtype == MNP.CHAT_HISTORY:
@@ -3172,7 +3222,14 @@ class WebRTCPeerSession:
def _group_ctx(self) -> dict:
if "groups" in self._ctx and self._group_id:
- return self._ctx["groups"][self._group_id]
+ # `.get`, not a bare subscript. A config reload removes a group
+ # from this map (daemon.py's reload does `groups_ctx.pop`) while
+ # sessions connected to it are still open, and the next request
+ # any of them made raised KeyError into _dispatch_message's
+ # catch-all. An absent group now reads the way an unconfigured
+ # one already does — the handlers all test for what they need —
+ # instead of failing every request the session has left.
+ return self._ctx["groups"].get(self._group_id) or {}
return self._ctx
def _indexing_status(self) -> dict:
@@ -4005,71 +4062,6 @@ class WebRTCPeerSession:
"director": director,
}
- def _do_stream_segment(self, msg: dict) -> None:
- self._spawn(self._do_stream_segment_async(msg))
-
- async def _do_stream_segment_async(self, msg: dict) -> None:
- """
- Legacy HLS segment extraction (superseded by stream_req/MSE).
-
- Finding H6: this ran subprocess.run(..., timeout=30) directly inside the
- event loop, so a single request stalled the whole daemon — every peer,
- every group — for up to thirty seconds. Now async and under the same
- transcode semaphore as _stream_video.
- """
- ctx = self._group_ctx()
- file_id = msg["file_id"]
- segment_index = msg["segment_index"]
- segment_duration = msg.get("segment_duration", 4)
-
- entry = ctx["index"].get_entry(file_id)
- if not entry:
- self._send({"type": "error", "detail": "File not found"})
- return
-
- file_path = entry_abs_path(ctx["roots"], entry)
- if not file_path.exists():
- self._send({"type": "error", "detail": "File not on disk"})
- return
-
- sem = self._transcode_semaphore()
-
- try:
- async with sem:
- proc = await asyncio.create_subprocess_exec(
- platform.ffmpeg_cmd(), "-hide_banner", "-loglevel", "error",
- "-ss", str(segment_index * segment_duration),
- "-i", str(file_path),
- "-t", str(segment_duration),
- "-c:v", "copy", "-c:a", "copy",
- "-f", "mpegts", "pipe:1",
- stdout=asyncio.subprocess.PIPE,
- stderr=asyncio.subprocess.DEVNULL,
- )
- try:
- stdout, _ = await asyncio.wait_for(proc.communicate(), timeout=30)
- except asyncio.TimeoutError:
- proc.kill()
- await proc.wait()
- self._send({"type": "error", "detail": "Segment extraction timed out"})
- return
- if proc.returncode != 0 or not stdout:
- self._send({"type": "error", "detail": "Segment extraction failed"})
- return
- segment_data = stdout
- except Exception:
- self._send({"type": "error", "detail": "Segment extraction failed"})
- return
-
- self._send({
- "type": MNP.STREAM_SEGMENT,
- "v": MNP_VERSION,
- "file_id": file_id,
- "segment_index": segment_index,
- "data_b64": base64.b64encode(segment_data).decode(),
- "size": len(segment_data),
- })
-
def _do_chat_message(self, msg: dict) -> None:
"""
Store one message and hand it to everyone else in this group.
@@ -4417,27 +4409,88 @@ class WebRTCPeerSession:
self._send(resp)
def _do_file_upload(self, msg: dict) -> None:
+ """
+ One chunk of an upload, sealed under the group key (MNP 2.0).
+
+ Sealing this direction is not symmetry for its own sake. Downloads have
+ been under a GEK-derived key since the beginning; uploads carried the
+ filename and the raw bytes in plain msgpack, so the same file was
+ ciphertext leaving a node and plaintext arriving at one. The node holds
+ the GEK for its own group, so it opens the payload here — before it
+ decides a destination, before it touches the disk — and refuses a chunk
+ that does not open.
+
+ `upload_id` is the correlation key and stays in clear; `filename`, `dir`
+ and `root` moved inside the seal, which is why every refusal below names
+ the upload rather than the file. A `code` says which refusal it is, and
+ the client already knows what it sent.
+ """
ctx = self._group_ctx()
- filename = msg.get("filename", "")
- chunk_index = msg.get("chunk_index", 0)
- total_chunks = msg.get("total_chunks", 1)
- data = msg.get("data")
+ upload_id = str(msg.get("upload_id") or "")[:64]
+
+ gek = ctx.get("gek")
+ if not gek:
+ self._send({"type": "error", "detail": "Group encryption not initialized",
+ "code": "no_group_key", "upload_id": upload_id})
+ return
- if not filename or data is None:
- self._send({"type": "error", "detail": "Missing filename or data",
- "filename": filename})
+ try:
+ payload = file_upload_payload(gek, self._group_id or "", msg)
+ except Exception:
+ # Deliberately one answer for "not sealed at all" and "sealed wrong":
+ # distinguishing them tells a peer which of the two it got right.
+ # An MNP 1.x client lands here, which is the whole of the upgrade
+ # story — everything else it does still works.
+ self._audit("upload_refused", "unsealed")
+ self._send({
+ "type": "error",
+ "detail": "This upload did not open under the group key — the "
+ "client may be running an older version",
+ "code": "upload_not_sealed",
+ "upload_id": upload_id,
+ })
return
+ filename = payload.get("filename") or ""
+ data = payload.get("data")
+ # From the clear part of the message, so peer-controlled and unchecked
+ # by the AEAD. Everything below compares and adds to them.
+ try:
+ chunk_index = int(msg.get("chunk_index", 0))
+ total_chunks = int(msg.get("total_chunks", 1))
+ except (TypeError, ValueError):
+ self._send({"type": "error", "detail": "Invalid chunk index",
+ "code": "bad_chunk_index", "upload_id": upload_id})
+ return
+
+ def _refuse(detail: str, code: str = "") -> None:
+ """A refusal names the upload, never the file: the name is sealed."""
+ out = {"type": "error", "detail": detail, "upload_id": upload_id}
+ if code:
+ out["code"] = code
+ self._send(out)
+
+ # Types first, and before any state is created. What comes out of a
+ # sealed payload is authenticated, not validated: it is msgpack a
+ # member wrote, and `SAFE_UPLOAD_NAME.match(123)` raises where a
+ # refusal was meant.
+ if not isinstance(filename, str) or not filename:
+ _refuse("Missing filename or data", "upload_incomplete")
+ return
+ # Bytes, always: base64 was the shape of the old plaintext `data` field
+ # and there is no sealed message that can carry a string here.
+ if not isinstance(data, (bytes, bytearray)):
+ _refuse("Invalid chunk encoding", "bad_chunk_encoding")
+ return
+ chunk_bytes = bytes(data)
+
if not SAFE_UPLOAD_NAME.match(filename):
- self._send({"type": "error", "detail": "Invalid filename",
- "filename": filename})
+ _refuse("Invalid filename", "invalid_filename")
return
roots: RootSet | None = ctx.get("roots")
if not roots:
- self._send({"type": "error",
- "detail": "No directories configured for this group",
- "filename": filename})
+ _refuse("No directories configured for this group", "no_roots")
return
# The client names the root it is uploading into — it is browsing one,
@@ -4448,47 +4501,32 @@ class WebRTCPeerSession:
# An unknown name is refused rather than falling back to a writable
# root, because "the file went somewhere else" is discovered weeks
# later — the same reason the old single upload root was never guessed.
- # A client that names nothing is an MNP 1.0 one, and there was exactly
- # one destination in its world: the first writable root.
# `dir` is the folder being browsed, as a virtual path
# (`Media/Films/1999`); `root` is the older, coarser form and is what
- # its first segment means on its own.
- target_rel = str(msg.get("dir") or "").strip().strip("/")
+ # its first segment means on its own. Both are sealed now, so a refusal
+ # below can no longer quote them back.
+ target_rel = str(payload.get("dir") or "").strip().strip("/")
target_root_name = (target_rel.split("/")[0] if target_rel
- else str(msg.get("root") or "").strip())
+ else str(payload.get("root") or "").strip())
upload_root = None
if target_root_name:
upload_root = roots.by_name(target_root_name)
if upload_root is None:
- self._send({"type": "error",
- "detail": f"No directory named "
- f"{target_root_name!r} in this group",
- "code": "no_such_root",
- "filename": filename})
+ _refuse("No such directory in this group", "no_such_root")
return
else:
writable = roots.writable_roots
upload_root = writable[0] if writable else None
if upload_root is None:
- self._send({"type": "error",
- "detail": "No writable directory in this group",
- "code": "no_writable_root",
- "filename": filename})
+ _refuse("No writable directory in this group", "no_writable_root")
return
if not upload_root.writable:
- self._send({"type": "error",
- "detail": f"Directory '{upload_root.name}' is read-only",
- "code": "root_read_only",
- "filename": filename})
+ _refuse("That directory is read-only", "root_read_only")
self._audit("upload_refused", filename[:64])
return
if not upload_root.available:
- self._send({"type": "error",
- "detail": f"Directory '{upload_root.name}' is "
- f"currently unavailable",
- "code": "root_unavailable",
- "filename": filename})
+ _refuse("That directory is currently unavailable", "root_unavailable")
return
# The folder the sender is looking at, and no subdirectory of the node's
@@ -4512,23 +4550,17 @@ class WebRTCPeerSession:
if target_rel:
target_dir = roots.resolve(target_rel)
if target_dir is None or not target_dir.is_dir():
- self._send({"type": "error",
- "detail": "Not a directory in this group",
- "code": "no_such_directory",
- "filename": filename})
+ _refuse("Not a directory in this group", "no_such_directory")
return
rel_dir = target_rel
else:
- # An MNP 1.0 client names nothing; the root itself is where its one
- # destination now is.
+ # A client that names nothing: the first writable root is where its
+ # one destination is.
target_dir = upload_root.path
rel_dir = upload_root.name
if not target_dir.is_dir():
- self._send({"type": "error",
- "detail": f"Directory '{upload_root.name}' is "
- f"currently unavailable",
- "code": "root_unavailable",
- "filename": filename})
+ _refuse("That directory is currently unavailable",
+ "root_unavailable")
return
upload_key = f"{rel_dir}/{filename}"
@@ -4544,33 +4576,24 @@ class WebRTCPeerSession:
# Backstop: _free_name already guarantees this, and it stays because
# it asserts the invariant where the write happens.
if final_path.exists():
- self._send({"type": "error", "detail": "File already exists",
- "filename": filename})
+ _refuse("File already exists", "already_exists")
return
state = {"next_index": 0, "bytes": 0, "stored_name": stored_name}
self._uploads[upload_key] = state
elif state is None:
- self._send({"type": "error", "detail": "Upload not started",
- "filename": filename})
+ _refuse("Upload not started", "not_started")
return
# Reject out-of-order or replayed chunks — otherwise chunk_index>0 appends
# blindly to whatever .part file is already on disk.
if chunk_index != state["next_index"]:
- self._send({"type": "error", "detail": "Unexpected chunk index",
- "filename": filename})
+ _refuse("Unexpected chunk index", "bad_chunk_index")
return
- if isinstance(data, str):
- chunk_bytes = base64.b64decode(data)
- else:
- chunk_bytes = bytes(data)
-
if state["bytes"] + len(chunk_bytes) > MAX_UPLOAD_BYTES:
self._uploads.pop(upload_key, None)
tmp_path.unlink(missing_ok=True)
- self._send({"type": "error", "detail": "Upload exceeds size limit",
- "filename": filename})
+ _refuse("Upload exceeds size limit", "too_large")
return
with open(tmp_path, "wb" if chunk_index == 0 else "ab") as f:
@@ -4578,16 +4601,16 @@ class WebRTCPeerSession:
state["next_index"] = chunk_index + 1
state["bytes"] += len(chunk_bytes)
- self._send({
- "type": MNP.FILE_UPLOAD_ACK,
- "v": MNP_VERSION,
- "chunk_index": chunk_index,
- "filename": filename,
+ self._send(file_upload_ack_wire(
+ gek, self._group_id or "",
+ upload_id=upload_id,
+ chunk_index=chunk_index,
+ filename=filename,
# What it is actually called on disk, which a chat attachment has to
# reference and the uploader deserves to be told.
- "stored_as": stored_name,
- "dir": rel_dir,
- })
+ stored_as=stored_name,
+ dir=rel_dir,
+ ))
if chunk_index + 1 >= total_chunks:
self._uploads.pop(upload_key, None)
@@ -5384,6 +5407,16 @@ class WebRTCPeerSession:
self._audit("stream_video", entry.name)
def _send(self, obj: dict) -> None:
+ # Stamp the reply with the id of the request being answered, so the
+ # caller never has to guess. Only for this session's own replies: a
+ # handler that also pushes to other peers (a chat broadcast, an index
+ # delta) reaches them through *their* _send, where the owner no longer
+ # matches and nothing is stamped — those messages answer no request.
+ # An explicit req_id already on the object wins, and an unsolicited
+ # push (no request in scope) carries none, exactly as before.
+ owner, req_id = _REPLY_TO.get()
+ if req_id is not None and owner is self and "req_id" not in obj:
+ obj = {**obj, "req_id": req_id}
if self._channel and self._channel.readyState == "open":
self._channel.send(_pack(obj))
else: