use std::borrow::Cow;
use std::future::Future;
use std::sync::Arc;
use rmcp::model::{
CacheScope, CallToolRequestParams, CallToolResult, CancelTaskParams, ContentBlock,
CreateTaskResult, DetailedTask, GetPromptRequestParams, GetPromptResult, GetTaskParams,
GetTaskResult, Implementation, InitializeResult, ListPromptsResult, ListResourcesResult,
ListToolsResult, PaginatedRequestParams, Prompt, PromptMessage, ProtocolVersion,
ReadResourceRequestParams, ReadResourceResult, Resource, ResourceContents, Role,
ServerCapabilities, Task, TaskPayload, TaskStatus, Tool,
};
use rmcp::model::{CallToolResponse, GetPromptResponse, ReadResourceResponse};
use rmcp::service::RequestContext;
use rmcp::{ErrorData as McpError, RoleServer, ServerHandler};
use crate::core::RunId;
use crate::core::{RuntimeError, SourceId, Tainted};
use crate::manifest::{Identity, Manifest};
use crate::runtime::{Admission, RunStatus, RunTerms, Runtime};
#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
pub enum ServeError {
#[error(
"agent '{agent}' provides '{capability}' and declares no `spec.input`, so there is \
no reviewed argument shape to offer a model. Declare one, or leave the agent off \
this catalogue"
)]
NoInputSchema { agent: String, capability: String },
#[error("'{name}' is offered by two agents, so a call to it names no agent in particular")]
Duplicate { name: String },
#[error(
"agent '{agent}' declares '{capability}' and nothing on this plane provides it, so \
the tool would be offered to a model and refused at every call"
)]
NoProvider { agent: String, capability: String },
#[error("this plane has no policy engine, so nothing would govern what a host may run")]
NoPolicy,
#[error("the policy set cannot evaluate a run a host admits: {problems}")]
PolicyUnevaluable { problems: String },
#[error(
"protected-resource metadata must name at least one authorization server whose \
tokens this plane verifies"
)]
NoAuthorizationServer,
}
const ADMISSION_SOURCE: &str = "mcp/client";
const CALLER_NAMESPACE: &str = "mcp/";
pub mod action {
pub const TOOL_LIST: &str = "mcp:tool.list";
pub const TOOL_CALL: &str = "mcp:tool.call";
pub const TASK_READ: &str = "mcp:task.read";
pub const TASK_CANCEL: &str = "mcp:task.cancel";
pub const PROMPT_READ: &str = "mcp:prompt.read";
pub const RESOURCE_READ: &str = "mcp:resource.read";
pub const ALL: &[&str] = &[
TOOL_LIST,
TOOL_CALL,
TASK_READ,
TASK_CANCEL,
PROMPT_READ,
RESOURCE_READ,
];
}
#[must_use]
pub fn policy_problems(engine: &dyn crate::core::PolicyEngine) -> Vec<String> {
let caller = Authenticated {
actor: "preflight".to_owned(),
source: "peer:preflight".to_owned(),
roles: vec!["preflight".to_owned()],
tenant: "preflight".to_owned(),
acting_as: None,
};
let context = caller.context();
let requests: Vec<crate::core::PolicyRequest<'_>> = action::ALL
.iter()
.map(|action| crate::core::PolicyRequest {
principal: &caller.actor,
principal_kind: crate::core::PrincipalKind::Subject,
action,
resource: "preflight.resource",
context: &context,
})
.collect();
engine.preflight(&requests)
}
#[derive(Debug, Clone)]
struct Authenticated {
actor: String,
source: String,
roles: Vec<String>,
tenant: String,
acting_as: Option<crate::core::Delegation>,
}
impl Authenticated {
fn context(&self) -> serde_json::Value {
serde_json::json!({ "roles": self.roles, "peer": self.actor, "tenant": self.tenant })
}
}
#[derive(Debug, Clone)]
struct Asker {
admission: String,
input: String,
caller: Option<Authenticated>,
}
pub const IDEMPOTENCY_META_KEY: &str = "io.agentplane/idempotencyKey";
#[derive(Debug, Clone)]
struct Served {
capability: String,
agent: String,
description: Option<String>,
input_schema: Arc<rmcp::model::JsonObject>,
output_schema: Option<Arc<rmcp::model::JsonObject>>,
prompt: Option<String>,
document: String,
digest: Option<String>,
may_suspend: bool,
}
#[derive(Debug)]
pub struct McpServer {
runtime: Arc<Runtime>,
served: Vec<Served>,
session: String,
authenticated: bool,
}
impl Clone for McpServer {
fn clone(&self) -> Self {
Self {
runtime: Arc::clone(&self.runtime),
served: self.served.clone(),
session: RunId::generate().to_string(),
authenticated: self.authenticated,
}
}
}
impl McpServer {
pub fn new(runtime: Arc<Runtime>, manifests: &[Manifest]) -> Result<Self, ServeError> {
if runtime.policy().is_none() {
return Err(ServeError::NoPolicy);
}
let problems = runtime.served_policy_problems();
if !problems.is_empty() {
return Err(ServeError::PolicyUnevaluable {
problems: problems.join("; "),
});
}
let mut served: Vec<Served> = Vec::new();
for manifest in manifests {
let identity = manifest.spec.identity.as_ref();
let description = identity.map(|i| i.role.clone());
let prompt = identity
.map(Identity::system_prompt)
.filter(|p| !p.trim().is_empty());
for capability in &manifest.spec.capabilities.provides {
let Some(schema) = manifest.input_schema() else {
return Err(ServeError::NoInputSchema {
agent: manifest.metadata.name.clone(),
capability: capability.clone(),
});
};
if served.iter().any(|s| &s.capability == capability) {
return Err(ServeError::Duplicate {
name: capability.clone(),
});
}
if !runtime.provides(capability) {
return Err(ServeError::NoProvider {
agent: manifest.metadata.name.clone(),
capability: capability.clone(),
});
}
served.push(Served {
document: serde_json::to_string_pretty(manifest)
.unwrap_or_else(|_| String::new()),
digest: manifest.digest().ok().map(crate::core::Digest::to_hex),
capability: capability.clone(),
agent: manifest.metadata.name.clone(),
description: description.clone(),
input_schema: Arc::new(object(schema)),
output_schema: manifest.output_schema().map(|s| Arc::new(object(s))),
prompt: prompt.clone(),
may_suspend: manifest.may_suspend(|server| runtime.wires_peer(server)),
});
}
}
Ok(Self {
runtime,
served,
session: RunId::generate().to_string(),
authenticated: false,
})
}
#[cfg(feature = "mcp-server-http")]
pub(crate) fn authenticated(mut self) -> Self {
self.authenticated = true;
self
}
fn asker(&self, context: &RequestContext<RoleServer>) -> Result<Asker, McpError> {
if let Some(caller) = authenticated_caller(context) {
return Ok(Asker {
admission: format!("{CALLER_NAMESPACE}{}", caller.source),
input: caller.source.clone(),
caller: Some(caller),
});
}
let over_http = context.extensions.get::<http::request::Parts>().is_some();
if self.authenticated || over_http {
return Err(McpError::invalid_request(
"this request was not authenticated",
None,
));
}
Ok(Asker {
admission: ADMISSION_SOURCE.to_owned(),
input: "mcp://client".to_owned(),
caller: None,
})
}
fn gate(&self, asker: &Asker, action: &str, resource: &str) -> Result<(), McpError> {
let Some(caller) = asker.caller.as_ref() else {
return Ok(());
};
let Some(policy) = self.runtime.policy() else {
return Err(McpError::invalid_request(
"this request was not permitted",
None,
));
};
let context = caller.context();
match policy.authorize(&crate::core::PolicyRequest {
principal: &caller.actor,
principal_kind: crate::core::PrincipalKind::Subject,
action,
resource,
context: &context,
}) {
crate::core::PolicyDecision::Permit => Ok(()),
crate::core::PolicyDecision::Deny { reason }
| crate::core::PolicyDecision::Malformed { reason } => {
tracing::warn!(
target: "agentplane::mcp",
action,
resource,
reason,
"MCP request denied by policy"
);
Err(McpError::invalid_request(
"this request was not permitted",
None,
))
}
}
}
#[cfg(feature = "mcp-server-http")]
pub(crate) fn runtime(&self) -> &Arc<Runtime> {
&self.runtime
}
fn find(&self, name: &str) -> Option<&Served> {
self.served.iter().find(|s| s.capability == name)
}
fn admission_key(
&self,
asker: &Asker,
capability: &str,
request: &CallToolRequestParams,
context: &RequestContext<RoleServer>,
) -> String {
let named = request
.meta
.as_ref()
.and_then(|meta| meta.get(IDEMPOTENCY_META_KEY))
.or_else(|| context.meta.get(IDEMPOTENCY_META_KEY))
.and_then(serde_json::Value::as_str)
.filter(|key| !key.trim().is_empty());
let id = match named {
None => format!("request:{}/{}", self.session, context.id),
Some(key) if asker.caller.is_some() => format!("host:{key}"),
Some(key) => format!("host:{}/{key}", self.session),
};
let id = crate::core::origin_key(capability, &id);
crate::core::origin_key(&asker.admission, &id)
}
async fn our_run(
&self,
asker: &Asker,
task_id: &str,
) -> Result<(RunId, Vec<crate::journal::Record>), McpError> {
let run =
RunId::parse(task_id).map_err(|_| McpError::invalid_params("no such task", None))?;
let records = self
.runtime
.journal()
.read(run, 1)
.await
.map_err(|e| internal("reading a task's journal", &e))?;
if records
.first()
.and_then(crate::journal::Record::admission_source)
!= Some(asker.admission.as_str())
{
return Err(McpError::invalid_params("no such task", None));
}
Ok((run, records))
}
}
#[cfg(feature = "mcp-server-http")]
fn authenticated_caller(context: &RequestContext<RoleServer>) -> Option<Authenticated> {
let caller = context
.extensions
.get::<axum::http::request::Parts>()?
.extensions
.get::<crate::api::Caller>()?;
Some(Authenticated {
actor: caller.actor.clone(),
source: crate::api::peer_source(&caller.actor),
roles: caller.roles.clone(),
tenant: caller.tenant.as_str().to_owned(),
acting_as: caller.acting_as.clone(),
})
}
#[cfg(not(feature = "mcp-server-http"))]
fn authenticated_caller(_context: &RequestContext<RoleServer>) -> Option<Authenticated> {
None
}
fn host_has_tasks(context: &RequestContext<RoleServer>) -> bool {
context.protocol_version() == Some(crate::tools::MCP_REVISION)
&& context
.client_capabilities()
.is_some_and(|c| c.supports_tasks())
}
fn waiting_without_tasks(run: RunId) -> CallToolResult {
CallToolResult::error(vec![ContentBlock::text(format!(
"run {run} is waiting on this plane and has not concluded; this host did not \
negotiate the {} extension, so the call cannot be handed back as a task",
rmcp::model::TASKS_EXTENSION_ID
))])
}
fn internal(doing: &str, error: &dyn std::fmt::Display) -> McpError {
McpError::internal_error(crate::core::withheld_fault("mcp", doing, error), None)
}
fn finished(
run: RunId,
status: &RunStatus,
output: Option<&Tainted<serde_json::Value>>,
) -> CallToolResult {
match status {
RunStatus::Succeeded => {
CallToolResult::structured(output.map_or(serde_json::Value::Null, |o| o.peek().clone()))
}
other => CallToolResult::error(vec![ContentBlock::text(format!(
"run {run} did not succeed: {}",
other.as_str()
))]),
}
}
fn declined() -> CallToolResult {
CallToolResult::error(vec![ContentBlock::text("this agent declined the request")])
}
fn manifest_uri(agent: &str) -> String {
format!("agentplane://manifest/{agent}")
}
fn error_object(message: &str) -> rmcp::model::JsonObject {
let mut out = serde_json::Map::new();
out.insert("code".to_owned(), serde_json::json!(-32603));
out.insert("message".to_owned(), serde_json::json!(message));
out
}
#[allow(clippy::disallowed_methods)]
fn protocol_now() -> String {
crate::core::format_timestamp(crate::core::Timestamp::now_utc())
}
const CACHE_SCOPE: CacheScope = CacheScope::Private;
const CACHE_TTL_MS: u64 = 0;
fn object(schema: &serde_json::Value) -> rmcp::model::JsonObject {
schema.as_object().cloned().unwrap_or_default()
}
impl ServerHandler for McpServer {
fn supported_protocol_versions(&self) -> Cow<'static, [ProtocolVersion]> {
Cow::Owned(crate::tools::McpClient::SPOKEN_REVISIONS.to_vec())
}
fn initialize(
&self,
request: rmcp::model::InitializeRequestParams,
context: RequestContext<RoleServer>,
) -> impl Future<Output = Result<InitializeResult, McpError>> + Send + '_ {
let supported = self.supported_protocol_versions();
std::future::ready(if supported.contains(&request.protocol_version) {
let mut info = self.get_info();
info.protocol_version = request.protocol_version.clone();
if info.protocol_version != crate::tools::MCP_REVISION
&& let Some(extensions) = info.capabilities.extensions.as_mut()
{
extensions.remove(rmcp::model::TASKS_EXTENSION_ID);
}
context.peer.set_peer_info(request);
Ok(info)
} else {
Err(McpError::unsupported_protocol_version(
request.protocol_version,
&supported,
))
})
}
fn get_info(&self) -> InitializeResult {
let mut info = InitializeResult::new(
ServerCapabilities::builder()
.enable_tools()
.enable_prompts()
.enable_resources()
.enable_tasks()
.build(),
)
.with_server_info(Implementation::new("agentplane", env!("CARGO_PKG_VERSION")))
.with_instructions(
"Each tool runs one governed agent: the call is admitted, journaled and \
dispatched under the agent's declared authority and budget. A refusal is \
an answer, not an outage.",
);
info.protocol_version = crate::tools::MCP_REVISION;
info
}
fn list_tools(
&self,
_request: Option<PaginatedRequestParams>,
context: RequestContext<RoleServer>,
) -> impl Future<Output = Result<ListToolsResult, McpError>> + Send + '_ {
let asked = self
.asker(&context)
.and_then(|asker| self.gate(&asker, action::TOOL_LIST, "catalogue"));
let tasks = host_has_tasks(&context);
std::future::ready(asked.map(|()| {
ListToolsResult {
tools: self
.served
.iter()
.filter(|s| tasks || !s.may_suspend)
.map(|s| {
let mut tool = Tool::new_with_raw(
Cow::Owned(s.capability.clone()),
s.description.clone().map(Cow::Owned),
Arc::clone(&s.input_schema),
);
tool.title = Some(s.agent.clone());
tool.output_schema.clone_from(&s.output_schema);
tool
})
.collect(),
ttl_ms: Some(CACHE_TTL_MS),
cache_scope: Some(CACHE_SCOPE),
..ListToolsResult::default()
}
}))
}
async fn call_tool(
&self,
request: CallToolRequestParams,
context: RequestContext<RoleServer>,
) -> Result<CallToolResponse, McpError> {
let asker = self.asker(&context)?;
if self.gate(&asker, action::TOOL_CALL, &request.name).is_err() {
return Ok(declined().into());
}
let Some(served) = self.find(&request.name) else {
return Err(McpError::invalid_params(
format!("no agent provides '{}'", request.name),
None,
));
};
let tasks = host_has_tasks(&context);
if served.may_suspend && !tasks {
return Err(McpError::invalid_params(
format!(
"'{}' may wait on a person or a timer, and this host did not negotiate \
the {} extension (MCP {}) that carries a waiting call back as a task",
request.name,
rmcp::model::TASKS_EXTENSION_ID,
crate::tools::MCP_REVISION
),
None,
));
}
let key = self.admission_key(&asker, &served.capability, &request, &context);
let input = Tainted::from_source(
serde_json::Value::Object(request.arguments.unwrap_or_default()),
SourceId::new(asker.input.clone()),
);
let terms = match &asker.caller {
Some(caller) => RunTerms::default()
.once(&key)
.served(caller.acting_as.clone())
.admitted_by(&caller.actor),
None => RunTerms::default().once(&key).served(None),
};
let admission = match self
.runtime
.run_under(&served.capability, input, terms)
.await
{
Ok(admission) => admission,
Err(RuntimeError::PolicyDenied(_) | RuntimeError::Delegation(_)) => {
return Ok(declined().into());
}
Err(e) => return Err(internal("admitting a tool call", &e)),
};
let outcome = match admission {
Admission::Replayed(outcome) if matches!(outcome.status, RunStatus::Succeeded) => self
.runtime
.replay(outcome.run_id, crate::runtime::Mode::Strict)
.await
.map_err(|e| internal("replaying a retried call", &e))?,
Admission::Fresh(outcome) | Admission::Replayed(outcome) => outcome,
Admission::InFlight(run) if !tasks => return Ok(waiting_without_tasks(run).into()),
Admission::InFlight(run) => {
let now = protocol_now();
return Ok(CreateTaskResult::new(
Task::new(run.to_string(), TaskStatus::Working, &now, &now)
.with_status_message("already in flight"),
)
.into());
}
};
Ok(match outcome.status {
RunStatus::Suspended(_) if !tasks => waiting_without_tasks(outcome.run_id).into(),
RunStatus::Suspended(ref why) => {
let now = protocol_now();
CreateTaskResult::new(
Task::new(outcome.run_id.to_string(), TaskStatus::Working, &now, &now)
.with_status_message(why.to_string()),
)
.into()
}
ref status => finished(outcome.run_id, status, outcome.output.as_ref()).into(),
})
}
fn list_prompts(
&self,
_request: Option<PaginatedRequestParams>,
context: RequestContext<RoleServer>,
) -> impl Future<Output = Result<ListPromptsResult, McpError>> + Send + '_ {
let asked = self
.asker(&context)
.and_then(|asker| self.gate(&asker, action::PROMPT_READ, "catalogue"));
let mut listed = std::collections::BTreeSet::new();
std::future::ready(asked.map(|()| {
ListPromptsResult {
prompts: self
.served
.iter()
.filter(|s| s.prompt.is_some())
.filter(|s| listed.insert(s.agent.clone()))
.map(|s| Prompt::new(s.agent.clone(), s.description.clone(), None))
.collect(),
ttl_ms: Some(CACHE_TTL_MS),
cache_scope: Some(CACHE_SCOPE),
..ListPromptsResult::default()
}
}))
}
fn list_resources(
&self,
_request: Option<PaginatedRequestParams>,
context: RequestContext<RoleServer>,
) -> impl Future<Output = Result<ListResourcesResult, McpError>> + Send + '_ {
let asked = self
.asker(&context)
.and_then(|asker| self.gate(&asker, action::RESOURCE_READ, "catalogue"));
let mut seen = std::collections::BTreeSet::new();
std::future::ready(asked.map(|()| {
ListResourcesResult {
resources: self
.served
.iter()
.filter(|s| seen.insert(s.agent.clone()))
.map(|s| {
let resource = Resource::new(manifest_uri(&s.agent), s.agent.clone())
.with_mime_type("application/json");
match &s.description {
Some(role) => resource.with_description(role.clone()),
None => resource,
}
})
.collect(),
ttl_ms: Some(CACHE_TTL_MS),
cache_scope: Some(CACHE_SCOPE),
..ListResourcesResult::default()
}
}))
}
fn read_resource(
&self,
request: ReadResourceRequestParams,
context: RequestContext<RoleServer>,
) -> impl Future<Output = Result<ReadResourceResponse, McpError>> + Send + '_ {
if let Err(refused) = self
.asker(&context)
.and_then(|asker| self.gate(&asker, action::RESOURCE_READ, &request.uri))
{
return std::future::ready(Err(refused));
}
std::future::ready(
self.served
.iter()
.find(|s| manifest_uri(&s.agent) == request.uri)
.map_or_else(
|| {
Err(McpError::invalid_params(
format!("no resource at '{}'", request.uri),
None,
))
},
|s| {
let mut contents =
ResourceContents::text(s.document.clone(), request.uri.clone())
.with_mime_type("application/json");
if let Some(digest) = &s.digest {
let mut meta = rmcp::model::MetaObject::new();
meta.insert("digest".to_owned(), serde_json::json!(digest));
contents = contents.with_meta(meta);
}
let mut result = ReadResourceResult::new(vec![contents]);
result.ttl_ms = Some(CACHE_TTL_MS);
result.cache_scope = Some(CACHE_SCOPE);
Ok(result.into())
},
),
)
}
async fn get_task(
&self,
request: GetTaskParams,
context: RequestContext<RoleServer>,
) -> Result<GetTaskResult, McpError> {
let asker = self.asker(&context)?;
let (run, records) = self.our_run(&asker, &request.task_id).await?;
self.gate(&asker, action::TASK_READ, &request.task_id)?;
let status = crate::runtime::observed_status(&records);
let payload = match &status {
Some(RunStatus::Succeeded) => {
let outcome = self
.runtime
.replay(run, crate::runtime::Mode::Strict)
.await
.map_err(|e| internal("replaying a completed task", &e))?;
let result =
serde_json::to_value(finished(run, &outcome.status, outcome.output.as_ref()))
.map_err(|e| internal("encoding a task's result", &e))?;
TaskPayload::Completed {
result: result.as_object().cloned().unwrap_or_default(),
}
}
Some(RunStatus::Cancelled { .. }) => TaskPayload::Cancelled,
Some(RunStatus::Suspended(_)) | None => TaskPayload::Working,
Some(other) => TaskPayload::Failed {
error: error_object(&format!("the run did not succeed: {}", other.as_str())),
},
};
let now = protocol_now();
Ok(GetTaskResult::new(DetailedTask::new(
Task::new(request.task_id, TaskStatus::Working, &now, &now),
payload,
)))
}
async fn cancel_task(
&self,
request: CancelTaskParams,
context: RequestContext<RoleServer>,
) -> Result<(), McpError> {
let asker = self.asker(&context)?;
let (run, _) = self.our_run(&asker, &request.task_id).await?;
self.gate(&asker, action::TASK_CANCEL, &request.task_id)?;
let operator = match &asker.caller {
Some(caller) => crate::core::Operator::authenticated(caller.actor.clone()),
None => crate::core::Operator::connected("mcp://client"),
}
.map_err(|e| McpError::internal_error(e.to_string(), None))?;
self.runtime
.request_cancel(run, &operator, "cancelled by the calling host")
.await
.map(|_| ())
.map_err(|e| match e {
RuntimeError::Store(_) => internal("requesting a cancellation", &e),
other => McpError::invalid_params(other.to_string(), None),
})
}
fn get_prompt(
&self,
request: GetPromptRequestParams,
context: RequestContext<RoleServer>,
) -> impl Future<Output = Result<GetPromptResponse, McpError>> + Send + '_ {
std::future::ready(
self.asker(&context)
.and_then(|asker| self.gate(&asker, action::PROMPT_READ, &request.name))
.and_then(|()| self.prompt_for(request)),
)
}
}
impl McpServer {
fn prompt_for(&self, request: GetPromptRequestParams) -> Result<GetPromptResponse, McpError> {
let Some(served) = self
.served
.iter()
.find(|s| s.agent == request.name && s.prompt.is_some())
else {
return Err(McpError::invalid_params(
format!("no agent named '{}' serves a prompt", request.name),
None,
));
};
if request.arguments.is_some_and(|a| !a.is_empty()) {
return Err(McpError::invalid_params(
"this prompt takes no arguments: it is a reviewed instruction, and a \
value spliced into it would be text nobody approved under a digest that \
covers text somebody did",
None,
));
}
let text = served.prompt.clone().unwrap_or_default();
let mut result = GetPromptResult::new(vec![PromptMessage::new_text(Role::User, text)]);
result.description.clone_from(&served.description);
Ok(result.into())
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn a_suspension_for_a_host_without_tasks_is_an_error_naming_the_run() {
let run = RunId::generate();
let answer = waiting_without_tasks(run);
assert_eq!(answer.is_error, Some(true), "a waiting run read as success");
assert!(answer.structured_content.is_none(), "{answer:?}");
let said = serde_json::to_string(&answer.content).expect("content");
assert!(
said.contains(&run.to_string()) && said.contains("waiting"),
"the answer does not name the run and that it is waiting: {said}"
);
}
}