use crate::{
cancellation::AgentCancellation,
config::{DEFAULT_MCP_TIMEOUT_SECONDS, McpHttpServerConfig},
mcp::{
McpError, McpResult,
jsonrpc::{JsonRpcMessage, JsonRpcNotification, JsonRpcRequest, RequestId},
oauth::{TokenProvider, parse_www_authenticate_resource_metadata},
protocol::{METHOD_INITIALIZE, METHOD_INITIALIZED, PROTOCOL_VERSION},
sse::{
DEFAULT_MAX_SSE_EVENT_BYTES, SseDecoder, jsonrpc_message_from_event, read_sse_event,
},
},
};
use reqwest::{
StatusCode,
blocking::{Client, Response},
header::{ACCEPT, AUTHORIZATION, CONTENT_TYPE, HeaderMap},
};
use serde_json::Value;
use std::{
collections::HashMap,
io::{BufReader, Read},
sync::{
Arc, Mutex,
atomic::{AtomicBool, Ordering},
mpsc::{self, SyncSender},
},
thread::{self, JoinHandle},
time::{Duration, Instant},
};
use tokio::sync::watch;
pub(crate) const MAX_MCP_HTTP_ERROR_BODY_BYTES: usize = 1024;
pub(crate) const MAX_MCP_HTTP_RESPONSE_BYTES: usize = 10 * 1024 * 1024;
const HEADER_MCP_SESSION_ID: &str = "Mcp-Session-Id";
const HEADER_MCP_PROTOCOL_VERSION: &str = "MCP-Protocol-Version";
const ACCEPT_POST: &str = "application/json, text/event-stream";
const ACCEPT_GET: &str = "text/event-stream, application/json";
const HTTP_SHUTDOWN_REQUEST_TIMEOUT: Duration = Duration::from_millis(500);
type PendingMap = Arc<Mutex<HashMap<RequestId, SyncSender<McpResult<Value>>>>>;
pub(crate) struct HttpConnection {
client: Client,
url: String,
display_url: String,
server_name: Option<String>,
headers: HeaderMap,
token_provider: Option<TokenProvider>,
session_id: Arc<Mutex<Option<String>>>,
protocol_version: Arc<Mutex<Option<String>>>,
timeout: Duration,
pending: PendingMap,
closed: Arc<AtomicBool>,
notification_stream_unavailable: Arc<Mutex<Option<String>>>,
last_event_id: Arc<Mutex<Option<String>>>,
shutdown_sender: watch::Sender<bool>,
get_thread: Mutex<Option<JoinHandle<()>>>,
shutdown_complete: Mutex<bool>,
}
impl std::fmt::Debug for HttpConnection {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("HttpConnection")
.field("url", &self.display_url)
.field("server_name", &self.server_name)
.field(
"headers",
&format_args!("<{} headers redacted>", self.headers.len()),
)
.field("timeout", &self.timeout)
.finish_non_exhaustive()
}
}
impl HttpConnection {
pub(crate) fn connect_named(
server_name: Option<&str>,
config: &McpHttpServerConfig,
mc_home: Option<&std::path::Path>,
) -> McpResult<Self> {
let headers = crate::mcp::headers::resolve_http_headers(&config.headers)?;
if config.oauth.is_some() && has_auth_header(&headers) {
return Err(McpError::Config(
"MCP HTTP OAuth cannot be combined with Authorization or Proxy-Authorization headers"
.to_string(),
));
}
let timeout = Duration::from_secs(config.timeout.unwrap_or(DEFAULT_MCP_TIMEOUT_SECONDS));
let client = Client::builder()
.redirect(reqwest::redirect::Policy::none())
.timeout(timeout)
.build()
.map_err(sanitize_reqwest_error)?;
let token_provider = match &config.oauth {
Some(oauth) => {
let server_name = server_name.ok_or_else(|| {
McpError::Config("MCP OAuth HTTP connection requires server name".to_string())
})?;
let mc_home = mc_home.ok_or_else(|| {
McpError::Config("MCP OAuth HTTP connection requires MC_HOME paths".to_string())
})?;
Some(TokenProvider::new(
mc_home.to_path_buf(),
server_name.to_string(),
config.url.clone(),
oauth.clone(),
client.clone(),
))
}
None => None,
};
let (shutdown_sender, _) = watch::channel(false);
Ok(Self {
client,
url: config.url.clone(),
display_url: sanitize_url(&config.url),
server_name: server_name.map(ToString::to_string),
headers,
token_provider,
session_id: Arc::new(Mutex::new(None)),
protocol_version: Arc::new(Mutex::new(None)),
timeout,
pending: Arc::new(Mutex::new(HashMap::new())),
closed: Arc::new(AtomicBool::new(false)),
notification_stream_unavailable: Arc::new(Mutex::new(None)),
last_event_id: Arc::new(Mutex::new(None)),
shutdown_sender,
get_thread: Mutex::new(None),
shutdown_complete: Mutex::new(false),
})
}
pub(crate) fn send_request(
&self,
id: RequestId,
method: &str,
params: Option<Value>,
cancellation: Option<&AgentCancellation>,
) -> McpResult<Value> {
if let Some(cancellation) = cancellation
&& let Err(error) = cancellation.check()
{
return Err(McpError::Transport(error.to_string()));
}
let message = serde_json::to_value(JsonRpcRequest::new(id.clone(), method, params))
.map_err(McpError::transport)?;
let (sender, receiver) = mpsc::sync_channel(1);
{
let mut pending = self
.pending
.lock()
.map_err(|_| McpError::Transport("pending request map poisoned".to_string()))?;
if self.closed.load(Ordering::SeqCst) {
return Err(McpError::Transport("MCP HTTP server is closed".to_string()));
}
pending.insert(id.clone(), sender);
}
let started = Instant::now();
let posted = self.post_json(message, true, Some(&id));
match posted {
Ok(Some(value)) => {
let _ = self.pending.lock().map(|mut pending| pending.remove(&id));
if method == METHOD_INITIALIZE {
self.capture_initialize_result(&value);
}
if let Some(cancellation) = cancellation
&& let Err(error) = cancellation.check()
{
self.close();
return Err(McpError::Transport(error.to_string()));
}
if self.closed.load(Ordering::SeqCst) {
return Err(McpError::Transport("MCP HTTP server is closed".to_string()));
}
Ok(value)
}
Ok(None) => self.wait_for_pending(id, receiver, cancellation),
Err(error) => {
let _ = self.pending.lock().map(|mut pending| pending.remove(&id));
if matches!(error, McpError::Timeout { .. })
|| started.elapsed() >= self.timeout
|| cancellation.is_some_and(AgentCancellation::is_canceled)
{
self.close();
}
Err(error)
}
}
}
pub(crate) fn send_notification(&self, method: &str, params: Option<Value>) -> McpResult<()> {
if self.closed.load(Ordering::SeqCst) {
return Err(McpError::Transport("MCP HTTP server is closed".to_string()));
}
let message = serde_json::to_value(JsonRpcNotification::new(method, params))
.map_err(McpError::transport)?;
match self.post_json(message, false, None)? {
None => {
if method == METHOD_INITIALIZED {
self.start_notification_stream();
}
Ok(())
}
Some(_) => {
if method == METHOD_INITIALIZED {
self.start_notification_stream();
}
Ok(())
}
}
}
fn close(&self) {
self.closed.store(true, Ordering::SeqCst);
self.shutdown_sender.send_replace(true);
fail_pending(
&self.pending,
McpError::Transport("MCP HTTP server closed".to_string()),
);
}
pub(crate) fn request_shutdown(&self) {
self.close();
}
pub(crate) fn shutdown(&self) {
let mut complete = self
.shutdown_complete
.lock()
.unwrap_or_else(|error| error.into_inner());
if *complete {
return;
}
self.close();
if let Ok(mut handle) = self.get_thread.lock()
&& let Some(handle) = handle.take()
{
let _ = handle.join();
}
self.delete_session_best_effort();
*complete = true;
}
fn post_json(
&self,
body: Value,
expect_response: bool,
request_id: Option<&RequestId>,
) -> McpResult<Option<Value>> {
let mut response = self.send_post_json(body.clone(), false)?;
if response.status() == StatusCode::UNAUTHORIZED && self.token_provider.is_some() {
response = self.send_post_json(body, true)?;
if response.status() == StatusCode::UNAUTHORIZED {
return Err(McpError::Config(format!(
"authentication failed after token refresh; run magi-code mcp login {}",
self.server_name.as_deref().unwrap_or("<server>")
)));
}
}
self.capture_session_id(response.headers());
match response.status() {
StatusCode::OK if expect_response => self.handle_ok_response(response, request_id),
StatusCode::OK => Err(McpError::Transport(format!(
"MCP HTTP notification expected 202 or 204 from {}",
self.display_url
))),
StatusCode::ACCEPTED | StatusCode::NO_CONTENT if !expect_response => Ok(None),
StatusCode::ACCEPTED | StatusCode::NO_CONTENT => Ok(None),
status if status.is_client_error() || status.is_server_error() => {
Err(http_status_error(
status,
response,
&self.display_url,
&self.headers,
self.server_name.as_deref(),
))
}
status => Err(McpError::Transport(format!(
"MCP HTTP unexpected status {status} from {}",
self.display_url
))),
}
}
fn handle_ok_response(
&self,
response: Response,
request_id: Option<&RequestId>,
) -> McpResult<Option<Value>> {
let content_type = response
.headers()
.get(CONTENT_TYPE)
.and_then(|value| value.to_str().ok())
.unwrap_or("")
.to_ascii_lowercase();
if content_type.contains("text/event-stream") {
return self.read_sse_response(response, request_id);
}
if content_type.contains("application/json") || content_type.is_empty() {
let mut bytes = Vec::new();
let mut limited = response.take(MAX_MCP_HTTP_RESPONSE_BYTES as u64 + 1);
limited.read_to_end(&mut bytes).map_err(|_| {
McpError::Transport("MCP HTTP response body read failed".to_string())
})?;
if bytes.len() > MAX_MCP_HTTP_RESPONSE_BYTES {
return Err(McpError::Transport(format!(
"MCP HTTP response body from {} exceeded {} bytes",
self.display_url, MAX_MCP_HTTP_RESPONSE_BYTES
)));
}
let value: Value = serde_json::from_slice(&bytes).map_err(McpError::transport)?;
return self.value_from_jsonrpc_message(value, request_id);
}
Err(McpError::Transport(format!(
"MCP HTTP unsupported content type '{}' from {}",
content_type, self.display_url
)))
}
fn read_sse_response(
&self,
response: Response,
request_id: Option<&RequestId>,
) -> McpResult<Option<Value>> {
let mut reader = BufReader::new(response);
loop {
let event = match read_sse_event(&mut reader, DEFAULT_MAX_SSE_EVENT_BYTES)? {
Some(event) => event,
None => return Ok(None),
};
if let Some(id) = &event.id {
let _ = self
.last_event_id
.lock()
.map(|mut last| *last = Some(id.clone()));
}
let Some(message) = jsonrpc_message_from_event(&event)? else {
continue;
};
if let Some(value) = self.handle_stream_message(message, request_id)? {
return Ok(Some(value));
}
}
}
fn value_from_jsonrpc_message(
&self,
value: Value,
request_id: Option<&RequestId>,
) -> McpResult<Option<Value>> {
let message: JsonRpcMessage = serde_json::from_value(value).map_err(McpError::transport)?;
self.handle_stream_message(message, request_id)
}
fn handle_stream_message(
&self,
message: JsonRpcMessage,
request_id: Option<&RequestId>,
) -> McpResult<Option<Value>> {
match message {
JsonRpcMessage::Response(response) => {
if request_id.is_none_or(|id| *id == response.id) {
Ok(Some(response.result))
} else {
dispatch_pending(&self.pending, response.id, Ok(response.result));
Ok(None)
}
}
JsonRpcMessage::Error(error) => {
if request_id.is_none_or(|id| Some(id) == error.id.as_ref()) {
let data = error.error;
Err(McpError::Protocol {
code: data.code,
message: data.message,
})
} else if let Some(id) = error.id {
let data = error.error;
dispatch_pending(
&self.pending,
id,
Err(McpError::Protocol {
code: data.code,
message: data.message,
}),
);
Ok(None)
} else {
Ok(None)
}
}
JsonRpcMessage::Notification(_) | JsonRpcMessage::Request(_) => Ok(None),
}
}
fn wait_for_pending(
&self,
id: RequestId,
receiver: mpsc::Receiver<McpResult<Value>>,
cancellation: Option<&AgentCancellation>,
) -> McpResult<Value> {
let started = Instant::now();
loop {
if let Some(cancellation) = cancellation
&& let Err(error) = cancellation.check()
{
let _ = self.pending.lock().map(|mut pending| pending.remove(&id));
self.close();
return Err(McpError::Transport(error.to_string()));
}
if let Some(reason) = self
.notification_stream_unavailable
.lock()
.ok()
.and_then(|reason| reason.clone())
{
let _ = self.pending.lock().map(|mut pending| pending.remove(&id));
return Err(McpError::Transport(format!(
"MCP HTTP GET SSE unavailable while waiting for response: {reason}"
)));
}
let remaining = self.timeout.saturating_sub(started.elapsed());
if remaining.is_zero() {
let _ = self.pending.lock().map(|mut pending| pending.remove(&id));
self.close();
return Err(McpError::Timeout {
seconds: self.timeout.as_secs(),
});
}
match receiver.recv_timeout(remaining.min(Duration::from_millis(100))) {
Ok(result) => return result,
Err(mpsc::RecvTimeoutError::Timeout) => continue,
Err(mpsc::RecvTimeoutError::Disconnected) => {
let _ = self.pending.lock().map(|mut pending| pending.remove(&id));
return Err(McpError::Transport(
"MCP HTTP pending response channel disconnected".to_string(),
));
}
}
}
}
fn send_post_json(&self, body: Value, force_refresh: bool) -> McpResult<Response> {
let builder = self
.request_with_session_headers(self.client.post(&self.url), force_refresh)?
.header(CONTENT_TYPE, "application/json")
.header(ACCEPT, ACCEPT_POST)
.json(&body);
builder
.send()
.map_err(|error| sanitize_reqwest_error_with_timeout(error, self.timeout))
}
fn request_with_session_headers(
&self,
builder: reqwest::blocking::RequestBuilder,
force_refresh: bool,
) -> McpResult<reqwest::blocking::RequestBuilder> {
let mut builder = builder.headers(self.headers.clone());
if let Some(provider) = &self.token_provider {
let token = if force_refresh {
provider.force_refresh_access_token()?
} else {
provider.access_token()?
};
builder = builder.header(AUTHORIZATION, format!("Bearer {token}"));
}
if let Some(session_id) = self.session_id.lock().ok().and_then(|id| id.clone()) {
builder = builder.header(HEADER_MCP_SESSION_ID, session_id);
}
if let Some(protocol_version) = self.protocol_version.lock().ok().and_then(|v| v.clone()) {
builder = builder.header(HEADER_MCP_PROTOCOL_VERSION, protocol_version);
}
Ok(builder)
}
fn capture_session_id(&self, headers: &HeaderMap) {
if let Some(value) = headers
.get(HEADER_MCP_SESSION_ID)
.and_then(|value| value.to_str().ok())
&& let Ok(mut session_id) = self.session_id.lock()
{
*session_id = Some(value.to_string());
}
}
fn capture_initialize_result(&self, value: &Value) {
if let Some(protocol_version) = value.get("protocolVersion").and_then(Value::as_str)
&& protocol_version == PROTOCOL_VERSION
&& let Ok(mut stored) = self.protocol_version.lock()
{
*stored = Some(protocol_version.to_string());
}
}
fn start_notification_stream(&self) {
let Ok(mut slot) = self.get_thread.lock() else {
return;
};
if self.closed.load(Ordering::SeqCst) || slot.is_some() {
return;
}
let url = self.url.clone();
let display_url = self.display_url.clone();
let headers = self.headers.clone();
let token_provider = self.token_provider.clone();
let session_id = Arc::clone(&self.session_id);
let protocol_version = Arc::clone(&self.protocol_version);
let pending = Arc::clone(&self.pending);
let unavailable = Arc::clone(&self.notification_stream_unavailable);
let timeout = self.timeout;
let last_event_id = Arc::clone(&self.last_event_id);
let shutdown = self.shutdown_sender.subscribe();
let handle = thread::spawn(move || {
if *shutdown.borrow() {
return;
}
let access_token = token_provider
.as_ref()
.map(TokenProvider::access_token)
.transpose();
drop(token_provider);
let access_token = match access_token {
Ok(token) => token,
Err(error) => {
record_unavailable(&unavailable, error.to_string());
return;
}
};
if *shutdown.borrow() {
return;
}
let runtime = match tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
{
Ok(runtime) => runtime,
Err(error) => {
record_unavailable(
&unavailable,
format!("MCP HTTP GET SSE runtime failed: {error}"),
);
return;
}
};
let client = match reqwest::Client::builder()
.redirect(reqwest::redirect::Policy::none())
.connect_timeout(timeout)
.build()
{
Ok(client) => client,
Err(error) => {
record_unavailable(&unavailable, sanitize_reqwest_error(error).to_string());
return;
}
};
runtime.block_on(run_get_stream(GetStreamContext {
client,
url,
display_url,
headers,
access_token,
session_id,
protocol_version,
pending,
unavailable,
last_event_id,
shutdown,
}));
});
*slot = Some(handle);
}
fn delete_session_best_effort(&self) {
let Some(session_id) = self.session_id.lock().ok().and_then(|id| id.clone()) else {
return;
};
let mut builder = self
.client
.delete(&self.url)
.headers(self.headers.clone())
.header(HEADER_MCP_SESSION_ID, session_id)
.timeout(HTTP_SHUTDOWN_REQUEST_TIMEOUT);
if let Some(provider) = &self.token_provider
&& let Ok(token) = provider.access_token()
{
builder = builder.header(AUTHORIZATION, format!("Bearer {token}"));
}
if let Some(protocol_version) = self.protocol_version.lock().ok().and_then(|v| v.clone()) {
builder = builder.header(HEADER_MCP_PROTOCOL_VERSION, protocol_version);
}
let _ = builder.send();
}
}
impl Drop for HttpConnection {
fn drop(&mut self) {
self.shutdown();
}
}
fn has_auth_header(headers: &HeaderMap) -> bool {
headers.keys().any(|name| {
name.as_str().eq_ignore_ascii_case("authorization")
|| name.as_str().eq_ignore_ascii_case("proxy-authorization")
})
}
fn dispatch_pending(pending: &PendingMap, id: RequestId, result: McpResult<Value>) {
if let Ok(mut pending) = pending.lock()
&& let Some(sender) = pending.remove(&id)
{
let _ = sender.send(result);
}
}
struct GetStreamContext {
client: reqwest::Client,
url: String,
display_url: String,
headers: HeaderMap,
access_token: Option<String>,
session_id: Arc<Mutex<Option<String>>>,
protocol_version: Arc<Mutex<Option<String>>>,
pending: PendingMap,
unavailable: Arc<Mutex<Option<String>>>,
last_event_id: Arc<Mutex<Option<String>>>,
shutdown: watch::Receiver<bool>,
}
async fn run_get_stream(context: GetStreamContext) {
let GetStreamContext {
client,
url,
display_url,
headers,
access_token,
session_id,
protocol_version,
pending,
unavailable,
last_event_id,
mut shutdown,
} = context;
let mut builder = client.get(url).headers(headers).header(ACCEPT, ACCEPT_GET);
if let Some(token) = access_token {
builder = builder.header(AUTHORIZATION, format!("Bearer {token}"));
}
if let Some(id) = session_id.lock().ok().and_then(|id| id.clone()) {
builder = builder.header(HEADER_MCP_SESSION_ID, id);
}
if let Some(version) = protocol_version.lock().ok().and_then(|v| v.clone()) {
builder = builder.header(HEADER_MCP_PROTOCOL_VERSION, version);
}
let response = tokio::select! {
result = builder.send() => match result {
Ok(response) => response,
Err(error) => {
if !*shutdown.borrow() {
record_unavailable(&unavailable, sanitize_reqwest_error(error).to_string());
}
return;
}
},
_ = shutdown.changed() => return,
};
if *shutdown.borrow() {
return;
}
if !response.status().is_success() {
record_unavailable(
&unavailable,
format!(
"MCP HTTP GET SSE unavailable: status {} from {display_url}",
response.status()
),
);
return;
}
let content_type = response
.headers()
.get(CONTENT_TYPE)
.and_then(|value| value.to_str().ok())
.unwrap_or("")
.to_ascii_lowercase();
if !content_type.contains("text/event-stream") {
record_unavailable(
&unavailable,
format!(
"MCP HTTP GET SSE unsupported content type '{content_type}' from {display_url}"
),
);
return;
}
let mut response = response;
let mut decoder = SseDecoder::default();
loop {
let chunk = tokio::select! {
result = response.chunk() => match result {
Ok(Some(chunk)) => chunk,
Ok(None) => {
match decoder.finish(DEFAULT_MAX_SSE_EVENT_BYTES) {
Ok(Some(event)) => {
if handle_get_event(event, &pending, &unavailable, &last_event_id) {
return;
}
}
Ok(None) => {}
Err(error) => {
record_unavailable(&unavailable, error.to_string());
return;
}
}
record_unavailable(&unavailable, "MCP HTTP GET SSE stream closed".to_string());
return;
}
Err(error) => {
if !*shutdown.borrow() {
record_unavailable(&unavailable, sanitize_reqwest_error(error).to_string());
}
return;
}
},
_ = shutdown.changed() => return,
};
match decoder.push_chunk(&chunk, DEFAULT_MAX_SSE_EVENT_BYTES) {
Ok(events) => {
for event in events {
if handle_get_event(event, &pending, &unavailable, &last_event_id) {
return;
}
}
}
Err(error) => {
record_unavailable(&unavailable, error.to_string());
return;
}
}
}
}
fn handle_get_event(
event: crate::mcp::sse::SseEvent,
pending: &PendingMap,
unavailable: &Arc<Mutex<Option<String>>>,
last_event_id: &Arc<Mutex<Option<String>>>,
) -> bool {
if let Some(id) = &event.id {
let _ = last_event_id
.lock()
.map(|mut last| *last = Some(id.clone()));
}
match jsonrpc_message_from_event(&event) {
Ok(Some(JsonRpcMessage::Response(response))) => {
dispatch_pending(pending, response.id, Ok(response.result));
false
}
Ok(Some(JsonRpcMessage::Error(error))) => {
if let Some(id) = error.id {
let data = error.error;
dispatch_pending(
pending,
id,
Err(McpError::Protocol {
code: data.code,
message: data.message,
}),
);
}
false
}
Ok(Some(JsonRpcMessage::Notification(_)) | Some(JsonRpcMessage::Request(_)) | None) => {
false
}
Err(error) => {
record_unavailable(unavailable, error.to_string());
true
}
}
}
fn fail_pending(pending: &PendingMap, error: McpError) {
if let Ok(mut pending) = pending.lock() {
for (_, sender) in pending.drain() {
let _ = sender.send(Err(error.clone()));
}
}
}
fn record_unavailable(slot: &Arc<Mutex<Option<String>>>, reason: String) {
let bounded = bound_text(&reason, MAX_MCP_HTTP_ERROR_BODY_BYTES);
if let Ok(mut slot) = slot.lock() {
*slot = Some(bounded);
}
}
fn http_status_error(
status: StatusCode,
mut response: Response,
display_url: &str,
headers: &HeaderMap,
server_name: Option<&str>,
) -> McpError {
if status == StatusCode::UNAUTHORIZED {
if let Some(value) = response
.headers()
.get(reqwest::header::WWW_AUTHENTICATE)
.and_then(|value| value.to_str().ok())
&& parse_www_authenticate_resource_metadata(value).is_some()
{
let server = server_name.unwrap_or("<server>");
return McpError::Transport(format!(
"server appears to require OAuth; add mcpServers.{server}.oauth to .mcp.json and run magi-code mcp login {server}"
));
}
return McpError::Transport(format!("MCP HTTP status {status} from {display_url}"));
}
if status == StatusCode::FORBIDDEN {
return McpError::Transport(format!("MCP HTTP status {status} from {display_url}"));
}
let mut bytes = Vec::new();
let mut limited = response
.by_ref()
.take(MAX_MCP_HTTP_ERROR_BODY_BYTES as u64 + 1);
let _ = limited.read_to_end(&mut bytes);
let truncated = bytes.len() > MAX_MCP_HTTP_ERROR_BODY_BYTES;
bytes.truncate(MAX_MCP_HTTP_ERROR_BODY_BYTES);
let raw_preview = String::from_utf8_lossy(&bytes).replace(['\r', '\n'], " ");
let mut preview = redact_http_error_preview(&raw_preview, headers);
if truncated {
preview.push_str("...");
}
if preview.is_empty() {
McpError::Transport(format!("MCP HTTP status {status} from {display_url}"))
} else {
McpError::Transport(format!(
"MCP HTTP status {status} from {display_url}: {preview}"
))
}
}
fn redact_http_error_preview(preview: &str, headers: &HeaderMap) -> String {
let mut redacted = preview.to_string();
for value in headers.values().filter_map(|value| value.to_str().ok()) {
if value.len() >= 4 {
redacted = redacted.replace(value, "[REDACTED]");
}
}
let words = redacted
.split_whitespace()
.map(redact_secret_word)
.collect::<Vec<_>>();
redact_bearer_tokens(&words.join(" "))
}
fn redact_bearer_tokens(text: &str) -> String {
let mut out = Vec::new();
let mut redact_next = false;
for word in text.split_whitespace() {
if redact_next {
out.push("[REDACTED]".to_string());
redact_next = false;
continue;
}
out.push(word.to_string());
if word.trim_end_matches(':').eq_ignore_ascii_case("bearer") {
redact_next = true;
}
}
out.join(" ")
}
fn redact_secret_word(word: &str) -> String {
let trimmed =
word.trim_matches(|c: char| !c.is_ascii_alphanumeric() && c != '-' && c != '_' && c != '.');
if looks_like_secret_token(trimmed) {
word.replace(trimmed, "[REDACTED]")
} else {
word.to_string()
}
}
fn looks_like_secret_token(token: &str) -> bool {
let lower = token.to_ascii_lowercase();
lower.starts_with("sk-")
|| lower.starts_with("key-")
|| lower.starts_with("api-key-")
|| (token.len() >= 32
&& token.bytes().all(|byte| {
byte.is_ascii_alphanumeric()
|| matches!(byte, b'-' | b'_' | b'.' | b'+' | b'/' | b'=')
}))
}
fn sanitize_reqwest_error_with_timeout(error: reqwest::Error, timeout: Duration) -> McpError {
if error.is_timeout() {
McpError::Timeout {
seconds: timeout.as_secs(),
}
} else {
sanitize_reqwest_error(error)
}
}
fn sanitize_reqwest_error(error: reqwest::Error) -> McpError {
let kind = if error.is_timeout() {
"request timed out"
} else if error.is_connect() {
"connection failed"
} else if error.is_decode() {
"decode failed"
} else if error.is_body() {
"body read failed"
} else {
"request failed"
};
McpError::Transport(format!("MCP HTTP {kind}"))
}
fn sanitize_url(url: &str) -> String {
match reqwest::Url::parse(url) {
Ok(parsed) => {
let host = parsed.host_str().unwrap_or("<unknown>");
let port = parsed
.port()
.map(|port| format!(":{port}"))
.unwrap_or_default();
format!("{}://{}{}{}", parsed.scheme(), host, port, parsed.path())
}
Err(_) => "<invalid-url>".to_string(),
}
}
fn bound_text(text: &str, max: usize) -> String {
if text.len() <= max {
return text.to_string();
}
let end = text
.char_indices()
.rev()
.find(|(i, _)| *i <= max)
.map(|(i, _)| i)
.unwrap_or(0);
let mut bounded = text[..end].to_string();
bounded.push_str("...");
bounded
}