| 1 | //! HTTP MCP transport. |
| 2 | //! |
| 3 | //! Speaks Streamable HTTP first and falls back to the legacy SSE endpoint |
| 4 | //! when the server rejects the newer protocol. Request-time authentication |
| 5 | //! and egress authority belong to the shared McpHttpClient session. |
| 6 | |
| 7 | use std::time::Duration; |
| 8 | |
| 9 | use anyhow::Result; |
| 10 | |
| 11 | use super::McpTransport; |
| 12 | use super::http_client::McpHttpClient; |
| 13 | use super::sse::SseTransport; |
| 14 | use super::streamable_http::{StreamableHttpTransport, StreamableSendError}; |
| 15 | use super::wire::McpSessionRejected; |
| 16 | pub(super) struct HttpTransport { |
| 17 | mode: HttpTransportMode, |
| 18 | client: McpHttpClient, |
| 19 | base_url: String, |
| 20 | cancel_token: tokio_util::sync::CancellationToken, |
| 21 | endpoint_timeout: Duration, |
| 22 | } |
| 23 | |
| 24 | enum HttpTransportMode { |
| 25 | Streamable(StreamableHttpTransport), |
| 26 | Sse(SseTransport), |
| 27 | } |
| 28 | |
| 29 | impl HttpTransport { |
| 30 | pub(super) fn new( |
| 31 | client: McpHttpClient, |
| 32 | url: String, |
| 33 | cancel_token: tokio_util::sync::CancellationToken, |
| 34 | endpoint_timeout: Duration, |
| 35 | ) -> Self { |
| 36 | Self { |
| 37 | mode: HttpTransportMode::Streamable(StreamableHttpTransport::new( |
| 38 | client.clone(), |
| 39 | url.clone(), |
| 40 | )), |
| 41 | client, |
| 42 | base_url: url, |
| 43 | cancel_token, |
| 44 | endpoint_timeout, |
| 45 | } |
| 46 | } |
| 47 | |
| 48 | async fn switch_to_sse_and_send(&mut self, msg: Vec<u8>) -> Result<()> { |
| 49 | let mut sse = SseTransport::connect( |
| 50 | self.client.clone(), |
| 51 | self.base_url.clone(), |
| 52 | self.cancel_token.clone(), |
| 53 | self.endpoint_timeout, |
| 54 | ) |
| 55 | .await?; |
| 56 | sse.send(msg).await?; |
| 57 | self.mode = HttpTransportMode::Sse(sse); |
| 58 | Ok(()) |
| 59 | } |
| 60 | |
| 61 | /// Best-effort session-establishment GET preflight. |
| 62 | /// |
| 63 | /// Per the Streamable HTTP spec, the server may return an |
| 64 | /// `Mcp-Session-Id` header on the `initialize` response (the normal |
| 65 | /// path handled inside [`StreamableHttpTransport::send`] above). |
| 66 | /// However some servers (e.g. Hindsight, #1629) **require** a session |
| 67 | /// ID on every POST including `initialize`, creating a chicken-and-egg |
| 68 | /// problem. For those servers we send a short-lived GET before the |
| 69 | /// first POST: if the server returns a session ID in the GET response |
| 70 | /// it will be captured by the header-reading code in |
| 71 | /// [`StreamableHttpTransport::send`] just as if it came from a POST |
| 72 | /// response. |
| 73 | /// |
| 74 | /// This is intentionally best-effort: |
| 75 | /// * The GET uses a tight per-request inner timeout so it never |
| 76 | /// blocks connection startup for long. |
| 77 | /// * If the server doesn't support GET (405, 404, …) we log a debug |
| 78 | /// line and move on — the `initialize` POST will proceed without a |
| 79 | /// session ID. |
| 80 | /// * If the server opens an SSE stream in response (the GET from old |
| 81 | /// SSE transport), we read only the headers, then discard the body |
| 82 | /// so the SSE stream is torn down. The actual SSE path uses a |
| 83 | /// dedicated `SseTransport` and is triggered by the incompatible- |
| 84 | /// status fallback in [`HttpTransport::send`]. |
| 85 | pub(super) async fn try_establish_session(&mut self) -> Result<()> { |
| 86 | let cancel = self.cancel_token.clone(); |
| 87 | let transport = match &mut self.mode { |
| 88 | HttpTransportMode::Streamable(t) => t, |
| 89 | // Already on SSE — session is implicit via the long-lived GET. |
| 90 | HttpTransportMode::Sse(_) => return Ok(()), |
| 91 | }; |
| 92 | |
| 93 | let request = tokio::select! { |
| 94 | biased; |
| 95 | _ = cancel.cancelled() => { |
| 96 | anyhow::bail!("MCP session preflight cancelled after plugin authority changed") |
| 97 | } |
| 98 | request = transport.client.prepare_mcp_request( |
| 99 | transport.client.get(&transport.url), false, |
| 100 | ) => request?, |
| 101 | }; |
| 102 | let response = tokio::select! { |
| 103 | biased; |
| 104 | _ = cancel.cancelled() => { |
| 105 | anyhow::bail!("MCP session preflight cancelled after plugin authority changed") |
| 106 | } |
| 107 | response = tokio::time::timeout(Duration::from_secs(5), transport.client.send(request)) => { |
| 108 | response |
| 109 | .map_err(|_| anyhow::anyhow!("GET timeout"))? |
| 110 | .map_err(|e| anyhow::anyhow!("GET error: {e}"))? |
| 111 | } |
| 112 | }; |
| 113 | |
| 114 | // Capture session ID from the GET response so subsequent POSTs |
| 115 | // (including `initialize`) can include it. This is the same |
| 116 | // header-reading logic that would be hit inside |
| 117 | // `StreamableHttpTransport::send` for POST responses, but since |
| 118 | // the GET is sent before any POST we do it here directly. |
| 119 | if let Some(sid) = response |
| 120 | .headers() |
| 121 | .get("Mcp-Session-Id") |
| 122 | .and_then(|v| v.to_str().ok()) |
| 123 | && transport.session_id.as_deref() != Some(sid) |
| 124 | { |
| 125 | let session_ref = crate::utils::redacted_identifier_for_log(sid); |
| 126 | tracing::debug!(target: "mcp", session = %session_ref, "captured MCP session ID via GET preflight"); |
| 127 | transport.session_id = Some(sid.to_string()); |
| 128 | } |
| 129 | |
| 130 | // We only care about the response headers — discard the body. |
| 131 | // If the server opened an SSE stream in response (some servers |
| 132 | // do this on GET), it will be torn down when response is dropped. |
| 133 | drop(response); |
| 134 | |
| 135 | Ok(()) |
| 136 | } |
| 137 | } |
| 138 | |
| 139 | #[async_trait::async_trait] |
| 140 | impl McpTransport for HttpTransport { |
| 141 | fn set_protocol_version(&mut self, version: &str) { |
| 142 | // Only Streamable HTTP carries the MCP-Protocol-Version header; the |
| 143 | // legacy SSE transport predates it and ignores the negotiation result. |
| 144 | if let HttpTransportMode::Streamable(transport) = &mut self.mode { |
| 145 | transport.set_protocol_version(version); |
| 146 | } |
| 147 | } |
| 148 | |
| 149 | async fn send(&mut self, msg: Vec<u8>) -> Result<()> { |
| 150 | match &mut self.mode { |
| 151 | HttpTransportMode::Streamable(transport) => match transport.send(msg.clone()).await { |
| 152 | Ok(()) => Ok(()), |
| 153 | Err(StreamableSendError::Incompatible(detail)) => { |
| 154 | tracing::debug!( |
| 155 | "MCP Streamable HTTP unavailable; falling back to SSE endpoint discovery: {}", |
| 156 | detail |
| 157 | ); |
| 158 | self.switch_to_sse_and_send(msg).await |
| 159 | } |
| 160 | Err(StreamableSendError::StaleSession(detail)) => { |
| 161 | if let HttpTransportMode::Streamable(transport) = &mut self.mode { |
| 162 | tracing::debug!( |
| 163 | target: "mcp", |
| 164 | error = %detail, |
| 165 | "MCP Streamable HTTP session expired; clearing cached session ID" |
| 166 | ); |
| 167 | transport.session_id = None; |
| 168 | } |
| 169 | Err(McpSessionRejected(format!( |
| 170 | "MCP Streamable HTTP session expired; retry with a new session required ({detail})" |
| 171 | )) |
| 172 | .into()) |
| 173 | } |
| 174 | Err(StreamableSendError::Other(err)) => Err(err), |
| 175 | }, |
| 176 | HttpTransportMode::Sse(transport) => transport.send(msg).await, |
| 177 | } |
| 178 | } |
| 179 | |
| 180 | async fn recv(&mut self) -> Result<Vec<u8>> { |
| 181 | match &mut self.mode { |
| 182 | HttpTransportMode::Streamable(transport) => transport.recv().await, |
| 183 | HttpTransportMode::Sse(transport) => transport.recv().await, |
| 184 | } |
| 185 | } |
| 186 | |
| 187 | fn probe_dead(&self) -> bool { |
| 188 | match &self.mode { |
| 189 | HttpTransportMode::Streamable(_) => false, |
| 190 | HttpTransportMode::Sse(transport) => transport.probe_dead(), |
| 191 | } |
| 192 | } |
| 193 | |
| 194 | async fn shutdown(&mut self) { |
| 195 | if let HttpTransportMode::Sse(transport) = &mut self.mode { |
| 196 | transport.shutdown().await; |
| 197 | } |
| 198 | } |
| 199 | } |
| 200 |