Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 5 additions & 1 deletion CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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, 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

- **`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)

Expand Down
2 changes: 1 addition & 1 deletion CLAUDE.md
Original file line number Diff line number Diff line change
Expand Up @@ -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, 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)

Expand Down
29 changes: 29 additions & 0 deletions docs/guides/mcp.md
Original file line number Diff line number Diff line change
Expand Up @@ -83,6 +83,35 @@ 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 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.

`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:
Expand Down
6 changes: 3 additions & 3 deletions docs/guides/wasm-workers.md
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down
17 changes: 14 additions & 3 deletions src/agent.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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::*;
Expand Down Expand Up @@ -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<Self, McpError> {
let client = McpClient::connect_http(url).await?;
pub async fn with_mcp_server_http(self, url: &str) -> Result<Self, McpError> {
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<Self, McpError> {
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 {
Expand Down
8 changes: 7 additions & 1 deletion src/mcp/client.rs
Original file line number Diff line number Diff line change
Expand Up @@ -45,7 +45,13 @@ impl McpClient {

/// Connect to an MCP server via HTTP.
pub async fn connect_http(url: &str) -> Result<Self, McpError> {
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<Self, McpError> {
let mut client = Self {
transport: Arc::new(Mutex::new(Box::new(transport))),
server_info: None,
Expand Down
88 changes: 82 additions & 6 deletions src/mcp/transport.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<Option<String>>,
Expand Down Expand Up @@ -454,10 +459,74 @@ 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 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, 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.
///
/// ```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<str>) -> Result<Self, McpError> {
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 or HTTP client and cannot be overridden"
)));
}
// 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)
}

/// 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
Expand Down Expand Up @@ -767,8 +836,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")
Expand Down Expand Up @@ -830,8 +898,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(&notification);
if let Some(session) = self.session_id.lock().await.as_ref() {
Expand Down Expand Up @@ -870,8 +937,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
Expand All @@ -896,6 +962,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() {
Expand Down
Loading
Loading