use std::{collections::HashMap, sync::Arc};
use schemars::{JsonSchema, Schema, schema_for};
use rig_core::{
memory::ConversationMemory,
message::ToolChoice,
vector_store::{VectorSearchRequest, VectorStoreIndexDyn},
};
use crate::{
agent::hook::{
AgentHook, CompletionCall, CompletionCallAction, HookContext, HookStack, RequestPatch,
},
completion::{CompletionModel, Document},
tool::{
DynamicTool, PortableDynamicTool, Tool, ToolSet,
server::{ToolServer, ToolServerHandle},
},
};
#[cfg(all(feature = "rmcp", not(target_family = "wasm")))]
#[cfg_attr(docsrs, doc(cfg(feature = "rmcp")))]
use crate::tool::rmcp::McpTool as RmcpTool;
use super::{Agent, OutputMode};
struct DynamicContext<I> {
samples: usize,
index: I,
}
impl<I> AgentHook for DynamicContext<I>
where
I: VectorStoreIndexDyn,
{
async fn on_completion_call(
&self,
_ctx: &HookContext,
event: CompletionCall<'_>,
) -> CompletionCallAction {
let query = event.prompt.rag_text().or_else(|| {
event
.history
.iter()
.rev()
.find_map(|message| message.rag_text())
});
let Some(query) = query else {
return CompletionCallAction::continue_run();
};
let request = VectorSearchRequest::builder()
.query(query)
.samples(self.samples as u64)
.build();
match self.index.top_n(request).await {
Ok(results) => CompletionCallAction::patch(RequestPatch::new().extra_context(
results.into_iter().map(|(_, id, value)| Document {
id,
text:
serde_json::to_string_pretty(&value).unwrap_or_else(|_| value.to_string()),
additional_props: Default::default(),
}),
)),
Err(error) => {
CompletionCallAction::stop(format!("failed to retrieve dynamic context: {error}"))
}
}
}
}
#[cfg(all(feature = "rmcp", not(target_family = "wasm")))]
fn build_rmcp_tools(
tools: Vec<rmcp::model::Tool>,
client: rmcp::service::ServerSink,
timeout: Option<std::time::Duration>,
) -> Vec<(String, RmcpTool)> {
tools
.into_iter()
.map(|tool| {
let name = tool.name.to_string();
let rmcp_tool = RmcpTool::from_mcp_server(tool, client.clone()).with_timeout(timeout);
(name, rmcp_tool)
})
.collect()
}
#[derive(Default)]
pub struct NoToolConfig;
pub struct WithToolServerHandle {
handle: ToolServerHandle,
}
pub struct WithBuilderTools {
tools: ToolSet,
retrieval_indexes: Vec<(usize, Arc<dyn VectorStoreIndexDyn + Send + Sync>)>,
}
pub struct AgentBuilder<M, ToolState = NoToolConfig>
where
M: CompletionModel,
{
name: Option<String>,
description: Option<String>,
model: M,
preamble: Option<String>,
static_context: Vec<Document>,
additional_params: Option<serde_json::Value>,
record_telemetry_content: bool,
max_tokens: Option<u64>,
temperature: Option<f64>,
tool_choice: Option<ToolChoice>,
default_max_turns: Option<usize>,
tool_state: ToolState,
hooks: HookStack,
output_schema: Option<schemars::Schema>,
output_mode: OutputMode,
memory: Option<Arc<dyn ConversationMemory>>,
default_conversation_id: Option<String>,
}
impl<M, ToolState> AgentBuilder<M, ToolState>
where
M: CompletionModel,
{
pub fn name(mut self, name: &str) -> Self {
self.name = Some(name.into());
self
}
pub fn description(mut self, description: &str) -> Self {
self.description = Some(description.into());
self
}
pub fn preamble(mut self, preamble: &str) -> Self {
self.preamble = Some(preamble.into());
self
}
pub fn without_preamble(mut self) -> Self {
self.preamble = None;
self
}
pub fn append_preamble(mut self, doc: &str) -> Self {
self.preamble = Some(format!("{}\n{}", self.preamble.unwrap_or_default(), doc));
self
}
pub fn context(mut self, doc: &str) -> Self {
self.static_context.push(Document {
id: format!("static_doc_{}", self.static_context.len()),
text: doc.into(),
additional_props: HashMap::new(),
});
self
}
pub fn dynamic_context<I>(self, samples: usize, index: I) -> Self
where
I: VectorStoreIndexDyn + 'static,
{
self.add_hook(DynamicContext { samples, index })
}
pub fn tool_choice(mut self, tool_choice: ToolChoice) -> Self {
self.tool_choice = Some(tool_choice);
self
}
pub fn default_max_turns(mut self, default_max_turns: usize) -> Self {
self.default_max_turns = Some(default_max_turns);
self
}
pub fn temperature(mut self, temperature: f64) -> Self {
self.temperature = Some(temperature);
self
}
pub fn max_tokens(mut self, max_tokens: u64) -> Self {
self.max_tokens = Some(max_tokens);
self
}
pub fn additional_params(mut self, params: serde_json::Value) -> Self {
self.additional_params = Some(params);
self
}
pub fn record_content_telemetry(mut self, enabled: bool) -> Self {
self.record_telemetry_content = enabled;
self
}
pub fn output_schema<T>(mut self) -> Self
where
T: JsonSchema,
{
self.output_schema = Some(schema_for!(T));
self
}
pub fn output_schema_raw(mut self, schema: Schema) -> Self {
self.output_schema = Some(schema);
self
}
pub fn output_mode(mut self, mode: OutputMode) -> Self {
self.output_mode = mode;
self
}
pub fn memory<B>(mut self, memory: B) -> Self
where
B: ConversationMemory + 'static,
{
self.memory = Some(Arc::new(memory));
self
}
pub fn conversation(mut self, id: impl Into<String>) -> Self {
self.default_conversation_id = Some(id.into());
self
}
pub fn add_hook<H>(mut self, hook: H) -> Self
where
H: AgentHook + 'static,
{
self.hooks.push(hook);
self
}
}
impl<M> AgentBuilder<M, NoToolConfig>
where
M: CompletionModel,
{
pub fn new(model: M) -> Self {
Self {
name: None,
description: None,
model,
preamble: None,
static_context: vec![],
temperature: None,
max_tokens: None,
additional_params: None,
record_telemetry_content: false,
tool_choice: None,
default_max_turns: None,
tool_state: NoToolConfig,
hooks: HookStack::new(),
output_schema: None,
output_mode: OutputMode::default(),
memory: None,
default_conversation_id: None,
}
}
}
impl<M> AgentBuilder<M, NoToolConfig>
where
M: CompletionModel,
{
pub fn tool_server_handle(
self,
handle: ToolServerHandle,
) -> AgentBuilder<M, WithToolServerHandle> {
AgentBuilder {
name: self.name,
description: self.description,
model: self.model,
preamble: self.preamble,
static_context: self.static_context,
additional_params: self.additional_params,
record_telemetry_content: self.record_telemetry_content,
max_tokens: self.max_tokens,
temperature: self.temperature,
tool_choice: self.tool_choice,
default_max_turns: self.default_max_turns,
tool_state: WithToolServerHandle { handle },
hooks: self.hooks,
output_schema: self.output_schema,
output_mode: self.output_mode,
memory: self.memory,
default_conversation_id: self.default_conversation_id,
}
}
pub fn tool<T>(self, tool: T) -> AgentBuilder<M, WithBuilderTools>
where
T: Tool + 'static,
{
let mut tools = ToolSet::default();
tools.add_tool(tool);
AgentBuilder {
name: self.name,
description: self.description,
model: self.model,
preamble: self.preamble,
static_context: self.static_context,
additional_params: self.additional_params,
record_telemetry_content: self.record_telemetry_content,
max_tokens: self.max_tokens,
temperature: self.temperature,
tool_choice: self.tool_choice,
default_max_turns: self.default_max_turns,
tool_state: WithBuilderTools {
tools,
retrieval_indexes: vec![],
},
hooks: self.hooks,
output_schema: self.output_schema,
output_mode: self.output_mode,
memory: self.memory,
default_conversation_id: self.default_conversation_id,
}
}
pub fn dynamic_tool(self, tool: DynamicTool) -> AgentBuilder<M, WithBuilderTools> {
self.dynamic_tools(vec![tool])
}
pub fn portable_dynamic_tool(
self,
tool: PortableDynamicTool,
) -> AgentBuilder<M, WithBuilderTools> {
self.dynamic_tool(DynamicTool::from_portable(tool))
}
pub fn dynamic_tools(self, tools: Vec<DynamicTool>) -> AgentBuilder<M, WithBuilderTools> {
let tools = ToolSet::from_dynamic_tools(tools);
AgentBuilder {
name: self.name,
description: self.description,
model: self.model,
preamble: self.preamble,
static_context: self.static_context,
additional_params: self.additional_params,
record_telemetry_content: self.record_telemetry_content,
max_tokens: self.max_tokens,
temperature: self.temperature,
tool_choice: self.tool_choice,
default_max_turns: self.default_max_turns,
hooks: self.hooks,
output_schema: self.output_schema,
output_mode: self.output_mode,
memory: self.memory,
default_conversation_id: self.default_conversation_id,
tool_state: WithBuilderTools {
tools,
retrieval_indexes: vec![],
},
}
}
#[cfg(all(feature = "rmcp", not(target_family = "wasm")))]
#[cfg_attr(docsrs, doc(cfg(feature = "rmcp")))]
pub fn rmcp_tool(
self,
tool: rmcp::model::Tool,
client: rmcp::service::ServerSink,
) -> AgentBuilder<M, WithBuilderTools> {
self.rmcp_tool_with_timeout(tool, client, crate::tool::rmcp::DEFAULT_MCP_TOOL_TIMEOUT)
}
#[cfg(all(feature = "rmcp", not(target_family = "wasm")))]
#[cfg_attr(docsrs, doc(cfg(feature = "rmcp")))]
pub fn rmcp_tool_with_timeout(
self,
tool: rmcp::model::Tool,
client: rmcp::service::ServerSink,
timeout: impl Into<Option<std::time::Duration>>,
) -> AgentBuilder<M, WithBuilderTools> {
self.with_rmcp_toolset(build_rmcp_tools(vec![tool], client, timeout.into()))
}
#[cfg(all(feature = "rmcp", not(target_family = "wasm")))]
#[cfg_attr(docsrs, doc(cfg(feature = "rmcp")))]
pub fn rmcp_tools(
self,
tools: Vec<rmcp::model::Tool>,
client: rmcp::service::ServerSink,
) -> AgentBuilder<M, WithBuilderTools> {
self.rmcp_tools_with_timeout(tools, client, crate::tool::rmcp::DEFAULT_MCP_TOOL_TIMEOUT)
}
#[cfg(all(feature = "rmcp", not(target_family = "wasm")))]
#[cfg_attr(docsrs, doc(cfg(feature = "rmcp")))]
pub fn rmcp_tools_with_timeout(
self,
tools: Vec<rmcp::model::Tool>,
client: rmcp::service::ServerSink,
timeout: impl Into<Option<std::time::Duration>>,
) -> AgentBuilder<M, WithBuilderTools> {
self.with_rmcp_toolset(build_rmcp_tools(tools, client, timeout.into()))
}
#[cfg(all(feature = "rmcp", not(target_family = "wasm")))]
fn with_rmcp_toolset(
self,
built: Vec<(String, RmcpTool)>,
) -> AgentBuilder<M, WithBuilderTools> {
AgentBuilder {
name: self.name,
description: self.description,
model: self.model,
preamble: self.preamble,
static_context: self.static_context,
additional_params: self.additional_params,
record_telemetry_content: self.record_telemetry_content,
max_tokens: self.max_tokens,
temperature: self.temperature,
tool_choice: self.tool_choice,
default_max_turns: self.default_max_turns,
hooks: self.hooks,
output_schema: self.output_schema,
output_mode: self.output_mode,
memory: self.memory,
default_conversation_id: self.default_conversation_id,
tool_state: WithBuilderTools {
tools: {
let mut set = ToolSet::default();
for (_, tool) in built {
set.add_erased(std::sync::Arc::new(tool));
}
set
},
retrieval_indexes: vec![],
},
}
}
pub fn retrieved_tools(
self,
sample: usize,
index: impl VectorStoreIndexDyn + Send + Sync + 'static,
toolset: ToolSet,
) -> AgentBuilder<M, WithBuilderTools> {
let mut tools = ToolSet::default();
tools.add_retrievable_tools(toolset);
AgentBuilder {
name: self.name,
description: self.description,
model: self.model,
preamble: self.preamble,
static_context: self.static_context,
additional_params: self.additional_params,
record_telemetry_content: self.record_telemetry_content,
max_tokens: self.max_tokens,
temperature: self.temperature,
tool_choice: self.tool_choice,
default_max_turns: self.default_max_turns,
hooks: self.hooks,
output_schema: self.output_schema,
output_mode: self.output_mode,
memory: self.memory,
default_conversation_id: self.default_conversation_id,
tool_state: WithBuilderTools {
tools,
retrieval_indexes: vec![(sample, Arc::new(index))],
},
}
}
pub fn build(self) -> Agent<M> {
let tool_server_handle = ToolServer::new().run();
Agent {
name: self.name,
description: self.description,
model: Arc::new(self.model),
preamble: self.preamble,
static_context: self.static_context,
temperature: self.temperature,
max_tokens: self.max_tokens,
additional_params: self.additional_params,
record_telemetry_content: self.record_telemetry_content,
tool_choice: self.tool_choice,
tool_server_handle,
default_max_turns: self.default_max_turns,
hooks: self.hooks,
output_schema: self.output_schema,
output_mode: self.output_mode,
memory: self.memory,
default_conversation_id: self.default_conversation_id,
}
}
}
impl<M> AgentBuilder<M, WithToolServerHandle>
where
M: CompletionModel,
{
pub fn build(self) -> Agent<M> {
Agent {
name: self.name,
description: self.description,
model: Arc::new(self.model),
preamble: self.preamble,
static_context: self.static_context,
temperature: self.temperature,
max_tokens: self.max_tokens,
additional_params: self.additional_params,
record_telemetry_content: self.record_telemetry_content,
tool_choice: self.tool_choice,
tool_server_handle: self.tool_state.handle,
default_max_turns: self.default_max_turns,
hooks: self.hooks,
output_schema: self.output_schema,
output_mode: self.output_mode,
memory: self.memory,
default_conversation_id: self.default_conversation_id,
}
}
}
impl<M> AgentBuilder<M, WithBuilderTools>
where
M: CompletionModel,
{
pub fn tool<T>(mut self, tool: T) -> Self
where
T: Tool + 'static,
{
self.tool_state.tools.add_tool(tool);
self
}
pub fn dynamic_tool(mut self, tool: DynamicTool) -> Self {
self.tool_state.tools.add_dynamic_tool(tool);
self
}
pub fn portable_dynamic_tool(mut self, tool: PortableDynamicTool) -> Self {
self.tool_state.tools.add_portable_dynamic_tool(tool);
self
}
pub fn dynamic_tools(mut self, tools: Vec<DynamicTool>) -> Self {
let tools = ToolSet::from_dynamic_tools(tools);
self.tool_state.tools.add_tools(tools);
self
}
#[cfg(all(feature = "rmcp", not(target_family = "wasm")))]
#[cfg_attr(docsrs, doc(cfg(feature = "rmcp")))]
pub fn rmcp_tools(
self,
tools: Vec<rmcp::model::Tool>,
client: rmcp::service::ServerSink,
) -> Self {
self.rmcp_tools_with_timeout(tools, client, crate::tool::rmcp::DEFAULT_MCP_TOOL_TIMEOUT)
}
#[cfg(all(feature = "rmcp", not(target_family = "wasm")))]
#[cfg_attr(docsrs, doc(cfg(feature = "rmcp")))]
pub fn rmcp_tools_with_timeout(
self,
tools: Vec<rmcp::model::Tool>,
client: rmcp::service::ServerSink,
timeout: impl Into<Option<std::time::Duration>>,
) -> Self {
self.add_rmcp_tools(build_rmcp_tools(tools, client, timeout.into()))
}
#[cfg(all(feature = "rmcp", not(target_family = "wasm")))]
fn add_rmcp_tools(mut self, built: Vec<(String, RmcpTool)>) -> Self {
for (_, tool) in built {
self.tool_state.tools.add_erased(std::sync::Arc::new(tool));
}
self
}
pub fn retrieved_tools(
mut self,
sample: usize,
index: impl VectorStoreIndexDyn + Send + Sync + 'static,
toolset: ToolSet,
) -> Self {
self.tool_state
.retrieval_indexes
.push((sample, Arc::new(index)));
self.tool_state.tools.add_retrievable_tools(toolset);
self
}
pub fn build(self) -> Agent<M> {
let tool_server_handle = ToolServer::new()
.add_tools(self.tool_state.tools)
.add_retrieval_indexes(self.tool_state.retrieval_indexes)
.run();
Agent {
name: self.name,
description: self.description,
model: Arc::new(self.model),
preamble: self.preamble,
static_context: self.static_context,
temperature: self.temperature,
max_tokens: self.max_tokens,
additional_params: self.additional_params,
record_telemetry_content: self.record_telemetry_content,
tool_choice: self.tool_choice,
tool_server_handle,
default_max_turns: self.default_max_turns,
hooks: self.hooks,
output_schema: self.output_schema,
output_mode: self.output_mode,
memory: self.memory,
default_conversation_id: self.default_conversation_id,
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::test_utils::{MockAddTool, MockCompletionModel, MockSubtractTool, MockToolIndex};
use crate::tool::{ToolContext, ToolExecutionError};
#[derive(Clone)]
struct BuilderHook;
impl AgentHook for BuilderHook {}
#[test]
fn hook_can_be_set_after_tool_configuration() {
let _agent = AgentBuilder::new(MockCompletionModel::text("ok"))
.tool(MockAddTool)
.add_hook(BuilderHook)
.build();
}
struct NamedTool;
impl NamedTool {
fn new() -> Self {
Self
}
}
impl Tool for NamedTool {
const NAME: &'static str = "registered_named";
type Error = rig::tool::ToolExecutionError;
type Args = serde_json::Value;
type Output = String;
fn description(&self) -> String {
"uses its canonical name".to_string()
}
fn parameters(&self) -> serde_json::Value {
serde_json::json!({"type": "object", "properties": {}})
}
async fn call(
&self,
_context: &mut ToolContext,
_args: Self::Args,
) -> Result<Self::Output, ToolExecutionError> {
Ok("ok".to_string())
}
}
#[tokio::test]
async fn typed_tool_builder_paths_advertise_canonical_name() {
for agent in [
AgentBuilder::new(MockCompletionModel::text("ok"))
.tool(NamedTool::new())
.build(),
AgentBuilder::new(MockCompletionModel::text("ok"))
.tool(MockAddTool)
.tool(NamedTool::new())
.build(),
] {
let definitions = agent.tool_server_handle.get_tool_defs(None).await.unwrap();
assert!(
definitions
.iter()
.any(|definition| definition.name == NamedTool::NAME),
"the provider definitions dropped the canonical tool name"
);
let mut context = ToolContext::new();
let result = agent
.tool_server_handle
.execute(NamedTool::NAME, "{}", &mut context)
.await;
assert!(result.is_success());
assert_eq!(result.output().as_text(), Some("ok"));
}
}
#[tokio::test]
async fn retrieved_tools_are_exposed_only_for_prompted_retrieval() {
let retrieval_only = AgentBuilder::new(MockCompletionModel::text("ok"))
.retrieved_tools(
1,
MockToolIndex::new(["add"]),
ToolSet::from_tools(vec![MockAddTool]),
)
.build();
assert!(
retrieval_only
.tool_server_handle
.get_tool_defs(None)
.await
.unwrap()
.is_empty()
);
let agent = AgentBuilder::new(MockCompletionModel::text("ok"))
.tool(MockSubtractTool)
.retrieved_tools(
1,
MockToolIndex::new(["add"]),
ToolSet::from_tools(vec![MockAddTool]),
)
.build();
let always = agent.tool_server_handle.get_tool_defs(None).await.unwrap();
assert_eq!(
always
.iter()
.map(|definition| definition.name.as_str())
.collect::<Vec<_>>(),
vec!["subtract"]
);
let with_retrieval = agent
.tool_server_handle
.get_tool_defs(Some("add two numbers".to_string()))
.await
.unwrap();
assert_eq!(
with_retrieval
.iter()
.map(|definition| definition.name.as_str())
.collect::<Vec<_>>(),
vec!["add", "subtract"]
);
}
#[cfg(all(feature = "rmcp", not(target_family = "wasm")))]
#[tokio::test]
async fn build_rmcp_tools_threads_timeout_into_built_tools() {
use crate::tool::rmcp::DEFAULT_MCP_TOOL_TIMEOUT;
use crate::tool::{ToolContext, ToolErrorKind, server::ToolServer};
use rmcp::model::{
CallToolRequestParams, CallToolResult, ClientInfo, ErrorData, Implementation,
ProtocolVersion, ServerCapabilities, ServerInfo, Tool,
};
use rmcp::service::RequestContext;
use rmcp::{RoleServer, ServerHandler, ServiceExt};
use std::sync::Arc;
use std::time::Duration;
#[derive(Clone)]
struct HangingServer;
impl ServerHandler for HangingServer {
fn get_info(&self) -> ServerInfo {
ServerInfo::new(ServerCapabilities::builder().enable_tools().build())
.with_protocol_version(ProtocolVersion::LATEST)
.with_server_info(Implementation::new("builder-timeout-test", "0.1.0"))
}
async fn call_tool(
&self,
_request: CallToolRequestParams,
_context: RequestContext<RoleServer>,
) -> Result<CallToolResult, ErrorData> {
std::future::pending::<Result<CallToolResult, ErrorData>>().await
}
}
fn tool(name: &str) -> Tool {
Tool::new(
name.to_string(),
String::new(),
Arc::new(serde_json::Map::new()),
)
}
let (c2s, sfc) = tokio::io::duplex(8192);
let (s2c, cfs) = tokio::io::duplex(8192);
let server_task = tokio::spawn(async move {
let running = HangingServer.serve((sfc, s2c)).await.expect("server start");
running.waiting().await.expect("server error");
});
let client = ClientInfo::default()
.serve((cfs, c2s))
.await
.expect("client connect");
let peer = client.peer().clone();
let built_default = build_rmcp_tools(
vec![tool("a")],
peer.clone(),
Some(DEFAULT_MCP_TOOL_TIMEOUT),
);
assert_eq!(built_default[0].1.timeout(), Some(DEFAULT_MCP_TOOL_TIMEOUT));
let built_none = build_rmcp_tools(vec![tool("b")], peer.clone(), None);
assert_eq!(built_none[0].1.timeout(), None);
let built = build_rmcp_tools(
vec![tool("hang_forever")],
peer,
Some(Duration::from_millis(200)),
);
assert_eq!(built.len(), 1);
assert_eq!(built[0].0, "hang_forever");
let handle = ToolServer::new().run();
handle
.add_erased_tool(Arc::new(built.into_iter().next().unwrap().1))
.await;
let timed = tokio::time::timeout(Duration::from_secs(5), async {
let mut context = ToolContext::new();
handle.execute("hang_forever", "{}", &mut context).await
})
.await;
let result = timed.expect("built tool hung past the safety timeout");
assert!(result.is_error_kind(ToolErrorKind::Timeout));
assert!(result.output().render().contains("timed out"));
drop(client);
server_task.abort();
}
}