use rmcp::{
handler::server::{router::tool::ToolRouter, wrapper::Parameters},
model::{ServerCapabilities, ServerInfo},
schemars, tool, tool_handler, tool_router, ServerHandler,
};
use serde::Deserialize;
#[derive(Debug, Deserialize, schemars::JsonSchema)]
pub struct ConvergenceInput {}
#[derive(Clone)]
pub struct FleetMcp {
tool_router: ToolRouter<Self>,
}
#[tool_router]
impl FleetMcp {
#[must_use]
pub fn new() -> Self {
Self {
tool_router: Self::tool_router(),
}
}
#[tool(description = "Report this node's GitOps convergence: verdict \
(converged/behind/stopped/failing/unknown/notEnrolled), \
the deployed rev, the branch HEAD the reconciler last \
observed, and how long ago it ticked. Reads state the \
reconciler published to disk, so it answers correctly \
even when the daemon is dead — a stopped loop reports \
`stopped`, never silence.")]
async fn gitops_convergence(
&self,
Parameters(ConvergenceInput {}): Parameters<ConvergenceInput>,
) -> String {
let node =
super::utils::run_command_output(std::process::Command::new("hostname").arg("-s"))
.unwrap_or_else(|_| "unknown".to_owned());
let now = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map_or(0, |d| d.as_secs());
let doc = super::convergence::local(
std::path::Path::new(super::convergence::DEFAULT_STATE_DIR),
node,
now,
);
serde_json::to_string_pretty(&doc).unwrap_or_else(|e| {
let mut s = String::from("{\"error\":\"could not serialize convergence document: ");
s.push_str(&e.to_string());
s.push_str("\"}");
s
})
}
}
impl Default for FleetMcp {
fn default() -> Self {
Self::new()
}
}
#[tool_handler]
impl ServerHandler for FleetMcp {
fn get_info(&self) -> ServerInfo {
ServerInfo {
capabilities: ServerCapabilities::builder().enable_tools().build(),
instructions: Some(
"fleet — GitOps convergence for the node this server runs on. \
`gitops_convergence` reads the reconciler's published heartbeat \
and receipt chain; it never reaches a remote host, so a node \
that is unreachable is not thereby reported unhealthy, and a \
node whose daemon has died still reports its last known state \
plus how stale that state is. Absent evidence returns the \
verdict `unknown` — it is never rounded to `converged`."
.to_owned(),
),
..Default::default()
}
}
}
pub fn mcp() -> anyhow::Result<()> {
let rt = tokio::runtime::Builder::new_multi_thread()
.enable_all()
.build()?;
rt.block_on(async {
use rmcp::ServiceExt as _;
let service = FleetMcp::new()
.serve(rmcp::transport::stdio())
.await
.map_err(|e| anyhow::anyhow!("fleet mcp: serve failed: {e}"))?;
service
.waiting()
.await
.map_err(|e| anyhow::anyhow!("fleet mcp: {e}"))?;
Ok::<_, anyhow::Error>(())
})
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn the_server_advertises_a_tool_capability() {
let info = FleetMcp::new().get_info();
assert!(info.capabilities.tools.is_some(), "no tools advertised");
let instr = info.instructions.expect("instructions guide the agent");
assert!(
instr.contains("unknown"),
"the unknown verdict must be documented"
);
assert!(
instr.contains("never reaches a remote host"),
"an agent must not read this as a fleet-wide answer"
);
}
#[tokio::test]
async fn the_tool_returns_the_convergence_document() {
let out = FleetMcp::new()
.gitops_convergence(Parameters(ConvergenceInput {}))
.await;
let v: serde_json::Value =
serde_json::from_str(&out).expect("the tool must emit valid JSON");
assert!(v["node"].is_string(), "document must name the node: {out}");
assert!(
v["verdict"].is_string(),
"document must carry a verdict: {out}"
);
}
}