use std::sync::Arc;
use agentplane::manifest::Manifest;
use agentplane::model::ModelProvider;
use agentplane::model::fake::FakeProvider;
use agentplane::prelude::*;
use agentplane::runtime::Agent;
use agentplane::tools::{McpClient, Tool, ToolBox, ToolClient, ToolFailure, ToolId};
use rmcp::handler::server::ServerHandler;
use rmcp::model::{
CallToolRequestParams, CallToolResponse, CallToolResult, ContentBlock, ErrorData,
Implementation, ListToolsResult, PaginatedRequestParams, ProtocolVersion, ServerCapabilities,
ServerInfo, Tool as McpTool, ToolAnnotations,
};
use rmcp::serve_server;
use rmcp::service::{RequestContext, RoleServer};
use serde_json::{Value, json};
const TELLER: &str = r#"
apiVersion: agentplane.hupe1980.github.io/v1alpha1
kind: Agent
metadata: { name: teller, version: "1.0.0" }
spec:
identity:
role: "Answer questions using the ledger and the ticket system."
constraints: "Use the tools. Do not guess."
capabilities:
provides: [desk.ask]
models:
privileged: { provider: fake, model: teller-1 }
security:
max_sensitivity_egress: internal
tools:
# Implemented in this binary, as a typed Rust tool.
- ref: tool://ledger/read
mutates: false
max_sensitivity: internal
description: Read a ledger account's balance.
# Served by the MCP server below. Declared identically — the manifest states
# *which tool*, and the wiring states how it is reached.
- ref: tool://tickets/read
mutates: false
max_sensitivity: internal
description: Read a support ticket.
execution: { kind: tool-calling, max_turns: 4 }
budgets:
max_tokens: 100000
"#;
#[derive(Debug, serde::Deserialize, schemars::JsonSchema)]
struct ReadBalance {
account: String,
}
#[async_trait::async_trait]
impl Tool for ReadBalance {
const SERVER: &'static str = "ledger";
const NAME: &'static str = "read";
fn mutates() -> bool {
false
}
async fn call(self) -> Result<Value, ToolFailure> {
println!(" → ledger/read {} (in process)", self.account);
Ok(json!({ "account": self.account, "balance": 42 }))
}
}
#[derive(Debug, Clone)]
struct TicketServer;
#[allow(clippy::unused_async_trait_impl)]
impl ServerHandler for TicketServer {
fn get_info(&self) -> ServerInfo {
let mut me = Implementation::default();
me.name = "tickets".into();
me.version = "0.0.0".into();
let mut info = ServerInfo::default();
info.protocol_version = ProtocolVersion::V_2026_07_28;
info.capabilities = ServerCapabilities::builder().enable_tools().build();
info.server_info = me;
info
}
async fn list_tools(
&self,
_p: Option<PaginatedRequestParams>,
_cx: RequestContext<RoleServer>,
) -> Result<ListToolsResult, ErrorData> {
let schema = Arc::new(
json!({ "type": "object", "properties": { "id": { "type": "string" } } })
.as_object()
.cloned()
.unwrap_or_default(),
);
let mut generous = ToolAnnotations::default();
generous.read_only_hint = Some(true);
let read = McpTool::new("read", "read a ticket", Arc::clone(&schema));
let mut close = McpTool::new("close", "close a ticket", schema);
close.annotations = Some(generous);
Ok(ListToolsResult {
tools: vec![read, close],
..Default::default()
})
}
async fn call_tool(
&self,
p: CallToolRequestParams,
_cx: RequestContext<RoleServer>,
) -> Result<CallToolResponse, ErrorData> {
match p.name.as_ref() {
"read" => {
println!(" → tickets/read (over MCP)");
Ok(
CallToolResult::success(vec![ContentBlock::text("ticket 7: printer on fire")])
.into(),
)
}
other => Err(ErrorData::invalid_params(
format!("no such tool: {other}"),
None,
)),
}
}
}
async fn connect() -> McpClient {
let (client_side, server_side) = tokio::io::duplex(8 * 1024);
let (sr, sw) = tokio::io::split(server_side);
let (cr, cw) = tokio::io::split(client_side);
tokio::spawn(async move {
if let Ok(running) = serve_server(TicketServer, (sr, sw)).await {
let _ = running.waiting().await;
}
});
McpClient::connect(
"tickets",
(cr, cw),
agentplane::tools::Destination::Local,
)
.await
.expect("client initialises on a known version")
}
#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
let tickets = Arc::new(connect().await);
println!("1. what the server says about itself");
for (id, advertised) in tickets.discover().await? {
println!(" {id}: read_only_hint={:?}", advertised.read_only);
}
println!(
" → recorded for comparison. Nothing here entered a catalogue: a tool\n\
\x20 absent from the operator's grants cannot be called however the\n\
\x20 server describes it, and `close` is not granted at all."
);
let manifest = Manifest::parse(TELLER)?;
let provider = FakeProvider::new();
provider.will_call_tool("call_1", "ledger__read", json!({ "account": "AC-1" }));
provider.will_call_tool("call_2", "tickets__read", json!({ "id": "7" }));
provider.will_say("AC-1 holds 42, and ticket 7 says the printer is on fire.");
let store: Arc<dyn JournalStore> = Arc::new(RedbStore::open_in_memory()?);
let driver: Arc<dyn ModelProvider> = provider.clone();
let transport: Arc<dyn ToolClient> = tickets.clone();
let rt = Runtime::builder(Arc::clone(&store))
.provider("fake", driver)
.agent(Agent::new(&manifest))
.toolbox(ToolBox::new().with::<ReadBalance>())
.tool_server("tickets", transport)
.build();
println!("\n2. the model uses both, and neither knows about the other");
let out = rt
.run(
"desk.ask",
Tainted::trusted(json!({ "question": "AC-1 balance and ticket 7?" })),
)
.await?;
assert_eq!(out.status, RunStatus::Succeeded);
println!(" answered: {}", out.output.as_ref().unwrap().peek());
println!("\n3. a tool id belonging to the other server");
match tickets
.call(&ToolId::new("ledger", "read"), &json!({}), None)
.await
{
Err(e) => println!(" → refused: {e}"),
Ok(v) => panic!("the ticket server answered a ledger tool: {v}"),
}
let before = provider.calls();
let replayed = rt.replay(out.run_id, Mode::Strict).await?;
assert_eq!(replayed.output, out.output);
assert_eq!(
provider.calls(),
before,
"strict replay called the model again"
);
println!(
"\n4. strict replay reassembled the conversation with zero model calls\n \
and zero tool calls — including the one that went over MCP"
);
store.verify(out.run_id).await?;
println!(" and the chain verifies end to end");
Ok(())
}