From 69c88974d666064cfbd9dc10d8c192f7c6540865 Mon Sep 17 00:00:00 2001 From: Yuanhao Li Date: Fri, 9 Oct 2026 13:59:56 +0200 Subject: [PATCH 1/3] feat(mcp): headers for HTTP MCP servers HttpTransport::with_header(name, value) sends a header on every request of the session (initialize, the initialized notification, tool calls and the closing DELETE) through one request helper. Replace on repeat; values marked sensitive; invalid or reserved names (Accept, Content-Type, Mcp-Session-Id) refused before sending, without echoing the value. McpClient::connect_http_with and Agent::with_mcp_server_http_transport connect through a configured transport; connect_http and with_mcp_server_http are unchanged and now delegate to them. Closes #275. Co-Authored-By: Claude Opus 5.5 (1M context) Claude-Session: https://claude.ai/code/session_01T7iq5hpndSiHQcnAsywKuG --- CHANGELOG.md | 6 +- CLAUDE.md | 2 +- docs/guides/mcp.md | 26 ++++ docs/guides/wasm-workers.md | 6 +- src/agent.rs | 17 ++- src/mcp/client.rs | 8 +- src/mcp/transport.rs | 72 ++++++++++- tests/mcp_http_transport_test.rs | 205 +++++++++++++++++++++++++++++++ 8 files changed, 327 insertions(+), 15 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 0b84770..ef630a1 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -6,13 +6,17 @@ adheres to [Semantic Versioning](https://semver.org/). ## Unreleased +### Added + +- **Headers for HTTP MCP servers** (#275). `HttpTransport::with_header(name, value)` sends a header on every request of the session — `initialize`, the `initialized` notification, tool calls and the closing `DELETE` — so servers that need a token, an API key or a `User-Agent` can be reached. Setting a name again replaces it; values are marked sensitive; an invalid name or value, or one the transport sets itself (`Accept`, `Content-Type`, `Mcp-Session-Id`), is an error before anything is sent. Connect with `McpClient::connect_http_with(transport)` or `Agent::with_mcp_server_http_transport(transport)` (same handshake and timeouts as `connect_http` / `with_mcp_server_http`, which are unchanged). Works on wasm32. Tests: every request of a session carries the headers (DELETE included, mutation-checked), an agent connects through a configured transport, refusals, replacement, and a plain transport sends none. + ### Fixed - **`BashTool` kills what a command started** (#277, suggested by @shahidcodes). On Unix, `bash` now leads its own process group, and a timeout, a cancel or dropping the call kills the whole group: pipeline stages, `&&` lists and background jobs no longer keep running. Only a process that starts its own session (`setsid`) escapes. A command that finishes on its own leaves a background job alone when the job's output is redirected (`cmd >log 2>&1 &`); one still writing to the tool's output keeps the call open until the timeout, which now kills it. **Behaviour change:** the command is no longer in the terminal's foreground group, so the terminal's Ctrl+C does not reach it — cancel or drop the call on Ctrl+C (the `cli` example now aborts the run), and a command prompting on `/dev/tty` (`sudo`, SSH) stops until the timeout. Windows is unchanged (the `bash` process only). Adds `libc` as a Unix-only dependency (already in the tree through tokio). Tests: timeout, cancel, drop and normal completion in `tests/tools_test.rs`, mutation-checked. ### Tests -- **HTTP MCP on wasm32.** `tests/wasm32.rs` runs an agent that connects to an HTTP MCP server, discovers its tool and calls it from the loop, through the host's `fetch` (a scripted server replaces the global `fetch`): handshake, `Mcp-Session-Id` replay and the tool result. The Workers guide notes that a stalled MCP stream has no idle read timeout of its own on wasm32, and that custom headers are not supported yet (#275). +- **HTTP MCP on wasm32.** `tests/wasm32.rs` runs an agent that connects to an HTTP MCP server, discovers its tool and calls it from the loop, through the host's `fetch` (a scripted server replaces the global `fetch`): handshake, `Mcp-Session-Id` replay and the tool result. The Workers guide notes that a stalled MCP stream has no idle read timeout of its own on wasm32. ## 0.25.1 (2026-10-09) diff --git a/CLAUDE.md b/CLAUDE.md index 9bdf1af..329b32d 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -108,7 +108,7 @@ Behind the `openapi` Cargo feature. `OpenApiToolAdapter` parses an OpenAPI 3.0 s ### MCP Integration (`mcp/`) -`McpClient` communicates via `McpTransport` trait (stdio or HTTP; `notify` sends a `JsonRpcNotification`; the trait default = old best-effort request-and-wait (error logged, `Ok`) for custom transports; a built-in transport's failed `initialized` fails the connect; `connect_*` bound the handshake at 60 s). Stdio: reads until the response with the request's id (numeric or string form), skips server notifications, answers server requests (`ping` → `{}`, others → -32601, `warn!`), drops stale ids, skips non-JSON lines with `warn!` (counted into the close error); a null-id error is taken as the current call's, with `warn!`; `StdinWriter.interrupted` starts the next message on a fresh line after a cut-off write; stderr drained lossily in 8 KiB chunks to `debug!` (target `yoagent::mcp::stderr`) with a 4 KiB `StderrTail`; EOF → `McpError::Transport("Connection closed: … exited with …; last stderr: …")` (waits ≤500 ms for the drain); `kill_on_drop`; tested against bash-scripted servers in `tests/mcp_stdio_test.rs` (non-UTF-8 stderr flood, late answer, kill on drop, death with stderr). Call timeout lives on `McpClient` (`call_timeout`: 300 s stdio, None HTTP/custom) and is copied by `McpToolAdapter::from_client*`; `McpToolAdapter::new` has none; the clock starts after the client lock is held (only cancel ends the wait). `McpToolCallResult.content` is lenient: an unmodelled block becomes `McpContent::Text` (a text resource's own text, else JSON with `data`/`blob` over 1 KiB replaced by their size; `warn!` once per type). `McpToolAdapter` wraps MCP tools to implement `AgentTool`, making them transparent to the agent loop. Added via `Agent::with_mcp_server_stdio()` / `with_mcp_server_http()`. +`McpClient` communicates via `McpTransport` trait (stdio or HTTP; `notify` sends a `JsonRpcNotification`; the trait default = old best-effort request-and-wait (error logged, `Ok`) for custom transports; a built-in transport's failed `initialized` fails the connect; `connect_*` bound the handshake at 60 s). Stdio: reads until the response with the request's id (numeric or string form), skips server notifications, answers server requests (`ping` → `{}`, others → -32601, `warn!`), drops stale ids, skips non-JSON lines with `warn!` (counted into the close error); a null-id error is taken as the current call's, with `warn!`; `StdinWriter.interrupted` starts the next message on a fresh line after a cut-off write; stderr drained lossily in 8 KiB chunks to `debug!` (target `yoagent::mcp::stderr`) with a 4 KiB `StderrTail`; EOF → `McpError::Transport("Connection closed: … exited with …; last stderr: …")` (waits ≤500 ms for the drain); `kill_on_drop`; tested against bash-scripted servers in `tests/mcp_stdio_test.rs` (non-UTF-8 stderr flood, late answer, kill on drop, death with stderr). Call timeout lives on `McpClient` (`call_timeout`: 300 s stdio, None HTTP/custom) and is copied by `McpToolAdapter::from_client*`; `McpToolAdapter::new` has none; the clock starts after the client lock is held (only cancel ends the wait). `McpToolCallResult.content` is lenient: an unmodelled block becomes `McpContent::Text` (a text resource's own text, else JSON with `data`/`blob` over 1 KiB replaced by their size; `warn!` once per type). `McpToolAdapter` wraps MCP tools to implement `AgentTool`, making them transparent to the agent loop. Added via `Agent::with_mcp_server_stdio()` / `with_mcp_server_http()`. Caller headers (#275): `HttpTransport::with_header(name, value) -> Result` (stored `HeaderMap`, values `set_sensitive`, insert = replace; invalid name/value or a reserved name — `accept`, `content-type`, `mcp-session-id` — is a `Transport` error that never echoes the value) applied by the private `request(Method)` helper that every path (send, notify, close's DELETE) uses; `McpClient::connect_http_with(transport)` (= `connect_http`'s handshake, `call_timeout: None`) and `Agent::with_mcp_server_http_transport(transport)`; tests in `tests/mcp_http_transport_test.rs` (`mount_server_requiring`). ### GASP Bridge (`gasp.rs`, feature-gated) diff --git a/docs/guides/mcp.md b/docs/guides/mcp.md index 792d77d..e9f30e4 100644 --- a/docs/guides/mcp.md +++ b/docs/guides/mcp.md @@ -83,6 +83,32 @@ let agent = Agent::from_config(ModelConfig::anthropic("claude-sonnet-5", "Claude .await?; ``` +#### Servers that need headers + +A server that wants a token, an API key or a `User-Agent` gets them from +`HttpTransport::with_header`, then connects with +`Agent::with_mcp_server_http_transport` (or `McpClient::connect_http_with`): + +```rust +use yoagent::mcp::HttpTransport; + +let token = std::env::var("MCP_TOKEN")?; +let transport = HttpTransport::new("https://mcp.example.com/mcp")? + .with_header("authorization", format!("Bearer {token}"))? + .with_header("user-agent", "my-app/1.0")?; +let agent = agent.with_mcp_server_http_transport(transport).await?; +``` + +- The headers go on every request of the session: `initialize`, the + `initialized` notification, tool calls and the closing `DELETE`. +- Setting a name again replaces its value. Values are marked sensitive, so they + stay out of debug output, and a refused value is never echoed in the error. +- An invalid name or value is an error from `with_header`, before anything is + sent. So is a header the transport sets itself: `Accept`, `Content-Type`, + `Mcp-Session-Id`. +- The headers are fixed for the transport's lifetime. A token that expires + during a session (OAuth refresh) needs a new transport and client. + `HttpTransport` handles both the plain JSON-RPC-over-POST shape and the **request/response subset of Streamable HTTP** — servers that answer a POST with an SSE-framed response, whether or not they then close the stream: diff --git a/docs/guides/wasm-workers.md b/docs/guides/wasm-workers.md index a4f3686..c5b3071 100644 --- a/docs/guides/wasm-workers.md +++ b/docs/guides/wasm-workers.md @@ -44,9 +44,9 @@ Only `wasm32-unknown-unknown` is supported. WASI targets are not. HTTP MCP goes through the host's `fetch` like the providers: the handshake, the `Mcp-Session-Id` replay and tool calls are covered by `tests/wasm32.rs` against a scripted server. Unlike natively, a stalled MCP stream has no idle -read timeout of its own; the platform's request limits end it. Custom request -headers (an auth token for the server) are not supported yet -([#275](https://github.com/yologdev/yoagent/issues/275)). +read timeout of its own; the platform's request limits end it. A server that +needs a token gets it from `HttpTransport::with_header` (see the +[MCP guide](mcp.md#servers-that-need-headers)). ## API keys diff --git a/src/agent.rs b/src/agent.rs index 6b98d67..c5b3137 100644 --- a/src/agent.rs +++ b/src/agent.rs @@ -6,7 +6,7 @@ use crate::agent_loop::{ BeforeTurnFn, OnErrorFn, }; use crate::context::{CompactionStrategy, ContextConfig, ExecutionLimits}; -use crate::mcp::{McpClient, McpError, McpToolAdapter}; +use crate::mcp::{HttpTransport, McpClient, McpError, McpToolAdapter}; use crate::provider::{ModelConfig, StreamProvider}; use crate::rt::JoinHandle; use crate::types::*; @@ -862,8 +862,19 @@ impl Agent { } /// Connect to an MCP server via HTTP and add its tools to the agent. - pub async fn with_mcp_server_http(mut self, url: &str) -> Result { - let client = McpClient::connect_http(url).await?; + pub async fn with_mcp_server_http(self, url: &str) -> Result { + self.with_mcp_server_http_transport(HttpTransport::new(url)?) + .await + } + + /// Connect over a configured [`HttpTransport`] — one carrying an auth + /// header ([`HttpTransport::with_header`]), say — and add the server's + /// tools to the agent. + pub async fn with_mcp_server_http_transport( + mut self, + transport: HttpTransport, + ) -> Result { + let client = McpClient::connect_http_with(transport).await?; let client = Arc::new(tokio::sync::Mutex::new(client)); let adapters = McpToolAdapter::from_client(client).await?; for adapter in adapters { diff --git a/src/mcp/client.rs b/src/mcp/client.rs index 07d65fd..eccdf53 100644 --- a/src/mcp/client.rs +++ b/src/mcp/client.rs @@ -45,7 +45,13 @@ impl McpClient { /// Connect to an MCP server via HTTP. pub async fn connect_http(url: &str) -> Result { - let transport = HttpTransport::new(url)?; + Self::connect_http_with(HttpTransport::new(url)?).await + } + + /// Connect over a configured [`HttpTransport`] — with request headers + /// ([`HttpTransport::with_header`]), say. The handshake and timeouts are + /// those of [`connect_http`](Self::connect_http). + pub async fn connect_http_with(transport: HttpTransport) -> Result { let mut client = Self { transport: Arc::new(Mutex::new(Box::new(transport))), server_info: None, diff --git a/src/mcp/transport.rs b/src/mcp/transport.rs index fd3d3fa..f13ec45 100644 --- a/src/mcp/transport.rs +++ b/src/mcp/transport.rs @@ -421,9 +421,14 @@ impl McpTransport for StdioTransport { /// the call has already returned by then. A server that blocks awaiting a reply /// to a `sampling/createMessage` it sent on this stream will therefore time out /// rather than be answered. +/// +/// Extra request headers — an auth token, a `User-Agent` — are set with +/// [`with_header`](Self::with_header) and go on every request of the session. pub struct HttpTransport { client: reqwest::Client, base_url: String, + /// The caller's headers, on every request. Values are marked sensitive. + headers: reqwest::header::HeaderMap, /// Session assigned by the server on `initialize`, replayed on subsequent /// requests. `Mutex` because [`McpTransport::send`] takes `&self`. session_id: Mutex>, @@ -454,10 +459,58 @@ impl HttpTransport { Ok(Self { client, base_url: url.trim_end_matches('/').to_string(), + headers: reqwest::header::HeaderMap::new(), session_id: Mutex::new(None), }) } + /// Headers the transport sets itself; a caller's header of the same name + /// would break the protocol, so [`with_header`](Self::with_header) refuses + /// them. + const RESERVED_HEADERS: [&'static str; 3] = ["accept", "content-type", "mcp-session-id"]; + + /// Send `name: value` on every request of the session: `initialize`, the + /// `initialized` notification, tool calls and the closing `DELETE`. + /// + /// Setting the same name again replaces the value. The value is marked + /// sensitive, so it stays out of debug output. An invalid name or value, + /// or a name the transport sets itself (`Accept`, `Content-Type`, + /// `Mcp-Session-Id`), is an error here, before anything is sent. The + /// headers are fixed for the transport's lifetime; a token that changes + /// during a session needs a new transport. + /// + /// ```no_run + /// # async fn run(token: &str) -> Result<(), yoagent::mcp::McpError> { + /// use yoagent::mcp::{HttpTransport, McpClient}; + /// let transport = HttpTransport::new("https://mcp.example.com/mcp")? + /// .with_header("authorization", format!("Bearer {token}"))?; + /// let client = McpClient::connect_http_with(transport).await?; + /// # Ok(()) } + /// ``` + pub fn with_header(mut self, name: &str, value: impl AsRef) -> Result { + let header = reqwest::header::HeaderName::from_bytes(name.as_bytes()) + .map_err(|e| McpError::Transport(format!("invalid header name {name:?}: {e}")))?; + if Self::RESERVED_HEADERS.contains(&header.as_str()) { + return Err(McpError::Transport(format!( + "header {name:?} is set by the MCP transport itself and cannot be overridden" + ))); + } + let mut value = reqwest::header::HeaderValue::from_str(value.as_ref()) + // The value may be a secret: name the header, never echo the value. + .map_err(|_| McpError::Transport(format!("invalid value for header {name:?}")))?; + value.set_sensitive(true); + self.headers.insert(header, value); + Ok(self) + } + + /// A request to the server with the caller's headers; every request path + /// goes through here so none of them misses them. + fn request(&self, method: reqwest::Method) -> reqwest::RequestBuilder { + self.client + .request(method, &self.base_url) + .headers(self.headers.clone()) + } + /// Decide whether `payload` is *this request's* JSON-RPC response. /// /// Selection is structural rather than "does it deserialize". Every field @@ -767,8 +820,7 @@ impl McpTransport for HttpTransport { let method = request.method.clone(); let mut builder = self - .client - .post(&self.base_url) + .request(reqwest::Method::POST) // Streamable HTTP servers pick their framing from this; JSON-only // servers still match `application/json`. .header("Accept", "application/json, text/event-stream") @@ -830,8 +882,7 @@ impl McpTransport for HttpTransport { async fn notify(&self, notification: JsonRpcNotification) -> Result<(), McpError> { let method = notification.method.clone(); let mut builder = self - .client - .post(&self.base_url) + .request(reqwest::Method::POST) .header("Accept", "application/json, text/event-stream") .json(¬ification); if let Some(session) = self.session_id.lock().await.as_ref() { @@ -870,8 +921,7 @@ impl McpTransport for HttpTransport { // never reached the server leaks the session there, and that only // surfaces later as an unrelated connect failure. match self - .client - .delete(&self.base_url) + .request(reqwest::Method::DELETE) .header("Mcp-Session-Id", &session) .send() .await @@ -896,6 +946,16 @@ impl McpTransport for HttpTransport { mod tests { use super::*; + #[test] + fn caller_header_values_are_sensitive() { + let transport = HttpTransport::new("http://127.0.0.1:9/mcp") + .unwrap() + .with_header("x-api-key", "s3cret") + .unwrap(); + assert!(transport.headers["x-api-key"].is_sensitive()); + assert!(!format!("{:?}", transport.headers).contains("s3cret")); + } + #[cfg(feature = "native")] #[tokio::test] async fn test_stdio_transport_with_cat() { diff --git a/tests/mcp_http_transport_test.rs b/tests/mcp_http_transport_test.rs index 444c20e..55b8e27 100644 --- a/tests/mcp_http_transport_test.rs +++ b/tests/mcp_http_transport_test.rs @@ -723,3 +723,208 @@ async fn a_rejected_initialized_notification_fails_the_connect() { "{err}" ); } + +// --------------------------------------------------------------------------- +// Caller headers (#275) +// --------------------------------------------------------------------------- + +/// Answers every JSON-RPC request (`initialize`, `tools/list`, `tools/call`) +/// with a session id and a minimal result, `202` for notifications and `204` +/// for the closing `DELETE` — but only when both caller headers are present. +/// A request without them matches nothing and gets wiremock's 404. +async fn mount_server_requiring(server: &MockServer, auth: &'static str, agent: &'static str) { + Mock::given(method("POST")) + .and(header("authorization", auth)) + .and(header("user-agent", agent)) + .respond_with(|request: &Request| { + let body: serde_json::Value = serde_json::from_slice(&request.body).unwrap(); + if body.get("id").is_none() { + return ResponseTemplate::new(202); + } + let result = match body["method"].as_str().unwrap() { + "initialize" => serde_json::json!({ + "protocolVersion": "2024-11-05", "capabilities": {"tools": {}}, + "serverInfo": {"name": "auth-fixture", "version": "1"} + }), + "tools/list" => serde_json::json!({"tools": [ + {"name": "whoami", "inputSchema": {"type": "object"}} + ]}), + _ => serde_json::json!({"content": [{"type": "text", "text": "you"}]}), + }; + ResponseTemplate::new(200) + .insert_header("Mcp-Session-Id", "auth-session") + .set_body_json( + serde_json::json!({"jsonrpc": "2.0", "id": body["id"], "result": result}), + ) + }) + .mount(server) + .await; + Mock::given(method("DELETE")) + .and(header("authorization", auth)) + .and(header("user-agent", agent)) + .and(header("mcp-session-id", "auth-session")) + .respond_with(ResponseTemplate::new(204)) + .mount(server) + .await; +} + +/// The headers reach every request of a session: the handshake (initialize and +/// the initialized notification), discovery, a call and the closing DELETE. +#[tokio::test] +async fn caller_headers_go_on_every_request_of_the_session() { + let server = MockServer::start().await; + mount_server_requiring(&server, "Bearer s3cret", "my-app/1.0").await; + let transport = HttpTransport::new(&server.uri()) + .unwrap() + .with_header("Authorization", "Bearer s3cret") + .unwrap() + .with_header("user-agent", "my-app/1.0") + .unwrap(); + let client = McpClient::connect_http_with(transport).await.unwrap(); + let tools = client.list_tools().await.unwrap(); + assert_eq!(tools[0].name, "whoami"); + client + .call_tool("whoami", serde_json::json!({})) + .await + .unwrap(); + client.close().await.unwrap(); + + let requests = server.received_requests().await.unwrap(); + // initialize, notifications/initialized, tools/list, tools/call, DELETE — + // each matched the header-requiring mocks (a miss would have been a 404). + let methods: Vec<&str> = requests.iter().map(|r| r.method.as_str()).collect(); + assert_eq!(methods, ["POST", "POST", "POST", "POST", "DELETE"]); + // `close` treats a rejected DELETE as best-effort, so check it directly. + for r in &requests { + assert_eq!(r.headers.get("authorization").unwrap(), "Bearer s3cret"); + assert_eq!(r.headers.get("user-agent").unwrap(), "my-app/1.0"); + } +} + +/// An agent connects the same way, and discovers the server's tools. +#[tokio::test] +async fn an_agent_connects_through_a_configured_transport() { + let server = MockServer::start().await; + mount_server_requiring(&server, "Bearer s3cret", "my-app/1.0").await; + let transport = HttpTransport::new(&server.uri()) + .unwrap() + .with_header("authorization", "Bearer s3cret") + .unwrap() + .with_header("user-agent", "my-app/1.0") + .unwrap(); + use yoagent::provider::mock::{MockResponse, MockToolCall}; + let provider = yoagent::provider::MockProvider::new(vec![ + MockResponse::ToolCalls(vec![MockToolCall { + name: "whoami".into(), + arguments: serde_json::json!({}), + provider_metadata: None, + }]), + MockResponse::Text("done".into()), + ]); + let mut agent = yoagent::Agent::from_provider(provider, yoagent::provider::ModelConfig::mock()) + .with_mcp_server_http_transport(transport) + .await + .unwrap(); + let mut events = agent.prompt("who am I").await; + let mut called = false; + while let Some(event) = events.recv().await { + if let yoagent::AgentEvent::ToolExecutionEnd { + tool_name, + is_error, + .. + } = event + { + called = tool_name == "whoami" && !is_error; + } + } + agent.finish().await; + assert!(called, "the discovered tool ran with the caller headers"); +} + +/// Without the header the same server refuses the handshake — the header is +/// what the tests above depend on, not an accident of the fixture. +#[tokio::test] +async fn without_the_header_the_server_refuses_the_connect() { + let server = MockServer::start().await; + mount_server_requiring(&server, "Bearer s3cret", "my-app/1.0").await; + let err = McpClient::connect_http(&server.uri()).await.err().unwrap(); + assert!(err.to_string().contains("404"), "{err}"); +} + +#[test] +fn invalid_or_reserved_headers_are_refused_before_sending() { + let new = || HttpTransport::new("http://127.0.0.1:9/mcp").unwrap(); + assert!(new().with_header("bad header", "x").is_err()); + let err = new() + .with_header("authorization", "Bearer s3cret\r\nInjected: 1") + .err() + .unwrap() + .to_string(); + assert!( + !err.contains("s3cret"), + "a refused value is never echoed: {err}" + ); + for reserved in ["Accept", "content-type", "MCP-Session-Id"] { + let err = new().with_header(reserved, "x").err().unwrap().to_string(); + assert!( + err.contains("set by the MCP transport"), + "{reserved}: {err}" + ); + } +} + +/// Setting a header twice keeps only the second value. +#[tokio::test] +async fn setting_a_header_twice_replaces_it() { + let server = MockServer::start().await; + let request = JsonRpcRequest::new("ping", None); + Mock::given(method("POST")) + .and(header("authorization", "Bearer new")) + .respond_with( + ResponseTemplate::new(200).set_body_json( + serde_json::json!({"jsonrpc": "2.0", "id": request.id, "result": {}}), + ), + ) + .expect(1) + .mount(&server) + .await; + HttpTransport::new(&server.uri()) + .unwrap() + .with_header("authorization", "Bearer old") + .unwrap() + .with_header("Authorization", "Bearer new") + .unwrap() + .send(request) + .await + .unwrap(); + let received = server.received_requests().await.unwrap(); + assert_eq!( + received[0].headers.get_all("authorization").iter().count(), + 1 + ); + server.verify().await; +} + +/// No headers configured: nothing extra is sent (no User-Agent either). +#[tokio::test] +async fn a_plain_transport_sends_no_caller_headers() { + let server = MockServer::start().await; + let request = JsonRpcRequest::new("ping", None); + Mock::given(method("POST")) + .and(HeaderAbsent("authorization")) + .and(HeaderAbsent("user-agent")) + .respond_with( + ResponseTemplate::new(200).set_body_json( + serde_json::json!({"jsonrpc": "2.0", "id": request.id, "result": {}}), + ), + ) + .expect(1) + .mount(&server) + .await; + HttpTransport::new(&server.uri()) + .unwrap() + .send(request) + .await + .unwrap(); + server.verify().await; +} From a18b0a5ca27a30ab0efde018453b600e2376715f Mon Sep 17 00:00:00 2001 From: Yuanhao Li Date: Fri, 9 Oct 2026 14:03:31 +0200 Subject: [PATCH 2/3] fix(mcp): refuse non-ASCII header values and framing headers Review of #280: HeaderValue accepts bytes 0x80-0xFF, which wasm32's reqwest refuses on every request (and many servers reject), so with_header now requires visible ASCII. Content-Length, Transfer-Encoding, Connection and Host are reserved too. The wasm32 MCP test now sends an authorization header through the host's fetch and checks it on every request. Docs note that a browser may drop some headers. Co-Authored-By: Claude Opus 5.5 (1M context) Claude-Session: https://claude.ai/code/session_01T7iq5hpndSiHQcnAsywKuG --- CHANGELOG.md | 2 +- CLAUDE.md | 2 +- docs/guides/mcp.md | 9 +++++--- src/mcp/transport.rs | 37 +++++++++++++++++++++++--------- tests/mcp_http_transport_test.rs | 10 ++++++++- tests/wasm32.rs | 22 ++++++++++++------- 6 files changed, 58 insertions(+), 24 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index ef630a1..36c9df1 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -8,7 +8,7 @@ adheres to [Semantic Versioning](https://semver.org/). ### Added -- **Headers for HTTP MCP servers** (#275). `HttpTransport::with_header(name, value)` sends a header on every request of the session — `initialize`, the `initialized` notification, tool calls and the closing `DELETE` — so servers that need a token, an API key or a `User-Agent` can be reached. Setting a name again replaces it; values are marked sensitive; an invalid name or value, or one the transport sets itself (`Accept`, `Content-Type`, `Mcp-Session-Id`), is an error before anything is sent. Connect with `McpClient::connect_http_with(transport)` or `Agent::with_mcp_server_http_transport(transport)` (same handshake and timeouts as `connect_http` / `with_mcp_server_http`, which are unchanged). Works on wasm32. Tests: every request of a session carries the headers (DELETE included, mutation-checked), an agent connects through a configured transport, refusals, replacement, and a plain transport sends none. +- **Headers for HTTP MCP servers** (#275). `HttpTransport::with_header(name, value)` sends a header on every request of the session — `initialize`, the `initialized` notification, tool calls and the closing `DELETE` — so servers that need a token, an API key or a `User-Agent` can be reached. Setting a name again replaces it; values are marked sensitive; an invalid name, a value that is not visible ASCII, or a name the transport or HTTP client sets itself (`Accept`, `Content-Type`, `Mcp-Session-Id`, `Content-Length`, `Transfer-Encoding`, `Connection`, `Host`) is an error before anything is sent. Connect with `McpClient::connect_http_with(transport)` or `Agent::with_mcp_server_http_transport(transport)` (same handshake and timeouts as `connect_http` / `with_mcp_server_http`, which are unchanged). Works on wasm32. Tests: every request of a session carries the headers (DELETE included, mutation-checked), an agent connects through a configured transport, refusals, replacement, and a plain transport sends none. ### Fixed diff --git a/CLAUDE.md b/CLAUDE.md index 329b32d..f543382 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -108,7 +108,7 @@ Behind the `openapi` Cargo feature. `OpenApiToolAdapter` parses an OpenAPI 3.0 s ### MCP Integration (`mcp/`) -`McpClient` communicates via `McpTransport` trait (stdio or HTTP; `notify` sends a `JsonRpcNotification`; the trait default = old best-effort request-and-wait (error logged, `Ok`) for custom transports; a built-in transport's failed `initialized` fails the connect; `connect_*` bound the handshake at 60 s). Stdio: reads until the response with the request's id (numeric or string form), skips server notifications, answers server requests (`ping` → `{}`, others → -32601, `warn!`), drops stale ids, skips non-JSON lines with `warn!` (counted into the close error); a null-id error is taken as the current call's, with `warn!`; `StdinWriter.interrupted` starts the next message on a fresh line after a cut-off write; stderr drained lossily in 8 KiB chunks to `debug!` (target `yoagent::mcp::stderr`) with a 4 KiB `StderrTail`; EOF → `McpError::Transport("Connection closed: … exited with …; last stderr: …")` (waits ≤500 ms for the drain); `kill_on_drop`; tested against bash-scripted servers in `tests/mcp_stdio_test.rs` (non-UTF-8 stderr flood, late answer, kill on drop, death with stderr). Call timeout lives on `McpClient` (`call_timeout`: 300 s stdio, None HTTP/custom) and is copied by `McpToolAdapter::from_client*`; `McpToolAdapter::new` has none; the clock starts after the client lock is held (only cancel ends the wait). `McpToolCallResult.content` is lenient: an unmodelled block becomes `McpContent::Text` (a text resource's own text, else JSON with `data`/`blob` over 1 KiB replaced by their size; `warn!` once per type). `McpToolAdapter` wraps MCP tools to implement `AgentTool`, making them transparent to the agent loop. Added via `Agent::with_mcp_server_stdio()` / `with_mcp_server_http()`. Caller headers (#275): `HttpTransport::with_header(name, value) -> Result` (stored `HeaderMap`, values `set_sensitive`, insert = replace; invalid name/value or a reserved name — `accept`, `content-type`, `mcp-session-id` — is a `Transport` error that never echoes the value) applied by the private `request(Method)` helper that every path (send, notify, close's DELETE) uses; `McpClient::connect_http_with(transport)` (= `connect_http`'s handshake, `call_timeout: None`) and `Agent::with_mcp_server_http_transport(transport)`; tests in `tests/mcp_http_transport_test.rs` (`mount_server_requiring`). +`McpClient` communicates via `McpTransport` trait (stdio or HTTP; `notify` sends a `JsonRpcNotification`; the trait default = old best-effort request-and-wait (error logged, `Ok`) for custom transports; a built-in transport's failed `initialized` fails the connect; `connect_*` bound the handshake at 60 s). Stdio: reads until the response with the request's id (numeric or string form), skips server notifications, answers server requests (`ping` → `{}`, others → -32601, `warn!`), drops stale ids, skips non-JSON lines with `warn!` (counted into the close error); a null-id error is taken as the current call's, with `warn!`; `StdinWriter.interrupted` starts the next message on a fresh line after a cut-off write; stderr drained lossily in 8 KiB chunks to `debug!` (target `yoagent::mcp::stderr`) with a 4 KiB `StderrTail`; EOF → `McpError::Transport("Connection closed: … exited with …; last stderr: …")` (waits ≤500 ms for the drain); `kill_on_drop`; tested against bash-scripted servers in `tests/mcp_stdio_test.rs` (non-UTF-8 stderr flood, late answer, kill on drop, death with stderr). Call timeout lives on `McpClient` (`call_timeout`: 300 s stdio, None HTTP/custom) and is copied by `McpToolAdapter::from_client*`; `McpToolAdapter::new` has none; the clock starts after the client lock is held (only cancel ends the wait). `McpToolCallResult.content` is lenient: an unmodelled block becomes `McpContent::Text` (a text resource's own text, else JSON with `data`/`blob` over 1 KiB replaced by their size; `warn!` once per type). `McpToolAdapter` wraps MCP tools to implement `AgentTool`, making them transparent to the agent loop. Added via `Agent::with_mcp_server_stdio()` / `with_mcp_server_http()`. Caller headers (#275): `HttpTransport::with_header(name, value) -> Result` (stored `HeaderMap`, values `set_sensitive`, insert = replace; invalid name, a value that is not visible ASCII (`to_str` fails — wasm32's reqwest refuses those on every request), or a reserved name — `accept`, `content-type`, `mcp-session-id`, `content-length`, `transfer-encoding`, `connection`, `host` — is a `Transport` error that never echoes the value) applied by the private `request(Method)` helper that every path (send, notify, close's DELETE) uses; `McpClient::connect_http_with(transport)` (= `connect_http`'s handshake, `call_timeout: None`) and `Agent::with_mcp_server_http_transport(transport)`; tests in `tests/mcp_http_transport_test.rs` (`mount_server_requiring`). ### GASP Bridge (`gasp.rs`, feature-gated) diff --git a/docs/guides/mcp.md b/docs/guides/mcp.md index e9f30e4..ba1a2a9 100644 --- a/docs/guides/mcp.md +++ b/docs/guides/mcp.md @@ -103,9 +103,12 @@ let agent = agent.with_mcp_server_http_transport(transport).await?; `initialized` notification, tool calls and the closing `DELETE`. - Setting a name again replaces its value. Values are marked sensitive, so they stay out of debug output, and a refused value is never echoed in the error. -- An invalid name or value is an error from `with_header`, before anything is - sent. So is a header the transport sets itself: `Accept`, `Content-Type`, - `Mcp-Session-Id`. +- An invalid name, or a value that is not visible ASCII, is an error from + `with_header`, before anything is sent. So is a header the transport or the + HTTP client sets itself: `Accept`, `Content-Type`, `Mcp-Session-Id`, + `Content-Length`, `Transfer-Encoding`, `Connection`, `Host`. +- On wasm32 the headers ride on the host's `fetch`. Cloudflare Workers send + them; a browser may drop some, such as `User-Agent`. - The headers are fixed for the transport's lifetime. A token that expires during a session (OAuth refresh) needs a new transport and client. diff --git a/src/mcp/transport.rs b/src/mcp/transport.rs index f13ec45..b7b4c81 100644 --- a/src/mcp/transport.rs +++ b/src/mcp/transport.rs @@ -464,18 +464,28 @@ impl HttpTransport { }) } - /// Headers the transport sets itself; a caller's header of the same name - /// would break the protocol, so [`with_header`](Self::with_header) refuses - /// them. - const RESERVED_HEADERS: [&'static str; 3] = ["accept", "content-type", "mcp-session-id"]; + /// Headers the transport or the HTTP client sets itself; a caller's header + /// of the same name would break the protocol or the request framing, so + /// [`with_header`](Self::with_header) refuses them. + const RESERVED_HEADERS: [&'static str; 7] = [ + "accept", + "content-type", + "mcp-session-id", + "content-length", + "transfer-encoding", + "connection", + "host", + ]; /// Send `name: value` on every request of the session: `initialize`, the /// `initialized` notification, tool calls and the closing `DELETE`. /// /// Setting the same name again replaces the value. The value is marked - /// sensitive, so it stays out of debug output. An invalid name or value, - /// or a name the transport sets itself (`Accept`, `Content-Type`, - /// `Mcp-Session-Id`), is an error here, before anything is sent. The + /// sensitive, so it stays out of debug output. An invalid name, a value + /// that is not visible ASCII, or a name the transport or HTTP client sets + /// itself (`Accept`, `Content-Type`, `Mcp-Session-Id`, `Content-Length`, + /// `Transfer-Encoding`, `Connection`, `Host`) is an error here, before + /// anything is sent. The /// headers are fixed for the transport's lifetime; a token that changes /// during a session needs a new transport. /// @@ -495,9 +505,16 @@ impl HttpTransport { "header {name:?} is set by the MCP transport itself and cannot be overridden" ))); } - let mut value = reqwest::header::HeaderValue::from_str(value.as_ref()) - // The value may be a secret: name the header, never echo the value. - .map_err(|_| McpError::Transport(format!("invalid value for header {name:?}")))?; + // The value may be a secret: name the header, never echo the value. + // Visible ASCII only: `HeaderValue` also takes bytes 0x80-0xFF, which + // many servers reject and wasm32's fetch-based client refuses on every + // request — refuse them here instead. + let invalid = || McpError::Transport(format!("invalid value for header {name:?}")); + let mut value = + reqwest::header::HeaderValue::from_str(value.as_ref()).map_err(|_| invalid())?; + if value.to_str().is_err() { + return Err(invalid()); + } value.set_sensitive(true); self.headers.insert(header, value); Ok(self) diff --git a/tests/mcp_http_transport_test.rs b/tests/mcp_http_transport_test.rs index 55b8e27..2658718 100644 --- a/tests/mcp_http_transport_test.rs +++ b/tests/mcp_http_transport_test.rs @@ -864,7 +864,15 @@ fn invalid_or_reserved_headers_are_refused_before_sending() { !err.contains("s3cret"), "a refused value is never echoed: {err}" ); - for reserved in ["Accept", "content-type", "MCP-Session-Id"] { + // Not visible ASCII: wasm32's client would refuse it on every request. + assert!(new().with_header("x-user", "café").is_err()); + for reserved in [ + "Accept", + "content-type", + "MCP-Session-Id", + "Content-Length", + "host", + ] { let err = new().with_header(reserved, "x").err().unwrap().to_string(); assert!( err.contains("set by the MCP transport"), diff --git a/tests/wasm32.rs b/tests/wasm32.rs index 655ddbd..4faa409 100644 --- a/tests/wasm32.rs +++ b/tests/wasm32.rs @@ -367,7 +367,7 @@ async fn an_extension_gates_and_observes_a_run_on_the_host() { /// test. reqwest's wasm client calls the global `fetch` on every request, as a /// Worker's does, so the whole transport (session header, JSON bodies, the /// response stream) runs against it. Requests are recorded in -/// `globalThis.__mcpRequests` as `" "`. +/// `globalThis.__mcpRequests` as `" "`. struct FakeMcpServer; impl FakeMcpServer { @@ -384,7 +384,8 @@ impl FakeMcpServer { const session = request.headers.get('mcp-session-id') ?? '-'; const text = request.method === 'POST' ? await request.text() : ''; const body = text ? JSON.parse(text) : {}; - globalThis.__mcpRequests.push(`${request.method} ${body.method ?? '-'} ${session}`); + const auth = request.headers.get('authorization') ?? '-'; + globalThis.__mcpRequests.push(`${request.method} ${body.method ?? '-'} ${session} ${auth}`); if (request.method === 'DELETE') return reply(null, { status: 204 }); if (body.id === undefined) return reply(null, { status: 202 }); const results = { @@ -447,8 +448,12 @@ async fn an_agent_calls_an_http_mcp_tool_through_the_hosts_fetch() { }]), MockResponse::Text("done".into()), ]); + let transport = yoagent::mcp::HttpTransport::new("https://mcp.example.test/mcp") + .unwrap() + .with_header("authorization", "Bearer edge-token") + .unwrap(); let mut agent = Agent::from_provider(provider, ModelConfig::mock()) - .with_mcp_server_http("https://mcp.example.test/mcp") + .with_mcp_server_http_transport(transport) .await .expect("connect and discover over fetch"); let events = drain(agent.prompt("shout it").await).await; @@ -470,11 +475,12 @@ async fn an_agent_calls_an_http_mcp_tool_through_the_hosts_fetch() { assert_eq!( server.requests(), [ - "POST initialize -", - "POST notifications/initialized edge-session", - "POST tools/list edge-session", - "POST tools/call edge-session", + "POST initialize - Bearer edge-token", + "POST notifications/initialized edge-session Bearer edge-token", + "POST tools/list edge-session Bearer edge-token", + "POST tools/call edge-session Bearer edge-token", ], - "the handshake, then the session replayed on every later request" + "the handshake, then the session replayed on every later request, \ + each carrying the caller's header" ); } From 98eeb4ed107d36ac233b53ca0b39f1bc66ceb598 Mon Sep 17 00:00:00 2001 From: Yuanhao Li Date: Fri, 9 Oct 2026 14:08:19 +0200 Subject: [PATCH 3/3] docs(mcp): reserved-header error names the HTTP client too; reflow Co-Authored-By: Claude Opus 5.5 (1M context) Claude-Session: https://claude.ai/code/session_01T7iq5hpndSiHQcnAsywKuG --- src/mcp/transport.rs | 7 +++---- 1 file changed, 3 insertions(+), 4 deletions(-) diff --git a/src/mcp/transport.rs b/src/mcp/transport.rs index b7b4c81..390a5f3 100644 --- a/src/mcp/transport.rs +++ b/src/mcp/transport.rs @@ -485,9 +485,8 @@ impl HttpTransport { /// that is not visible ASCII, or a name the transport or HTTP client sets /// itself (`Accept`, `Content-Type`, `Mcp-Session-Id`, `Content-Length`, /// `Transfer-Encoding`, `Connection`, `Host`) is an error here, before - /// anything is sent. The - /// headers are fixed for the transport's lifetime; a token that changes - /// during a session needs a new transport. + /// anything is sent. The headers are fixed for the transport's lifetime; + /// a token that changes during a session needs a new transport. /// /// ```no_run /// # async fn run(token: &str) -> Result<(), yoagent::mcp::McpError> { @@ -502,7 +501,7 @@ impl HttpTransport { .map_err(|e| McpError::Transport(format!("invalid header name {name:?}: {e}")))?; if Self::RESERVED_HEADERS.contains(&header.as_str()) { return Err(McpError::Transport(format!( - "header {name:?} is set by the MCP transport itself and cannot be overridden" + "header {name:?} is set by the MCP transport or HTTP client and cannot be overridden" ))); } // The value may be a secret: name the header, never echo the value.