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
8 changes: 2 additions & 6 deletions crates/agent/ARCHITECTURE.md
Original file line number Diff line number Diff line change
Expand Up @@ -4411,11 +4411,7 @@ response; pinned in `an_ordinary_request_is_answered_as_rmcps_own_client_would`)
and `McpAuthStore` `mcp-login` uses (so the new token is persisted), and the request is retried
**once**; a second 401 is returned as the server's answer, saying to run `agent mcp-login <server>`
again — on a tool call as at connect, since only a new login helps once a refreshed token is refused
too. **The server's reason survives the 401:** a 401 whose small body (≤ 64 KiB) is a JSON-RPC
error fails as rmcp's own bare-401 error does, `HTTP 401 Unauthorized: <message> (JSON-RPC error
<code>)` — still a 401 to the refresh, and for a server with no login (a static key it stopped
accepting, say) the tool error the model sees, with the server's "invalid API key" in it rather than a
bare "Auth required"; any other 401 is `AuthRequired` (`tests/mcp_unauthorized.rs`). That covers `tools/call`, `resources/*`,
too. **The server's reason survives the 401:** when its small body (≤ 64 KiB) is a JSON-RPC error, the message — untrusted text from an external system — is cut to 1 KiB (with `…`), its control characters made spaces, and fenced the way event payloads are (`<mcp_server_message untrusted>…`, with `<`/`>` — and their fullwidth, small-form and angle-bracket lookalikes — and quotes replaced so it cannot close its fence or a quoted-string, and bidi, zero-width and BOM characters dropped so it cannot reorder or hide what is shown). The same fence (`mcp_wire::fenced_server_message`) holds every server-supplied header or body that reaches model-visible text: a challenge header printed among a failed tool call's causes, the post-refresh `agent mcp-login` error, a non-JSON-RPC success's body, a legacy `server/discover` rejection's body. With a `WWW-Authenticate` challenge the 401 stays `AuthRequired`, carrying that challenge for whatever reads it (rmcp's auth client, `auth_challenge`) with the reason added as RFC 6750's `error_description`; without one it fails as rmcp's own bare-401 error does, `HTTP 401 Unauthorized: <fenced message> (JSON-RPC error <code>)`. Either way it is still a 401 to the refresh, and for a server with no login (a static key it stopped accepting, say) it is the tool error the model sees: a failed tool call prints its error's causes too, so a challenge and the reason in it are not hidden behind rmcp's "Auth required"; any other 401 is a bare `AuthRequired` (`tests/mcp_unauthorized.rs`). That covers `tools/call`, `resources/*`,
`prompts/*`, `skills/*`, MCP App view reads, the handshake, the standalone stream and `events/*`. A 403
(`InsufficientScope`) never refreshes.

Expand Down Expand Up @@ -4724,7 +4720,7 @@ mcp_events_subscribe (any session) ──► owned by that session
passes through the same `rescue`). Stateless
(`2026-07-28`) streamable-HTTP servers get `events/*` directly over HTTP (`MCP-Protocol-Version`,
`Mcp-Method`, per-request `_meta`, the server's resolved headers and OAuth bearer), with bodies
bounded by the same per-message cap as every MCP transport (`mcp_stdio::max_message_bytes`, `BEYOND_AI_AGENT_MCP_MAX_MESSAGE_BYTES`: a unary JSON body, a unary SSE answer's event, an `events/stream` event — an over-cap answer fails its request, never read whole). An over-cap **event notification** (on a push stream, or from a stdio server) is skipped, not reconnected into: its bounded head is read structurally (`mcp_stdio::oversized_stand_in`, the same top-level member walk as `scan_head`) for its routing, cursor and id — **in any key order**: `method`, `id` and `params` are found wherever they fall within the window, and inside `params` every scalar and `_meta` that appears whole is kept — the rest is skipped unread, and a small `$oversized` stand-in takes its place. When the head does not reach the method (a payload-first `params` fills it), the message is still taken for an event unless the head proves it is a request (an `id`) or a response; over stdio, a stand-in whose routing lay past the head is delivered to every push stream on that connection (each records the gap, since which one it belonged to is unknowable), and on a direct-HTTP stream it is that stream's. In either case the subscription keeps whatever cursor the head showed (or else the next event's or heartbeat's — never reconnecting into the same event), records an `oversized` gap (a `gap` frame, and a notice to the model that an event was dropped unread) and carries on (`tests/mcp_message_cap.rs`, both key orders over both transports) — JWKS documents at 64 KiB — and an SSE reader whose scan and drain
bounded by the same per-message cap as every MCP transport (`mcp_stdio::max_message_bytes`, `BEYOND_AI_AGENT_MCP_MAX_MESSAGE_BYTES`: a unary JSON body, a unary SSE answer's event, an `events/stream` event — an over-cap answer fails its request, never read whole). An over-cap **event notification** (on a push stream, or from a stdio server) is skipped, not reconnected into: its bounded head is read structurally (`mcp_stdio::oversized_stand_in`, the same top-level member walk as `scan_head`) for its routing, cursor and id — **in any key order**: `method`, `id` and `params` are found wherever they fall within the window, and inside `params` every scalar and `_meta` that appears whole is kept — the rest is skipped unread, and a small `$oversized` stand-in takes its place. When the head does not reach the method (a payload-first `params` fills it), the message is still taken for an event unless the head proves it is a request (an `id`) or a response — and its gap says it was _possibly not an event_ (to the client, `possibly_not_an_event`, and to the model), since a request whose id and method both lay past the window looks the same. Over stdio, a stand-in whose routing lay past the head is delivered to every push stream on that connection (each records the gap in its own state, since which one it belonged to is unknowable), stamped with one id for the one dropped message (a per-process sequence from a secret start), so a session holding several of those streams tells its model **once**; the host-reserved `$`-keys (`$oversized`, `$ambiguous`, `$dropped_id`) are stripped from every server notification at ingress (`mcp_stdio::host_params`, at the router and on direct-HTTP streams) unless it carries the per-process secret the host stamps into its own stand-ins (`$host`), so a server cannot forge a gap or pre-empt a real drop's notice; with no push stream open it is logged, not reported. On a direct-HTTP stream it is that stream's. In either case the subscription keeps whatever cursor the head showed (or else the next event's or heartbeat's — never reconnecting into the same event), records an `oversized` gap (a `gap` frame, and a notice to the model that an event was dropped unread) and carries on (`tests/mcp_message_cap.rs`, both key orders over both transports) — JWKS documents at 64 KiB — and an SSE reader whose scan and drain
are both linear (events are parsed in place; the consumed prefix is dropped once per read, not once
per event); an older,
session-bound HTTP server goes through rmcp's own connection, which carries its `Mcp-Session-Id`.
Expand Down
39 changes: 34 additions & 5 deletions crates/agent/src/tools/mcp.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1216,9 +1216,38 @@ impl McpServerHandle {
}
}

fn tool_call_err(server: &str, remote: &str, e: impl std::fmt::Display) -> ToolError {
fn tool_call_err(server: &str, remote: &str, e: &ServiceError) -> ToolError {
// The causes too, each once: a 401's challenge (and the server's reason in it) is in the
// transport error's source, not its own text ("Auth required").
let mut text = e.to_string();
// rmcp's `TransportSend` shows its transport error but does not chain it as a source.
let mut cause: Option<&(dyn std::error::Error + 'static)> = match e {
ServiceError::TransportSend(transport) => std::error::Error::source(transport),
_ => std::error::Error::source(e),
};
while let Some(c) = cause {
// A 401/403's challenge is the server's own header: fenced and cut short like any text a
// server supplies, not shown raw (rmcp's display of these errors prints it verbatim).
use rmcp::transport::streamable_http_client::{AuthRequiredError, InsufficientScopeError};
let fenced = crate::tools::mcp_wire::fenced_server_message;
let more = if let Some(a) = c.downcast_ref::<AuthRequiredError>() {
format!(
"authorization required: {}",
fenced(&a.www_authenticate_header)
)
} else if let Some(s) = c.downcast_ref::<InsufficientScopeError>() {
format!("insufficient scope: {}", fenced(&s.www_authenticate_header))
} else {
c.to_string()
};
if !text.contains(&more) {
text.push_str(": ");
text.push_str(&more);
}
cause = c.source();
}
ToolError::Execution(format!(
"mcp server `{server}` tool `{remote}` call failed: {e}"
"mcp server `{server}` tool `{remote}` call failed: {text}"
))
}

Expand All @@ -1240,7 +1269,7 @@ async fn drive_tool_call(
let host_arc = calling_host(&client.service().host);
match call_tool_tracked(&client, params.clone(), host_arc)
.await
.map_err(|e| tool_call_err(server_name, remote_name, e))?
.map_err(|e| tool_call_err(server_name, remote_name, &e))?
{
CallToolResponse::Complete(result) => return Ok(result),
CallToolResponse::InputRequired(required) => {
Expand Down Expand Up @@ -1686,7 +1715,7 @@ async fn await_task(
client = recover(conn, client, loss, losses).await;
continue;
}
None => return Err(tool_call_err(server_name, remote_name, e).into()),
None => return Err(tool_call_err(server_name, remote_name, &e).into()),
},
}
}
Expand Down Expand Up @@ -1726,7 +1755,7 @@ async fn await_task(
client = recover(conn, client, loss, losses).await;
continue;
}
None => return Err(tool_call_err(server_name, remote_name, e).into()),
None => return Err(tool_call_err(server_name, remote_name, &e).into()),
},
};
let detailed = info.task;
Expand Down
74 changes: 65 additions & 9 deletions crates/agent/src/tools/mcp_events/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -496,6 +496,10 @@ struct Hub {
/// `mcp_events_*` commands in flight — spawned so none ever blocks the session's command loop.
command_tasks: Mutex<tokio::task::JoinSet<()>>,
owns_configured: bool,
/// The last few over-cap messages a stdio connection reported to every push stream at once
/// (`$dropped_id`): each is told to the model once, however many of this session's
/// subscriptions it reached.
dropped_seen: Mutex<std::collections::VecDeque<u64>>,
}

/// A session's MCP Events client: its subscriptions, its coalescer, and the injection path into
Expand Down Expand Up @@ -553,6 +557,7 @@ impl McpEventsHub {
store,
command_tasks: Mutex::new(tokio::task::JoinSet::new()),
owns_configured: cfg.owns_configured,
dropped_seen: Mutex::new(std::collections::VecDeque::new()),
});
let weak_hub = Arc::downgrade(&hub);
let on_change: Arc<dyn Fn() + Send + Sync> = Arc::new(move || {
Expand Down Expand Up @@ -1390,10 +1395,39 @@ impl Hub {
if let Some(cursor) = stand_in.get("cursor") {
carrier["cursor"] = cursor.clone();
}
self.gap(spec, state, &carrier);
// Taken for an event without its method seen: said so, not asserted.
if stand_in.get("$ambiguous") == Some(&json!(true)) {
carrier["possibly_not_an_event"] = json!(true);
}
// One dropped message reported to several of this session's subscriptions: each keeps
// its own state, the model is told once.
let tell_model = match stand_in.get("$dropped_id").and_then(Value::as_u64) {
Some(id) => {
let mut seen = self
.dropped_seen
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
let first = !seen.contains(&id);
if first {
if seen.len() == 64 {
seen.pop_front();
}
seen.push_back(id);
}
first
}
None => true,
};
self.gap_with(spec, state, &carrier, tell_model);
}

fn gap(&self, spec: &SubSpec, state: &SubState, carrier: &Value) {
self.gap_with(spec, state, carrier, true);
}

/// Record a gap: the subscription's position, its status frame, and — when `tell_model` —
/// the notice queued for the model.
fn gap_with(&self, spec: &SubSpec, state: &SubState, carrier: &Value, tell_model: bool) {
if carrier
.as_object()
.is_some_and(|o| o.contains_key("cursor"))
Expand All @@ -1402,19 +1436,22 @@ impl Hub {
}
let reason = carrier.get("reason").and_then(Value::as_str);
let event_id = carrier.get("eventId").filter(|v| !v.is_null()).cloned();
self.status_event(
spec,
"gap",
json!({ "cursor": state.cursor(), "reason": reason, "event_id": event_id }),
);
if spec.sub.action != McpEventAction::Notify {
let possibly_not = carrier.get("possibly_not_an_event") == Some(&json!(true));
let mut status =
json!({ "cursor": state.cursor(), "reason": reason, "event_id": event_id });
if possibly_not {
status["possibly_not_an_event"] = json!(true);
}
self.status_event(spec, "gap", status);
if tell_model && spec.sub.action != McpEventAction::Notify {
let queued = self.store.push_pending(PendingEvent::new(
spec.sub.action,
spec.server.clone(),
spec.sub.name.clone(),
spec.arguments(),
spec.sub.instructions.clone(),
json!({ "gap": true, "cursor": state.cursor(), "reason": reason, "eventId": event_id }),
json!({ "gap": true, "cursor": state.cursor(), "reason": reason, "eventId": event_id,
"possibly_not_an_event": possibly_not }),
));
if !queued {
tracing::warn!("pending queue full; a gap notice was not queued for the model");
Expand Down Expand Up @@ -2285,8 +2322,13 @@ fn render_injection(batch: u64, events: &[PendingEvent]) -> String {
for e in events {
if e.event.get("gap") == Some(&json!(true)) {
if e.event.get("reason") == Some(&json!("oversized")) {
let what = if e.event.get("possibly_not_an_event") == Some(&json!(true)) {
"A message — possibly not an event — "
} else {
"An event "
};
out.push_str(&format!(
"\n[gap] An event for `{}` on `{}`{} was larger than the message-size limit and \
"\n[gap] {what}for `{}` on `{}`{} was larger than the message-size limit and \
was dropped unread. If it matters, re-check the authoritative state with tools.\n",
e.name,
e.server,
Expand Down Expand Up @@ -2482,6 +2524,20 @@ mod tests {
let _ = task.await;
}

/// An over-cap gap whose message was taken for an event without its method seen says it may
/// not have been one.
#[test]
fn an_ambiguous_oversized_gap_says_it_may_not_have_been_an_event() {
let gap = |possibly: bool| {
pending(
McpEventAction::Steer,
json!({ "gap": true, "reason": "oversized", "possibly_not_an_event": possibly }),
)
};
assert!(render_injection(1, &[gap(true)]).contains("possibly not an event"));
assert!(!render_injection(1, &[gap(false)]).contains("possibly not an event"));
}

/// Every injection's text names its batch on its first line — for anyone reading the
/// transcript; delivery itself is tracked by tag (`AgentEvent::Steered`), not by text.
#[test]
Expand Down
Loading
Loading