use std::collections::BTreeMap;
use std::sync::Arc;
use std::time::Duration;
use async_trait::async_trait;
use rmcp::model::{
CallToolRequestParams, CallToolResponse, CancelTaskParams, ClientInfo, ErrorCode,
ExtensionCapabilities, GetPromptRequestParams, GetPromptResponse, GetTaskParams,
InputResponses, JsonObject, ProtocolVersion, ReadResourceRequestParams, ReadResourceResponse,
TASKS_EXTENSION_ID, UpdateTaskParams,
};
use rmcp::service::{
ClientLifecycleMode, ClientServiceExt as _, RoleClient, RunningService, ServiceError,
};
use rmcp::transport::IntoTransport;
use serde_json::Value;
use super::{Advertised, ToolClient, ToolError, ToolId};
use crate::core::{Effect, EffectDescriptor, EffectError, Recovery, Sensitivity, Trust};
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct McpDataSafety {
pub max_input_sensitivity: Sensitivity,
pub output_sensitivity: Sensitivity,
}
impl McpDataSafety {
#[must_use]
pub const fn public() -> Self {
Self {
max_input_sensitivity: Sensitivity::Public,
output_sensitivity: Sensitivity::Public,
}
}
#[must_use]
pub const fn max_input(mut self, sensitivity: Sensitivity) -> Self {
self.max_input_sensitivity = sensitivity;
self
}
#[must_use]
pub const fn output(mut self, sensitivity: Sensitivity) -> Self {
self.output_sensitivity = sensitivity;
self
}
}
#[derive(Debug, Clone, Default)]
pub struct McpAccess {
prompts: BTreeMap<String, McpDataSafety>,
resources: BTreeMap<String, McpDataSafety>,
task_input: Option<McpDataSafety>,
}
impl McpAccess {
#[must_use]
pub fn new() -> Self {
Self::default()
}
#[must_use]
pub fn prompt(mut self, name: impl Into<String>, safety: McpDataSafety) -> Self {
self.prompts.insert(name.into(), safety);
self
}
#[must_use]
pub fn resource(mut self, uri: impl Into<String>, safety: McpDataSafety) -> Self {
self.resources.insert(uri.into(), safety);
self
}
#[must_use]
pub fn task_input(mut self, safety: McpDataSafety) -> Self {
self.task_input = Some(safety);
self
}
#[cfg(feature = "manifest")]
#[must_use]
pub fn from_manifest(server: &str, manifest: &crate::manifest::Manifest) -> Self {
let mut access = Self::new();
for grant in &manifest.spec.context.prompts {
if grant.server == server {
access = access.prompt(
grant.name.clone(),
McpDataSafety::public()
.max_input(grant.max_input_sensitivity)
.output(grant.output_sensitivity),
);
}
}
for grant in &manifest.spec.context.resources {
if grant.server == server {
access = access.resource(
grant.uri.clone(),
McpDataSafety::public().output(grant.output_sensitivity),
);
}
}
if let Some(grant) = manifest.task_input_grant(server) {
access =
access.task_input(McpDataSafety::public().max_input(grant.max_input_sensitivity));
}
access
}
}
#[derive(Debug)]
pub struct McpClient {
server: String,
service: Arc<RunningService<RoleClient, ClientInfo>>,
access: McpAccess,
timeout: Duration,
destination: crate::tools::Destination,
}
impl McpClient {
#[must_use]
pub fn host_info() -> ClientInfo {
let mut info = ClientInfo::default();
info.protocol_version = ProtocolVersion::V_2026_07_28;
info.client_info =
rmcp::model::Implementation::new(env!("CARGO_PKG_NAME"), env!("CARGO_PKG_VERSION"));
let mut extensions = ExtensionCapabilities::new();
extensions.insert(TASKS_EXTENSION_ID.to_owned(), JsonObject::new());
info.capabilities.extensions = Some(extensions);
info
}
pub const KNOWN_VERSIONS: [ProtocolVersion; 5] = [
ProtocolVersion::V_2024_11_05,
ProtocolVersion::V_2025_03_26,
ProtocolVersion::V_2025_06_18,
ProtocolVersion::V_2025_11_25,
ProtocolVersion::V_2026_07_28,
];
pub async fn connect<T, E, A>(
server: impl Into<String>,
transport: T,
destination: crate::tools::Destination,
) -> Result<Self, ToolError>
where
T: IntoTransport<RoleClient, E, A>,
E: std::error::Error + Send + Sync + 'static,
{
let server = server.into();
let service = Self::host_info()
.serve_with_lifecycle(
transport,
ClientLifecycleMode::Auto {
preferred_versions: vec![ProtocolVersion::V_2026_07_28],
legacy_version: Some(ProtocolVersion::V_2025_11_25),
},
)
.await
.map_err(|e| ToolError::Unreachable {
tool: ToolId::new(&server, "initialize"),
detail: format!("the MCP server did not initialise: {e}"),
})?;
Self::new(server, Arc::new(service), destination)
}
pub fn new(
server: impl Into<String>,
service: Arc<RunningService<RoleClient, ClientInfo>>,
destination: crate::tools::Destination,
) -> Result<Self, ToolError> {
let server = server.into();
if let Some(info) = service.peer_info() {
let negotiated = &info.protocol_version;
if !Self::KNOWN_VERSIONS.contains(negotiated) {
return Err(ToolError::Unreachable {
tool: ToolId::new(&server, "initialize"),
detail: format!(
"the server negotiated MCP protocol version '{}', which this \
host does not speak — proceeding would issue requests in a \
dialect nobody here implements",
negotiated.as_str()
),
});
}
}
Ok(Self {
server,
service,
access: McpAccess::default(),
timeout: Self::DEFAULT_TIMEOUT,
destination,
})
}
pub const DEFAULT_TIMEOUT: Duration = Duration::from_secs(60);
#[must_use]
pub fn with_timeout(mut self, timeout: Duration) -> Self {
self.timeout = timeout;
self
}
#[must_use]
pub fn negotiated_version(&self) -> Option<String> {
self.service
.peer_info()
.map(|info| info.protocol_version.as_str().to_owned())
}
#[must_use]
pub fn with_access(mut self, access: McpAccess) -> Self {
self.access = access;
self
}
pub fn prompt(
&self,
name: impl Into<String>,
arguments: Value,
) -> Result<McpPrompt, ToolError> {
let name = name.into();
let Some(safety) = self.access.prompts.get(&name).copied() else {
return Err(ToolError::Refused {
tool: ToolId::new(&self.server, format!("prompt/{name}")),
detail: "the operator did not grant this MCP prompt".to_owned(),
});
};
if !arguments.is_object() && !arguments.is_null() {
return Err(ToolError::Refused {
tool: ToolId::new(&self.server, format!("prompt/{name}")),
detail: "MCP prompt arguments must be an object or null".to_owned(),
});
}
Ok(McpPrompt {
server: self.server.clone(),
name,
arguments,
safety,
service: Arc::clone(&self.service),
timeout: self.timeout,
})
}
pub fn resource(&self, uri: impl Into<String>) -> Result<McpResource, ToolError> {
let uri = uri.into();
let Some(safety) = self.access.resources.get(&uri).copied() else {
return Err(ToolError::Refused {
tool: ToolId::new(&self.server, format!("resource/{uri}")),
detail: "the operator did not grant this MCP resource".to_owned(),
});
};
Ok(McpResource {
server: self.server.clone(),
uri,
safety,
service: Arc::clone(&self.service),
timeout: self.timeout,
})
}
pub fn task(
&self,
task: McpTask,
output_sensitivity: Sensitivity,
) -> Result<McpTaskPoll, ToolError> {
self.check_task_server(&task)?;
Ok(McpTaskPoll {
task,
service: Arc::clone(&self.service),
timeout: self.timeout,
output_sensitivity,
})
}
pub fn update_task(
&self,
task: McpTask,
input_responses: InputResponses,
) -> Result<McpTaskUpdate, ToolError> {
self.check_task_server(&task)?;
let Some(safety) = self.access.task_input else {
return Err(ToolError::Refused {
tool: ToolId::new(&self.server, format!("task/{}", task.id)),
detail: "the operator did not grant this server input responses — an \
elicitation is a server asking this plane for data, and nothing \
about the server raising one says it may have an answer"
.to_owned(),
});
};
let arguments =
serde_json::to_value(&input_responses).map_err(|error| ToolError::Refused {
tool: ToolId::new(&self.server, format!("task/{}", task.id)),
detail: format!(
"the input responses could not be serialized for policy \
inspection, so they were not sent: {error}"
),
})?;
Ok(McpTaskUpdate {
task,
input_responses,
arguments,
safety,
service: Arc::clone(&self.service),
timeout: self.timeout,
})
}
pub fn cancel_task(&self, task: McpTask) -> Result<McpTaskCancel, ToolError> {
self.check_task_server(&task)?;
Ok(McpTaskCancel {
task,
service: Arc::clone(&self.service),
timeout: self.timeout,
})
}
fn check_task_server(&self, task: &McpTask) -> Result<(), ToolError> {
if task.server == self.server {
return Ok(());
}
Err(ToolError::Unreachable {
tool: ToolId::new(&task.server, format!("task/{}", task.id)),
detail: format!(
"this client is connected to MCP server '{}', not '{}'",
self.server, task.server
),
})
}
pub async fn discover(&self) -> Result<Vec<(ToolId, Advertised)>, ToolError> {
let listed = bounded(self.timeout, self.service.list_all_tools())
.await
.map_err(|e| Self::classify(&ToolId::new(&self.server, "tools/list"), &e))?;
Ok(listed
.into_iter()
.map(|t| {
let annotations = t.annotations.as_ref();
(
ToolId::new(&self.server, t.name.to_string()),
Advertised {
read_only: annotations.and_then(|a| a.read_only_hint),
destructive: annotations.and_then(|a| a.destructive_hint),
idempotent: annotations.and_then(|a| a.idempotent_hint),
},
)
})
.collect())
}
#[allow(clippy::match_same_arms)]
fn classify(tool: &ToolId, e: &ServiceError) -> ToolError {
let detail = e.to_string();
match e {
ServiceError::McpError(err)
if matches!(
err.code,
ErrorCode::METHOD_NOT_FOUND
| ErrorCode::INVALID_PARAMS
| ErrorCode::INVALID_REQUEST
| ErrorCode::PARSE_ERROR
| ErrorCode::RESOURCE_NOT_FOUND
) =>
{
ToolError::Refused {
tool: tool.clone(),
detail,
}
}
ServiceError::McpError(_) => ToolError::TimedOut {
tool: tool.clone(),
detail,
},
ServiceError::Timeout { .. } => ToolError::TimedOut {
tool: tool.clone(),
detail,
},
ServiceError::Cancelled { .. } => ToolError::TimedOut {
tool: tool.clone(),
detail,
},
ServiceError::TransportSend(_) | ServiceError::TransportClosed => ToolError::TimedOut {
tool: tool.clone(),
detail,
},
ServiceError::UnexpectedResponse => ToolError::Malformed {
tool: tool.clone(),
detail,
},
_ => ToolError::TimedOut {
tool: tool.clone(),
detail,
},
}
}
}
#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
pub struct McpTask {
server: String,
id: String,
}
impl McpTask {
pub fn from_result(server: impl Into<String>, value: &Value) -> Result<Self, String> {
if value.get("resultType").and_then(Value::as_str) != Some("task") {
return Err("MCP result is not a task handle".to_owned());
}
let id = value
.get("taskId")
.and_then(Value::as_str)
.ok_or_else(|| "MCP task handle has no taskId".to_owned())?;
Ok(Self {
server: server.into(),
id: id.to_owned(),
})
}
#[must_use]
pub fn id(&self) -> &str {
&self.id
}
#[must_use]
pub fn server(&self) -> &str {
&self.server
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
#[serde(rename_all = "snake_case")]
#[non_exhaustive]
pub enum McpTaskState {
Working,
InputRequired,
Completed,
Failed,
Cancelled,
}
impl McpTaskState {
fn parse(value: &Value) -> Result<Self, String> {
match value.get("status").and_then(Value::as_str) {
Some("working") => Ok(Self::Working),
Some("input_required") => Ok(Self::InputRequired),
Some("completed") => Ok(Self::Completed),
Some("failed") => Ok(Self::Failed),
Some("cancelled") => Ok(Self::Cancelled),
Some(other) => Err(format!("unknown MCP task state '{other}'")),
None => Err("MCP task result has no status".to_owned()),
}
}
#[must_use]
pub const fn is_terminal(self) -> bool {
matches!(self, Self::Completed | Self::Failed | Self::Cancelled)
}
}
#[derive(Debug, Clone, PartialEq, serde::Serialize, serde::Deserialize)]
pub struct McpTaskSnapshot {
pub task: McpTask,
pub state: McpTaskState,
pub value: Value,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum TaskRetention {
Until { ms: u64 },
Unlimited,
Unstated,
}
impl McpTaskSnapshot {
#[must_use]
pub fn retention(&self) -> TaskRetention {
match self.value.get("ttlMs") {
None => TaskRetention::Unstated,
Some(Value::Null) => TaskRetention::Unlimited,
Some(value) => value
.as_u64()
.map_or(TaskRetention::Unstated, |ms| TaskRetention::Until { ms }),
}
}
#[must_use]
pub fn poll_interval_ms(&self) -> Option<u64> {
self.value.get("pollIntervalMs").and_then(Value::as_u64)
}
}
#[cfg(test)]
mod task_codec_tests {
use super::*;
#[test]
fn a_task_handle_reads_the_fields_the_extension_spells() {
let create = serde_json::json!({
"resultType": "task",
"taskId": "task-7",
"status": "working",
"createdAt": "2026-07-28T09:00:00Z",
"lastUpdatedAt": "2026-07-28T09:00:00Z",
"ttlMs": 60_000,
"pollIntervalMs": 2_500
});
let task = McpTask::from_result("tickets", &create).expect("a task handle");
assert_eq!(task.id(), "task-7");
assert_eq!(
McpTaskState::parse(&create).expect("a status"),
McpTaskState::Working
);
let snapshot = McpTaskSnapshot {
task,
state: McpTaskState::Working,
value: create,
};
assert_eq!(snapshot.retention(), TaskRetention::Until { ms: 60_000 });
assert_eq!(snapshot.poll_interval_ms(), Some(2_500));
}
#[test]
fn an_unlimited_retention_is_not_an_unstated_one() {
let snapshot = |ttl: Value| McpTaskSnapshot {
task: McpTask {
server: "tickets".to_owned(),
id: "task-7".to_owned(),
},
state: McpTaskState::Working,
value: serde_json::json!({ "taskId": "task-7", "status": "working", "ttlMs": ttl }),
};
assert_eq!(snapshot(Value::Null).retention(), TaskRetention::Unlimited);
assert_eq!(
snapshot(serde_json::json!(1)).retention(),
TaskRetention::Until { ms: 1 }
);
let absent = McpTaskSnapshot {
task: McpTask {
server: "tickets".to_owned(),
id: "task-7".to_owned(),
},
state: McpTaskState::Working,
value: serde_json::json!({ "taskId": "task-7", "status": "working" }),
};
assert_eq!(absent.retention(), TaskRetention::Unstated);
}
}
#[derive(Debug)]
pub struct McpPrompt {
server: String,
name: String,
arguments: Value,
safety: McpDataSafety,
service: Arc<RunningService<RoleClient, ClientInfo>>,
timeout: Duration,
}
#[async_trait]
impl Effect for McpPrompt {
type Output = Value;
fn descriptor(&self) -> EffectDescriptor {
EffectDescriptor::new(
"mcp.prompt/get",
serde_json::json!({
"server": self.server,
"name": self.name,
"arguments": self.arguments,
}),
)
}
fn mutates(&self) -> bool {
false
}
fn recovery(&self) -> Recovery {
Recovery::Retry
}
fn max_sensitivity(&self) -> Sensitivity {
self.safety.max_input_sensitivity
}
fn output_sensitivity(&self) -> Sensitivity {
self.safety.output_sensitivity
}
fn trust(&self) -> Trust {
Trust::Untrusted
}
fn sink_arguments(&self) -> Option<&Value> {
Some(&self.arguments)
}
async fn perform(&self) -> Result<Value, EffectError> {
let mut params = GetPromptRequestParams::new(&self.name);
if let Value::Object(arguments) = &self.arguments {
params = params.with_arguments(arguments.clone());
}
match bounded(self.timeout, self.service.get_prompt_once(params)).await {
Ok(GetPromptResponse::Complete(result)) => serde_json::to_value(result)
.map_err(|error| EffectError::Other(error.to_string())),
Ok(GetPromptResponse::InputRequired(_)) => Err(EffectError::Interrupted {
driver: self.server.clone(),
detail: "MCP prompt retrieval requested elicitation; this host does not allow a server to open an ungoverned human-input loop"
.to_owned(),
}),
Ok(_) => Err(EffectError::Interrupted {
driver: self.server.clone(),
detail: "MCP prompt retrieval returned an unknown response variant".to_owned(),
}),
Err(error) => Err(mcp_effect_error(&self.server, &self.name, &error)),
}
}
}
#[derive(Debug)]
pub struct McpResource {
server: String,
uri: String,
safety: McpDataSafety,
service: Arc<RunningService<RoleClient, ClientInfo>>,
timeout: Duration,
}
#[async_trait]
impl Effect for McpResource {
type Output = Value;
fn descriptor(&self) -> EffectDescriptor {
EffectDescriptor::new(
"mcp.resource/read",
serde_json::json!({
"server": self.server,
"uri": self.uri,
}),
)
}
fn mutates(&self) -> bool {
false
}
fn recovery(&self) -> Recovery {
Recovery::Retry
}
fn output_sensitivity(&self) -> Sensitivity {
self.safety.output_sensitivity
}
fn trust(&self) -> Trust {
Trust::Untrusted
}
async fn perform(&self) -> Result<Value, EffectError> {
match bounded(
self.timeout,
self.service
.read_resource_once(ReadResourceRequestParams::new(&self.uri)),
)
.await
{
Ok(ReadResourceResponse::Complete(result)) => serde_json::to_value(result)
.map_err(|error| EffectError::Other(error.to_string())),
Ok(ReadResourceResponse::InputRequired(_)) => Err(EffectError::Interrupted {
driver: self.server.clone(),
detail: "MCP resource retrieval requested elicitation; this host does not allow a server to open an ungoverned human-input loop"
.to_owned(),
}),
Ok(_) => Err(EffectError::Interrupted {
driver: self.server.clone(),
detail: "MCP resource retrieval returned an unknown response variant".to_owned(),
}),
Err(error) => Err(mcp_effect_error(&self.server, &self.uri, &error)),
}
}
}
#[derive(Debug)]
pub struct McpTaskPoll {
task: McpTask,
service: Arc<RunningService<RoleClient, ClientInfo>>,
timeout: Duration,
output_sensitivity: Sensitivity,
}
#[derive(Debug)]
pub struct McpTaskUpdate {
task: McpTask,
input_responses: InputResponses,
arguments: Value,
safety: McpDataSafety,
service: Arc<RunningService<RoleClient, ClientInfo>>,
timeout: Duration,
}
#[async_trait]
impl Effect for McpTaskUpdate {
type Output = ();
fn descriptor(&self) -> EffectDescriptor {
EffectDescriptor::new(
"mcp.task/update",
serde_json::json!({
"server": self.task.server,
"task_id": self.task.id,
"input_responses": self.input_responses,
}),
)
}
fn mutates(&self) -> bool {
true
}
fn recovery(&self) -> Recovery {
Recovery::RequiresOperator
}
fn max_sensitivity(&self) -> Sensitivity {
self.safety.max_input_sensitivity
}
fn sink_arguments(&self) -> Option<&Value> {
Some(&self.arguments)
}
async fn perform(&self) -> Result<(), EffectError> {
bounded(
self.timeout,
self.service.update_task(UpdateTaskParams::new(
&self.task.id,
self.input_responses.clone(),
)),
)
.await
.map_err(|error| mcp_effect_error(&self.task.server, &self.task.id, &error))
}
}
#[derive(Debug)]
pub struct McpTaskCancel {
task: McpTask,
service: Arc<RunningService<RoleClient, ClientInfo>>,
timeout: Duration,
}
#[async_trait]
impl Effect for McpTaskCancel {
type Output = ();
fn descriptor(&self) -> EffectDescriptor {
EffectDescriptor::new(
"mcp.task/cancel",
serde_json::json!({"server": self.task.server, "task_id": self.task.id}),
)
}
fn mutates(&self) -> bool {
true
}
fn recovery(&self) -> Recovery {
Recovery::Retry
}
async fn perform(&self) -> Result<(), EffectError> {
bounded(
self.timeout,
self.service
.cancel_task(CancelTaskParams::new(&self.task.id)),
)
.await
.map_err(|error| mcp_effect_error(&self.task.server, &self.task.id, &error))
}
}
#[async_trait]
impl Effect for McpTaskPoll {
type Output = McpTaskSnapshot;
fn descriptor(&self) -> EffectDescriptor {
EffectDescriptor::new(
"mcp.task/get",
serde_json::json!({"server": self.task.server, "task_id": self.task.id}),
)
}
fn mutates(&self) -> bool {
false
}
fn recovery(&self) -> Recovery {
Recovery::Retry
}
fn trust(&self) -> Trust {
Trust::Untrusted
}
fn output_sensitivity(&self) -> Sensitivity {
self.output_sensitivity
}
async fn perform(&self) -> Result<McpTaskSnapshot, EffectError> {
let result = bounded(
self.timeout,
self.service.get_task(GetTaskParams::new(&self.task.id)),
)
.await
.map_err(|error| mcp_effect_error(&self.task.server, &self.task.id, &error))?;
let value =
serde_json::to_value(result).map_err(|error| EffectError::Other(error.to_string()))?;
let state = McpTaskState::parse(&value).map_err(EffectError::Other)?;
Ok(McpTaskSnapshot {
task: self.task.clone(),
state,
value,
})
}
}
async fn bounded<T>(
timeout: Duration,
call: impl std::future::Future<Output = Result<T, ServiceError>> + Send,
) -> Result<T, ServiceError> {
match tokio::time::timeout(timeout, call).await {
Ok(result) => result,
Err(_) => Err(ServiceError::Timeout { timeout }),
}
}
fn mcp_effect_error(server: &str, operation: &str, error: &ServiceError) -> EffectError {
let tool = ToolId::new(server, operation);
let error = McpClient::classify(&tool, error);
match error {
ToolError::Unreachable { .. } | ToolError::Refused { .. } => {
EffectError::Rejected(error.to_string())
}
ToolError::TimedOut { .. } => EffectError::Interrupted {
driver: server.to_owned(),
detail: error.to_string(),
},
ToolError::Malformed { .. } | ToolError::ToolFailed { .. } => {
EffectError::Performed(error.to_string())
}
}
}
#[async_trait]
impl ToolClient for McpClient {
async fn call(
&self,
tool: &ToolId,
arguments: &Value,
provenance: Option<&crate::core::Provenance>,
) -> Result<Value, ToolError> {
if tool.server != self.server {
return Err(ToolError::Unreachable {
tool: tool.clone(),
detail: format!(
"this client is connected to MCP server '{}', not '{}'",
self.server, tool.server
),
});
}
let object = match arguments {
Value::Object(map) => Some(map.clone()),
Value::Null => None,
other => {
return Err(ToolError::Refused {
tool: tool.clone(),
detail: format!(
"MCP tool arguments must be a JSON object, got {}",
kind_of(other)
),
});
}
};
let mut params = match object {
Some(args) => CallToolRequestParams::new(tool.tool.clone()).with_arguments(args),
None => CallToolRequestParams::new(tool.tool.clone()),
};
if let Some(p) = provenance {
use rmcp::model::RequestParamsMeta;
params.set_meta(rmcp::model::RequestMetaObject(rmcp::model::MetaObject(
p.to_meta(),
)));
}
let result = match bounded(self.timeout, self.service.call_tool_once(params))
.await
.map_err(|e| Self::classify(tool, &e))?
{
CallToolResponse::Complete(result) => result,
CallToolResponse::Task(task) => {
return serde_json::to_value(task).map_err(|error| ToolError::Malformed {
tool: tool.clone(),
detail: format!("MCP task handle could not be represented: {error}"),
});
}
CallToolResponse::InputRequired(_) => {
return Err(ToolError::TimedOut {
tool: tool.clone(),
detail: "the MCP server requested elicitation, which this host does not advertise or answer inside a tool call"
.to_owned(),
});
}
_ => {
return Err(ToolError::TimedOut {
tool: tool.clone(),
detail: "the MCP server returned an unknown tool response variant".to_owned(),
});
}
};
if result.is_error == Some(true) {
return Err(ToolError::ToolFailed {
tool: tool.clone(),
detail: render(&result),
});
}
if let Some(structured) = result.structured_content {
return Ok(structured);
}
serde_json::to_value(&result.content).map_err(|error| ToolError::Malformed {
tool: tool.clone(),
detail: format!("MCP tool result content could not be represented: {error}"),
})
}
fn destination(&self, _tool: &ToolId) -> crate::tools::Destination {
self.destination.clone()
}
}
fn kind_of(v: &Value) -> &'static str {
match v {
Value::Null => "null",
Value::Bool(_) => "a boolean",
Value::Number(_) => "a number",
Value::String(_) => "a string",
Value::Array(_) => "an array",
Value::Object(_) => "an object",
}
}
fn render(result: &rmcp::model::CallToolResult) -> String {
let text = result
.content
.iter()
.filter_map(|c| c.as_text().map(|t| t.text.clone()))
.collect::<Vec<_>>()
.join("\n");
if !text.is_empty() {
return text;
}
if let Some(structured) = &result.structured_content {
return structured.to_string();
}
serde_json::to_string(&result.content).unwrap_or_default()
}
#[cfg(test)]
mod classify_tests {
use super::*;
use crate::core::Disposition;
#[test]
fn a_legacy_resource_not_found_is_a_refusal_not_an_unknown_outcome() {
let tool = ToolId::new("kb", "resource/read");
let error = ServiceError::McpError(rmcp::model::ErrorData::resource_not_found(
"no such resource",
None,
));
let classified = McpClient::classify(&tool, &error);
assert!(
matches!(classified, ToolError::Refused { .. }),
"-32002 classified as: {classified}"
);
assert_eq!(classified.disposition(), Disposition::DidNotHappen);
}
}