mcp_trace_validator/checks/draft/transport/
stream.rs1use std::collections::BTreeSet;
14
15use serde_json::Value;
16
17use super::super::super::FindingSink;
18use crate::context::TraceContext;
19use mcp_conformance_core::trace::{Direction, EventBody, LifecycleEvent, TransportKind};
20
21#[cfg(test)]
22mod tests;
23
24pub(in crate::checks) fn client_no_responses(context: &TraceContext<'_>, sink: &mut FindingSink) {
34 for (event, _, _) in context.messages() {
35 if event.direction != Direction::ClientToServer {
36 continue;
37 }
38 let Some(payload) = event.message_payload() else {
39 continue;
40 };
41 sink.examined();
42 let is_response = payload.get("id").is_some()
43 && payload.get("method").is_none()
44 && (payload.get("result").is_some() || payload.get("error").is_some());
45 if is_response {
46 sink.push(
47 Some(event.seq),
48 "client sent a JSON-RPC response; 2026-07-28 removed server-initiated \
49 requests, so there is nothing for one to answer"
50 .to_owned(),
51 );
52 }
53 }
54}
55
56pub(in crate::checks) fn no_independent_server_requests(
63 context: &TraceContext<'_>,
64 sink: &mut FindingSink,
65) {
66 for (event, _, _) in context.messages() {
67 if event.direction != Direction::ServerToClient {
68 continue;
69 }
70 let Some(payload) = event.message_payload() else {
71 continue;
72 };
73 sink.examined();
74 if let Some(method) = payload.get("method").and_then(Value::as_str)
75 && payload.get("id").is_some_and(|id| !id.is_null())
76 {
77 sink.push(
78 Some(event.seq),
79 format!(
80 "server sent an independent request `{method}`; 2026-07-28 replaces \
81 server-initiated requests with MRTR input requests"
82 ),
83 );
84 }
85 }
86}
87
88pub(in crate::checks) fn accel_buffering_header(
90 context: &TraceContext<'_>,
91 sink: &mut FindingSink,
92) {
93 for event in context.events() {
94 if event.direction != Direction::ServerToClient {
95 continue;
96 }
97 let EventBody::Http { headers, .. } = &event.body else {
98 continue;
99 };
100 let is_sse = headers
101 .get("content-type")
102 .is_some_and(|value| value.starts_with("text/event-stream"));
103 if !is_sse {
104 continue; }
106 sink.examined();
107 if headers.get("x-accel-buffering").map(String::as_str) != Some("no") {
108 sink.push(
109 Some(event.seq),
110 "SSE response does not carry `X-Accel-Buffering: no`".to_owned(),
111 );
112 }
113 }
114}
115
116pub(in crate::checks) fn no_messages_after_cancellation(
128 context: &TraceContext<'_>,
129 sink: &mut FindingSink,
130) {
131 let mut outstanding: BTreeSet<String> = BTreeSet::new();
132 let mut closed_at: Option<u64> = None;
133 for event in context.events() {
134 if let Some(closed_at) = closed_at {
135 report_after_close(event, &outstanding, closed_at, sink);
136 } else if is_cancellation(event) {
137 closed_at = Some(event.seq);
138 } else {
139 track_outstanding(event, &mut outstanding);
140 }
141 }
142}
143
144fn is_cancellation(event: &mcp_conformance_core::trace::TraceEvent) -> bool {
146 let closed = matches!(
147 event.body,
148 EventBody::Lifecycle {
149 event: LifecycleEvent::TransportClose | LifecycleEvent::TransportAbort
150 }
151 );
152 closed && event.transport == TransportKind::StreamableHttp
153}
154
155fn track_outstanding(
158 event: &mcp_conformance_core::trace::TraceEvent,
159 outstanding: &mut BTreeSet<String>,
160) {
161 let Some(payload) = event.message_payload() else {
162 return;
163 };
164 let Some(id) = payload.get("id").filter(|id| !id.is_null()) else {
165 return;
166 };
167 if payload.get("method").is_some() {
168 outstanding.insert(id.to_string());
169 } else {
170 outstanding.remove(&id.to_string());
171 }
172}
173
174fn report_after_close(
176 event: &mcp_conformance_core::trace::TraceEvent,
177 outstanding: &BTreeSet<String>,
178 closed_at: u64,
179 sink: &mut FindingSink,
180) {
181 if event.direction != Direction::ServerToClient {
182 return;
183 }
184 let Some(id) = event
185 .message_payload()
186 .and_then(|payload| payload.get("id"))
187 .filter(|id| !id.is_null())
188 else {
189 return;
190 };
191 sink.examined();
194 if outstanding.contains(&id.to_string()) {
195 sink.push(
196 Some(event.seq),
197 format!(
198 "server sent a further message for request id {id}, whose response \
199 stream closed at seq {closed_at}; a close is cancellation at this revision"
200 ),
201 );
202 }
203}