#[cfg(feature = "mcp")]
mod mcp;
mod typed;
#[cfg(feature = "mcp")]
pub use mcp::{
McpAccess, McpClient, McpDataSafety, McpPrompt, McpResource, McpTask, McpTaskCancel,
McpTaskPoll, McpTaskSnapshot, McpTaskState, McpTaskUpdate, TaskRetention,
};
#[cfg(feature = "mcp")]
pub use rmcp;
pub use typed::{Tool, ToolBox, ToolFailure};
use std::collections::BTreeMap;
use std::fmt::Debug;
use std::sync::Arc;
use async_trait::async_trait;
use serde::{Deserialize, Serialize};
use serde_json::Value;
use crate::core::{
Disposition, Effect, EffectDescriptor, EffectError, ProtectedField, Recovery, RetryPolicy,
Sensitivity, Trust,
};
#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize)]
pub struct ToolId {
pub server: String,
pub tool: String,
}
pub const TOOL_SCHEME: &str = "tool://";
pub const AGENT_SERVER: &str = "agent";
impl ToolId {
pub fn new(server: impl Into<String>, tool: impl Into<String>) -> Self {
Self {
server: server.into(),
tool: tool.into(),
}
}
#[must_use]
pub fn reference(&self) -> String {
format!("{TOOL_SCHEME}{}/{}", self.server, self.tool)
}
#[must_use]
pub fn parse(reference: &str) -> Option<Self> {
let rest = reference.strip_prefix(TOOL_SCHEME)?;
let (server, tool) = rest.split_once('/')?;
(valid_component(server) && valid_component(tool) && !tool.contains('/'))
.then(|| Self::new(server, tool))
}
#[must_use]
pub fn wire_name(&self) -> String {
format!(
"{}__{}",
wire_component(&self.server),
wire_component(&self.tool)
)
}
}
fn wire_component(value: &str) -> String {
value.replace('.', "-")
}
fn valid_component(value: &str) -> bool {
!value.is_empty()
&& value
.chars()
.all(|c| c.is_ascii_alphanumeric() || c == '_' || c == '.')
&& !value.contains("__")
&& !value.starts_with('_')
&& !value.ends_with('_')
}
impl std::fmt::Display for ToolId {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(f, "{}/{}", self.server, self.tool)
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct ToolSafety {
pub mutates: bool,
pub recovery: Recovery,
pub max_sensitivity: Sensitivity,
pub output_sensitivity: Sensitivity,
pub retry: RetryPolicy,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub protected_fields: Vec<ProtectedField>,
}
impl Default for ToolSafety {
fn default() -> Self {
Self {
mutates: true,
recovery: Recovery::RequiresOperator,
max_sensitivity: Sensitivity::Public,
output_sensitivity: Sensitivity::Public,
retry: RetryPolicy::never(),
protected_fields: Vec::new(),
}
}
}
pub(crate) fn sorted_fields(fields: &[ProtectedField]) -> Vec<ProtectedField> {
let mut out = fields.to_vec();
out.sort_by(|left, right| left.path().cmp(right.path()));
out
}
impl ToolSafety {
#[cfg(feature = "manifest")]
#[must_use]
pub fn from_grant(grant: &crate::manifest::ToolGrant) -> Self {
let defaults = Self::default();
Self {
mutates: grant.mutates,
protected_fields: grant.protected_fields.clone(),
max_sensitivity: grant.max_sensitivity.unwrap_or(defaults.max_sensitivity),
..defaults
}
}
#[must_use]
pub fn read_only() -> Self {
Self {
mutates: false,
recovery: Recovery::Retry,
..Self::default()
}
}
#[must_use]
pub fn recovery(mut self, r: Recovery) -> Self {
self.recovery = r;
self
}
#[must_use]
pub const fn max_sensitivity(mut self, s: Sensitivity) -> Self {
self.max_sensitivity = s;
self
}
#[must_use]
pub const fn output_sensitivity(mut self, s: Sensitivity) -> Self {
self.output_sensitivity = s;
self
}
#[must_use]
pub fn retry(mut self, r: RetryPolicy) -> Self {
self.retry = r;
self
}
#[must_use]
pub fn protect(mut self, field: ProtectedField) -> Self {
assert!(
!self
.protected_fields
.iter()
.any(|existing| existing.path() == field.path()),
"a protected tool argument may be declared only once"
);
self.protected_fields.push(field);
self.protected_fields
.sort_by(|left, right| left.path().cmp(right.path()));
self
}
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
pub struct Advertised {
pub read_only: Option<bool>,
pub destructive: Option<bool>,
pub idempotent: Option<bool>,
}
impl Advertised {
#[must_use]
pub fn overclaims(&self, safety: &ToolSafety) -> bool {
if !safety.mutates {
return false;
}
self.read_only == Some(true)
|| self.destructive == Some(false)
|| (self.idempotent == Some(true)
&& matches!(safety.recovery, Recovery::RequiresOperator))
}
}
#[derive(Debug, thiserror::Error)]
pub enum ToolError {
#[error("could not reach '{tool}': {detail}")]
Unreachable { tool: ToolId, detail: String },
#[error("'{tool}' refused the request: {detail}")]
Refused { tool: ToolId, detail: String },
#[error("'{tool}' did not answer in time: {detail}")]
TimedOut { tool: ToolId, detail: String },
#[error("'{tool}' reported an error: {detail}")]
ToolFailed { tool: ToolId, detail: String },
#[error("'{tool}' returned a malformed response: {detail}")]
Malformed { tool: ToolId, detail: String },
}
impl ToolError {
#[must_use]
pub const fn disposition(&self) -> Disposition {
match self {
Self::Unreachable { .. } | Self::Refused { .. } => Disposition::DidNotHappen,
Self::TimedOut { .. } => Disposition::InDoubt,
Self::ToolFailed { .. } | Self::Malformed { .. } => Disposition::Landed,
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum Destination {
Local,
Remote(String),
}
impl Destination {
pub fn remote(host: impl Into<String>) -> Self {
Self::Remote(host.into())
}
#[must_use]
pub fn host(&self) -> Option<&str> {
match self {
Self::Local => None,
Self::Remote(host) => Some(host),
}
}
}
#[async_trait]
pub trait ToolClient: Send + Sync + Debug {
async fn call(
&self,
tool: &ToolId,
arguments: &Value,
provenance: Option<&crate::core::Provenance>,
) -> Result<Value, ToolError>;
fn destination(&self, tool: &ToolId) -> Destination;
}
#[derive(Debug, Default, Clone)]
pub struct ToolCatalog {
entries: BTreeMap<ToolId, (ToolSafety, Advertised)>,
#[cfg(feature = "manifest")]
declarations: BTreeMap<ToolId, (String, Value)>,
}
impl ToolCatalog {
#[must_use]
pub fn new() -> Self {
Self::default()
}
#[cfg(feature = "manifest")]
#[must_use]
pub fn from_manifest(manifest: &crate::manifest::Manifest) -> Self {
let mut catalog = Self::new();
for grant in &manifest.spec.tools {
let Some(id) = ToolId::parse(&grant.reference) else {
continue;
};
catalog = catalog.allow(id.clone(), ToolSafety::from_grant(grant));
if grant.description.is_some() || grant.arguments.is_some() {
catalog = catalog.declare(
id,
grant.description.clone().unwrap_or_default(),
grant
.arguments
.clone()
.unwrap_or_else(|| serde_json::json!({ "type": "object" })),
);
}
}
catalog
}
#[must_use]
pub fn allow(mut self, id: ToolId, safety: ToolSafety) -> Self {
assert!(
valid_component(&id.server) && valid_component(&id.tool),
"tool id '{id}' cannot be rendered as a wire name: server and tool may \
hold letters, digits, `_` and `.`, with no `-`, no `__`, and no leading \
or trailing `_` — a name outside that set would collide with, or split \
on, the `__` separator and the `.`→`-` rendering"
);
self.entries.insert(id, (safety, Advertised::default()));
self
}
#[cfg(feature = "manifest")]
pub(crate) fn declare(mut self, id: ToolId, description: String, arguments: Value) -> Self {
self.declarations.insert(id, (description, arguments));
self
}
#[cfg(feature = "manifest")]
pub(crate) fn declaration(&self, id: &ToolId) -> Option<(&str, &Value)> {
self.declarations
.get(id)
.map(|(description, arguments)| (description.as_str(), arguments))
}
#[must_use]
pub fn entries(self) -> Vec<(ToolId, ToolSafety)> {
self.entries
.into_iter()
.map(|(id, (safety, _))| (id, safety))
.collect()
}
#[must_use]
pub fn resolve(&self, name: &str) -> Option<ToolId> {
self.entries
.keys()
.find(|id| id.wire_name() == name)
.cloned()
}
#[must_use]
pub fn resolve_reference(&self, reference: &str) -> Option<ToolId> {
self.entries
.keys()
.find(|id| id.reference() == reference)
.cloned()
}
pub fn granted(&self) -> impl Iterator<Item = &ToolId> {
self.entries.keys()
}
#[must_use]
pub fn observed(mut self, id: &ToolId, advertised: Advertised) -> Self {
if let Some(entry) = self.entries.get_mut(id) {
if advertised.overclaims(&entry.0) {
tracing::warn!(
tool = %id,
"this tool's server now advertises more safety than the \
operator granted; the grant still rules, but the \
advertisement changed"
);
}
entry.1 = advertised;
}
self
}
#[must_use]
pub fn safety(&self, id: &ToolId) -> Option<&ToolSafety> {
self.entries.get(id).map(|(s, _)| s)
}
#[must_use]
pub fn advertised(&self, id: &ToolId) -> Option<&Advertised> {
self.entries.get(id).map(|(_, a)| a)
}
pub fn overclaiming(&self) -> impl Iterator<Item = &ToolId> {
self.entries
.iter()
.filter(|(_, (safety, adv))| adv.overclaims(safety))
.map(|(id, _)| id)
}
#[must_use]
pub fn len(&self) -> usize {
self.entries.len()
}
#[must_use]
pub fn is_empty(&self) -> bool {
self.entries.is_empty()
}
}
#[derive(Debug, Default)]
pub struct ToolRouter {
routes: BTreeMap<String, Arc<dyn ToolClient>>,
}
impl ToolRouter {
#[must_use]
pub fn new() -> Self {
Self::default()
}
#[must_use]
pub fn server(mut self, name: impl Into<String>, client: Arc<dyn ToolClient>) -> Self {
let name = name.into();
assert!(
!self.routes.contains_key(&name),
"tool server '{name}' is routed twice — one of the two transports would \
silently never be called"
);
self.routes.insert(name, client);
self
}
#[must_use]
pub fn toolbox(self, tools: &Arc<ToolBox>) -> Self {
tools
.servers()
.map(ToOwned::to_owned)
.collect::<Vec<String>>()
.into_iter()
.fold(self, |router, name| {
router.server(name, Arc::clone(tools) as Arc<dyn ToolClient>)
})
}
pub fn servers(&self) -> impl Iterator<Item = &str> {
self.routes.keys().map(String::as_str)
}
}
impl ToolRouter {
fn route(&self, tool: &ToolId) -> Option<&Arc<dyn ToolClient>> {
self.routes.get(&tool.server)
}
}
#[async_trait]
impl ToolClient for ToolRouter {
async fn call(
&self,
tool: &ToolId,
arguments: &Value,
provenance: Option<&crate::core::Provenance>,
) -> Result<Value, ToolError> {
let Some(client) = self.route(tool) else {
return Err(ToolError::Unreachable {
tool: tool.clone(),
detail: format!(
"no transport is wired for tool server '{}'; this plane routes {:?}",
tool.server,
self.routes.keys().collect::<Vec<_>>()
),
});
};
client.call(tool, arguments, provenance).await
}
fn destination(&self, tool: &ToolId) -> Destination {
self.route(tool)
.map_or(Destination::Local, |client| client.destination(tool))
}
}
#[derive(Debug)]
pub struct ToolCall {
id: ToolId,
arguments: Value,
safety: ToolSafety,
client: std::sync::Arc<dyn ToolClient>,
provenance: Option<crate::core::Provenance>,
}
impl ToolCall {
pub fn prepare(
catalog: &ToolCatalog,
client: std::sync::Arc<dyn ToolClient>,
id: ToolId,
arguments: Value,
egress: Option<&crate::core::Egress>,
) -> Result<Self, ToolError> {
let Some(safety) = catalog.safety(&id) else {
return Err(ToolError::Unreachable {
detail: "this tool is not in the catalogue; a tool nobody declared is a \
tool nobody has reasoned about"
.into(),
tool: id,
});
};
if let Some(egress) = egress
&& let Some(host) = client.destination(&id).host().map(ToOwned::to_owned)
&& let Err(refused) = egress.permits(Some(&host))
{
return Err(ToolError::Unreachable {
detail: format!(
"the transport for server '{}' was refused — {refused}",
id.server
),
tool: id,
});
}
Ok(Self {
safety: safety.clone(),
id,
arguments,
client,
provenance: None,
})
}
}
#[async_trait]
impl Effect for ToolCall {
fn gen_ai_operation(&self) -> Option<&'static str> {
Some(crate::runtime::telemetry::GEN_AI_EXECUTE_TOOL)
}
fn gen_ai_request(&self) -> Option<crate::core::GenAiRequest> {
Some(crate::core::GenAiRequest {
provider: None,
name: self.id.reference(),
})
}
type Output = Value;
fn descriptor(&self) -> EffectDescriptor {
EffectDescriptor::new(
"tool.call",
serde_json::json!({
"server": self.id.server,
"tool": self.id.tool,
"arguments": self.arguments,
"protected_fields": sorted_fields(&self.safety.protected_fields),
}),
)
}
fn mutates(&self) -> bool {
self.safety.mutates
}
fn recovery(&self) -> Recovery {
self.safety.recovery.clone()
}
fn retry(&self) -> RetryPolicy {
self.safety.retry
}
fn max_sensitivity(&self) -> Sensitivity {
self.safety.max_sensitivity
}
fn sink_arguments(&self) -> Option<&Value> {
Some(&self.arguments)
}
fn protected_fields(&self) -> &[ProtectedField] {
&self.safety.protected_fields
}
fn output_sensitivity(&self) -> Sensitivity {
self.safety.output_sensitivity
}
fn trust(&self) -> Trust {
Trust::Untrusted
}
fn source(&self) -> crate::core::SourceId {
crate::core::SourceId::new(self.id.reference())
}
fn attach(&mut self, provenance: &crate::core::Provenance) {
self.provenance = Some(provenance.clone());
}
async fn perform(&self) -> Result<Value, EffectError> {
self.client
.call(&self.id, &self.arguments, self.provenance.as_ref())
.await
.map_err(|e| {
let detail = e.to_string();
match e.disposition() {
Disposition::DidNotHappen => EffectError::Rejected(detail),
Disposition::InDoubt => EffectError::Interrupted {
driver: self.id.to_string(),
detail,
},
Disposition::Landed => EffectError::Performed(detail),
}
})
}
}
#[cfg(test)]
mod overclaim_tests {
use super::*;
#[test]
fn a_server_claiming_more_safety_than_granted_is_flagged_per_clause() {
let mutating = ToolSafety::default(); let operator_needs_operator = ToolSafety::default().recovery(Recovery::RequiresOperator);
assert!(
Advertised {
read_only: Some(true),
..Advertised::default()
}
.overclaims(&mutating)
);
assert!(
Advertised {
destructive: Some(false),
..Advertised::default()
}
.overclaims(&mutating)
);
assert!(
Advertised {
idempotent: Some(true),
..Advertised::default()
}
.overclaims(&operator_needs_operator)
);
assert!(
!Advertised {
destructive: Some(true),
read_only: Some(false),
idempotent: Some(false),
}
.overclaims(&mutating)
);
assert!(!Advertised::default().overclaims(&mutating));
assert!(
!Advertised {
read_only: Some(true),
destructive: Some(false),
idempotent: Some(true),
}
.overclaims(&ToolSafety::read_only())
);
}
}