""" A reply names the request it answers. MNP carried no correlation id until 2026-09-07. A reply named its own type and nothing else, so a client with more than one request outstanding had to work out which one a message answered from the message itself — and for the replies that name nothing, it could not. This module sends `{"type": "error"}` from 240 places and two of them say what they are about; `_dispatch_message`'s catch-all is one of the 238. Such a refusal reached no caller at all: the browser handed it to whichever request happened to be waiting, and the request it belonged to sat until its own 30s timeout. Live symptom (2026-09-06): the Chat composer is disabled while a send is in flight, so a chat message whose refusal went astray froze the tab for thirty seconds. `req_id` is the client's own pending-map key, put on the wire and stamped back onto the reply by `_send`. What matters here, and what the browser cannot check for itself: * a reply carries it, including the refusals that name nothing else; * a *broadcast* does not — it answers no request, and stamping it would hand another peer's client a reply to a request it never made; * work handed to a background task still answers under the right id, which is why this is a ContextVar and not an attribute on the session. """ import asyncio import base64 import msgpack import pytest from meshbay_common.protocol import MNP from meshbay_node.transport.webrtc_server import WebRTCPeerSession pytestmark = pytest.mark.asyncio # A device key's raw bytes. Only its length and its identity with the # connection's own pin matter here; nothing verifies a signature over it. _DEVICE = b"\x07" * 32 class _Channel: readyState = "open" def __init__(self): self.sent = [] def send(self, framed): # Skip the 4-byte length prefix _pack writes. self.sent.append(msgpack.unpackb(framed[4:], raw=False)) def _session(peer_id="p"): s = WebRTCPeerSession.__new__(WebRTCPeerSession) s._ctx = {} s._peer_id = peer_id s._user_id = None s._group_id = "" s._channel = _Channel() s._tasks = set() return s async def test_a_refusal_that_names_nothing_else_names_the_request(): """The reply at the root of the defect: no type of its own to match on.""" s = _session() # No handshake yet, so any other message is refused — a bare `error`, the # same shape the catch-all sends and the same shape a browser could not # route. s._handle_message({"type": MNP.CHAT_MESSAGE, "req_id": 41}) (reply,) = s._channel.sent assert reply["type"] == "error" assert reply["req_id"] == 41, ( "a refusal that names neither the request nor a type of its own is a " "reply no caller can claim") async def test_the_catch_all_refusal_names_the_request_too(): """Every failure in the dispatch loop funnels into one generic reply.""" s = _session() s._user_id = "u" def _boom(msg): raise RuntimeError("filesystem path that must not reach the peer") s._do_chat_message = _boom s._handle_message({"type": MNP.CHAT_MESSAGE, "req_id": 7}) (reply,) = s._channel.sent assert reply == {"type": "error", "detail": "Request failed", "req_id": 7}, ( "the catch-all is where an unforeseen failure ends up, so it is exactly " "the reply that must still be routable") async def test_a_request_without_an_id_is_answered_without_one(): """An older client sends none; nothing may be invented for it.""" s = _session() s._handle_message({"type": MNP.CHAT_MESSAGE}) (reply,) = s._channel.sent assert "req_id" not in reply async def test_a_broadcast_to_another_peer_is_not_stamped(): """The reply goes to the asker; the broadcast goes to everyone else. They travel out of the same handler, and only the first answers anything. Stamping the second would hand another browser a reply keyed to a pending request of its own that it never sent — the very confusion this fixes. """ asker, other = _session("asker"), _session("other") asker._user_id, other._user_id = "a", "b" registry = {"ka": asker, "kb": other} asker._peer_registry = lambda: registry asker._user_names = lambda: {} asker._audit = lambda *a, **k: None asker._group_ctx = lambda: {} asker._spawn = lambda coro: coro.close() asker._registry_key, other._registry_key = "ka", "kb" asker._pinned_pk = base64.b64encode(_DEVICE).decode() asker._device_confirmed = True # A sealed message, because MNP 2.0 has no plaintext chat and the node # refuses one — the bytes need not decrypt, since nothing here opens them. # What this test is about is unchanged: which of the two messages leaving # this handler carries the id. asker._handle_message({ "type": MNP.CHAT_MESSAGE, "req_id": 3, "format": 1, "epoch": 1, "device": _DEVICE, "ct": b"ciphertext", "nonce": b"\x02" * 12, "sig": b"\x03" * 64, }) (ack,) = asker._channel.sent assert ack["type"] == "ack" and ack["req_id"] == 3 (broadcast,) = other._channel.sent assert broadcast["type"] == MNP.CHAT_MESSAGE assert "req_id" not in broadcast, ( "a broadcast answers no request and must not look like a reply") async def test_work_handed_to_a_task_still_answers_under_the_right_id(): """Most handlers `_spawn` their real work, and the reply leaves long after the dispatch call that started it has returned. This is the reason the id lives in a ContextVar: asyncio copies the current context into a task, so the answer keeps the id even though nothing passed it along. An attribute on the session would have been overwritten by the next message to arrive in the meantime. """ s = _session() s._user_id = "u" async def _late(reply): await asyncio.sleep(0.01) s._send({"type": "roster_read_resp", "detail": reply}) s._do_chat_message = lambda msg: s._spawn(_late(msg["payload"])) s._handle_message({"type": MNP.CHAT_MESSAGE, "payload": "first", "req_id": 11}) # A second request arrives while the first one's task is still asleep. s._handle_message({"type": MNP.CHAT_MESSAGE, "payload": "second", "req_id": 12}) await asyncio.gather(*list(s._tasks)) by_id = {m["req_id"]: m["detail"] for m in s._channel.sent} assert by_id == {11: "first", 12: "second"}, ( "a late reply answered under whichever request arrived most recently")