diff options
Diffstat (limited to 'packages/meshbay-node/src')
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: |