From f453cd39ef0fec606b3d1a67b3fae5fa96e719f9 Mon Sep 17 00:00:00 2001 From: Christophe Besson Date: Thu, 24 Sep 2026 12:12:15 +0200 Subject: refactor(node): move the session core out of webrtc_server SessionCore in transport/webrtc/core.py, last among the bases: state, the data channel, task ownership, the peer registry, sending and teardown. The facade keeps _dispatch_message and WebRTCTransport. Co-Authored-By: Claude Opus 5.5 --- CLAUDE.md | 4 +- docs/MESHBAY_NODE_PROTOCOL.md | 11 +- docs/transfers-v1.md | 2 +- .../src/meshbay_node/transport/webrtc/core.py | 338 +++++++++++++++++++++ .../src/meshbay_node/transport/webrtc_server.py | 320 +------------------ .../tests/test_security_regressions.py | 6 +- 6 files changed, 351 insertions(+), 330 deletions(-) create mode 100644 packages/meshbay-node/src/meshbay_node/transport/webrtc/core.py diff --git a/CLAUDE.md b/CLAUDE.md index a062a5e..fee4b73 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -885,8 +885,8 @@ here are kept only where they are a rule about *editing* the code. | Uploads, resume state | `meshbay_node/uploads.py` | §6.4 | | WebRTC transport and every MNP handler | `meshbay_node/transport/webrtc_server.py` | the big one; framing, shared limits and one-shot ffmpeg jobs are in `transport/webrtc/`, handlers for one app only in `transport/webrtc/apps/` | | QUIC transport | `meshbay_node/transport/quic_server.py` | off by default; **must stay at parity with the WebRTC handlers** | -| Background tasks (node) | `webrtc_server.py` — `_spawn()` | the only way to start one; a bare `ensure_future` can be collected | -| Stream handover, backpressure | `transport/webrtc/apps/streaming.py` — `_replace_stream`; `webrtc/files.py` — `DOWNLOAD_BUFFER_HIGH`; `webrtc_server.py` — `shutdown_tasks` | §8.5 | +| Background tasks (node) | `transport/webrtc/core.py` — `_spawn()` | the only way to start one; a bare `ensure_future` can be collected | +| Stream handover, backpressure | `transport/webrtc/apps/streaming.py` — `_replace_stream`; `webrtc/files.py` — `DOWNLOAD_BUFFER_HIGH`; `webrtc/core.py` — `shutdown_tasks` | §8.5 | | Stream diagnosis | `webrtc_server.py` — `client_diag` at DEBUG | the player's own view in the node's log; the only window into a phone | | Chat store and paging | `meshbay_node/chat/store.py` | `get_recent`/`get_before`/`has_before`; `get_messages` pages *forwards* and is not what a chat opens with | | Chat epochs | `meshbay_node/ops.py` — `open_chat_epoch`, `ensure_chat_epoch`, `chat_epoch_keys` | §4.5 | diff --git a/docs/MESHBAY_NODE_PROTOCOL.md b/docs/MESHBAY_NODE_PROTOCOL.md index 35366f0..a219b50 100644 --- a/docs/MESHBAY_NODE_PROTOCOL.md +++ b/docs/MESHBAY_NODE_PROTOCOL.md @@ -4,8 +4,8 @@ **Oldest peer accepted:** `3.0` — `handshake.py` (`MNP_MIN_SUPPORTED`) **Normative implementation:** `meshbay-common` (`protocol.py`, `handshake.py`, `groupbox.py`, `chatbox.py`, `adminop.py`, `join.py`, `device.py`, `crypto.py`, -`webcrypto.py`), `meshbay-node` (`transport/wire.py`, `transport/webrtc_server.py`, -`transport/quic_server.py`, `transfers.py`, `uploads.py`), browser client +`webcrypto.py`), `meshbay-node` (`transport/wire.py`, `transport/webrtc_server.py` +and `transport/webrtc/`, `transport/quic_server.py`, `transfers.py`, `uploads.py`), browser client (`meshbay-hub/static/transport.js`, `static/crypto.js`). **Document status:** descriptive specification of the protocol as implemented on 2026-09-10. It describes the protocol as it stands. Where the code and this document @@ -2291,7 +2291,7 @@ LP(x) = uint32be(len(x)) || x every field, no exceptions | Device attempts | 5 per connection | ” | | `MAX_DEVICES_PER_USER` | 5 | `roster.py` | | `MAX_LINK_INVITES_PER_GROUP` | 20 unredeemed invitation links | `roster.py` | -| `PRE_HANDSHAKE_MAX_MSG` / `MAX_MSG` | 64 KiB / 64 MiB | `webrtc_server.py` / `webrtc/limits.py` | +| `PRE_HANDSHAKE_MAX_MSG` / `MAX_MSG` | 64 KiB / 64 MiB | `webrtc/core.py` / `webrtc/limits.py` | | `CHUNK_SIZE` | 1 MiB | `webrtc/limits.py` | | `DOWNLOAD_BUFFER_HIGH` | 2 MiB | `webrtc/files.py` | | `MAX_UPLOAD_BYTES` | 8 GiB, default only — `max_upload_gb` overrides it per node | `webrtc/upload_handlers.py` | @@ -2301,7 +2301,7 @@ LP(x) = uint32be(len(x)) || x every field, no exceptions | `UPLOAD_PROBE_INDEX` / probe timeout | `-1` / 5 s | `protocol.py`, `transport.js` | | `PART_SUFFIX` / `ORPHAN_AFTER_SECS` | `.part` / 24 h | `uploads.py` | | `DEFAULT_MAX_CONCURRENT` (node-wide, per kind) | 8 | `transfers.py` | -| `DEFAULT_MAX_PER_MEMBER` (per group, per kind) | 2, settable 1–32 | `transfers.py`, `webrtc_server.py` | +| `DEFAULT_MAX_PER_MEMBER` (per group, per kind) | 2, settable 1–32 | `transfers.py`, `webrtc/node_ops.py` | | `GRANT_DEADLINE_SECS` / `IDLE_TIMEOUT_SECS` | 30 s / 120 s | `transfers.py` | | `MAX_QUEUED_PER_MEMBER` / `MAX_MISSED_GRANTS` | 32 / 3 | ” | | `MAX_LEASELESS_IN_FLIGHT` / `LEASELESS_IDLE_SECS` | 12 files / 60 s | ” | @@ -2340,7 +2340,8 @@ meshbay-common/ protocol.py message types, chunk and upload codecs crypto.py GEK, chunk keys, ECIES wrap, BLAKE3 ids webcrypto.py the AES variants the browser can also compute -meshbay-node/ transport/webrtc_server.py the reference implementation of MNP +meshbay-node/ transport/webrtc_server.py the reference implementation of MNP, + transport/webrtc/ assembled from these modules transport/quic_server.py the QUIC transport, in development (§5.2) transport/wire.py the one index encoder transfers.py leases, queues, caps, leaseless reads diff --git a/docs/transfers-v1.md b/docs/transfers-v1.md index f1002a9..47e82a6 100644 --- a/docs/transfers-v1.md +++ b/docs/transfers-v1.md @@ -181,7 +181,7 @@ node-wide slot they would then hold while a second member has none. Reversed, one member arriving first takes all eight. **"Per member" means per account, summed across their devices**, resolved with -`_sessions_of(user_id)` (`webrtc_server.py:3325`) — which exists for exactly +`_sessions_of(user_id)` (`webrtc/core.py`) — which exists for exactly this reason, since device linking landed. Two browsers and a desktop client signed in as the same person share the two slots. Anything else makes the cap a function of how many tabs someone opens. diff --git a/packages/meshbay-node/src/meshbay_node/transport/webrtc/core.py b/packages/meshbay-node/src/meshbay_node/transport/webrtc/core.py new file mode 100644 index 0000000..c3fb623 --- /dev/null +++ b/packages/meshbay-node/src/meshbay_node/transport/webrtc/core.py @@ -0,0 +1,338 @@ +"""The session itself: its state, the data channel it reads from and writes to, +the tasks it owns, the group it belongs to, and how it ends.""" + +import asyncio +import logging +import os +import time +import uuid +from typing import TYPE_CHECKING + +from aiortc import RTCDataChannel, RTCPeerConnection +from meshbay_common import MNP_VERSION +from meshbay_common.protocol import MNP + +from meshbay_node import ops +from meshbay_node import transfers as transfers_mod +from meshbay_node.transport.webrtc.channel import ( + _REPLY_TO, + _DataChannelBuffer, + _get_remote_ip, + _pack, +) + +log = logging.getLogger("meshbay_node.transport.webrtc_server") + +if TYPE_CHECKING: # annotations only: the facade imports this module + from meshbay_node.transport.webrtc_server import WebRTCPeerSession + + +# Budget for an unauthenticated peer: enough for a handshake and a bundle fetch, +# nowhere near enough to be a memory-exhaustion primitive (H6). +PRE_HANDSHAKE_MAX_MSG = 64 * 1024 + + +# Opt-in, off by default: a per-session heartbeat log (message count, time +# since the last message, ICE state) and ICE-state-change logging, on top of +# the connectionstatechange logging that already runs unconditionally. Added +# while chasing a report of the browser side going unresponsive after a +# mobile screen lock; --log-level DEBUG was not the right knob for this, +# since it is already used for the per-message request/response tracing +# every group index lookup produces, and turning that on for days of normal +# operation just to catch one intermittent session is not viable. Set +# MESHBAY_WEBRTC_TRACE=1 in the node's environment for the duration of a +# debugging session. +_WEBRTC_TRACE = os.environ.get("MESHBAY_WEBRTC_TRACE") == "1" +_WEBRTC_TRACE_INTERVAL_S = 30.0 + + +class SessionCore: + def __init__(self, pc: RTCPeerConnection, node_ctx: dict, peer_id: str = ""): + self._pc = pc + self._ctx = node_ctx + # Every background task this session starts. asyncio keeps only a *weak* + # reference to a task, so one that is merely fired and forgotten can be + # collected while it is still running — "Task was destroyed but it is + # pending!" in the log. For _stream_video that meant its `async with + # sem` never reached __aexit__ and the transcode slot was gone for good. + # There are two slots: after two abandoned streams the node answered + # "Server busy" to everything and no video would start at all. + self._tasks: set[asyncio.Task] = set() + self._channel: RTCDataChannel | None = None + self._buffer = _DataChannelBuffer(max_message=PRE_HANDSHAKE_MAX_MSG) + self._pre_proof_fetches = 0 + self._user_id: str | None = None + self._group_id: str | None = None + self._peer_id: str = peer_id + self._remote_ip: str = "" + self._username: str = "" + # This connection's key in the group's peer registry. **Per connection, + # never per account**: one person may hold several devices here, and + # keying the registry by user_id makes the second evict the first, and + # the symptom is invisible: two devices of one account cannot both be + # connected, and whichever disconnects takes the other's chat delivery + # with it. + self._registry_key: str = uuid.uuid4().hex + # Set from the roster: the key this node pinned for this account. Never + # from the JWT — the hub picks what goes in there. + # + # This is the account's *oldest* live device unless `device_hello` has + # told us better — see _do_device_hello. Treat it as "a device of this + # account", not "the device on this connection", anywhere that has not + # checked `_device_confirmed`. + self._pinned_pk: str = "" + # True once this connection proved which device it is. Until then the + # node knows the account and not the key, which is all it ever knew + # before device linking existed. + self._device_confirmed: bool = False + # Flow control for video: how many segments the client says it can take. + self._stream_credit = 0 + self._stream_credit_evt = asyncio.Event() + self._stream_stopped = False + # When the peer last said anything about this stream. See + # _await_stream_credit: silence is what ends a stream, not stinginess. + self._stream_heard_at = 0.0 + # Diagnostics: how many `stream_more n=0` the peer sent. See + # _grant_stream_credit — it tells a paced client from an unpaced one. + self._stream_keepalives = 0 + # The stream this session currently owns. One viewer plays one film at + # a time, so a second request means the first is over — see + # _replace_stream for why waiting for it to time out is not an option. + self._stream_task: asyncio.Task | None = None + # Diagnostics only: when the current stream began and how far it got. + self._stream_started_at: float = 0.0 + self._stream_segments: int = 0 + self._gek_challenge: bytes | None = None + # Same value as the GEK challenge, but kept for the life of the connection: + # a join_request is signed over it, and it must stay verifiable after the + # handshake clears the challenge (an operator pairs while already connected). + self._nonce_node: bytes = b"" + self._join_attempts = 0 + self._nonce_client: bytes = b"" + self._admin_ops: dict[str, dict] = {} # op_id → pending admin operation + # Uploads in progress live in the group context, not here: see + # `_partial_uploads` and `uploads.py`. + # + # Leaseless reads, though, *are* this connection's: the bound is on what + # one session may do while claiming to be browsing, not a pool shared + # between them. Three tabs open is browsing in three tabs. + self._leaseless = transfers_mod.LeaselessReads() + # Whether this session has already been noted as transferring under a + # lease the node does not have (see `_note_unleased`). One line per + # connection, not per chunk. + self._unleased_noted = False + # Diagnostics only (_WEBRTC_TRACE): when the last DataChannel message + # arrived, so the heartbeat can report silence duration. + self._last_msg_at: float = 0.0 + + def _setup_channel(self, channel: RTCDataChannel) -> None: + self._channel = channel + self._msg_count = 0 + + @channel.on("message") + def on_message(message): + if isinstance(message, str): + message = message.encode() + self._msg_count += 1 + self._last_msg_at = time.monotonic() + if self._msg_count <= 3: + 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(): + self._handle_message(msg) + + if _WEBRTC_TRACE: + self._spawn(self._trace_heartbeat()) + + async def _trace_heartbeat(self) -> None: + """Diagnostics only (_WEBRTC_TRACE): periodic proof-of-life for this + session, so a gap in these lines pinpoints when the node stopped + hearing from a peer that (from its own side) may still look connected.""" + while True: + await asyncio.sleep(_WEBRTC_TRACE_INTERVAL_S) + silence = time.monotonic() - self._last_msg_at if self._last_msg_at else -1 + log.info( + "WebRTC heartbeat peer=%s msgs=%d silence=%.0fs pc=%s ice=%s", + self._peer_id, self._msg_count, silence, + self._pc.connectionState, self._pc.iceConnectionState, + ) + + 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 in `_dispatch_message`, 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 _audit(self, event: str, detail: str = "") -> None: + audit = self._ctx.get("audit_store") + if audit and self._user_id: + if not self._remote_ip: + self._remote_ip = _get_remote_ip(self._pc) + self._spawn(audit.log_event( + user_id=self._user_id, + event=event, + ip=self._remote_ip, + username=self._username, + group_id=self._group_id or "", + detail=detail, + )) + + def _broadcast_to_group(self, notice: dict) -> None: + """ + Tell everyone connected to this group about a setting that changed. + + Enforcement never depends on this reaching them — the node is what + refuses — but a control that stays on screen until the next + reconnection is a control people use. + """ + for _uid, session in list(self._peer_registry().items()): + try: + session._send(notice) + except Exception: + pass + + async def _run_op(self, fn, *args, **kwargs): + """ + Call an operation from `meshbay_node.ops` with the daemon's own view. + + The transport carries its own context and the loopback API carries the + daemon state; they overlap but are not the same dict. Handing the MNP + path a *second* set of lookups is exactly how two implementations of one + operation start disagreeing — C1 and C6 one size down — so the daemon + publishes its state here and both adapters call the same function. + """ + state = self._ctx.get("daemon_state") + if state is None: + raise ops.OpError("Node state not available", status=503) + return await fn(state, *args, **kwargs) + + def _spawn(self, coro) -> asyncio.Task: + """Run a coroutine in the background and hold on to it. + + The reference is what keeps the task alive; the done callback is what + stops the set growing. Anything that owns a resource for its lifetime — + a transcode slot, an ffmpeg process — must go through here rather than + `asyncio.ensure_future`. + """ + task = asyncio.ensure_future(coro) + self._tasks.add(task) + def _on_done(t): + self._tasks.discard(t) + if not t.cancelled() and t.exception(): + log.error("Spawned task failed: %s", t.exception(), exc_info=t.exception()) + task.add_done_callback(_on_done) + return task + + def _group_ctx(self) -> dict: + if "groups" in self._ctx and 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 _register_peer(self) -> None: + """Add this connection to its group's peer set. + + One place decides the key, and it is `_registry_key` — per connection, + never per account. Written as a method so a test drives the real + registration rather than a second copy of this line that agrees with it + by construction. + """ + self._peer_registry()[self._registry_key] = self + + def _unregister_peer(self) -> None: + self._peer_registry().pop(self._registry_key, None) + + def _sessions_of(self, user_id: str) -> list["WebRTCPeerSession"]: + """Every live connection this account holds in this group. + + Never "the" connection: with device linking a person may be connected + from a laptop and a phone at once, and an operation that acts on one of + them at random is a revocation that leaves a session running. + """ + return [s for s in list(self._peer_registry().values()) + if s._user_id == user_id] + + def _peer_registry(self) -> dict: + """ + Connected peers for THIS group only. + + Finding H1: this used to live on the shared transport context, so a chat + message was broadcast to every peer on the node regardless of which group + they had authenticated to. + """ + return self._group_ctx().setdefault("_peers", {}) + + def _user_names(self) -> dict: + """Display-name cache, per group — same leak as _peer_registry (H1).""" + return self._group_ctx().setdefault("_user_names", {}) + + def _do_ping(self, msg: dict) -> None: + """Answer a liveness probe on an open channel, echoing the caller's token. + + Echoed rather than bare so a client can match the answer to the probe it + sent and measure a round trip, instead of being reassured by a reply to + some earlier one. + """ + self._send({"type": MNP.PONG, "v": MNP_VERSION, "token": msg.get("token")}) + + 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: + log.warning("WebRTC send skipped: channel=%s", + self._channel.readyState if self._channel else "none") + + async def shutdown_tasks(self) -> None: + """Stop everything this session is doing and give back what it holds. + + Separate from close() because the connection-state handler runs while + aiortc is already tearing the peer connection down — calling pc.close() + from in there would re-enter it. What matters for the transcode slot is + here: cancelling the task runs the exit of its `async with sem`. + """ + self._stop_stream() + # Before the tasks are cancelled: a lease is not held by a task, so + # nothing else would give it back, and this hook is the one place every + # way of walking away arrives at (see the connectionstatechange handler, + # which calls it for a closed tab, a quit browser and a dead network + # alike). + self._release_transfers() + for task in list(self._tasks): + task.cancel() + if self._tasks: + await asyncio.gather(*self._tasks, return_exceptions=True) + + async def close(self) -> None: + self._audit("disconnect") + self._release_transfers() + if self._user_id: + self._unregister_peer() + await self.shutdown_tasks() + await self._pc.close() 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 6a34b11..92e1cd4 100644 --- a/packages/meshbay-node/src/meshbay_node/transport/webrtc_server.py +++ b/packages/meshbay-node/src/meshbay_node/transport/webrtc_server.py @@ -24,21 +24,17 @@ Signaling flow (handled externally by the hub): import asyncio import logging -import os import time -import uuid from typing import Any from aiortc import RTCDataChannel, RTCPeerConnection, RTCSessionDescription from cryptography.hazmat.primitives.asymmetric.ed25519 import ( Ed25519PrivateKey, ) -from meshbay_common import MNP_VERSION from meshbay_common.protocol import ( MNP, ) -from meshbay_node import ops from meshbay_node import transfers as transfers_mod from meshbay_node.indexer import GroupIndex @@ -56,13 +52,8 @@ from meshbay_node.transport.webrtc.apps.streaming import StreamingMixin from meshbay_node.transport.webrtc.apps.subtitles import SubtitlesMixin from meshbay_node.transport.webrtc.apps.video_meta import VideoMetaMixin from meshbay_node.transport.webrtc.blobs import BlobsMixin -from meshbay_node.transport.webrtc.channel import ( - _REPLY_TO, - _DataChannelBuffer, - _get_remote_ip, - _pack, -) from meshbay_node.transport.webrtc.chat import ChatMixin +from meshbay_node.transport.webrtc.core import _WEBRTC_TRACE, SessionCore from meshbay_node.transport.webrtc.files import FilesMixin from meshbay_node.transport.webrtc.group_ops import GroupOpsMixin from meshbay_node.transport.webrtc.handshake import HandshakeMixin @@ -73,10 +64,6 @@ from meshbay_node.transport.webrtc.upload_handlers import UploadMixin log = logging.getLogger(__name__) -# Budget for an unauthenticated peer: enough for a handshake and a bundle fetch, -# nowhere near enough to be a memory-exhaustion primitive (H6). -PRE_HANDSHAKE_MAX_MSG = 64 * 1024 - # How many peer connections this node holds at once, and how long one may stay # without completing the MNP handshake. The budget above bounds what *one* # unauthenticated peer costs; these bound how many there may be and how long @@ -93,154 +80,14 @@ UNAUTHENTICATED_SESSION_TIMEOUT = 60 # seconds MAX_PRE_PROOF_FETCHES = 4 -# Opt-in, off by default: a per-session heartbeat log (message count, time -# since the last message, ICE state) and ICE-state-change logging, on top of -# the connectionstatechange logging that already runs unconditionally. Added -# while chasing a report of the browser side going unresponsive after a -# mobile screen lock; --log-level DEBUG was not the right knob for this, -# since it is already used for the per-message request/response tracing -# every group index lookup produces, and turning that on for days of normal -# operation just to catch one intermittent session is not viable. Set -# MESHBAY_WEBRTC_TRACE=1 in the node's environment for the duration of a -# debugging session. -_WEBRTC_TRACE = os.environ.get("MESHBAY_WEBRTC_TRACE") == "1" -_WEBRTC_TRACE_INTERVAL_S = 30.0 - - class WebRTCPeerSession( AdminMixin, AdmissionMixin, BlobsMixin, ChatMixin, FilesMixin, GroupOpsMixin, HandshakeMixin, NodeOpsMixin, TransferMixin, UploadMixin, StreamingMixin, VideoMetaMixin, MusicMixin, SubtitlesMixin, + SessionCore, ): """One WebRTC peer connection, handling MNP over a DataChannel.""" - def __init__(self, pc: RTCPeerConnection, node_ctx: dict, peer_id: str = ""): - self._pc = pc - self._ctx = node_ctx - # Every background task this session starts. asyncio keeps only a *weak* - # reference to a task, so one that is merely fired and forgotten can be - # collected while it is still running — "Task was destroyed but it is - # pending!" in the log. For _stream_video that meant its `async with - # sem` never reached __aexit__ and the transcode slot was gone for good. - # There are two slots: after two abandoned streams the node answered - # "Server busy" to everything and no video would start at all. - self._tasks: set[asyncio.Task] = set() - self._channel: RTCDataChannel | None = None - self._buffer = _DataChannelBuffer(max_message=PRE_HANDSHAKE_MAX_MSG) - self._pre_proof_fetches = 0 - self._user_id: str | None = None - self._group_id: str | None = None - self._peer_id: str = peer_id - self._remote_ip: str = "" - self._username: str = "" - # This connection's key in the group's peer registry. **Per connection, - # never per account**: one person may hold several devices here, and - # keying the registry by user_id makes the second evict the first, and - # the symptom is invisible: two devices of one account cannot both be - # connected, and whichever disconnects takes the other's chat delivery - # with it. - self._registry_key: str = uuid.uuid4().hex - # Set from the roster: the key this node pinned for this account. Never - # from the JWT — the hub picks what goes in there. - # - # This is the account's *oldest* live device unless `device_hello` has - # told us better — see _do_device_hello. Treat it as "a device of this - # account", not "the device on this connection", anywhere that has not - # checked `_device_confirmed`. - self._pinned_pk: str = "" - # True once this connection proved which device it is. Until then the - # node knows the account and not the key, which is all it ever knew - # before device linking existed. - self._device_confirmed: bool = False - # Flow control for video: how many segments the client says it can take. - self._stream_credit = 0 - self._stream_credit_evt = asyncio.Event() - self._stream_stopped = False - # When the peer last said anything about this stream. See - # _await_stream_credit: silence is what ends a stream, not stinginess. - self._stream_heard_at = 0.0 - # Diagnostics: how many `stream_more n=0` the peer sent. See - # _grant_stream_credit — it tells a paced client from an unpaced one. - self._stream_keepalives = 0 - # The stream this session currently owns. One viewer plays one film at - # a time, so a second request means the first is over — see - # _replace_stream for why waiting for it to time out is not an option. - self._stream_task: asyncio.Task | None = None - # Diagnostics only: when the current stream began and how far it got. - self._stream_started_at: float = 0.0 - self._stream_segments: int = 0 - self._gek_challenge: bytes | None = None - # Same value as the GEK challenge, but kept for the life of the connection: - # a join_request is signed over it, and it must stay verifiable after the - # handshake clears the challenge (an operator pairs while already connected). - self._nonce_node: bytes = b"" - self._join_attempts = 0 - self._nonce_client: bytes = b"" - self._admin_ops: dict[str, dict] = {} # op_id → pending admin operation - # Uploads in progress live in the group context, not here: see - # `_partial_uploads` and `uploads.py`. - # - # Leaseless reads, though, *are* this connection's: the bound is on what - # one session may do while claiming to be browsing, not a pool shared - # between them. Three tabs open is browsing in three tabs. - self._leaseless = transfers_mod.LeaselessReads() - # Whether this session has already been noted as transferring under a - # lease the node does not have (see `_note_unleased`). One line per - # connection, not per chunk. - self._unleased_noted = False - # Diagnostics only (_WEBRTC_TRACE): when the last DataChannel message - # arrived, so the heartbeat can report silence duration. - self._last_msg_at: float = 0.0 - - def _setup_channel(self, channel: RTCDataChannel) -> None: - self._channel = channel - self._msg_count = 0 - - @channel.on("message") - def on_message(message): - if isinstance(message, str): - message = message.encode() - self._msg_count += 1 - self._last_msg_at = time.monotonic() - if self._msg_count <= 3: - 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(): - self._handle_message(msg) - - if _WEBRTC_TRACE: - self._spawn(self._trace_heartbeat()) - - async def _trace_heartbeat(self) -> None: - """Diagnostics only (_WEBRTC_TRACE): periodic proof-of-life for this - session, so a gap in these lines pinpoints when the node stopped - hearing from a peer that (from its own side) may still look connected.""" - while True: - await asyncio.sleep(_WEBRTC_TRACE_INTERVAL_S) - silence = time.monotonic() - self._last_msg_at if self._last_msg_at else -1 - log.info( - "WebRTC heartbeat peer=%s msgs=%d silence=%.0fs pc=%s ice=%s", - self._peer_id, self._msg_count, silence, - self._pc.connectionState, self._pc.iceConnectionState, - ) - - 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) @@ -476,169 +323,6 @@ class WebRTCPeerSession( log.error("Error handling %s on DataChannel: %s", mtype, e, exc_info=True) self._send({"type": "error", "detail": "Request failed"}) - def _audit(self, event: str, detail: str = "") -> None: - audit = self._ctx.get("audit_store") - if audit and self._user_id: - if not self._remote_ip: - self._remote_ip = _get_remote_ip(self._pc) - self._spawn(audit.log_event( - user_id=self._user_id, - event=event, - ip=self._remote_ip, - username=self._username, - group_id=self._group_id or "", - detail=detail, - )) - - def _broadcast_to_group(self, notice: dict) -> None: - """ - Tell everyone connected to this group about a setting that changed. - - Enforcement never depends on this reaching them — the node is what - refuses — but a control that stays on screen until the next - reconnection is a control people use. - """ - for _uid, session in list(self._peer_registry().items()): - try: - session._send(notice) - except Exception: - pass - - async def _run_op(self, fn, *args, **kwargs): - """ - Call an operation from `meshbay_node.ops` with the daemon's own view. - - The transport carries its own context and the loopback API carries the - daemon state; they overlap but are not the same dict. Handing the MNP - path a *second* set of lookups is exactly how two implementations of one - operation start disagreeing — C1 and C6 one size down — so the daemon - publishes its state here and both adapters call the same function. - """ - state = self._ctx.get("daemon_state") - if state is None: - raise ops.OpError("Node state not available", status=503) - return await fn(state, *args, **kwargs) - - def _spawn(self, coro) -> asyncio.Task: - """Run a coroutine in the background and hold on to it. - - The reference is what keeps the task alive; the done callback is what - stops the set growing. Anything that owns a resource for its lifetime — - a transcode slot, an ffmpeg process — must go through here rather than - `asyncio.ensure_future`. - """ - task = asyncio.ensure_future(coro) - self._tasks.add(task) - def _on_done(t): - self._tasks.discard(t) - if not t.cancelled() and t.exception(): - log.error("Spawned task failed: %s", t.exception(), exc_info=t.exception()) - task.add_done_callback(_on_done) - return task - - def _group_ctx(self) -> dict: - if "groups" in self._ctx and 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 _register_peer(self) -> None: - """Add this connection to its group's peer set. - - One place decides the key, and it is `_registry_key` — per connection, - never per account. Written as a method so a test drives the real - registration rather than a second copy of this line that agrees with it - by construction. - """ - self._peer_registry()[self._registry_key] = self - - def _unregister_peer(self) -> None: - self._peer_registry().pop(self._registry_key, None) - - def _sessions_of(self, user_id: str) -> list["WebRTCPeerSession"]: - """Every live connection this account holds in this group. - - Never "the" connection: with device linking a person may be connected - from a laptop and a phone at once, and an operation that acts on one of - them at random is a revocation that leaves a session running. - """ - return [s for s in list(self._peer_registry().values()) - if s._user_id == user_id] - - def _peer_registry(self) -> dict: - """ - Connected peers for THIS group only. - - Finding H1: this used to live on the shared transport context, so a chat - message was broadcast to every peer on the node regardless of which group - they had authenticated to. - """ - return self._group_ctx().setdefault("_peers", {}) - - def _user_names(self) -> dict: - """Display-name cache, per group — same leak as _peer_registry (H1).""" - return self._group_ctx().setdefault("_user_names", {}) - - def _do_ping(self, msg: dict) -> None: - """Answer a liveness probe on an open channel, echoing the caller's token. - - Echoed rather than bare so a client can match the answer to the probe it - sent and measure a round trip, instead of being reassured by a reply to - some earlier one. - """ - self._send({"type": MNP.PONG, "v": MNP_VERSION, "token": msg.get("token")}) - - 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: - log.warning("WebRTC send skipped: channel=%s", - self._channel.readyState if self._channel else "none") - - async def shutdown_tasks(self) -> None: - """Stop everything this session is doing and give back what it holds. - - Separate from close() because the connection-state handler runs while - aiortc is already tearing the peer connection down — calling pc.close() - from in there would re-enter it. What matters for the transcode slot is - here: cancelling the task runs the exit of its `async with sem`. - """ - self._stop_stream() - # Before the tasks are cancelled: a lease is not held by a task, so - # nothing else would give it back, and this hook is the one place every - # way of walking away arrives at (see the connectionstatechange handler, - # which calls it for a closed tab, a quit browser and a dead network - # alike). - self._release_transfers() - for task in list(self._tasks): - task.cancel() - if self._tasks: - await asyncio.gather(*self._tasks, return_exceptions=True) - - async def close(self) -> None: - self._audit("disconnect") - self._release_transfers() - if self._user_id: - self._unregister_peer() - await self.shutdown_tasks() - await self._pc.close() - class WebRTCTransport: """ diff --git a/packages/meshbay-node/tests/test_security_regressions.py b/packages/meshbay-node/tests/test_security_regressions.py index 321d22d..b3f373d 100644 --- a/packages/meshbay-node/tests/test_security_regressions.py +++ b/packages/meshbay-node/tests/test_security_regressions.py @@ -720,11 +720,9 @@ def test_pre_handshake_message_budget_is_small(): H6: the frame limit was a flat 64 MB applied before authentication, so an unauthenticated peer could announce a huge frame and dribble bytes into it. """ + from meshbay_node.transport.webrtc.channel import _DataChannelBuffer + from meshbay_node.transport.webrtc.core import PRE_HANDSHAKE_MAX_MSG from meshbay_node.transport.webrtc.limits import MAX_MSG - from meshbay_node.transport.webrtc_server import ( - PRE_HANDSHAKE_MAX_MSG, - _DataChannelBuffer, - ) assert PRE_HANDSHAKE_MAX_MSG <= 1024 * 1024 assert PRE_HANDSHAKE_MAX_MSG < MAX_MSG -- cgit v1.2.3