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}