Skip to main content

bugwarden/
stdio.rs

1//! What wraps rmcp's stdio transport: a frame bound ([`BoundedLines`])
2//! and a `server/discover` answer ([`DiscoverAnswering`]).
3//!
4//! rmcp's `AsyncRwTransport` reads a newline-delimited frame with
5//! `read_until(b'\n', &mut self.line_buf)` into a `Vec<u8>` it never
6//! bounds, so a peer that writes bytes and no `\n` grows that buffer until
7//! the allocator gives up (#234). HTTP has the derived POST cap and a
8//! `413`; stdio had nothing.
9//!
10//! [`BoundedLines`] is that missing half: an [`AsyncRead`] wrapper that
11//! counts the bytes since the last `\n` and fails the read past the cap.
12//! The obvious alternatives do not work — `JsonRpcMessageCodec`'s
13//! `new_with_max_length` bounds only the `FramedWrite` half, which is
14//! bugwarden's own output, and `AsyncReadExt::take` truncates the whole
15//! stream rather than one frame.
16//!
17//! An over-cap frame is fatal, not skippable: the request id lives inside
18//! the unparsed frame, so no response can name it, and rmcp clients set no
19//! default request timeout — a peer that resumed would hang forever, which
20//! is worse than a closed transport. rmcp maps the read error to
21//! `receive() -> None`, i.e. a silent close indistinguishable from a clean
22//! peer hangup, so the trip is also published through [`BoundedLines::over_cap`]
23//! for `main` to turn into a non-zero exit.
24//!
25//! [`DiscoverAnswering`] is the other half: rmcp reads the stdio lifecycle
26//! off the first frame, so a `server/discover` probe committed the session
27//! to the handshake-free lifecycle before it was answered (#267). Both are
28//! transport-level because both are: the frame never reaches a handler.
29
30use std::future::Future;
31use std::io;
32use std::pin::Pin;
33use std::sync::atomic::{AtomicBool, Ordering};
34use std::sync::Arc;
35use std::task::{ready, Context, Poll};
36
37use rmcp::model::{
38    ClientJsonRpcMessage, ClientRequest, DiscoverResult, ErrorData as McpError, GetMeta as _,
39    ProtocolVersion, RequestId, ServerJsonRpcMessage, ServerResult,
40};
41use rmcp::transport::async_rw::AsyncRwTransport;
42use rmcp::transport::Transport;
43use rmcp::{RoleServer, ServerHandler as _};
44use tokio::io::{AsyncRead, AsyncWrite, ReadBuf};
45
46use crate::server::{BugWarden, SUPPORTED_PROTOCOL_VERSIONS};
47
48/// An [`AsyncRead`] that fails once a newline-delimited frame exceeds
49/// `cap` bytes.
50///
51/// Wraps stdin under rmcp's own `BufReader`, so the counter sees the
52/// delivered chunks and not the frames: the cap is enforced as soon as the
53/// running frame passes it, without waiting for the `\n` that may never
54/// come. A frame of exactly `cap` bytes is accepted, mirroring the http
55/// side's `<= max_bytes`.
56///
57/// The count lives here rather than in a per-`receive` future because rmcp
58/// keeps partial bytes across cancelled polls: `receive` is polled inside a
59/// `select!` and an in-progress line read is dropped whenever an outgoing
60/// response wins, so a per-call counter would restart mid-frame and the cap
61/// would bound nothing.
62#[derive(Debug)]
63pub struct BoundedLines<R> {
64    inner: R,
65    /// Largest frame accepted, in bytes before the delimiter.
66    cap: usize,
67    /// Bytes delivered since the last `\n` — the length of the frame in
68    /// progress.
69    since_newline: usize,
70    /// Set once, never cleared: the read stays failed, and `main` reads it
71    /// after `waiting()` to tell an over-cap close from a clean one.
72    over_cap: Arc<AtomicBool>,
73}
74
75impl<R> BoundedLines<R> {
76    /// Bound `inner` to frames of at most `cap` bytes before the `\n`.
77    pub fn new(inner: R, cap: usize) -> Self {
78        Self {
79            inner,
80            cap,
81            since_newline: 0,
82            over_cap: Arc::new(AtomicBool::new(false)),
83        }
84    }
85
86    /// The frame cap this reader enforces, in bytes before the delimiter.
87    ///
88    /// Exposed so a test can assert that the stdio and http transports were
89    /// sized from the same derivation rather than from two numbers that
90    /// happen to agree today.
91    pub fn cap(&self) -> usize {
92        self.cap
93    }
94
95    /// The trip flag, shared with whoever outlives the transport.
96    ///
97    /// rmcp turns the read error into `receive() -> None`, which the
98    /// service reports as an ordinary close (`QuitReason::Closed`) — so
99    /// without this flag a refused frame and a peer hangup are the same
100    /// `Ok` exit. `main` bails on it after `waiting()` so both stages exit
101    /// non-zero.
102    pub fn over_cap(&self) -> Arc<AtomicBool> {
103        self.over_cap.clone()
104    }
105
106    /// The error a tripped reader returns, on the trip and on every read
107    /// after it. `InvalidData` matches rmcp's own mapping of its codec's
108    /// `MaxLineLengthExceeded`.
109    fn refusal(&self) -> io::Error {
110        io::Error::new(
111            io::ErrorKind::InvalidData,
112            format!("stdio frame exceeds the {}-byte request cap", self.cap),
113        )
114    }
115
116    /// Trip the flag, log once from bugwarden's side, and return the error.
117    ///
118    /// The log matters: this is the only line that names the byte bound it
119    /// applied. rmcp echoes this error's `Display` (`Error reading from
120    /// stream: …`), which stops before the tail added here and names
121    /// neither the bound nor the transport; before `initialize` `main` adds
122    /// a line saying the session is over, which names the refusal but not
123    /// the number.
124    ///
125    /// WARN and not ERROR (#272): the line IS the refusal, the cap doing
126    /// what it exists for, so no operator has anything to act on and a
127    /// client repeats it as often as it can reconnect. ERROR stays with
128    /// `main`'s statement that the process ended — which only the
129    /// pre-`initialize` half writes, the other bailing out of `waiting()`
130    /// silently, so a refused frame costs bugwarden one ERROR or none.
131    fn trip(&mut self) -> io::Error {
132        self.over_cap.store(true, Ordering::Release);
133        tracing::warn!(
134            "stdio frame exceeds the {}-byte request cap; closing the transport",
135            self.cap
136        );
137        self.refusal()
138    }
139}
140
141impl<R: AsyncRead + Unpin> AsyncRead for BoundedLines<R> {
142    /// Delegates the read, then measures what arrived.
143    ///
144    /// Every frame in the chunk is checked at its delimiter, not only the
145    /// bytes after the chunk's last newline: a frame that both passes the
146    /// cap and *ends* inside the tripping chunk would otherwise be seen as a
147    /// short tail and let through. (Chunk length is not the measure either —
148    /// many short frames may share one chunk.) Overshoot is bounded by one chunk
149    /// (8 KiB under rmcp's `BufReader`), since the cap can only be observed
150    /// on bytes already delivered.
151    ///
152    /// Bytes the tripping read delivered stay in `buf`; the caller drops
153    /// them with the error, and the transport is closed either way.
154    fn poll_read(
155        self: Pin<&mut Self>,
156        cx: &mut Context<'_>,
157        buf: &mut ReadBuf<'_>,
158    ) -> Poll<io::Result<()>> {
159        let this = self.get_mut();
160        if this.over_cap.load(Ordering::Acquire) {
161            return Poll::Ready(Err(this.refusal()));
162        }
163        let before = buf.filled().len();
164        ready!(Pin::new(&mut this.inner).poll_read(cx, buf))?;
165        let mut rest = &buf.filled()[before..];
166        while let Some(newline) = rest.iter().position(|&byte| byte == b'\n') {
167            if this.since_newline.saturating_add(newline) > this.cap {
168                return Poll::Ready(Err(this.trip()));
169            }
170            this.since_newline = 0;
171            rest = &rest[newline + 1..];
172        }
173        this.since_newline = this.since_newline.saturating_add(rest.len());
174        if this.since_newline > this.cap {
175            return Poll::Ready(Err(this.trip()));
176        }
177        Poll::Ready(Ok(()))
178    }
179}
180
181/// One reply's write, boxed so it can outlive the `receive` that queued
182/// it. See [`DiscoverAnswering::receive`].
183type QueuedSend<E> = Pin<Box<dyn Future<Output = Result<(), E>> + Send>>;
184
185/// A [`Transport`] that answers `server/discover` itself and hands rmcp
186/// every other frame untouched.
187///
188/// `server/discover` is a probe, but over stdio rmcp lets it CHOOSE the
189/// session lifecycle before anyone answers it: `serve_server_with_ct_inner`
190/// (rmcp 3.1.4 `service/server.rs:510-577`) takes any first non-`ping`
191/// frame that is not `initialize` as a commitment to the handshake-free
192/// lifecycle and calls `Peer::require_request_metadata()` (`:562`) — a
193/// sticky `AtomicBool` that no later `initialize` clears and no public API
194/// can reset, both methods being `pub(crate)`. Every subsequent request
195/// but `initialize` that carries no `_meta` is then refused -32602
196/// (`handler/server.rs:78-100`), so a client that probes and then opens a
197/// session is dead (#267). A probe rmcp refuses is worse: it answers and
198/// then ENDS the process with `ExpectedInitializeRequest`.
199///
200/// Answering here keeps discover off that path entirely. rmcp's first frame
201/// is then the first one that IS a commitment — `initialize` for a session,
202/// or a `_meta`-carrying request that sets the flag exactly as before. The
203/// per-request gate stays rmcp's and no other frame is intercepted — but a
204/// probe no longer moves rmcp past its pre-`initialize` loop, so a `ping`
205/// after one is answered `{}` there rather than refused -32601 by the
206/// per-request handler (`handler/server.rs:112-118`), which is where a
207/// probe used to leave it.
208///
209/// http needs none of this: there `serve_negotiated_request_directly`
210/// answers a discover per POST and sets no flag.
211pub struct DiscoverAnswering<T: Transport<RoleServer>> {
212    inner: T,
213    /// Cloned, not snapshotted: `get_info()` is read per discover, as
214    /// rmcp's default handler reads it.
215    server: BugWarden,
216    /// A discover reply whose write has not finished, parked across
217    /// `receive` calls because `receive` is one arm of rmcp's `select!`
218    /// (`service.rs:1395`) and is dropped whenever another arm wins. A
219    /// reply awaited on that stack dies with the frame that asked for it —
220    /// consumed, never answered, the client hung on that id — while the
221    /// inner transport survives the same cancellation by keeping its
222    /// partial line in `line_buf` (`transport/async_rw.rs:126-131`).
223    pending: Option<QueuedSend<T::Error>>,
224    /// Set once a probe's reply has reached the peer, served or refused.
225    /// See [`DiscoverAnswering::answered`].
226    probed: Arc<AtomicBool>,
227}
228
229impl<T: Transport<RoleServer>> DiscoverAnswering<T> {
230    /// Answer `server/discover` on `inner` from `server`.
231    pub fn new(inner: T, server: BugWarden) -> Self {
232        Self {
233            inner,
234            server,
235            pending: None,
236            probed: Arc::new(AtomicBool::new(false)),
237        }
238    }
239
240    /// Whether a probe has been answered, shared with whoever outlives the
241    /// transport.
242    ///
243    /// A probe answered here leaves rmcp still waiting for its first
244    /// committing frame, so a probe-only client's hangup surfaces as
245    /// `ServerInitializeError::ConnectionClosed` rather than as the clean
246    /// close it is. `main` reads this to tell the two apart (#267).
247    pub fn answered(&self) -> Arc<AtomicBool> {
248        self.probed.clone()
249    }
250
251    /// rmcp's own answer for one discover request, in rmcp's own order:
252    /// the pre-`initialize` metadata check (`service/server.rs:541-551`,
253    /// its text at `:486`), then the declared-revision check
254    /// (`handler/server.rs:64-72`), then that file's default result (`:347`).
255    ///
256    /// That order is the pre-`initialize` path's, which is the one this
257    /// replaces; rmcp's in-session handler runs the two the other way
258    /// round, so a `_meta` malformed BOTH ways — an unserved revision and
259    /// no `clientCapabilities` — now draws -32602 where a mid-session
260    /// probe drew -32022. Refused either way, and a `_meta` missing a
261    /// required key declares no lifecycle worth reading a version out of.
262    ///
263    /// Not stripped for a legacy peer: `strip_result_type_for_legacy_peer`
264    /// (`model.rs:4596`) has no `DiscoverResult` arm, so `resultType`
265    /// stays on the wire whatever revision the probe declares.
266    fn discover_reply(&self, request: &ClientRequest, id: RequestId) -> ServerJsonRpcMessage {
267        let meta = request.get_meta();
268        let missing = meta.missing_required_keys(&ProtocolVersion::V_2026_07_28);
269        if !missing.is_empty() {
270            return ServerJsonRpcMessage::error(
271                McpError::invalid_params(
272                    format!(
273                        "request _meta is missing or has malformed required fields: {}",
274                        missing.join(", ")
275                    ),
276                    None,
277                ),
278                Some(id),
279            );
280        }
281        if let Some(requested) = meta
282            .protocol_version()
283            .filter(|version| !SUPPORTED_PROTOCOL_VERSIONS.contains(version))
284        {
285            return ServerJsonRpcMessage::error(
286                McpError::unsupported_protocol_version(requested, SUPPORTED_PROTOCOL_VERSIONS),
287                Some(id),
288            );
289        }
290        ServerJsonRpcMessage::response(
291            ServerResult::DiscoverResult(DiscoverResult::from_server_info(
292                SUPPORTED_PROTOCOL_VERSIONS.to_vec(),
293                self.server.get_info(),
294            )),
295            id,
296        )
297    }
298}
299
300impl<R, W> DiscoverAnswering<AsyncRwTransport<RoleServer, R, W>>
301where
302    R: AsyncRead + Send + Unpin,
303    W: AsyncWrite + Send + Unpin + 'static,
304{
305    /// The same wrapper over a newline-framed pair, so `main` frames stdio
306    /// through rmcp rather than naming rmcp's transport itself.
307    pub fn framed(read: R, write: W, server: BugWarden) -> Self {
308        Self::new(AsyncRwTransport::new_server(read, write), server)
309    }
310}
311
312impl<T: Transport<RoleServer>> Transport<RoleServer> for DiscoverAnswering<T> {
313    type Error = T::Error;
314
315    fn send(
316        &mut self,
317        item: ServerJsonRpcMessage,
318    ) -> impl Future<Output = Result<(), Self::Error>> + Send + 'static {
319        self.inner.send(item)
320    }
321
322    /// Loops instead of returning: an answered discover is consumed here
323    /// and the next frame awaited, so rmcp never learns a lifecycle from
324    /// one. A failed send closes the transport, like every other write
325    /// failure on this path.
326    ///
327    /// The reply is queued in `pending` and driven at the top of the loop,
328    /// never awaited on this stack: this future is cancelled whenever
329    /// another arm of rmcp's `select!` wins, which under any pipelining is
330    /// most of the time.
331    async fn receive(&mut self) -> Option<ClientJsonRpcMessage> {
332        loop {
333            if let Some(pending) = self.pending.as_mut() {
334                let sent = pending.await;
335                self.pending = None;
336                sent.ok()?;
337                self.probed.store(true, Ordering::Release);
338            }
339            match self.inner.receive().await? {
340                ClientJsonRpcMessage::Request(request)
341                    if matches!(request.request, ClientRequest::DiscoverRequest(_)) =>
342                {
343                    let reply = self.discover_reply(&request.request, request.id);
344                    self.pending = Some(Box::pin(self.inner.send(reply)));
345                }
346                other => return Some(other),
347            }
348        }
349    }
350
351    /// Drains `pending` first: rmcp closes the transport on its way out of
352    /// the serve loop, and the last `receive` may have been cancelled with
353    /// a probe's answer still queued.
354    async fn close(&mut self) -> Result<(), Self::Error> {
355        if let Some(pending) = self.pending.take() {
356            let _ = pending.await;
357        }
358        self.inner.close().await
359    }
360}
361
362#[cfg(test)]
363mod tests {
364    use std::time::Duration;
365
366    use rmcp::transport::async_rw::AsyncRwTransport;
367    use rmcp::transport::Transport as _;
368    use rmcp::RoleServer;
369    use tokio::io::{AsyncReadExt as _, AsyncWriteExt as _, DuplexStream};
370
371    use super::BoundedLines;
372
373    /// Read everything `bytes` yields through a reader capped at `cap`,
374    /// returning the first error if there was one.
375    async fn drain(bytes: &[u8], cap: usize) -> std::io::Result<()> {
376        let mut reader = BoundedLines::new(bytes, cap);
377        let mut sink = [0u8; 512];
378        while reader.read(&mut sink).await? != 0 {}
379        Ok(())
380    }
381
382    fn duplex(cap: usize) -> (DuplexStream, BoundedLines<DuplexStream>) {
383        let (peer, ours) = tokio::io::duplex(64 * 1024);
384        (peer, BoundedLines::new(ours, cap))
385    }
386
387    #[tokio::test]
388    async fn a_frame_of_exactly_the_cap_is_accepted() {
389        // `<= cap`, mirroring http's `<= max_bytes`. A `>=` here would
390        // refuse the largest frame the cap names.
391        drain(b"0123456789abcdef\n", 16)
392            .await
393            .expect("16 bytes is the cap, not past it");
394    }
395
396    #[tokio::test]
397    async fn a_frame_of_exactly_the_cap_waits_for_its_delimiter() {
398        // The chunked half of the same `<=`: a 4 MiB frame arrives in 8 KiB
399        // pieces, so the cap is reached with the `\n` still in flight. A
400        // `>=` in the no-delimiter-yet branch refuses the largest legal
401        // frame on every real stream while the whole-chunk test above still
402        // passes.
403        let (mut peer, mut reader) = duplex(16);
404        let mut buf = [0u8; 512];
405        peer.write_all(b"0123456789abcdef")
406            .await
407            .expect("duplex write");
408        assert_eq!(reader.read(&mut buf).await.expect("exactly the cap"), 16);
409        peer.write_all(b"\n").await.expect("duplex write");
410        assert_eq!(reader.read(&mut buf).await.expect("the delimiter"), 1);
411    }
412
413    #[tokio::test]
414    async fn one_byte_past_the_cap_is_refused() {
415        let err = drain(b"0123456789abcdefg\n", 16)
416            .await
417            .expect_err("17 bytes is past a 16-byte cap");
418        assert_eq!(err.kind(), std::io::ErrorKind::InvalidData);
419        assert!(
420            err.to_string().contains("16-byte request cap"),
421            "the refusal must name the cap: {err}"
422        );
423    }
424
425    #[tokio::test]
426    async fn the_count_is_per_frame_and_not_a_running_total() {
427        // Chunked on purpose: within a single chunk the loop starts from
428        // zero whether or not it resets, so a whole-frames-in-one-read test
429        // passes even with the reset deleted. Only a frame that STRADDLES a
430        // chunk boundary leaves a stale count for the next one to inherit —
431        // and then an ordinary session is refused after a few requests.
432        let (mut peer, mut reader) = duplex(16);
433        let mut buf = [0u8; 512];
434        peer.write_all(b"0123456789").await.expect("duplex write");
435        assert_eq!(reader.read(&mut buf).await.expect("10 bytes"), 10);
436        // Closes that frame at 12 bytes and opens the next one.
437        peer.write_all(b"ab\ncdef").await.expect("duplex write");
438        assert_eq!(reader.read(&mut buf).await.expect("12 then 4"), 7);
439        // 10 bytes into the second frame, but 20 into the session.
440        peer.write_all(b"ghijkl\n").await.expect("duplex write");
441        assert_eq!(reader.read(&mut buf).await.expect("10 bytes, not 20"), 7);
442    }
443
444    #[tokio::test]
445    async fn a_frame_that_ends_inside_the_tripping_chunk_is_still_refused() {
446        // One chunk carrying an over-cap frame AND its delimiter. Checking
447        // only the bytes after the chunk's last newline lets this through
448        // as a short tail.
449        let err = drain(b"0123456789abcdefg\nshort\n", 16)
450            .await
451            .expect_err("the delimiter does not excuse the frame before it");
452        assert_eq!(err.kind(), std::io::ErrorKind::InvalidData);
453    }
454
455    #[tokio::test]
456    async fn an_over_cap_second_frame_inside_one_chunk_is_refused() {
457        // The under-cap first frame must not reset the scan so far that the
458        // second one, entirely inside the same chunk, escapes its check.
459        let err = drain(b"short\n0123456789abcdefg\n", 16)
460            .await
461            .expect_err("the second frame of the chunk is over the cap");
462        assert_eq!(err.kind(), std::io::ErrorKind::InvalidData);
463    }
464
465    #[tokio::test]
466    async fn several_under_cap_frames_in_one_chunk_are_accepted() {
467        // Four 8-byte frames in a 32-byte chunk at cap 16: a check that
468        // measured the chunk, or forgot to reset at each delimiter, would
469        // refuse a stream of frames that individually fit.
470        drain(b"aaaaaaa\nbbbbbbb\nccccccc\nddddddd\n", 16)
471            .await
472            .expect("every frame is half the cap");
473    }
474
475    #[tokio::test]
476    async fn the_count_survives_a_cancelled_read() {
477        // rmcp polls `receive` inside a `select!` and keeps the partial
478        // bytes when another branch wins, so the counter has to live in the
479        // reader. A per-future count would restart here and accept 20 bytes
480        // under a 16-byte cap.
481        let (mut peer, mut reader) = duplex(16);
482        let mut buf = [0u8; 512];
483        peer.write_all(b"0123456789").await.expect("duplex write");
484        assert_eq!(reader.read(&mut buf).await.expect("first chunk"), 10);
485        // A read with nothing to deliver, dropped mid-flight.
486        assert!(
487            tokio::time::timeout(Duration::from_millis(50), reader.read(&mut buf))
488                .await
489                .is_err(),
490            "an empty duplex must leave the read pending"
491        );
492        peer.write_all(b"0123456789").await.expect("duplex write");
493        let err = reader
494            .read(&mut buf)
495            .await
496            .expect_err("20 bytes of one frame is past a 16-byte cap");
497        assert_eq!(err.kind(), std::io::ErrorKind::InvalidData);
498    }
499
500    #[tokio::test]
501    async fn the_refusal_is_sticky() {
502        // rmcp stops reading after the first error, but nothing in the
503        // `AsyncRead` contract makes it: a reader that recovered would hand
504        // the tail of a refused frame to the parser as a fresh one.
505        //
506        // The refused frame carries its own delimiter, so the count is back
507        // at zero and the next frame is legal on its face — without the
508        // sticky check the second read succeeds. A refusal with no
509        // delimiter would leave the count over the cap and trip again
510        // whether or not the reader remembered anything.
511        let (mut peer, mut reader) = duplex(16);
512        let mut buf = [0u8; 512];
513        peer.write_all(b"0123456789abcdefg\n")
514            .await
515            .expect("duplex write");
516        reader.read(&mut buf).await.expect_err("over the cap");
517        peer.write_all(b"{}\n").await.expect("duplex write");
518        let err = reader
519            .read(&mut buf)
520            .await
521            .expect_err("a tripped reader stays failed");
522        assert_eq!(err.kind(), std::io::ErrorKind::InvalidData);
523    }
524
525    #[tokio::test]
526    async fn an_under_cap_stream_reaches_the_reader_unchanged() {
527        // The wrapper is transparent below the cap: a `\r\n` frame, an
528        // empty frame and a trailing fragment all pass through byte for
529        // byte, so nothing it does can corrupt the framing rmcp parses.
530        let payload = b"{\"a\":1}\r\n\n{\"b\":2}\ntrailing";
531        let mut reader = BoundedLines::new(&payload[..], 16);
532        let mut read = Vec::new();
533        reader.read_to_end(&mut read).await.expect("under the cap");
534        assert_eq!(read, payload);
535    }
536
537    #[tokio::test]
538    async fn rmcp_turns_the_refusal_into_a_closed_transport() {
539        // The contract this whole design rests on, and the one thing an
540        // rmcp bump can silently change: `receive` maps a read error to
541        // `None`, so the service quits instead of retrying. If a future
542        // rmcp propagated the error or resumed the read, this fails and the
543        // exit-code path in `main` needs rethinking.
544        //
545        // Bounded: a reader that does not refuse leaves rmcp blocked in
546        // `read_until` on a duplex nothing else will write to, and this
547        // would hang the suite instead of failing.
548        let (mut peer, reader) = duplex(16);
549        let mut transport =
550            AsyncRwTransport::<RoleServer, _, _>::new_server(reader, tokio::io::sink());
551        peer.write_all(b"0123456789abcdefghij\n")
552            .await
553            .expect("duplex write");
554        let received = tokio::time::timeout(Duration::from_secs(5), transport.receive())
555            .await
556            .expect("the transport must not block on a frame it will never accept");
557        assert!(
558            received.is_none(),
559            "an over-cap frame must close the transport, not yield a message"
560        );
561    }
562}
563
564/// [`DiscoverAnswering`]: what the wrapper answers, what it lets past, and
565/// that the lifecycle rmcp ends up on is the one the client's FIRST
566/// committing request chose (#267).
567#[cfg(test)]
568mod discover {
569    use std::sync::Arc;
570    use std::time::Duration;
571
572    use bugwarden_core::guard::Guard;
573    use bugwarden_core::policy::Policy;
574    use rmcp::model::{
575        ClientJsonRpcMessage, ClientRequest, DiscoverResult, ProtocolVersion, ServerJsonRpcMessage,
576    };
577    use rmcp::transport::Transport as _;
578    use rmcp::{RoleServer, ServerHandler as _, ServiceExt as _};
579    use serde_json::{json, Value};
580    use tokio::io::{
581        AsyncBufReadExt as _, AsyncWriteExt as _, BufReader, DuplexStream, ReadHalf, WriteHalf,
582    };
583
584    use super::DiscoverAnswering;
585    use crate::pinned_cli::pinned;
586    use crate::server::{bugzilla_client, BugWarden, SUPPORTED_PROTOCOL_VERSIONS};
587
588    /// Bounded so a wrapper that swallows a frame fails rather than hanging
589    /// the suite until CI's own timeout kills it.
590    const REPLY_BUDGET: Duration = Duration::from_secs(10);
591
592    /// A revision this build serves, so a probe carrying it is answered.
593    const SERVED: &str = "2026-07-28";
594
595    /// A revision no build serves, standing in for the next one a client
596    /// adopts before bugwarden does — the reporter's case.
597    const UNSERVED: &str = "2027-01-01";
598
599    /// The narrow pipe [`Session::open_narrow`] serves over: wide enough
600    /// for a whole burst of requests, far too narrow for the answers.
601    const PIPE: usize = 8 * 1024;
602
603    fn test_server() -> BugWarden {
604        let cfg = Arc::new(pinned(&[
605            "bugwarden",
606            "--bugzilla-server",
607            "https://bugzilla.example.invalid",
608            "--transport",
609            "stdio",
610            "--api-key",
611            "test-key",
612        ]));
613        let guard = Arc::new(Guard {
614            policy: Policy::from_toml_str("").expect("the empty policy must parse"),
615        });
616        let bz = Arc::new(bugzilla_client(&cfg).expect("client must build"));
617        BugWarden::new(cfg, guard, bz).expect("server must build")
618    }
619
620    /// The `_meta` a 2026-07-28 client sends, with the two required keys.
621    fn meta(version: &str) -> Value {
622        json!({
623            "io.modelcontextprotocol/protocolVersion": version,
624            "io.modelcontextprotocol/clientCapabilities": {},
625        })
626    }
627
628    fn discover(id: u32, meta: Value) -> Value {
629        json!({ "jsonrpc": "2.0", "id": id, "method": "server/discover",
630                "params": { "_meta": meta } })
631    }
632
633    /// The one frame rmcp answers before `initialize`, and this file's
634    /// filler for pipelining probes against.
635    fn ping(id: u32) -> Value {
636        json!({ "jsonrpc": "2.0", "id": id, "method": "ping" })
637    }
638
639    /// The client end of a served stdio session: raw JSON lines in, raw
640    /// JSON lines out, exactly what a subprocess client writes.
641    struct Session {
642        write: WriteHalf<DuplexStream>,
643        lines: tokio::io::Lines<BufReader<ReadHalf<DuplexStream>>>,
644    }
645
646    impl Session {
647        /// Serve `test_server()` over a duplex pair, wrapped the way `main`
648        /// wraps stdio.
649        fn open() -> Self {
650            Self::serving(true, 256 * 1024)
651        }
652
653        /// The same server with rmcp answering the probe itself: the
654        /// oracle the wrapper is diffed against.
655        fn open_bare() -> Self {
656            Self::serving(false, 256 * 1024)
657        }
658
659        /// A pipe too small for the replies a burst produces, so the
660        /// writes really do park — which is what makes rmcp's `select!`
661        /// cancel a `receive` mid-send.
662        fn open_narrow() -> Self {
663            Self::serving(true, PIPE)
664        }
665
666        fn serving(wrapped: bool, capacity: usize) -> Self {
667            let (theirs, ours) = tokio::io::duplex(capacity);
668            let (read, write) = tokio::io::split(ours);
669            let server = test_server();
670            tokio::spawn(async move {
671                let service = if wrapped {
672                    let transport = DiscoverAnswering::framed(read, write, server.clone());
673                    server.serve(transport).await?
674                } else {
675                    server.serve((read, write)).await?
676                };
677                service.waiting().await?;
678                Ok::<(), anyhow::Error>(())
679            });
680            let (read, write) = tokio::io::split(theirs);
681            Session {
682                write,
683                lines: BufReader::new(read).lines(),
684            }
685        }
686
687        async fn send(&mut self, message: Value) {
688            self.write
689                .write_all(format!("{message}\n").as_bytes())
690                .await
691                .expect("the session must accept input");
692        }
693
694        /// Write every frame in one go, so replies are still being written
695        /// when the next frames land — the contention that makes rmcp's
696        /// `select!` cancel `receive`.
697        async fn burst(&mut self, frames: &[Value]) {
698            let mut buffer = String::new();
699            for frame in frames {
700                buffer.push_str(&format!("{frame}\n"));
701            }
702            self.write
703                .write_all(buffer.as_bytes())
704                .await
705                .expect("the session must accept input");
706        }
707
708        /// The next frame the server wrote, parsed.
709        async fn recv(&mut self) -> Value {
710            let line = tokio::time::timeout(REPLY_BUDGET, self.lines.next_line())
711                .await
712                .expect("the server must answer within the budget")
713                .expect("the session must be readable")
714                .expect("the server must not close before answering");
715            serde_json::from_str(&line).unwrap_or_else(|e| panic!("{line:?}: {e}"))
716        }
717
718        /// Send `message` and read its reply.
719        async fn call(&mut self, message: Value) -> Value {
720            self.send(message).await;
721            self.recv().await
722        }
723
724        /// Handshake at 2025-11-25 and announce it, the legacy lifecycle.
725        async fn initialize(&mut self, id: u32) -> Value {
726            let reply = self
727                .call(json!({
728                    "jsonrpc": "2.0", "id": id, "method": "initialize",
729                    "params": {
730                        "protocolVersion": "2025-11-25",
731                        "capabilities": {},
732                        "clientInfo": { "name": "discover-test", "version": "0" },
733                    }
734                }))
735                .await;
736            self.send(json!({ "jsonrpc": "2.0", "method": "notifications/initialized" }))
737                .await;
738            reply
739        }
740    }
741
742    /// The `tools` array of a served listing.
743    fn tools_of(reply: &Value) -> &Vec<Value> {
744        reply["result"]["tools"]
745            .as_array()
746            .unwrap_or_else(|| panic!("a served listing carries a tools array: {reply}"))
747    }
748
749    /// The JSON-RPC error code of a refusal.
750    fn error_code(reply: &Value) -> i64 {
751        reply["error"]["code"]
752            .as_i64()
753            .unwrap_or_else(|| panic!("a refusal carries an error code: {reply}"))
754    }
755
756    /// rmcp's per-request metadata refusal, spelled out rather than read
757    /// off the code under test.
758    const MISSING_META: &str = "request _meta is missing or has malformed required fields: ";
759
760    #[tokio::test]
761    async fn the_reporter_chain_reaches_the_tool_list() {
762        // Issue #267 exactly: a Go SDK client probes with a revision this
763        // build does not serve, is told so, falls back to the legacy
764        // handshake — and used to find every later request refused -32602,
765        // because rmcp had already committed the session to the
766        // handshake-free lifecycle on the probe alone.
767        let mut session = Session::open();
768        let probe = session.call(discover(1, meta(UNSERVED))).await;
769        assert_eq!(error_code(&probe), -32022, "{probe}");
770        let init = session.initialize(2).await;
771        assert_eq!(init["result"]["protocolVersion"], "2025-11-25", "{init}");
772        let listed = session
773            .call(json!({ "jsonrpc": "2.0", "id": 3, "method": "tools/list", "params": {} }))
774            .await;
775        assert!(!tools_of(&listed).is_empty(), "{listed}");
776    }
777
778    #[tokio::test]
779    async fn a_served_probe_then_a_session_reaches_the_tool_list() {
780        // The other half of the same defect, and the commoner one: the
781        // probe SUCCEEDS and the client opens a session anyway. Nothing
782        // about the answer differs; the lifecycle must still be the
783        // session's.
784        let mut session = Session::open();
785        let probe = session.call(discover(1, meta(SERVED))).await;
786        assert!(probe["result"]["supportedVersions"].is_array(), "{probe}");
787        session.initialize(2).await;
788        let listed = session
789            .call(json!({ "jsonrpc": "2.0", "id": 3, "method": "tools/list", "params": {} }))
790            .await;
791        assert!(!tools_of(&listed).is_empty(), "{listed}");
792    }
793
794    #[tokio::test]
795    async fn a_probe_with_no_meta_is_refused_and_the_session_survives() {
796        // On main this frame was answered and then ENDED the process
797        // (`ExpectedInitializeRequest`), so the client's `initialize`
798        // reached a closed pipe. The refusal is unchanged; what changes is
799        // that the session outlives it.
800        let mut session = Session::open();
801        let probe = session
802            .call(
803                json!({ "jsonrpc": "2.0", "id": 1, "method": "server/discover",
804                          "params": {} }),
805            )
806            .await;
807        assert_eq!(error_code(&probe), -32602, "{probe}");
808        assert_eq!(
809            probe["error"]["message"],
810            json!(format!(
811                "{MISSING_META}io.modelcontextprotocol/protocolVersion, \
812                 io.modelcontextprotocol/clientCapabilities"
813            )),
814            "{probe}"
815        );
816        let init = session.initialize(2).await;
817        assert_eq!(init["result"]["protocolVersion"], "2025-11-25", "{init}");
818        let listed = session
819            .call(json!({ "jsonrpc": "2.0", "id": 3, "method": "tools/list", "params": {} }))
820            .await;
821        assert!(!tools_of(&listed).is_empty(), "{listed}");
822    }
823
824    #[tokio::test]
825    async fn a_probe_mid_session_is_answered_and_changes_nothing() {
826        // A probe is not only an opener: a client may re-probe a live
827        // session. The wrapper intercepts that one too, so the answer must
828        // still be the full result and the session must keep serving
829        // `_meta`-free requests afterwards.
830        let mut session = Session::open();
831        session.initialize(1).await;
832        let probe = session.call(discover(2, meta(SERVED))).await;
833        assert_eq!(probe["result"]["resultType"], "complete", "{probe}");
834        let listed = session
835            .call(json!({ "jsonrpc": "2.0", "id": 3, "method": "tools/list", "params": {} }))
836            .await;
837        assert!(!tools_of(&listed).is_empty(), "{listed}");
838    }
839
840    #[tokio::test]
841    async fn a_pipelined_probe_is_answered_under_write_contention() {
842        // The one thing this wrapper does that the inner transport does
843        // not: hold a reply. rmcp polls `receive` as one arm of a
844        // `select!` (rmcp 3.1.4 `service.rs:1395`) and drops the future
845        // whenever another arm wins, so a reply awaited on `receive`'s own
846        // stack dies with the frame that asked for it — consumed, never
847        // answered, the client hung on that id while every other reply
848        // arrives. Every other row here is lock-step, where nothing else
849        // is ever ready and the drop never happens. A session first:
850        // before `initialize` rmcp awaits `receive` outside its `select!`
851        // (`service/server.rs:511`), so a probe-only chain never shows
852        // this either. The whole burst fits `PIPE`, so only the SERVER's
853        // writes park.
854        const FRAMES: u32 = 40;
855        let mut session = Session::open_narrow();
856        session.initialize(1).await;
857        let frames: Vec<Value> = (2..=FRAMES + 1)
858            .map(|id| {
859                if id % 2 == 0 {
860                    discover(id, meta(SERVED))
861                } else {
862                    ping(id)
863                }
864            })
865            .collect();
866        session.burst(&frames).await;
867        let mut answered = std::collections::HashSet::new();
868        for _ in 0..FRAMES {
869            let reply = session.recv().await;
870            answered.insert(
871                reply["id"]
872                    .as_u64()
873                    .unwrap_or_else(|| panic!("every reply names its request: {reply}")),
874            );
875        }
876        let unanswered: Vec<u64> = (2..=u64::from(FRAMES) + 1)
877            .filter(|id| !answered.contains(id))
878            .collect();
879        assert!(unanswered.is_empty(), "unanswered ids: {unanswered:?}");
880    }
881
882    #[tokio::test]
883    async fn the_per_request_lifecycle_is_still_rmcps_to_choose() {
884        // The invariant the fix must not weaken. A probe commits nothing,
885        // so the FIRST `_meta`-carrying request is what puts rmcp on the
886        // handshake-free lifecycle — and rmcp's gate then refuses a
887        // request that drops the `_meta`, exactly as before. A wrapper
888        // that had answered or relaxed anything past discover would serve
889        // the second listing.
890        let mut session = Session::open();
891        let probe = session.call(discover(1, meta(SERVED))).await;
892        assert!(probe["result"]["supportedVersions"].is_array(), "{probe}");
893        let declared = session
894            .call(json!({ "jsonrpc": "2.0", "id": 2, "method": "tools/list",
895                          "params": { "_meta": meta(SERVED) } }))
896            .await;
897        assert!(!tools_of(&declared).is_empty(), "{declared}");
898        let bare = session
899            .call(json!({ "jsonrpc": "2.0", "id": 3, "method": "tools/list", "params": {} }))
900            .await;
901        assert_eq!(error_code(&bare), -32602, "{bare}");
902        assert!(
903            bare["error"]["message"]
904                .as_str()
905                .is_some_and(|message| message.starts_with(MISSING_META)),
906            "{bare}"
907        );
908    }
909
910    #[tokio::test]
911    async fn a_ping_before_initialize_is_still_answered_by_rmcp() {
912        // The one other frame rmcp accepts before `initialize`. The
913        // wrapper must pass it through, not eat it, and the probe that
914        // follows must still be answered here.
915        let mut session = Session::open();
916        let pong = session.call(ping(1)).await;
917        assert_eq!(pong["result"], json!({}), "{pong}");
918        let probe = session.call(discover(2, meta(SERVED))).await;
919        assert_eq!(probe["result"]["resultType"], "complete", "{probe}");
920        let init = session.initialize(3).await;
921        assert_eq!(init["result"]["protocolVersion"], "2025-11-25", "{init}");
922    }
923
924    #[tokio::test]
925    async fn a_malformed_meta_names_only_the_key_that_is_wrong() {
926        // `missing_required_keys` counts a key present but undecodable as
927        // missing, and the refusal names it. A number where the revision
928        // belongs must not be read as a revision, nor drag the well-formed
929        // capabilities key into the message.
930        let mut session = Session::open();
931        let probe = session
932            .call(discover(
933                1,
934                json!({
935                    "io.modelcontextprotocol/protocolVersion": 7,
936                    "io.modelcontextprotocol/clientCapabilities": {},
937                }),
938            ))
939            .await;
940        assert_eq!(error_code(&probe), -32602, "{probe}");
941        assert_eq!(
942            probe["error"]["message"],
943            json!(format!(
944                "{MISSING_META}io.modelcontextprotocol/protocolVersion"
945            )),
946            "{probe}"
947        );
948    }
949
950    #[tokio::test]
951    async fn an_in_session_probe_wrong_two_ways_is_refused_on_the_metadata() {
952        // The one shape whose refusal CODE this fix changes, pinned as a
953        // decision. rmcp's in-session handler checks the declared revision
954        // before the required keys (-32022); its pre-`initialize` path —
955        // the one the wrapper replaces, and the one a probe normally takes
956        // — checks the keys first (-32602). A `_meta` missing a required
957        // key declares no lifecycle worth reading a revision out of, so
958        // the pre-`initialize` order is the one kept on both.
959        let mut session = Session::open();
960        session.initialize(1).await;
961        let probe = session
962            .call(discover(
963                2,
964                json!({ "io.modelcontextprotocol/protocolVersion": UNSERVED }),
965            ))
966            .await;
967        assert_eq!(error_code(&probe), -32602, "{probe}");
968        assert_eq!(
969            probe["error"]["message"],
970            json!(format!(
971                "{MISSING_META}io.modelcontextprotocol/clientCapabilities"
972            )),
973            "{probe}"
974        );
975    }
976
977    #[tokio::test]
978    async fn the_answer_is_the_sdk_default_discover_result() {
979        // The parity oracle: the wrapper replaces rmcp's default
980        // `ServerHandler::discover`, so what reaches the wire must be that
981        // default's result, built from this server's own `get_info()` —
982        // serverInfo, instructions, capabilities and the SEP-2549 cache
983        // hints included. Compared whole: a field this wrapper forgot, or
984        // one it added, is exactly how the identity surface drifts.
985        let mut session = Session::open();
986        let probe = session.call(discover(1, meta(SERVED))).await;
987        let expected = serde_json::to_value(DiscoverResult::from_server_info(
988            SUPPORTED_PROTOCOL_VERSIONS.to_vec(),
989            test_server().get_info(),
990        ))
991        .expect("the SDK result must serialize");
992        assert_eq!(probe["result"], expected, "{probe}");
993        assert_eq!(
994            probe["result"]["_meta"]["io.modelcontextprotocol/serverInfo"],
995            json!({ "name": "bugwarden", "version": env!("CARGO_PKG_VERSION") }),
996            "the probe must name this build, and nothing else: {probe}"
997        );
998        // Present whatever revision the probe declares:
999        // `strip_result_type_for_legacy_peer` (rmcp 3.1.4 model.rs) has no
1000        // `DiscoverResult` arm, so rmcp never stripped it here either.
1001        assert_eq!(probe["result"]["resultType"], "complete", "{probe}");
1002    }
1003
1004    /// The `_meta` shape whose refusal CODE the wrapper moves, named once
1005    /// so the differential and its exception cannot drift apart.
1006    const UNSERVED_WITHOUT_CAPABILITIES: &str = "an unserved revision, no capabilities";
1007
1008    /// Every `_meta` shape the differential drives: what a probe carries,
1009    /// well formed and not.
1010    fn probe_shapes() -> Vec<(&'static str, Value)> {
1011        vec![
1012            ("both keys, a served revision", meta(SERVED)),
1013            ("both keys, a legacy revision", meta("2025-11-25")),
1014            ("both keys, the oldest served revision", meta("2024-11-05")),
1015            ("both keys, an unserved revision", meta(UNSERVED)),
1016            (
1017                "the Go SDK shape, clientInfo included",
1018                json!({
1019                    "io.modelcontextprotocol/protocolVersion": SERVED,
1020                    "io.modelcontextprotocol/clientCapabilities": {},
1021                    "io.modelcontextprotocol/clientInfo": { "name": "agy", "version": "1.7.0" },
1022                }),
1023            ),
1024            ("an empty _meta", json!({})),
1025            (
1026                "a served revision, no capabilities",
1027                json!({ "io.modelcontextprotocol/protocolVersion": SERVED }),
1028            ),
1029            (
1030                "capabilities, no revision",
1031                json!({ "io.modelcontextprotocol/clientCapabilities": {} }),
1032            ),
1033            (
1034                UNSERVED_WITHOUT_CAPABILITIES,
1035                json!({ "io.modelcontextprotocol/protocolVersion": UNSERVED }),
1036            ),
1037            (
1038                "a revision that is not a string",
1039                json!({
1040                    "io.modelcontextprotocol/protocolVersion": 7,
1041                    "io.modelcontextprotocol/clientCapabilities": {},
1042                }),
1043            ),
1044            (
1045                "capabilities that are not an object",
1046                json!({
1047                    "io.modelcontextprotocol/protocolVersion": SERVED,
1048                    "io.modelcontextprotocol/clientCapabilities": "yes",
1049                }),
1050            ),
1051        ]
1052    }
1053
1054    /// One probe on a session opened for it alone: rmcp ends a bare
1055    /// session on the shapes it refuses pre-`initialize`, so no two rows
1056    /// may share one.
1057    async fn probe_reply(mut session: Session, in_session: bool, meta: Value) -> Value {
1058        let id = if in_session {
1059            session.initialize(1).await;
1060            2
1061        } else {
1062            1
1063        };
1064        session.call(discover(id, meta)).await
1065    }
1066
1067    #[tokio::test]
1068    async fn the_wrapper_answers_what_rmcp_answered() {
1069        // The oracle `the_answer_is_the_sdk_default_discover_result`
1070        // cannot be: that one compares the wrapper against the constructor
1071        // the wrapper itself calls, so it sees nothing rmcp does BETWEEN
1072        // handler and wire (`handler/server.rs:246-259`, the legacy-peer
1073        // strip). This serves the same server bare and diffs its
1074        // replies over every probe shape, on both paths a probe can take —
1075        // so an rmcp bump that changes its default `discover`, or starts
1076        // stripping a `DiscoverResult`, fails here.
1077        for (name, meta) in probe_shapes() {
1078            for in_session in [false, true] {
1079                let bare = probe_reply(Session::open_bare(), in_session, meta.clone()).await;
1080                let wrapped = probe_reply(Session::open(), in_session, meta.clone()).await;
1081                if in_session && name == UNSERVED_WITHOUT_CAPABILITIES {
1082                    // The only row that moves, pinned in both directions:
1083                    // rmcp's in-session handler checks the declared
1084                    // revision before the required keys, the
1085                    // pre-`initialize` path the wrapper replaces checks
1086                    // the keys first.
1087                    assert_eq!(error_code(&bare), -32022, "{name}: {bare}");
1088                    assert_eq!(error_code(&wrapped), -32602, "{name}: {wrapped}");
1089                    continue;
1090                }
1091                assert_eq!(bare, wrapped, "{name}, in a session: {in_session}");
1092            }
1093        }
1094    }
1095
1096    #[tokio::test]
1097    async fn a_legacy_declaration_is_answered_unstripped() {
1098        // The other side of that no-op: a pre-2026 revision is a served
1099        // one, so the probe is answered — and the answer is byte-identical
1100        // to the 2026-07-28 one, `resultType` included.
1101        let mut session = Session::open();
1102        let legacy = session.call(discover(1, meta("2025-11-25"))).await;
1103        let modern = session.call(discover(2, meta(SERVED))).await;
1104        assert_eq!(legacy["result"], modern["result"], "{legacy}");
1105        assert_eq!(legacy["result"]["resultType"], "complete", "{legacy}");
1106    }
1107
1108    #[tokio::test]
1109    async fn the_refusal_lists_exactly_the_revisions_this_build_serves() {
1110        // What the Go SDK's retry loop reads to pick its second attempt.
1111        // A list narrower than the served set makes a client give up on a
1112        // revision that works; a wider one sends it back with a revision
1113        // that does not.
1114        let mut session = Session::open();
1115        let probe = session.call(discover(1, meta(UNSERVED))).await;
1116        assert_eq!(error_code(&probe), -32022, "{probe}");
1117        assert_eq!(probe["error"]["data"]["requested"], UNSERVED, "{probe}");
1118        assert_eq!(
1119            probe["error"]["data"]["supported"],
1120            serde_json::to_value(SUPPORTED_PROTOCOL_VERSIONS).expect("versions serialize"),
1121            "{probe}"
1122        );
1123    }
1124
1125    /// A [`Transport`] over queued frames, so the wrapper's own dispatch is
1126    /// testable without a service behind it.
1127    ///
1128    /// [`Transport`]: rmcp::transport::Transport
1129    struct Queued {
1130        inbound: std::collections::VecDeque<ClientJsonRpcMessage>,
1131        sent: Arc<std::sync::Mutex<Vec<ServerJsonRpcMessage>>>,
1132        /// Fails every send, standing in for a peer that closed its end.
1133        broken: bool,
1134    }
1135
1136    #[derive(Debug, thiserror::Error)]
1137    #[error("send refused")]
1138    struct SendRefused;
1139
1140    impl rmcp::transport::Transport<RoleServer> for Queued {
1141        type Error = SendRefused;
1142
1143        fn send(
1144            &mut self,
1145            item: ServerJsonRpcMessage,
1146        ) -> impl std::future::Future<Output = Result<(), Self::Error>> + Send + 'static {
1147            let broken = self.broken;
1148            let sent = self.sent.clone();
1149            async move {
1150                sent.lock().expect("sent lock poisoned").push(item);
1151                if broken {
1152                    Err(SendRefused)
1153                } else {
1154                    Ok(())
1155                }
1156            }
1157        }
1158
1159        async fn receive(&mut self) -> Option<ClientJsonRpcMessage> {
1160            self.inbound.pop_front()
1161        }
1162
1163        async fn close(&mut self) -> Result<(), Self::Error> {
1164            self.inbound.clear();
1165            Ok(())
1166        }
1167    }
1168
1169    /// Parse `frames` as inbound client messages and wrap them.
1170    fn queued(frames: &[Value], broken: bool) -> DiscoverAnswering<Queued> {
1171        let sent = Arc::new(std::sync::Mutex::new(Vec::new()));
1172        DiscoverAnswering::new(
1173            Queued {
1174                inbound: frames
1175                    .iter()
1176                    .map(|frame| {
1177                        serde_json::from_value(frame.clone())
1178                            .unwrap_or_else(|e| panic!("{frame}: {e}"))
1179                    })
1180                    .collect(),
1181                sent,
1182                broken,
1183            },
1184            test_server(),
1185        )
1186    }
1187
1188    /// What the inner transport was asked to write.
1189    fn written(wrapper: &DiscoverAnswering<Queued>) -> Vec<Value> {
1190        wrapper
1191            .inner
1192            .sent
1193            .lock()
1194            .expect("sent lock poisoned")
1195            .iter()
1196            .map(|message| serde_json::to_value(message).expect("a reply serializes"))
1197            .collect()
1198    }
1199
1200    #[tokio::test]
1201    async fn every_frame_but_a_discover_request_reaches_rmcp() {
1202        // The pass-through half, at the one place it can be observed
1203        // directly: a notification, a response, an error and an ordinary
1204        // request all come back out of `receive`, and none of them makes
1205        // the wrapper write anything.
1206        let passed = [
1207            json!({ "jsonrpc": "2.0", "method": "notifications/initialized" }),
1208            json!({ "jsonrpc": "2.0", "id": 1, "result": {} }),
1209            json!({ "jsonrpc": "2.0", "id": 2, "error": { "code": -1, "message": "no" } }),
1210            json!({ "jsonrpc": "2.0", "id": 3, "method": "ping" }),
1211            json!({ "jsonrpc": "2.0", "id": 4, "method": "tools/list", "params": {} }),
1212        ];
1213        let mut wrapper = queued(&passed, false);
1214        for frame in &passed {
1215            let received = wrapper
1216                .receive()
1217                .await
1218                .unwrap_or_else(|| panic!("{frame} must reach rmcp"));
1219            assert_eq!(
1220                serde_json::to_value(&received).expect("a frame serializes"),
1221                *frame
1222            );
1223        }
1224        assert!(wrapper.receive().await.is_none(), "the queue is drained");
1225        assert!(written(&wrapper).is_empty(), "{:?}", written(&wrapper));
1226    }
1227
1228    #[tokio::test]
1229    async fn a_discover_is_answered_and_the_next_frame_is_returned() {
1230        // The interception half: the probe never leaves `receive`, its
1231        // answer is written to the inner transport, and the frame after it
1232        // is what rmcp gets — so rmcp learns the lifecycle from THAT one.
1233        let ping = json!({ "jsonrpc": "2.0", "id": 2, "method": "ping" });
1234        let mut wrapper = queued(&[discover(1, meta(SERVED)), ping.clone()], false);
1235        let received = wrapper.receive().await.expect("the ping must come through");
1236        assert_eq!(
1237            serde_json::to_value(&received).expect("a frame serializes"),
1238            ping
1239        );
1240        let written = written(&wrapper);
1241        assert_eq!(written.len(), 1, "{written:?}");
1242        assert_eq!(written[0]["id"], json!(1), "{written:?}");
1243        assert_eq!(
1244            written[0]["result"]["resultType"], "complete",
1245            "{written:?}"
1246        );
1247    }
1248
1249    #[tokio::test]
1250    async fn a_reply_that_cannot_be_written_closes_the_transport() {
1251        // A wrapper that ignored the send error would loop on to the next
1252        // frame and hand rmcp a session whose peer is already gone, with
1253        // the probe silently unanswered. `None` is rmcp's own reading of a
1254        // dead transport.
1255        let mut wrapper = queued(
1256            &[
1257                discover(1, meta(SERVED)),
1258                json!({ "jsonrpc": "2.0", "id": 2, "method": "ping" }),
1259            ],
1260            true,
1261        );
1262        assert!(
1263            wrapper.receive().await.is_none(),
1264            "a failed reply must close the transport, not skip to the next frame"
1265        );
1266    }
1267
1268    #[tokio::test]
1269    async fn send_and_close_are_the_inner_transport_s() {
1270        // The two delegating methods: `send` must reach the inner
1271        // transport (and carry its error back), `close` must reach it too.
1272        let mut wrapper = queued(
1273            &[json!({ "jsonrpc": "2.0", "id": 1, "method": "ping" })],
1274            false,
1275        );
1276        let pong = ServerJsonRpcMessage::response(
1277            rmcp::model::ServerResult::EmptyResult(rmcp::model::EmptyObject {}),
1278            rmcp::model::RequestId::Number(1),
1279        );
1280        wrapper.send(pong).await.expect("the inner send must run");
1281        assert_eq!(written(&wrapper).len(), 1);
1282        wrapper.close().await.expect("the inner close must run");
1283        assert!(
1284            wrapper.receive().await.is_none(),
1285            "close must reach the inner transport"
1286        );
1287    }
1288
1289    #[test]
1290    fn a_discover_frame_is_recognised_by_its_method() {
1291        // The guard the interception arm turns on, spelled out: rmcp's
1292        // deserializer routes `server/discover` to `DiscoverRequest` and
1293        // everything else elsewhere, so `matches!` on that variant is a
1294        // method test. If a bump renamed the method, this fails here
1295        // rather than by every probe silently reaching rmcp again.
1296        let frame: ClientJsonRpcMessage =
1297            serde_json::from_value(discover(1, meta(SERVED))).expect("the probe must parse");
1298        let ClientJsonRpcMessage::Request(request) = frame else {
1299            panic!("a probe is a request");
1300        };
1301        assert!(matches!(request.request, ClientRequest::DiscoverRequest(_)));
1302        assert_eq!(request.request.method(), "server/discover");
1303    }
1304
1305    #[test]
1306    fn the_wrapper_reads_the_list_the_handler_advertises() {
1307        // The wrapper answers from the constant while rmcp's default
1308        // handler answers from `supported_protocol_versions()`. They are
1309        // the same list today; if they ever diverge the probe and the
1310        // handshake would tell a client two different things.
1311        assert_eq!(
1312            test_server().supported_protocol_versions().as_ref(),
1313            SUPPORTED_PROTOCOL_VERSIONS
1314        );
1315        assert!(SUPPORTED_PROTOCOL_VERSIONS.contains(&ProtocolVersion::V_2026_07_28));
1316        assert!(!SUPPORTED_PROTOCOL_VERSIONS
1317            .iter()
1318            .any(|version| version.as_str() == UNSERVED));
1319    }
1320}