use crate::mcp::auth::McpCredential;
use crate::mcp::elicitation::{
ElicitationAction, UrlElicitation, UrlElicitationHandler, UrlElicitationPending,
validate_elicitation_url,
};
use crate::mcp::form_elicitation::{
FormElicitation, FormElicitationHandler, FormOutcome, parse_requested_schema,
};
use crate::mcp::protocol::{self, ClientCapabilities, Negotiated};
use crate::mcp::result::extract_json_from_response;
use crate::mcp::transport::{McpConnection, McpEndpoint, McpTransport};
use crate::{
EgressRequest, EgressRequestKind, EgressService, McpProtocolMode, McpToolCallResponse,
McpToolCallResult, McpToolDefinition, McpToolsListResponse, normalize_mcp_error_code,
};
use anyhow::{Result, anyhow};
use async_trait::async_trait;
use everruns_contracts::url_validation::validate_url_dns_pinned;
use serde_json::Value;
use std::collections::{BTreeMap, HashMap, hash_map::DefaultHasher};
use std::hash::{Hash, Hasher};
use std::sync::{Arc, Mutex};
use std::time::{Duration, Instant};
const DISCOVERY_TIMEOUT: Duration = Duration::from_secs(30);
const CALL_TIMEOUT: Duration = Duration::from_secs(60);
const NEGOTIATION_TTL: Duration = Duration::from_secs(300);
const MAX_INPUT_REQUIRED_ROUNDS: usize = 2;
#[derive(Debug, thiserror::Error)]
#[error("MCP server returned error: {status} - {body}")]
pub struct McpHttpStatusError {
pub status: u16,
pub body: String,
}
impl McpHttpStatusError {
pub fn is_unauthorized(&self) -> bool {
self.status == 401
}
}
pub struct RawMcpResponse {
pub status: u16,
pub headers: BTreeMap<String, String>,
pub body: String,
}
async fn send_raw(
egress: &dyn EgressService,
url: &str,
headers: &HashMap<String, String>,
extra_headers: &[(String, String)],
credential: Option<&McpCredential>,
body: Vec<u8>,
timeout: Duration,
) -> Result<RawMcpResponse> {
let (validated_url, resolved_addrs) = validate_url_dns_pinned(url).await.map_err(|e| {
tracing::warn!(url = %url, error = %e, "Blocked MCP request: URL failed SSRF validation");
anyhow!("MCP server URL blocked: {}", e)
})?;
let pin_host = validated_url.host_str().unwrap_or("").to_string();
let mut request = EgressRequest::new("POST", url, EgressRequestKind::Mcp)
.pinned_addrs(pin_host, resolved_addrs)
.header("Content-Type", "application/json")
.header("Accept", "application/json, text/event-stream")
.timeout_ms(timeout.as_millis() as u64)
.body(body);
for (name, value) in headers {
request = request.header(name, value);
}
for (name, value) in extra_headers {
request = request.header(name, value);
}
if let Some(credential) = credential {
if let Some(auth) = &credential.authorization {
request = request.header("Authorization", auth.clone());
}
for (name, value) in &credential.headers {
request = request.header(name, value.clone());
}
}
let response = egress
.send(request)
.await
.map_err(|e| anyhow!("Failed to call MCP server: {}", e))?;
let body = String::from_utf8(response.body)
.map_err(|e| anyhow!("Failed to read MCP response body: {}", e))?;
Ok(RawMcpResponse {
status: response.status,
headers: response.headers,
body,
})
}
pub async fn http_send_rpc(
egress: &dyn EgressService,
url: &str,
headers: &HashMap<String, String>,
credential: Option<&McpCredential>,
body: Vec<u8>,
timeout: Duration,
) -> Result<String> {
let response = send_raw(egress, url, headers, &[], credential, body, timeout).await?;
if !(200..300).contains(&response.status) {
return Err(McpHttpStatusError {
status: response.status,
body: response.body,
}
.into());
}
Ok(response.body)
}
async fn do_handshake(
egress: &dyn EgressService,
url: &str,
headers: &HashMap<String, String>,
credential: Option<&McpCredential>,
preferred_version: &str,
timeout: Duration,
) -> Result<Negotiated> {
let body = serde_json::to_vec(&protocol::initialize_body(0, preferred_version))?;
let extra = protocol::routable_headers(preferred_version, "initialize", None);
let response = send_raw(egress, url, headers, &extra, credential, body, timeout).await?;
if !(200..300).contains(&response.status) {
return Err(McpHttpStatusError {
status: response.status,
body: response.body,
}
.into());
}
let init_json = extract_json_from_response(&response.body).unwrap_or(&response.body);
let version = protocol::protocol_version_from_initialize(init_json)
.unwrap_or_else(|| preferred_version.to_string());
let session_id = protocol::session_id_from_headers(&response.headers);
if let Ok(note) = serde_json::to_vec(&protocol::initialized_notification()) {
let mut note_extra =
protocol::routable_headers(&version, "notifications/initialized", None);
if let Some(session_id) = &session_id {
note_extra.push((protocol::HEADER_SESSION_ID.to_string(), session_id.clone()));
}
let _ = send_raw(egress, url, headers, ¬e_extra, credential, note, timeout).await;
}
Ok(Negotiated {
version,
stateful: true,
session_id,
})
}
#[allow(clippy::too_many_arguments)]
async fn send_op(
egress: &dyn EgressService,
url: &str,
headers: &HashMap<String, String>,
credential: Option<&McpCredential>,
negotiated: &Negotiated,
method: &str,
tool_name: Option<&str>,
body: Vec<u8>,
timeout: Duration,
) -> Result<RawMcpResponse> {
let mut extra = protocol::routable_headers(&negotiated.version, method, tool_name);
if let Some(session_id) = &negotiated.session_id {
extra.push((protocol::HEADER_SESSION_ID.to_string(), session_id.clone()));
}
send_raw(egress, url, headers, &extra, credential, body, timeout).await
}
#[allow(clippy::too_many_arguments)]
async fn negotiate_and_send(
egress: &dyn EgressService,
url: &str,
headers: &HashMap<String, String>,
credential: Option<&McpCredential>,
mode: McpProtocolMode,
method: &str,
tool_name: Option<&str>,
build_body: &(dyn Fn(&str) -> Value + Send + Sync),
cached: Option<Negotiated>,
timeout: Duration,
) -> Result<(String, Negotiated)> {
let mut negotiated = cached.unwrap_or_else(|| Negotiated::initial_for_mode(mode));
if negotiated.stateful && negotiated.session_id.is_none() {
negotiated = do_handshake(
egress,
url,
headers,
credential,
&negotiated.version,
timeout,
)
.await?;
}
let response = send_op(
egress,
url,
headers,
credential,
&negotiated,
method,
tool_name,
serde_json::to_vec(&build_body(&negotiated.version))?,
timeout,
)
.await?;
let mut rejected_probe = None;
let response = if mode == McpProtocolMode::Auto
&& !negotiated.stateful
&& !(200..300).contains(&response.status)
&& protocol::looks_like_handshake_required(response.status, &response.body)
{
tracing::debug!(
url = %url,
"MCP stateless attempt rejected; falling back to stateful handshake"
);
rejected_probe = Some((response.status, response.body.clone()));
negotiated = do_handshake(
egress,
url,
headers,
credential,
protocol::DEFAULT_STATEFUL_VERSION,
timeout,
)
.await
.map_err(|fallback_error| {
fallback_error.context(format!(
"MCP 2026-07-28 probe failed: {} - {}; stateful fallback failed",
response.status, response.body
))
})?;
send_op(
egress,
url,
headers,
credential,
&negotiated,
method,
tool_name,
serde_json::to_vec(&build_body(&negotiated.version))?,
timeout,
)
.await
.map_err(|fallback_error| {
anyhow!(
"MCP 2026-07-28 probe failed: {} - {}; stateful fallback request failed: {}",
response.status,
response.body,
fallback_error
)
})?
} else {
response
};
if !(200..300).contains(&response.status) {
if let Some((probe_status, probe_body)) = rejected_probe {
let fallback_status = response.status;
let fallback_body = response.body;
return Err(anyhow!(McpHttpStatusError {
status: fallback_status,
body: fallback_body.clone(),
})
.context(format!(
"MCP 2026-07-28 probe failed: {probe_status} - {probe_body}; \
stateful fallback failed: {fallback_status} - {fallback_body}"
)));
}
return Err(McpHttpStatusError {
status: response.status,
body: response.body,
}
.into());
}
Ok((response.body, negotiated))
}
fn parse_tools_list(text: &str) -> Result<Vec<McpToolDefinition>> {
let json_str = extract_json_from_response(text)
.ok_or_else(|| anyhow!("SSE response missing data line"))?;
let response: McpToolsListResponse = serde_json::from_str(json_str)?;
if let Some(error) = response.error {
return Err(anyhow!(
"MCP server error: {} ({})",
error.message,
normalize_mcp_error_code(error.code)
));
}
Ok(response
.result
.ok_or_else(|| anyhow!("MCP server returned empty result"))?
.tools)
}
fn parse_tool_call(text: &str) -> Result<McpToolCallResult> {
let json_str = extract_json_from_response(text)
.ok_or_else(|| anyhow!("SSE response missing data line"))?;
let response: McpToolCallResponse = serde_json::from_str(json_str)?;
if let Some(error) = response.error {
return Err(anyhow!(
"MCP tool error: {} (code: {})",
error.message,
normalize_mcp_error_code(error.code)
));
}
response
.result
.ok_or_else(|| anyhow!("MCP server returned empty result"))
}
#[derive(Debug, Clone)]
pub struct HttpToolsList {
pub tools: Vec<McpToolDefinition>,
pub cache_hints: Option<protocol::CacheHints>,
}
pub async fn http_list_tools_with_cache_hints(
egress: &dyn EgressService,
url: &str,
headers: &HashMap<String, String>,
credential: Option<&McpCredential>,
) -> Result<HttpToolsList> {
let capabilities = ClientCapabilities::none();
let (text, _negotiated) = negotiate_and_send(
egress,
url,
headers,
credential,
McpProtocolMode::Auto,
"tools/list",
None,
&|version| protocol::tools_list_body(1, version, capabilities),
None,
DISCOVERY_TIMEOUT,
)
.await?;
let cache_hints = extract_json_from_response(&text).and_then(protocol::cache_hints_from_result);
Ok(HttpToolsList {
tools: parse_tools_list(&text)?,
cache_hints,
})
}
pub async fn http_list_tools(
egress: &dyn EgressService,
url: &str,
headers: &HashMap<String, String>,
credential: Option<&McpCredential>,
) -> Result<Vec<McpToolDefinition>> {
Ok(
http_list_tools_with_cache_hints(egress, url, headers, credential)
.await?
.tools,
)
}
#[derive(Debug, Clone, thiserror::Error)]
#[error("MCP server error: {message} (code {code})")]
pub struct McpRpcError {
pub code: i64,
pub message: String,
pub data: Option<Value>,
}
pub async fn http_request(
egress: &dyn EgressService,
url: &str,
headers: &HashMap<String, String>,
credential: Option<&McpCredential>,
method: &str,
params: &Value,
) -> Result<Value> {
let capabilities = ClientCapabilities::none();
let (text, _negotiated) = negotiate_and_send(
egress,
url,
headers,
credential,
McpProtocolMode::Auto,
method,
None,
&|version| protocol::request_body(1, method, params, version, capabilities),
None,
CALL_TIMEOUT,
)
.await?;
let json_str = extract_json_from_response(&text)
.ok_or_else(|| anyhow!("SSE response missing data line"))?;
let mut response: Value = serde_json::from_str(json_str)?;
if let Some(error) = response.get("error").filter(|error| !error.is_null()) {
return Err(McpRpcError {
code: error.get("code").and_then(Value::as_i64).unwrap_or(-32603),
message: error
.get("message")
.and_then(Value::as_str)
.unwrap_or("unknown error")
.to_string(),
data: error.get("data").cloned(),
}
.into());
}
response
.get_mut("result")
.map(Value::take)
.ok_or_else(|| anyhow!("MCP server returned empty result"))
}
#[allow(clippy::too_many_arguments)]
pub async fn http_call_tool(
egress: &dyn EgressService,
url: &str,
headers: &HashMap<String, String>,
server_name: &str,
tool_name: &str,
arguments: Value,
credential: Option<&McpCredential>,
elicitation: Option<&dyn UrlElicitationHandler>,
) -> Result<McpToolCallResult> {
let capabilities = ClientCapabilities {
url_elicitation: elicitation.is_some(),
form_elicitation: false,
};
let (text, negotiated) = negotiate_and_send(
egress,
url,
headers,
credential,
McpProtocolMode::Auto,
"tools/call",
Some(tool_name),
&|version| protocol::tools_call_body(1, tool_name, &arguments, version, capabilities),
None,
CALL_TIMEOUT,
)
.await?;
let text = resolve_input_required(
egress,
url,
headers,
credential,
&negotiated,
capabilities,
server_name,
tool_name,
&arguments,
InputHandlers {
url: elicitation,
form: None,
},
text,
)
.await?;
parse_tool_call(&text)
}
#[derive(Clone, Copy)]
struct InputHandlers<'a> {
url: Option<&'a dyn UrlElicitationHandler>,
form: Option<&'a dyn FormElicitationHandler>,
}
#[allow(clippy::too_many_arguments)]
async fn resolve_input_required(
egress: &dyn EgressService,
url: &str,
headers: &HashMap<String, String>,
credential: Option<&McpCredential>,
negotiated: &Negotiated,
capabilities: ClientCapabilities,
server_name: &str,
tool_name: &str,
arguments: &Value,
handlers: InputHandlers<'_>,
text: String,
) -> Result<String> {
let mut text = text;
for round in 0..MAX_INPUT_REQUIRED_ROUNDS {
let Some(json) = extract_json_from_response(&text) else {
return Ok(text);
};
let Some(input_required) = protocol::input_required_from_result(json) else {
return Ok(text);
};
let input_responses =
gather_input_responses(&input_required, server_name, tool_name, handlers, url).await?;
tracing::debug!(
url = %url,
tool = %tool_name,
round = round,
inputs = input_responses.len(),
"MCP tools/call returned input_required; retrying"
);
let body = serde_json::to_vec(&protocol::tools_call_retry_body(
round as i64 + 2,
tool_name,
arguments,
&negotiated.version,
capabilities,
input_required.request_state.as_deref(),
&input_responses,
))?;
let response = send_op(
egress,
url,
headers,
credential,
negotiated,
"tools/call",
Some(tool_name),
body,
CALL_TIMEOUT,
)
.await?;
if !(200..300).contains(&response.status) {
return Err(McpHttpStatusError {
status: response.status,
body: response.body,
}
.into());
}
text = response.body;
}
if extract_json_from_response(&text)
.and_then(protocol::input_required_from_result)
.is_some()
{
return Err(anyhow!(
"MCP tool '{tool_name}' still requires input after {MAX_INPUT_REQUIRED_ROUNDS} \
rounds; if you were asked to complete something in your browser, finish it and \
run the tool again"
));
}
Ok(text)
}
async fn gather_input_responses(
input_required: &protocol::InputRequired,
server_name: &str,
tool_name: &str,
handlers: InputHandlers<'_>,
url: &str,
) -> Result<BTreeMap<String, Value>> {
let forms = input_required
.requests
.iter()
.filter(|request| request.form_elicitation().is_some())
.count();
if forms > 1 {
return Err(anyhow!(
"MCP server '{server_name}' sent {forms} form elicitations in one round for tool \
'{tool_name}'; Everruns answers at most one at a time"
));
}
let mut responses = BTreeMap::new();
for request in &input_required.requests {
if let Some((message, requested_schema)) = request.form_elicitation() {
let response = answer_form_elicitation(
request.key.clone(),
message,
&requested_schema,
server_name,
tool_name,
handlers.form,
url,
)
.await?;
responses.insert(request.key.clone(), response);
continue;
}
let Some((message, elicitation_url)) = request.url_elicitation() else {
return Err(anyhow!(
"MCP server '{server_name}' requested input '{}' ({}) for tool '{tool_name}' \
that this client does not support; everruns answers only URL and form mode \
elicitation",
request.key,
request.method
));
};
let Some(handler) = handlers.url else {
return Err(anyhow!(
"MCP server '{server_name}' sent a URL elicitation for tool '{tool_name}', \
but this host declared no elicitation capability and has no way to ask a user"
));
};
let (host, punycode) = validate_elicitation_url(&elicitation_url).map_err(|e| {
anyhow!("MCP server '{server_name}' sent an unusable elicitation URL: {e}")
})?;
let elicitation_request = UrlElicitation {
server_name: server_name.to_string(),
tool_name: tool_name.to_string(),
key: request.key.clone(),
message,
url: elicitation_url,
host,
punycode,
};
let action = handler.request_url_consent(&elicitation_request).await?;
tracing::info!(
url = %url,
server = %server_name,
tool = %tool_name,
elicitation_host = %elicitation_request.host,
action = action.as_str(),
"MCP URL elicitation resolved"
);
match action {
ElicitationAction::Accept => {
responses.insert(
request.key.clone(),
serde_json::json!({ "action": "accept" }),
);
}
ElicitationAction::Decline | ElicitationAction::Cancel => {
return Err(anyhow!(UrlElicitationPending {
server_name: elicitation_request.server_name,
tool_name: elicitation_request.tool_name,
message: elicitation_request.message,
url: elicitation_request.url,
host: elicitation_request.host,
punycode: elicitation_request.punycode,
action,
}));
}
}
}
Ok(responses)
}
async fn answer_form_elicitation(
key: String,
message: String,
requested_schema: &Value,
server_name: &str,
tool_name: &str,
handler: Option<&dyn FormElicitationHandler>,
url: &str,
) -> Result<Value> {
let Some(handler) = handler else {
return Err(anyhow!(
"MCP server '{server_name}' sent a form elicitation for tool '{tool_name}', but \
this client did not declare form mode to it; everruns declares form mode only to \
servers whose elicitation_policy is url_and_form"
));
};
let schema = parse_requested_schema(requested_schema).map_err(|refusal| {
anyhow!(
"MCP server '{server_name}' asked questions for tool '{tool_name}' that Everruns \
will not put to a person: {refusal}"
)
})?;
if message.chars().count() > crate::mcp::form_elicitation::MAX_FORM_TEXT_CHARS {
return Err(anyhow!(
"MCP server '{server_name}' sent a form elicitation message longer than {} \
characters for tool '{tool_name}'",
crate::mcp::form_elicitation::MAX_FORM_TEXT_CHARS
));
}
let elicitation = FormElicitation {
server_name: server_name.to_string(),
tool_name: tool_name.to_string(),
key,
message,
schema,
};
let outcome = handler.answer_form(&elicitation).await?;
tracing::info!(
url = %url,
server = %server_name,
tool = %tool_name,
outcome = match &outcome {
FormOutcome::Accept(_) => "accept",
FormOutcome::Decline => "decline",
FormOutcome::Ask => "ask",
},
"MCP form elicitation resolved"
);
match outcome {
FormOutcome::Accept(content) => Ok(serde_json::json!({
"action": "accept",
"content": content,
})),
FormOutcome::Decline => Ok(serde_json::json!({ "action": "decline" })),
FormOutcome::Ask => Err(anyhow!(elicitation.pending())),
}
}
#[derive(Debug, Clone, PartialEq, Eq, Hash)]
struct NegotiationCacheKey {
url: String,
server_name: String,
protocol_mode: String,
headers_hash: u64,
credential_hash: u64,
}
impl NegotiationCacheKey {
fn new(
connection: &McpConnection,
url: &str,
headers: &HashMap<String, String>,
credential: Option<&McpCredential>,
) -> Self {
Self {
url: url.to_string(),
server_name: connection.name.clone(),
protocol_mode: connection.protocol_mode.to_string(),
headers_hash: hash_headers(headers),
credential_hash: hash_credential(credential),
}
}
}
fn hash_headers(headers: &HashMap<String, String>) -> u64 {
let mut hasher = DefaultHasher::new();
let sorted: BTreeMap<_, _> = headers.iter().collect();
sorted.hash(&mut hasher);
hasher.finish()
}
fn hash_credential(credential: Option<&McpCredential>) -> u64 {
let mut hasher = DefaultHasher::new();
credential
.and_then(|credential| credential.authorization.as_deref())
.hash(&mut hasher);
if let Some(credential) = credential {
let sorted: BTreeMap<_, _> = credential.headers.iter().collect();
sorted.hash(&mut hasher);
}
hasher.finish()
}
#[derive(Debug, Clone, PartialEq, Eq, Hash)]
enum ToolsCacheScope {
Shared,
Credential(u64),
}
#[derive(Debug, Clone, PartialEq, Eq, Hash)]
struct ToolsCacheKey {
url: String,
server_name: String,
protocol_mode: String,
headers_hash: u64,
scope: ToolsCacheScope,
}
impl ToolsCacheKey {
fn new(key: &NegotiationCacheKey, scope: ToolsCacheScope) -> Self {
Self {
url: key.url.clone(),
server_name: key.server_name.clone(),
protocol_mode: key.protocol_mode.clone(),
headers_hash: key.headers_hash,
scope,
}
}
fn shared(key: &NegotiationCacheKey) -> Self {
Self::new(key, ToolsCacheScope::Shared)
}
fn credential(key: &NegotiationCacheKey) -> Self {
Self::new(key, ToolsCacheScope::Credential(key.credential_hash))
}
}
struct CachedTools {
tools: Vec<McpToolDefinition>,
fetched_at: Instant,
ttl: Duration,
}
impl CachedTools {
fn is_fresh(&self) -> bool {
self.fetched_at.elapsed() < self.ttl
}
}
pub struct HttpTransport {
egress: Arc<dyn EgressService>,
negotiations: Mutex<HashMap<NegotiationCacheKey, (Negotiated, Instant)>>,
tools: Mutex<HashMap<ToolsCacheKey, CachedTools>>,
elicitation: Option<Arc<dyn UrlElicitationHandler>>,
form_elicitation: Option<Arc<dyn FormElicitationHandler>>,
}
impl HttpTransport {
pub fn new(egress: Arc<dyn EgressService>) -> Self {
Self {
egress,
negotiations: Mutex::new(HashMap::new()),
tools: Mutex::new(HashMap::new()),
elicitation: None,
form_elicitation: None,
}
}
pub fn with_elicitation_handler(mut self, handler: Arc<dyn UrlElicitationHandler>) -> Self {
self.elicitation = Some(handler);
self
}
pub fn with_form_elicitation_handler(
mut self,
handler: Arc<dyn FormElicitationHandler>,
) -> Self {
self.form_elicitation = Some(handler);
self
}
fn handlers<'a>(&'a self, connection: &McpConnection) -> InputHandlers<'a> {
let policy = connection.elicitation_policy;
InputHandlers {
url: self.elicitation.as_deref().filter(|_| policy.allows_url()),
form: self
.form_elicitation
.as_deref()
.filter(|_| policy.allows_form()),
}
}
fn capabilities(&self, connection: &McpConnection) -> ClientCapabilities {
let handlers = self.handlers(connection);
ClientCapabilities {
url_elicitation: handlers.url.is_some(),
form_elicitation: handlers.form.is_some(),
}
}
fn http_parts(connection: &McpConnection) -> Result<(&str, &HashMap<String, String>)> {
match &connection.endpoint {
McpEndpoint::Http { url, headers } => Ok((url.as_str(), headers)),
#[cfg(feature = "mcp-stdio")]
_ => Err(anyhow!(
"HttpTransport received a non-HTTP endpoint for server '{}'",
connection.name
)),
}
}
fn cached_negotiation(&self, key: &NegotiationCacheKey) -> Option<Negotiated> {
let cache = self.negotiations.lock().ok()?;
let (negotiated, at) = cache.get(key)?;
(at.elapsed() < NEGOTIATION_TTL).then(|| negotiated.clone())
}
fn store_negotiation(&self, key: NegotiationCacheKey, negotiated: Negotiated) {
if let Ok(mut cache) = self.negotiations.lock() {
cache.insert(key, (negotiated, Instant::now()));
}
}
fn cached_tools(&self, key: &NegotiationCacheKey) -> Option<Vec<McpToolDefinition>> {
let cache = self.tools.lock().ok()?;
for candidate in [ToolsCacheKey::shared(key), ToolsCacheKey::credential(key)] {
if let Some(entry) = cache.get(&candidate)
&& entry.is_fresh()
{
return Some(entry.tools.clone());
}
}
None
}
fn store_tools(
&self,
key: &NegotiationCacheKey,
hints: Option<protocol::CacheHints>,
tools: &[McpToolDefinition],
) {
let Some(hints) = hints else {
return;
};
let cache_key = match hints.scope {
protocol::CacheScope::Public => ToolsCacheKey::shared(key),
protocol::CacheScope::Private => ToolsCacheKey::credential(key),
};
if let Ok(mut cache) = self.tools.lock() {
cache.retain(|_, entry| entry.is_fresh());
cache.insert(
cache_key,
CachedTools {
tools: tools.to_vec(),
fetched_at: Instant::now(),
ttl: hints.ttl,
},
);
}
}
}
#[async_trait]
impl McpTransport for HttpTransport {
async fn list_tools(
&self,
connection: &McpConnection,
credential: Option<&McpCredential>,
) -> Result<Vec<McpToolDefinition>> {
let (url, headers) = Self::http_parts(connection)?;
let cache_key = NegotiationCacheKey::new(connection, url, headers, credential);
if let Some(tools) = self.cached_tools(&cache_key) {
return Ok(tools);
}
let capabilities = self.capabilities(connection);
let cached = self.cached_negotiation(&cache_key);
let (text, negotiated) = negotiate_and_send(
self.egress.as_ref(),
url,
headers,
credential,
connection.protocol_mode,
"tools/list",
None,
&|version| protocol::tools_list_body(1, version, capabilities),
cached,
DISCOVERY_TIMEOUT,
)
.await?;
self.store_negotiation(cache_key.clone(), negotiated);
let tools = parse_tools_list(&text)?;
let hints = extract_json_from_response(&text).and_then(protocol::cache_hints_from_result);
self.store_tools(&cache_key, hints, &tools);
Ok(tools)
}
async fn call_tool(
&self,
connection: &McpConnection,
tool_name: &str,
arguments: Value,
credential: Option<&McpCredential>,
) -> Result<McpToolCallResult> {
let (url, headers) = Self::http_parts(connection)?;
let cache_key = NegotiationCacheKey::new(connection, url, headers, credential);
let capabilities = self.capabilities(connection);
let cached = self.cached_negotiation(&cache_key);
let (text, negotiated) = negotiate_and_send(
self.egress.as_ref(),
url,
headers,
credential,
connection.protocol_mode,
"tools/call",
Some(tool_name),
&|version| protocol::tools_call_body(1, tool_name, &arguments, version, capabilities),
cached,
CALL_TIMEOUT,
)
.await?;
let text = resolve_input_required(
self.egress.as_ref(),
url,
headers,
credential,
&negotiated,
capabilities,
&connection.name,
tool_name,
&arguments,
self.handlers(connection),
text,
)
.await?;
self.store_negotiation(cache_key, negotiated);
parse_tool_call(&text)
}
}