| 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: ¶ms.ticket, |
| 157 | kind: TicketKind::McpLaunch, |
| 158 | tier: HostTier::Builtin, |
| 159 | host_generation: generation, |
| 160 | owner: ¶ms.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(¶ms.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, ¶ms.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(¶ms.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, ¶ms.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, ¶ms.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, ¶ms.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 | ¶ms.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 |