use std::process::ExitCode;
use std::sync::Arc;
use agentplane::core::Tainted;
use agentplane::journal::JournalStore;
use agentplane::manifest::Manifest;
use agentplane::model::ModelProvider;
use agentplane::runtime::{Mode, RunStatus, Runtime, RuntimeBuilder};
use agentplane::store::RedbStore;
fn install_tracing() {
use tracing_subscriber::{EnvFilter, fmt};
let filter = EnvFilter::try_from_default_env()
.unwrap_or_else(|_| EnvFilter::new("warn,agentplane=info"));
let _ = fmt()
.with_env_filter(filter)
.with_writer(std::io::stderr)
.try_init();
}
#[derive(clap::Parser, Debug)]
#[command(
name = "agentplane",
version,
about = "Run an agent that is only a file",
long_about = "Run, host and pin agents declared entirely in YAML.\n\n\
A file may hold several manifests separated by `---` (the \
Kubernetes convention), so a whole multi-agent room deploys as \
one file. Each document keeps its own digest — the file is \
packaging, not identity.",
disable_help_subcommand = true
)]
struct Cli {
#[command(subcommand)]
verb: Verb,
}
#[derive(clap::Subcommand, Debug)]
enum Verb {
Run(RunArgs),
Serve(Box<ServeArgs>),
Validate(FileArgs),
Digest(FileArgs),
}
#[derive(clap::Args, Debug)]
struct FileArgs {
manifest: String,
}
#[derive(clap::Args, Debug)]
struct RunArgs {
manifest: String,
#[arg(long, conflicts_with = "input_file")]
input: Option<String>,
#[arg(long)]
input_file: Option<String>,
#[arg(long)]
capability: Option<String>,
#[arg(long, env = "AGENTPLANE_STORE")]
store: Option<String>,
#[arg(long)]
replay: Option<String>,
#[arg(long, requires = "replay")]
strict: bool,
#[arg(long, value_name = "NAME=COMMAND")]
mcp: Vec<String>,
}
#[derive(clap::Args, Debug)]
struct ServeArgs {
manifest: String,
#[arg(long, env = "AGENTPLANE_URL")]
url: Option<String>,
#[arg(long, env = "AGENTPLANE_ADDR", default_value = "127.0.0.1:8080")]
addr: String,
#[arg(long, env = "AGENTPLANE_POLICY")]
policy: Option<String>,
#[arg(long, env = "AGENTPLANE_TOKENS")]
tokens: Option<String>,
#[arg(long, env = "AGENTPLANE_STORE")]
store: Option<String>,
#[arg(long, env = "AGENTPLANE_OPERATOR_ADDR")]
operator_addr: Option<String>,
#[arg(long, value_name = "SECS", env = "AGENTPLANE_SWEEP_EVERY")]
sweep_every: Option<u32>,
#[arg(long, value_name = "HOST")]
push_host: Vec<String>,
#[arg(long, value_name = "NAME=COMMAND")]
mcp: Vec<String>,
}
fn manifests_at(path: &str) -> Result<Vec<Manifest>, String> {
let text = std::fs::read_to_string(path).map_err(|e| format!("reading {path}: {e}"))?;
Manifest::parse_all(&text).map_err(|e| e.to_string())
}
fn dispatch(cli: Cli) -> Result<ExitCode, String> {
match cli.verb {
Verb::Validate(a) => {
for m in &manifests_at(&a.manifest)? {
println!("ok: {} {}", m.metadata.name, m.metadata.version);
}
Ok(ExitCode::SUCCESS)
}
Verb::Digest(a) => {
let manifests = manifests_at(&a.manifest)?;
if let [only] = manifests.as_slice() {
println!("{}", only.digest().map_err(|e| e.to_string())?.to_hex());
} else {
for m in &manifests {
println!(
"{} {} {}",
m.digest().map_err(|e| e.to_string())?.to_hex(),
m.metadata.name,
m.metadata.version
);
}
}
Ok(ExitCode::SUCCESS)
}
Verb::Run(a) => {
let manifests = manifests_at(&a.manifest)?;
execute(&manifests, &a)
}
Verb::Serve(a) => {
let manifests = manifests_at(&a.manifest)?;
serve(&manifests, &a)
}
}
}
fn main() -> ExitCode {
install_tracing();
let cli = <Cli as clap::Parser>::parse();
match dispatch(cli) {
Ok(code) => code,
Err(e) => {
eprintln!("agentplane: {e}");
ExitCode::FAILURE
}
}
}
impl RunArgs {
fn read_input(&self) -> Result<serde_json::Value, String> {
let text = match (&self.input, &self.input_file) {
(Some(s), _) => s.clone(),
(_, Some(p)) => std::fs::read_to_string(p).map_err(|e| format!("reading {p}: {e}"))?,
_ => "{}".into(),
};
serde_json::from_str(&text).map_err(|e| format!("the input is not valid JSON: {e}"))
}
}
#[cfg(all(feature = "a2a-server", feature = "cedar"))]
const PUSH_BATCH: usize = 64;
#[cfg(all(feature = "a2a-server", feature = "cedar"))]
const DEFAULT_SWEEP_SECONDS: u32 = 30;
#[cfg(all(feature = "a2a-server", feature = "cedar"))]
fn serve(manifests: &[Manifest], opts: &ServeArgs) -> Result<ExitCode, String> {
use agentplane::api::a2a::A2aServer;
use agentplane::api::tokens::TokenAuthenticator;
let [manifest] = manifests else {
return Err(format!(
"`serve` hosts one agent and this file holds {}. A2A's card path is \
well-known and singular, so a room would have to advertise one \
document and quietly not serve the rest — split the file, or run \
one process per agent",
manifests.len()
));
};
let url = opts.url.as_deref().ok_or(
"`serve` needs --url: the address callers reach this plane on. It goes on the \
Agent Card, so it is the public URL rather than what you bind — an agent's \
declaration must not change when its address does",
)?;
let policy_path = opts.policy.as_deref().ok_or(
"`serve` needs --policy: a Cedar policy set. There is deliberately no default — \
a permissive engine and no engine are the same behaviour, and only one of them \
looks governed",
)?;
let tokens_path = opts.tokens.as_deref().ok_or(
"`serve` needs --tokens: bearer tokens naming the callers this plane accepts. \
There is deliberately no default — a server that authenticates nobody has no \
actor to record a decision against",
)?;
let addr = opts.addr.as_str();
let policy_src = std::fs::read_to_string(policy_path)
.map_err(|e| format!("reading the policy set {policy_path}: {e}"))?;
let policy = agentplane::policy::CedarEngine::new(&policy_src)
.map_err(|e| format!("the policy set {policy_path} was refused: {e}"))?;
let tokens_src = std::fs::read_to_string(tokens_path)
.map_err(|e| format!("reading the token file {tokens_path}: {e}"))?;
let auth: Arc<dyn agentplane::api::Authenticator> = Arc::new(
TokenAuthenticator::from_yaml(&tokens_src)
.map_err(|e| format!("the token file {tokens_path} was refused: {e}"))?,
);
let operator_auth = Arc::clone(&auth);
let rt = tokio::runtime::Builder::new_multi_thread()
.enable_all()
.build()
.map_err(|e| format!("could not start the async runtime: {e}"))?;
rt.block_on(async move {
let path = opts.store.as_deref().ok_or(
"`serve` needs --store: a served task's id is a promise that it can be \
fetched again, and an in-memory journal breaks that promise at the next \
restart. `run` may journal to memory because it exits with its answer",
)?;
let store = Arc::new(RedbStore::open(path).map_err(|e| e.to_string())?);
let mut builder = with_providers(
Runtime::builder(Arc::clone(&store) as Arc<dyn JournalStore>),
std::slice::from_ref(manifest),
)
.await?;
for (name, client) in connect_mcp_servers(&opts.mcp).await? {
builder = builder.tool_server(name, client);
}
builder = builder
.cases(Arc::clone(&store) as Arc<dyn agentplane::case::CaseStore>)
.tasks(Arc::clone(&store) as Arc<dyn agentplane::case::TaskStore>)
.events(Arc::clone(&store) as Arc<dyn agentplane::case::EventStore>)
.timers(Arc::clone(&store) as Arc<dyn agentplane::case::TimerStore>)
.memory(Arc::clone(&store) as Arc<dyn agentplane::memory::MemoryStore>)
.policy(Arc::new(policy) as Arc<dyn agentplane::core::PolicyEngine>)
.agent(agentplane::runtime::Agent::new(manifest));
let runtime = builder.try_build().map_err(|e| e.to_string())?;
let security = agentplane::peers::CardSecurity::bearer("bearer", Vec::<String>::new());
let mut server = A2aServer::new(Arc::clone(&runtime), auth, &security, manifest, url)
.map_err(|e| e.to_string())?;
server = wire_push(server, &opts.push_host, &store)?;
if let Some(worker) = server.push_worker() {
spawn_push_worker(worker, opts.sweep_every.unwrap_or(DEFAULT_SWEEP_SECONDS));
}
spawn_sweeper(&runtime, opts.sweep_every.unwrap_or(DEFAULT_SWEEP_SECONDS));
if let Some(operator_addr) = opts.operator_addr.as_deref() {
spawn_operator_surface(&runtime, operator_auth, operator_addr).await?;
}
let listener = tokio::net::TcpListener::bind(addr)
.await
.map_err(|e| format!("could not bind {addr}: {e}"))?;
eprintln!(
"serving {} {} on {addr} as {url}",
manifest.metadata.name, manifest.metadata.version
);
eprintln!(" card: {url}/.well-known/agent-card.json");
axum::serve(listener, server.router())
.await
.map_err(|e| format!("the server stopped: {e}"))?;
Ok(ExitCode::SUCCESS)
})
}
#[cfg(all(feature = "a2a-server", feature = "cedar"))]
fn spawn_sweeper(runtime: &Arc<Runtime>, every: u32) {
if every == 0 {
return;
}
let sweeper = Arc::clone(runtime);
tokio::spawn(async move {
let mut tick = tokio::time::interval(std::time::Duration::from_secs(u64::from(every)));
loop {
tick.tick().await;
#[allow(clippy::disallowed_methods)]
let now = time::OffsetDateTime::now_utc();
match sweeper.fire_timers(now).await {
Ok(n) if n > 0 => tracing::info!(fired = n, "timers fired"),
Ok(_) => {}
Err(error) => tracing::error!(%error, "firing timers failed"),
}
match sweeper.sweep(now, time::Duration::hours(1)).await {
Ok(report) if report.needs_attention() => {
tracing::warn!(?report, "the sweep needs attention");
}
Ok(report) if !report.is_quiet() => tracing::info!(?report, "swept"),
Ok(_) => {}
Err(error) => tracing::error!(%error, "the sweep failed"),
}
}
});
}
#[cfg(all(feature = "a2a-server", feature = "cedar"))]
async fn spawn_operator_surface(
runtime: &Arc<Runtime>,
auth: Arc<dyn agentplane::api::Authenticator>,
addr: &str,
) -> Result<(), String> {
let api = agentplane::api::Api::new(Arc::clone(runtime), auth).map_err(|e| e.to_string())?;
let listener = tokio::net::TcpListener::bind(addr)
.await
.map_err(|e| format!("could not bind the operator surface {addr}: {e}"))?;
eprintln!(" operator: http://{addr}/runs?outcome=failed");
tokio::spawn(async move {
if let Err(error) = axum::serve(listener, api.router()).await {
tracing::error!(%error, "the operator surface stopped");
}
});
Ok(())
}
#[cfg(feature = "mcp-stdio")]
async fn connect_mcp_servers(
specs: &[String],
) -> Result<Vec<(String, Arc<dyn agentplane::tools::ToolClient>)>, String> {
use rmcp::ServiceExt as _;
let mut wired = Vec::with_capacity(specs.len());
for spec in specs {
let (name, command) = spec.split_once('=').ok_or_else(|| {
format!(
"--mcp wants `<server>=<command>`, got `{spec}`. The server name is the \
one your manifest's grants use: a grant `tool://tickets/read` needs \
`--mcp tickets=...`"
)
})?;
if name.trim().is_empty() {
return Err(format!("--mcp `{spec}` names no server"));
}
let mut parts = command.split_whitespace();
let program = parts
.next()
.ok_or_else(|| format!("--mcp `{spec}` names server `{name}` but no command"))?;
let mut process = tokio::process::Command::new(program);
process.args(parts);
let transport = rmcp::transport::TokioChildProcess::new(process)
.map_err(|e| format!("could not start the MCP server `{name}` (`{command}`): {e}"))?;
let service = agentplane::tools::McpClient::host_info()
.serve(transport)
.await
.map_err(|e| format!("the MCP server `{name}` did not initialise: {e}"))?;
eprintln!(" mcp: {name} <- {command}");
wired.push((
name.to_owned(),
Arc::new(agentplane::tools::McpClient::new(name, Arc::new(service)))
as Arc<dyn agentplane::tools::ToolClient>,
));
}
Ok(wired)
}
#[cfg(not(feature = "mcp-stdio"))]
#[allow(clippy::unused_async)]
async fn connect_mcp_servers(
specs: &[String],
) -> Result<Vec<(String, Arc<dyn agentplane::tools::ToolClient>)>, String> {
if specs.is_empty() {
return Ok(Vec::new());
}
Err(
"this build cannot run an MCP server: `--mcp` needs the `mcp-stdio` feature. \
Reinstall with `--features cli,mcp-stdio`, or use the `:full` container \
image, which is built with it"
.to_owned(),
)
}
#[cfg(all(feature = "a2a-server", feature = "cedar"))]
fn wire_push(
server: agentplane::api::a2a::A2aServer,
hosts: &[String],
store: &Arc<RedbStore>,
) -> Result<agentplane::api::a2a::A2aServer, String> {
if hosts.is_empty() {
return Ok(server);
}
let policy = hosts
.iter()
.fold(agentplane::push::PushPolicy::new(), |policy, host| {
policy.allow_host(host)
});
let server = server
.with_push(
Arc::clone(store) as Arc<dyn agentplane::push::PushStore>,
Arc::new(agentplane::push::PushSender::new(policy))
as Arc<dyn agentplane::push::PushTransport>,
)
.map_err(|e| e.to_string())?;
for host in hosts {
eprintln!(" push: https://{host}");
}
Ok(server)
}
#[cfg(all(feature = "a2a-server", feature = "cedar"))]
fn spawn_push_worker(worker: agentplane::api::a2a::A2aPushWorker, every: u32) {
if every == 0 {
return;
}
tokio::spawn(async move {
let mut tick = tokio::time::interval(std::time::Duration::from_secs(u64::from(every)));
loop {
tick.tick().await;
#[allow(clippy::disallowed_methods)]
let at = time::OffsetDateTime::now_utc().unix_timestamp();
let Ok(at) = u64::try_from(at) else { continue };
match worker.run_once(at, PUSH_BATCH).await {
Ok(report) if report.saturated => {
tracing::warn!(?report, "push delivery is saturated");
}
Ok(report) if report.registrations > 0 => {
tracing::info!(?report, "push delivered");
}
Ok(_) => {}
Err(error) => tracing::error!(%error, "push delivery failed"),
}
}
});
}
#[cfg(not(all(feature = "a2a-server", feature = "cedar")))]
#[allow(clippy::unnecessary_wraps)]
fn serve(_manifests: &[Manifest], _opts: &ServeArgs) -> Result<ExitCode, String> {
Err(
"this build cannot serve: `serve` needs the `a2a-server` and `cedar` features. \
Reinstall with `--features cli,a2a-server,cedar`, or use the `:full` \
container image, which is built with them"
.to_owned(),
)
}
fn execute(manifests: &[Manifest], opts: &RunArgs) -> Result<ExitCode, String> {
for manifest in manifests {
if manifest.spec.execution.is_none() {
return Err(format!(
"manifest '{}' declares no `spec.execution`, so its behaviour is a skill somebody \
wrote and there is nothing here for this binary to run. Register it in your own \
binary with `RuntimeBuilder::agent(Agent::new(&manifest).skill(YourSkill))` instead",
manifest.metadata.name
));
}
}
let rt = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.map_err(|e| format!("could not start the async runtime: {e}"))?;
rt.block_on(async {
let store: Arc<dyn JournalStore> = if let Some(path) = &opts.store {
Arc::new(RedbStore::open(path).map_err(|e| e.to_string())?)
} else {
eprintln!("note: journaling to memory; this run will not survive the process");
Arc::new(RedbStore::open_in_memory().map_err(|e| e.to_string())?)
};
let mut builder = with_providers(Runtime::builder(Arc::clone(&store)), manifests).await?;
for (name, client) in connect_mcp_servers(&opts.mcp).await? {
builder = builder.tool_server(name, client);
}
for manifest in manifests {
builder = builder.agent(agentplane::runtime::Agent::new(manifest));
}
let agent = builder.try_build().map_err(|e| e.to_string())?;
let outcome = if let Some(id) = &opts.replay {
let run = agentplane::core::RunId::parse(id)
.map_err(|e| format!("`{id}` is not a run id: {e}"))?;
let mode = if opts.strict {
Mode::Strict
} else {
Mode::Resume
};
agent.replay(run, mode).await
} else {
let capability = entry_capability(manifests, opts.capability.as_deref())?;
agent
.run(&capability, Tainted::trusted(opts.read_input()?))
.await
}
.map_err(|e| e.to_string())?;
eprintln!("run {} — {:?}", outcome.run_id, outcome.status);
if let Some(output) = &outcome.output {
println!("{}", output.peek());
}
Ok(if matches!(outcome.status, RunStatus::Succeeded) {
ExitCode::SUCCESS
} else {
ExitCode::FAILURE
})
})
}
async fn with_providers(
builder: RuntimeBuilder,
manifests: &[Manifest],
) -> Result<RuntimeBuilder, String> {
let mut builder = builder;
let mut seen: Vec<String> = Vec::new();
for manifest in manifests {
let Some(models) = &manifest.spec.models else {
continue;
};
for m in [models.privileged.as_ref(), models.quarantined.as_ref()]
.into_iter()
.flatten()
{
if seen.contains(&m.provider) {
continue;
}
seen.push(m.provider.clone());
builder = builder.provider(m.provider.clone(), driver(&m.provider).await?);
}
}
Ok(builder)
}
fn entry_capability(manifests: &[Manifest], asked: Option<&str>) -> Result<String, String> {
let all: Vec<(&str, &str)> = manifests
.iter()
.flat_map(|m| {
m.spec
.capabilities
.provides
.iter()
.map(move |c| (m.metadata.name.as_str(), c.as_str()))
})
.collect();
if let Some(asked) = asked {
if all.iter().any(|(_, c)| *c == asked) {
return Ok(asked.to_owned());
}
return Err(format!(
"no agent in this file provides '{asked}'. It provides: {}",
all.iter().map(|(_, c)| *c).collect::<Vec<_>>().join(", ")
));
}
if let [(_, only)] = all.as_slice() {
return Ok((*only).to_owned());
}
let orchestrators: Vec<&Manifest> = manifests
.iter()
.filter(|m| {
m.spec
.topology
.as_ref()
.is_some_and(|t| t.role == agentplane::manifest::Role::Orchestrator)
})
.collect();
if let [desk] = orchestrators.as_slice()
&& let [only] = desk.spec.capabilities.provides.as_slice()
{
return Ok(only.clone());
}
Err(format!(
"this file provides several capabilities and no single orchestrator to \
start at — say which one with --capability. It provides: {}",
all.iter()
.map(|(agent, c)| format!("{c} ({agent})"))
.collect::<Vec<_>>()
.join(", ")
))
}
async fn driver(name: &str) -> Result<Arc<dyn ModelProvider>, String> {
match name {
#[cfg(feature = "providers")]
"anthropic" => Ok(Arc::new(
agentplane::model::anthropic::Anthropic::new(key("ANTHROPIC_API_KEY")?)
.map_err(|e| e.to_string())?,
)),
#[cfg(feature = "bedrock")]
"bedrock" => Ok(Arc::new(
agentplane::model::bedrock::Bedrock::from_env(
std::env::var("AWS_REGION").map_err(|_| {
"AWS_REGION is not set, and the manifest names Bedrock".to_owned()
})?,
)
.await?,
)),
#[cfg(feature = "providers")]
"gemini" => Ok(Arc::new(
agentplane::model::gemini::Gemini::from_env().map_err(|e| e.to_string())?,
)),
#[cfg(feature = "providers")]
"openai" => Ok(Arc::new(
agentplane::model::openai::OpenAi::new(key("OPENAI_API_KEY")?)
.map_err(|e| e.to_string())?,
)),
#[cfg(feature = "providers")]
"chat-completions" => {
let base = key("CHAT_COMPLETIONS_BASE_URL").map_err(|_| {
"CHAT_COMPLETIONS_BASE_URL is not set, and the manifest names the \
chat-completions provider. Point it at the server: Ollama is \
http://localhost:11434, vLLM http://localhost:8000, TGI \
http://localhost:8080, Hugging Face's router \
https://router.huggingface.co/v1"
.to_owned()
})?;
let mut driver = agentplane::model::chat_completions::ChatCompletions::new(base)
.map_err(|e| e.to_string())?;
if let Ok(token) = std::env::var("CHAT_COMPLETIONS_API_KEY") {
driver = driver.bearer(token);
}
Ok(Arc::new(driver))
}
#[cfg(feature = "testkit")]
"fake" => Ok(agentplane::testkit::FakeProvider::new()),
other => Err(format!(
"no driver for provider '{other}'. This binary ships {}; anything else is an \
embedder's own driver, registered through RuntimeBuilder::provider",
shipped_providers().join(", "),
)),
}
}
fn shipped_providers() -> Vec<&'static str> {
#[allow(unused_mut)]
let mut names: Vec<&'static str> = Vec::new();
#[cfg(feature = "providers")]
names.extend(["anthropic", "chat-completions", "gemini", "openai"]);
#[cfg(feature = "bedrock")]
names.push("bedrock");
#[cfg(feature = "testkit")]
names.push("fake");
names.sort_unstable();
names
}
#[cfg(feature = "providers")]
fn key(var: &str) -> Result<String, String> {
std::env::var(var)
.map_err(|_| format!("{var} is not set, and the manifest names a provider that needs it"))
}