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::{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 },
}
#[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>,
}
#[derive(Debug, Clone)]
pub struct McpServer {
runtime: Arc<Runtime>,
served: Vec<Served>,
}
impl McpServer {
pub fn new(runtime: Arc<Runtime>, manifests: &[Manifest]) -> Result<Self, ServeError> {
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(),
});
}
}
Ok(Self { runtime, served })
}
fn find(&self, name: &str) -> Option<&Served> {
self.served.iter().find(|s| s.capability == name)
}
}
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(vec![crate::tools::MCP_REVISION])
}
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) {
context.peer.set_peer_info(request);
Ok(self.get_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 + '_ {
std::future::ready(Ok(ListToolsResult {
tools: self
.served
.iter()
.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 Some(served) = self.find(&request.name) else {
return Err(McpError::invalid_params(
format!("no agent provides '{}'", request.name),
None,
));
};
let input = Tainted::from_source(
serde_json::Value::Object(request.arguments.unwrap_or_default()),
SourceId::new("mcp://client"),
);
let admission = self
.runtime
.run_under(&served.capability, input, RunTerms::default())
.await
.map_err(|e| McpError::internal_error(e.to_string(), None))?;
let outcome = match admission {
Admission::Fresh(outcome) | Admission::Replayed(outcome) => outcome,
Admission::InFlight(run) => {
return Ok(CallToolResult::error(vec![ContentBlock::text(format!(
"run {run} is already in flight for this call"
))])
.into());
}
};
Ok(match outcome.status {
RunStatus::Succeeded => {
let value = outcome
.output
.map_or(serde_json::Value::Null, |o| o.peek().clone());
CallToolResult::structured(value).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(format!("{why:?}")),
)
.into()
}
other => CallToolResult::error(vec![ContentBlock::text(format!(
"run {} did not succeed: {other:?}",
outcome.run_id
))])
.into(),
})
}
fn list_prompts(
&self,
_request: Option<PaginatedRequestParams>,
_context: RequestContext<RoleServer>,
) -> impl Future<Output = Result<ListPromptsResult, McpError>> + Send + '_ {
std::future::ready(Ok(ListPromptsResult {
prompts: self
.served
.iter()
.filter(|s| s.prompt.is_some())
.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 mut seen = std::collections::BTreeSet::new();
std::future::ready(Ok(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 + '_ {
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 run = RunId::parse(&request.task_id)
.map_err(|_| McpError::invalid_params("no such task", None))?;
let records = self
.runtime
.journal()
.read(run, 1)
.await
.map_err(|e| McpError::internal_error(e.to_string(), None))?;
if records.is_empty() {
return Err(McpError::invalid_params("no such task", None));
}
let status = crate::runtime::observed_status(&records);
let payload = match &status {
Some(RunStatus::Succeeded) => TaskPayload::Completed {
result: serde_json::Map::new(),
},
Some(RunStatus::Cancelled { .. }) => TaskPayload::Failed {
error: error_object("the run was cancelled"),
},
Some(RunStatus::Suspended(_)) | None => TaskPayload::Working,
Some(other) => TaskPayload::Failed {
error: error_object(&format!("{other:?}")),
},
};
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 run = RunId::parse(&request.task_id)
.map_err(|_| McpError::invalid_params("no such task", None))?;
self.runtime
.request_cancel(
run,
&crate::core::Operator::connected("mcp://client")
.map_err(|e| McpError::internal_error(e.to_string(), None))?,
"cancelled by the calling host",
)
.await
.map(|_| ())
.map_err(|e| McpError::invalid_params(e.to_string(), None))
}
fn get_prompt(
&self,
request: GetPromptRequestParams,
_context: RequestContext<RoleServer>,
) -> impl Future<Output = Result<GetPromptResponse, McpError>> + Send + '_ {
std::future::ready(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())
}
}