返回 CodeWhale
mcp_http.rs
根目录 / crates / tui / src / extension_host / mcp_http.rs
1 //! The official SDK's Fetch implementation; all credentials and network
2 //! authority stay in the existing Rust McpHttpClient. URLs on the wire are
3 //! opaque session selectors, never configured URL/query/credential material.
4 use super::*;
5 use crate::mcp::http_client::McpHttpClient;
6 use crate::mcp::wire::{
7 MAX_SSE_FRAME_BYTES, McpSessionRejected, find_sse_event_separator_bytes,
8 is_streamable_http_incompatible_status, is_streamable_http_stale_session_status,
9 resolve_sse_endpoint_url, sse_field_value,
10 };
11 use std::sync::atomic::{AtomicBool, Ordering};
12
13 const MAX_RESPONSES: usize = 4;
14 const CHUNK: usize = 32 * 1024;
15 pub(super) struct HttpSession {
16 client: McpHttpClient,
17 url: String,
18 legacy: AtomicBool,
19 started: Mutex<bool>,
20 get_opened: Mutex<bool>,
21 endpoint: Mutex<Option<String>>,
22 server_session: Mutex<Option<String>>,
23 responses: Mutex<HashMap<String, Arc<AsyncMutex<Body>>>>,
24 slots: Arc<std::sync::atomic::AtomicUsize>,
25 failure: Mutex<Option<HttpFailure>>,
26 }
27 impl HttpSession {
28 pub(super) fn new(client: McpHttpClient, url: &str, legacy: bool) -> Self {
29 Self {
30 client,
31 url: url.to_string(),
32 legacy: AtomicBool::new(legacy),
33 started: Mutex::new(false),
34 get_opened: Mutex::new(false),
35 endpoint: Mutex::new(None),
36 server_session: Mutex::new(None),
37 responses: Mutex::new(HashMap::new()),
38 slots: Arc::new(std::sync::atomic::AtomicUsize::new(0)),
39 failure: Mutex::new(None),
40 }
41 }
42 /// The same best-effort preflight as the Rust default. It runs in the
43 /// Rust connection factory, before the SDK initialize deadline begins.
44 /// A slow GET therefore cannot consume a healthy handshake's whole budget.
45 pub(super) async fn preflight(&self, session: &Session, shared: &ManagerShared) -> Result<()> {
46 if self.legacy() {
47 return Ok(());
48 }
49 let run = async {
50 let request = self.client.send_mcp_request(
51 self.client.get(&self.url),
52 false,
53 false,
54 false,
55 || session.validate(shared, session.host_generation, &session.owner),
56 |_| Ok(()),
57 );
58 let response = tokio::time::timeout(Duration::from_secs(5), request).await??;
59 session.validate(shared, session.host_generation, &session.owner)?;
60 if let Some(sid) = response
61 .headers()
62 .get("mcp-session-id")
63 .and_then(|value| value.to_str().ok())
64 {
65 if sid.len() > 8192 {
66 bail!("MCP response framing header exceeds bound");
67 }
68 *self.server_session.lock().expect("MCP session lock") = Some(sid.to_string());
69 }
70 // Observe only headers. Dropping the response closes quiet SSE.
71 Ok::<(), anyhow::Error>(())
72 };
73 tokio::select! { biased;
74 _ = session.cancel.cancelled() => bail!("MCP preflight revoked"),
75 _ = run => {}, // ordinary preflight failures are best effort
76 }
77 session.validate(shared, session.host_generation, &session.owner)
78 }
79 pub(super) fn close(&self) {
80 self.responses.lock().expect("MCP response lock").clear();
81 }
82 pub(super) fn failure(&self) -> Option<anyhow::Error> {
83 self.failure
84 .lock()
85 .expect("MCP HTTP failure lock")
86 .as_ref()
87 .map(|failure| match failure {
88 HttpFailure::Detail(detail) => anyhow::Error::msg(detail.clone()),
89 HttpFailure::Stale(detail) => McpSessionRejected(detail.clone()).into(),
90 })
91 }
92 fn legacy(&self) -> bool {
93 self.legacy.load(Ordering::SeqCst)
94 }
95 }
96 #[derive(Clone)]
97 enum HttpFailure {
98 Detail(String),
99 Stale(String),
100 }
101 struct Permit(Arc<std::sync::atomic::AtomicUsize>);
102 impl Drop for Permit {
103 fn drop(&mut self) {
104 self.0.fetch_sub(1, std::sync::atomic::Ordering::SeqCst);
105 }
106 }
107 struct Body {
108 _permit: Permit,
109 response: Option<reqwest::Response>,
110 sse: bool,
111 buffer: Vec<u8>,
112 ready: Vec<u8>,
113 offset: usize,
114 finished: bool,
115 }
116 fn opaque(session: &Session) -> String {
117 format!("https://mcp-proxy.invalid/{}", session.session_id)
118 }
119
120 pub(super) async fn serve(
121 broker: &Broker,
122 shared: &ManagerShared,
123 generation: u64,
124 request: HostRequest,
125 cx: &HostRequestContext,
126 ) -> Result<Value> {
127 let (id, owner) = match &request {
128 HostRequest::NetStart(p) => (&p.session_id, &p.owner),
129 HostRequest::NetFetch(p) => (&p.session_id, &p.owner),
130 HostRequest::NetRead(p) | HostRequest::NetRelease(p) => (&p.session_id, &p.owner),
131 HostRequest::NetClose(p) => (&p.session_id, &p.owner),
132 _ => bail!("not a FetchProxy request"),
133 };
134 let session = broker.session(id)?;
135 if matches!(
136 &request,
137 HostRequest::NetClose(_) | HostRequest::NetRelease(_)
138 ) {
139 if &session.owner != owner || session.host_generation != generation {
140 bail!("MCP HTTP cleanup has wrong owner");
141 }
142 } else {
143 session.validate(shared, generation, owner)?;
144 }
145 let http = session.http.as_ref().context("not an HTTP session")?;
146 match request {
147 HostRequest::NetStart(params) => {
148 let mut guard = CancelGuard {
149 cancel: session.cancel.clone(),
150 armed: true,
151 };
152 shared
153 .core_calls
154 .tickets
155 .redeem(&Presented {
156 ticket: &params.ticket,
157 kind: TicketKind::McpLaunch,
158 tier: HostTier::Builtin,
159 host_generation: generation,
160 owner: &params.owner,
161 method: "net/start",
162 target: Some(&json!({"session_id": params.session_id})),
163 })
164 .map_err(|bad| {
165 if bad.violation {
166 cx.violation("too many invalid MCP HTTP tickets".into());
167 }
168 anyhow::anyhow!("MCP HTTP start grant refused")
169 })?;
170 session
171 .launch_ticket
172 .lock()
173 .expect("MCP launch ticket lock")
174 .take();
175 {
176 let mut started = http.started.lock().expect("MCP HTTP state lock");
177 if *started || cx.cancel.is_cancelled() {
178 bail!("MCP HTTP start cancelled or replayed");
179 }
180 *started = true;
181 }
182 guard.armed = false;
183 Ok(json!({}))
184 }
185 HostRequest::NetFetch(params) => {
186 fetch(broker, shared, generation, &session, http, params, cx).await
187 }
188 HostRequest::NetRead(params) => {
189 let mut guard = CancelGuard {
190 cancel: session.cancel.clone(),
191 armed: true,
192 };
193 let body = http
194 .responses
195 .lock()
196 .expect("MCP response lock")
197 .get(&params.response_id)
198 .cloned()
199 .context("MCP HTTP response is closed")?;
200 let mut body = tokio::select! { biased; _ = cx.cancel.cancelled() => bail!("MCP HTTP read cancelled"), _ = session.cancel.cancelled() => bail!("MCP HTTP withdrawn"), body = body.lock() => body };
201 let run = read(&session, http, &mut body);
202 let result = tokio::select! { biased; _ = cx.cancel.cancelled() => bail!("MCP HTTP read cancelled"), _ = session.cancel.cancelled() => bail!("MCP HTTP withdrawn"), result = run => result };
203 if result.is_err() {
204 session.cancel();
205 }
206 let (data, done) = result?;
207 session.validate(shared, generation, &params.owner)?;
208 guard.armed = false;
209 Ok(json!({"data":data,"done":done}))
210 }
211 HostRequest::NetRelease(params) => {
212 http.responses
213 .lock()
214 .expect("MCP response lock")
215 .remove(&params.response_id);
216 Ok(json!({}))
217 }
218 HostRequest::NetClose(_) => {
219 http.close();
220 broker.remove(shared, &session.session_id);
221 Ok(json!({}))
222 }
223 _ => bail!("not a FetchProxy request"),
224 }
225 }
226 async fn fetch(
227 broker: &Broker,
228 shared: &ManagerShared,
229 generation: u64,
230 session: &Arc<Session>,
231 http: &HttpSession,
232 params: NetFetchParams,
233 cx: &HostRequestContext,
234 ) -> Result<Value> {
235 let mut guard = CancelGuard {
236 cancel: session.cancel.clone(),
237 armed: true,
238 };
239 if !*http.started.lock().expect("MCP HTTP state lock")
240 || http.responses.lock().expect("MCP response lock").len() >= MAX_RESPONSES
241 {
242 bail!("MCP HTTP response bound or stale start");
243 }
244 let slots = Arc::clone(&http.slots);
245 // Match the existing admission loops: try_update exceeds our Rust MSRV.
246 let mut count = slots.load(std::sync::atomic::Ordering::SeqCst);
247 loop {
248 if count >= MAX_RESPONSES {
249 bail!("MCP HTTP active response cap exceeded");
250 }
251 match slots.compare_exchange_weak(
252 count,
253 count + 1,
254 std::sync::atomic::Ordering::SeqCst,
255 std::sync::atomic::Ordering::SeqCst,
256 ) {
257 Ok(_) => break,
258 Err(current) => count = current,
259 }
260 }
261 let permit = Permit(slots);
262 let base = opaque(session);
263 let endpoint_selector = format!("{base}/endpoint");
264 let url = if params.url == base {
265 http.url.clone()
266 } else if http.legacy() && params.url == endpoint_selector {
267 http.endpoint
268 .lock()
269 .expect("MCP endpoint lock")
270 .clone()
271 .context("SSE endpoint not admitted")?
272 } else {
273 bail!("MCP HTTP URL has no Rust selector");
274 };
275 let headers = [
276 ("accept", params.headers.accept.as_deref()),
277 ("content-type", params.headers.content_type.as_deref()),
278 ("mcp-session-id", params.headers.mcp_session_id.as_deref()),
279 (
280 "mcp-protocol-version",
281 params.headers.mcp_protocol_version.as_deref(),
282 ),
283 ];
284 if headers
285 .iter()
286 .any(|(_, value)| value.is_some_and(|value| value.len() > 8192))
287 {
288 bail!("MCP HTTP framing header exceeds bound");
289 }
290 if let Some(value) = params.headers.mcp_session_id.as_deref()
291 && http
292 .server_session
293 .lock()
294 .expect("MCP session lock")
295 .as_deref()
296 != Some(value)
297 {
298 bail!("MCP server session ID not observed");
299 }
300 if let Some(value) = params.headers.mcp_protocol_version.as_deref()
301 && !crate::mcp::MCP_CLIENT_ACCEPTED_PROTOCOL_VERSIONS.contains(&value)
302 {
303 bail!("MCP protocol header refused");
304 }
305 let mut deadline = Instant::now() + Duration::from_secs(30);
306 let original = params.operation_id.as_ref().and_then(|id| {
307 session
308 .operations
309 .lock()
310 .expect("MCP operation lock")
311 .get(id)
312 .cloned()
313 });
314 let request = match params.method.as_str() {
315 "POST" => {
316 let write = ProcWriteParams {
317 owner: params.owner.clone(),
318 session_id: params.session_id.clone(),
319 frame: params.frame.clone().context("MCP HTTP frame absent")?,
320 ticket: params.ticket.clone(),
321 operation_id: params.operation_id.clone(),
322 };
323 let (frame, expires) =
324 broker.authorize_frame(shared, generation, session, write, cx, "net/fetch")?;
325 deadline = expires;
326 http.client.post(&url).body(serde_json::to_vec(&frame)?)
327 }
328 "GET" => {
329 if params.frame.is_some()
330 || params.ticket.is_some()
331 || params.operation_id.is_some()
332 || params.url != base
333 {
334 bail!("MCP HTTP GET is not a channel open");
335 }
336 let mut opened = http.get_opened.lock().expect("MCP GET state lock");
337 if *opened {
338 bail!("MCP HTTP channel cannot reconnect or replay");
339 }
340 *opened = true;
341 http.client.get(&url)
342 }
343 _ => bail!("MCP HTTP method has no authority"),
344 };
345 // A preflight session header is an observed Rust value. The SDK removes
346 // it from initialize by spec, so apply it here for strict legacy servers.
347 let mut request = request;
348 for (name, value) in headers {
349 if let Some(value) = value {
350 request = request.header(name, value);
351 }
352 }
353 if params.method == "POST"
354 && !http.legacy()
355 && let Some(sid) = http
356 .server_session
357 .lock()
358 .expect("MCP session lock")
359 .as_ref()
360 {
361 request = request.header("mcp-session-id", sid);
362 }
363 let send = http.client.send_mcp_request(
364 request,
365 params.method == "POST",
366 params.method == "GET",
367 params.method == "POST" && !http.legacy(),
368 || {
369 session.validate(shared, generation, &params.owner)?;
370 if cx.cancel.is_cancelled() || Instant::now() >= deadline {
371 bail!("MCP HTTP operation expired or cancelled");
372 }
373 Ok(())
374 },
375 |response| {
376 if let Some(sid) = response
377 .headers()
378 .get("mcp-session-id")
379 .and_then(|v| v.to_str().ok())
380 {
381 if sid.len() > 8192 {
382 bail!("MCP response framing header exceeds bound");
383 }
384 *http.server_session.lock().expect("MCP session lock") = Some(sid.to_string());
385 }
386 Ok(())
387 },
388 );
389 let result = tokio::select! { biased;
390 _ = cx.cancel.cancelled() => bail!("MCP HTTP cancelled"),
391 _ = session.cancel.cancelled() => bail!("MCP HTTP revoked"),
392 result = tokio::time::timeout_at(tokio::time::Instant::from_std(deadline), send) => result.context("MCP HTTP operation expired")?,
393 };
394 let response = match result {
395 Ok(response) => response,
396 Err(error) => {
397 // Retain a Rust-only diagnostic; net/* RPC never exposes configured
398 // origins, OAuth metadata, credentials or reflected provider text.
399 *http.failure.lock().expect("MCP HTTP failure lock") =
400 Some(HttpFailure::Detail(format!("{error:#}")));
401 return Err(error);
402 }
403 };
404 session.validate(shared, generation, &params.owner)?;
405 let status = response.status().as_u16();
406 let mut headers = HashMap::new();
407 for name in ["content-type", "mcp-session-id"] {
408 if let Some(value) = response
409 .headers()
410 .get(name)
411 .and_then(|value| value.to_str().ok())
412 {
413 if value.len() > 8192 {
414 bail!("MCP response framing header exceeds bound");
415 }
416 headers.insert(name.to_string(), value.to_string());
417 }
418 }
419 if let Some(sid) = headers.get("mcp-session-id") {
420 *http.server_session.lock().expect("MCP session lock") = Some(sid.clone());
421 }
422 if params.method == "POST"
423 && !http.legacy()
424 && !headers.contains_key("mcp-session-id")
425 && let Some(sid) = http
426 .server_session
427 .lock()
428 .expect("MCP session lock")
429 .as_ref()
430 {
431 headers.insert("mcp-session-id".to_string(), sid.clone());
432 }
433 let sse = response
434 .headers()
435 .get("content-type")
436 .and_then(|v| v.to_str().ok())
437 .is_some_and(|v| {
438 v.split(';')
439 .next()
440 .unwrap_or("")
441 .trim()
442 .eq_ignore_ascii_case("text/event-stream")
443 });
444 if !sse
445 && response
446 .content_length()
447 .is_some_and(|size| size > MAX_MCP_RESPONSE_BYTES as u64)
448 {
449 bail!("MCP HTTP body exceeds bound");
450 }
451 // Only a parsed, explicit transport refusal admits a fresh exact grant.
452 // Never negotiate on auth, cancellation, partial writes, malformed replies
453 // or JSON-RPC errors. A stale session is a typed Rust refusal instead.
454 if params.method == "POST" && !(200..300).contains(&status) {
455 let had_session = http
456 .server_session
457 .lock()
458 .expect("MCP session lock")
459 .is_some();
460 let excerpt = tokio::select! { biased;
461 _ = cx.cancel.cancelled() => bail!("MCP refusal read cancelled"),
462 _ = session.cancel.cancelled() => bail!("MCP refusal read revoked"),
463 excerpt = tokio::time::timeout_at(tokio::time::Instant::from_std(deadline),
464 crate::mcp::bounded_body_excerpt(response, crate::mcp::ERROR_BODY_PREVIEW_BYTES)) => excerpt.context("MCP refusal read expired")?,
465 };
466 session.validate(shared, generation, &params.owner)?;
467 if cx.cancel.is_cancelled() || Instant::now() >= deadline {
468 bail!("MCP refusal expired or cancelled");
469 }
470 let safe_excerpt = http.client.server_error_preview(&excerpt);
471 let detail = format!(
472 "status={} body={safe_excerpt}",
473 reqwest::StatusCode::from_u16(status)?
474 );
475 let stale = if http.legacy() {
476 crate::mcp::wire::is_mcp_stale_session_body(&excerpt)
477 } else {
478 had_session
479 && is_streamable_http_stale_session_status(
480 reqwest::StatusCode::from_u16(status)?,
481 &excerpt,
482 )
483 };
484 if stale {
485 let detail = if http.legacy() {
486 format!(
487 "MCP session expired (transport=sse endpoint={} status={}): {safe_excerpt}",
488 crate::mcp::mask_url_secrets(&url),
489 reqwest::StatusCode::from_u16(status)?
490 )
491 } else {
492 format!(
493 "MCP Streamable HTTP session expired; retry with a new session required ({detail})"
494 )
495 };
496 *http.failure.lock().expect("MCP HTTP failure lock") = Some(HttpFailure::Stale(detail));
497 } else if !http.legacy()
498 && is_streamable_http_incompatible_status(reqwest::StatusCode::from_u16(status)?)
499 {
500 let operation = original.context("transport refusal has no semantic operation")?;
501 let grant = renegotiate(
502 shared,
503 generation,
504 session,
505 http,
506 &params.owner,
507 operation,
508 deadline,
509 )?;
510 guard.armed = false;
511 return Ok(json!({"status":status,"headers":headers,"legacy_grant":grant}));
512 } else {
513 let detail = if http.legacy() {
514 format!(
515 "MCP SSE POST rejected (transport=sse endpoint={} status={}): {safe_excerpt}",
516 crate::mcp::mask_url_secrets(&url),
517 reqwest::StatusCode::from_u16(status)?
518 )
519 } else {
520 format!(
521 "MCP Streamable HTTP rejected (transport=http url={} status={}): {safe_excerpt}",
522 crate::mcp::mask_url_secrets(&url),
523 reqwest::StatusCode::from_u16(status)?
524 )
525 };
526 *http.failure.lock().expect("MCP HTTP failure lock") =
527 Some(HttpFailure::Detail(detail));
528 }
529 guard.armed = false;
530 return Ok(json!({"status":status,"headers":headers}));
531 }
532 // Never expose reflected provider errors or authentication headers to TS.
533 if !(200..300).contains(&status) {
534 guard.armed = false;
535 return Ok(json!({"status":status,"headers":headers}));
536 }
537 if status == 204 || status == 202 {
538 guard.armed = false;
539 return Ok(json!({"status": status,"headers":headers}));
540 }
541 if params.method == "GET" && !sse {
542 bail!("MCP channel response is not SSE");
543 }
544 let response_id = uuid::Uuid::new_v4().to_string();
545 let body = Body {
546 _permit: permit,
547 response: Some(response),
548 sse,
549 buffer: Vec::new(),
550 ready: Vec::new(),
551 offset: 0,
552 finished: false,
553 };
554 http.responses
555 .lock()
556 .expect("MCP response lock")
557 .insert(response_id.clone(), Arc::new(AsyncMutex::new(body)));
558 guard.armed = false;
559 Ok(json!({"response_id":response_id,"status":status,"headers":headers}))
560 }
561 fn renegotiate(
562 shared: &ManagerShared,
563 generation: u64,
564 session: &Session,
565 http: &HttpSession,
566 owner: &OwnerRef,
567 mut operation: Operation,
568 deadline: Instant,
569 ) -> Result<McpOperationGrant> {
570 session.validate(shared, generation, owner)?;
571 let ttl = deadline
572 .checked_duration_since(Instant::now())
573 .context("MCP negotiation expired")?;
574 if ttl.is_zero() || http.legacy.swap(true, Ordering::SeqCst) {
575 bail!("MCP negotiation unavailable");
576 }
577 let operation_id = operation.target["operation_id"]
578 .as_str()
579 .context("MCP operation id absent")?
580 .to_string();
581 let method = operation.target["method"]
582 .as_str()
583 .context("MCP operation method absent")?
584 .to_string();
585 let params = operation.target["params"].clone();
586 let id = operation.target["wire_id"].as_str().map(str::to_string);
587 if let Some(id) = &id {
588 let mut pending = session.pending.lock().expect("MCP pending lock");
589 let key = wire_id(&json!(id))?;
590 if pending
591 .get(&key)
592 .is_none_or(|binding| binding.operation_id != operation_id)
593 {
594 bail!("MCP negotiation target changed");
595 }
596 pending.remove(&key);
597 }
598 let ticket = shared.core_calls.tickets.mint(Grant {
599 kind: TicketKind::McpOperation,
600 tier: HostTier::Builtin,
601 host_generation: generation,
602 owner: owner.clone(),
603 method: "net/fetch",
604 target: operation.target.clone(),
605 ttl,
606 uses: 1,
607 });
608 operation.ticket = ticket.clone();
609 let mut operations = session.operations.lock().expect("MCP operation lock");
610 if operations.len() >= MAX_PENDING || operations.contains_key(&operation_id) {
611 shared.core_calls.tickets.revoke(&ticket);
612 bail!("MCP negotiation operation unavailable");
613 }
614 operations.insert(operation_id.clone(), operation);
615 *http.get_opened.lock().expect("MCP GET state lock") = false;
616 *http.server_session.lock().expect("MCP session lock") = None;
617 http.close();
618 Ok(McpOperationGrant {
619 ticket: ticket.expose().to_string(),
620 operation_id,
621 method,
622 wire_id: id,
623 params,
624 })
625 }
626
627 async fn read(session: &Session, http: &HttpSession, body: &mut Body) -> Result<(Vec<u8>, bool)> {
628 loop {
629 if body.offset < body.ready.len() {
630 let end = (body.offset + CHUNK).min(body.ready.len());
631 let chunk = body.ready[body.offset..end].to_vec();
632 body.offset = end;
633 if body.offset == body.ready.len() {
634 body.ready.clear();
635 body.offset = 0;
636 }
637 return Ok((chunk, body.finished && body.ready.is_empty()));
638 }
639 if body.finished {
640 return Ok((Vec::new(), true));
641 }
642 if body.sse
643 && let Some((at, separator)) = find_sse_event_separator_bytes(&body.buffer)
644 {
645 if at > MAX_SSE_FRAME_BYTES {
646 bail!("MCP SSE event exceeds bound");
647 }
648 let mut event = body.buffer.drain(..at + separator).collect::<Vec<_>>();
649 let text = std::str::from_utf8(&event)?;
650 let mut kind = "message";
651 let mut data = String::new();
652 for line in text.lines() {
653 if let Some(value) = sse_field_value(line, "event:") {
654 kind = value;
655 } else if let Some(value) = sse_field_value(line, "data:") {
656 if !data.is_empty() {
657 data.push('\n');
658 }
659 data.push_str(value);
660 }
661 }
662 if kind == "endpoint" {
663 if !http.legacy() || http.endpoint.lock().expect("MCP endpoint lock").is_some() {
664 bail!("unexpected SSE endpoint");
665 }
666 let endpoint = resolve_sse_endpoint_url(&http.url, &data)?;
667 *http.endpoint.lock().expect("MCP endpoint lock") = Some(endpoint);
668 event =
669 format!("event: endpoint\ndata: {}/endpoint\n\n", opaque(session)).into_bytes();
670 } else if kind == "message" && !data.trim().is_empty() {
671 if http.legacy() && http.endpoint.lock().expect("MCP endpoint lock").is_none() {
672 bail!("SSE message before endpoint");
673 }
674 session.observe(&serde_json::from_str(&data)?)?;
675 }
676 body.ready = event;
677 continue;
678 }
679 let response = body.response.as_mut().context("MCP body is closed")?;
680 match response.chunk().await? {
681 Some(chunk) => {
682 let max = MAX_MCP_RESPONSE_BYTES;
683 if body
684 .buffer
685 .len()
686 .checked_add(chunk.len())
687 .is_none_or(|n| n > max)
688 {
689 bail!("MCP response assembly exceeds bound");
690 }
691 body.buffer.extend_from_slice(&chunk);
692 }
693 None => {
694 body.response.take();
695 body.finished = true;
696 if body.sse {
697 if !body.buffer.is_empty() {
698 bail!("MCP SSE event ended without separator");
699 }
700 } else if !body.buffer.is_empty() {
701 let value: Value = serde_json::from_slice(&body.buffer)?;
702 if let Some(values) = value.as_array() {
703 if values.len() > MAX_PENDING {
704 bail!("MCP batch exceeds bound");
705 }
706 for frame in values {
707 session.observe(frame)?;
708 }
709 } else {
710 session.observe(&value)?;
711 }
712 body.ready = std::mem::take(&mut body.buffer);
713 }
714 }
715 }
716 }
717 }
718 #[cfg(test)]
719 #[path = "mcp_http_tests.rs"]
720 mod tests;
721
721 lines RUST