返回 CodeWhale
http.rs
根目录 / crates / tui / src / mcp / http.rs
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
200 lines RUST