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),
Replay(ReplayArgs),
Card(CardArgs),
Serve(Box<ServeArgs>),
Validate(FileArgs),
Schema,
Digest(FileArgs),
Audit(AuditArgs),
Export(StoreArgs),
Verify(VerifyArgs),
Restore(RestoreArgs),
Drill(DrillArgs),
ForgetAdmissions(ForgetArgs),
Retain(RetainArgs),
Halt(HaltArgs),
Halts(HaltsArgs),
}
#[derive(clap::Args, Debug)]
struct RetainArgs {
#[arg(long, env = "AGENTPLANE_STORE")]
store: String,
#[arg(long, env = "AGENTPLANE_TENANT")]
tenant: Option<String>,
#[arg(long)]
older_than_days: u32,
#[arg(long)]
reason: String,
#[arg(long)]
dry_run: bool,
}
#[derive(clap::Args, Debug)]
struct HaltArgs {
#[arg(long, env = "AGENTPLANE_STORE")]
store: String,
#[arg(long, env = "AGENTPLANE_TENANT")]
tenant: Option<String>,
#[arg(long, default_value = "tenant")]
scope: String,
#[arg(long)]
reason: Option<String>,
#[arg(long, conflicts_with = "reason")]
lift: bool,
}
#[derive(clap::Args, Debug)]
struct HaltsArgs {
#[arg(long, env = "AGENTPLANE_STORE")]
store: String,
#[arg(long, env = "AGENTPLANE_TENANT")]
tenant: Option<String>,
}
#[derive(clap::Args, Debug)]
struct ForgetArgs {
#[arg(long, env = "AGENTPLANE_STORE")]
store: String,
#[arg(long)]
older_than_days: u32,
}
#[derive(clap::Args, Debug)]
struct DrillArgs {
#[arg(long, env = "AGENTPLANE_STORE")]
store: String,
}
#[derive(clap::Args, Debug)]
struct RestoreArgs {
file: String,
#[arg(long, env = "AGENTPLANE_STORE")]
store: String,
}
#[derive(clap::Args, Debug)]
struct VerifyArgs {
file: String,
#[arg(long)]
key: Vec<String>,
#[arg(long)]
checkpoint: Option<String>,
}
#[derive(clap::Args, Debug)]
struct AuditArgs {
#[command(flatten)]
store: StoreArgs,
#[arg(long)]
key: Vec<String>,
#[arg(long)]
prior: Option<String>,
#[arg(long)]
require_signatures: bool,
}
#[derive(clap::Args, Debug)]
struct StoreArgs {
#[arg(long, env = "AGENTPLANE_STORE")]
store: String,
#[arg(long)]
outcome: Vec<String>,
#[arg(long, default_value_t = 1000)]
limit: usize,
}
fn verifier_from(keys: &[String]) -> Result<Option<agentplane::policy::Ed25519Verifier>, String> {
if keys.is_empty() {
return Ok(None);
}
let mut verifier = agentplane::policy::Ed25519Verifier::new();
for entry in keys {
let Some((id, hex_key)) = entry.split_once('=') else {
return Err(format!(
"--key takes <key-id>=<64 hex chars>, got '{entry}' — the id is what records \
name as their signer, and the hex is the Ed25519 public key"
));
};
let mut bytes = [0u8; 32];
hex::decode_to_slice(hex_key, &mut bytes)
.map_err(|e| format!("--key {id}: not 64 hex characters: {e}"))?;
verifier = verifier
.trust(id, &bytes)
.map_err(|e| format!("--key {id}: not a valid Ed25519 public key: {e}"))?;
}
Ok(Some(verifier))
}
#[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, value_name = "NAME=COMMAND")]
mcp: Vec<String>,
#[arg(long, value_name = "NAME=URL")]
peer: Vec<String>,
}
#[derive(clap::Args, Debug)]
struct ReplayArgs {
run_id: String,
#[arg(long, env = "AGENTPLANE_STORE")]
store: String,
#[arg(long)]
manifest: String,
#[arg(long)]
strict: bool,
#[arg(long, value_name = "NAME=COMMAND")]
mcp: Vec<String>,
#[arg(long, value_name = "NAME=URL")]
peer: Vec<String>,
}
#[derive(clap::Args, Debug)]
struct CardArgs {
manifest: String,
#[arg(long, env = "AGENTPLANE_URL")]
url: 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 = "SECS", env = "AGENTPLANE_DRILL_EVERY")]
drill_every: Option<u32>,
#[arg(long, value_name = "HOST")]
push_host: Vec<String>,
#[arg(long, value_name = "NAME=COMMAND")]
mcp: Vec<String>,
#[arg(long, value_name = "NAME=URL")]
peer: Vec<String>,
}
fn journal_verb(opts: &StoreArgs, audit: Option<&AuditArgs>) -> Result<ExitCode, String> {
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 redb = Arc::new(RedbStore::open(&opts.store).map_err(|e| e.to_string())?);
let store: Arc<dyn JournalStore> = redb.clone();
let cases: Arc<dyn agentplane::case::CaseStore> = redb;
let wanted: Vec<String> = if opts.outcome.is_empty() {
agentplane::runtime::OUTCOMES_OF_RECORD
.iter()
.map(|s| (*s).to_owned())
.collect()
} else {
opts.outcome.clone()
};
let mut runs = Vec::new();
let mut truncated = Vec::new();
for outcome in &wanted {
let found = store
.runs_by_outcome(outcome, opts.limit + 1)
.await
.map_err(|e| e.to_string())?;
if found.len() > opts.limit {
truncated.push(outcome.clone());
}
runs.extend(found.into_iter().take(opts.limit));
}
if !truncated.is_empty() {
eprintln!(
"warning: --limit {} was reached for: {}. This is a partial view; \
raise --limit or narrow --outcome",
opts.limit,
truncated.join(", ")
);
}
let Some(audit) = audit else {
let stdout = std::io::stdout();
let trailer = agentplane::export::to_jsonl(
&store,
Some(&cases),
&runs,
std::io::BufWriter::new(stdout.lock()),
)
.await
.map_err(|e| e.to_string())?;
eprintln!(
"exported {} record(s) from {}/{} run(s) and {} case(s)",
trailer.records, trailer.runs_exported, trailer.runs_requested, trailer.cases
);
if !trailer.unreadable.is_empty() {
for u in &trailer.unreadable {
eprintln!("unreadable: {} — {}", u.run, u.reason);
}
return Ok(ExitCode::FAILURE);
}
return Ok(ExitCode::SUCCESS);
};
let verifier = verifier_from(&audit.key)?;
let prior: Option<agentplane::journal::Checkpoint> = match &audit.prior {
Some(path) => Some(
std::fs::read_to_string(path)
.map_err(|e| format!("reading --prior {path}: {e}"))
.and_then(|text| {
serde_json::from_str(&text).map_err(|e| {
format!(
"--prior {path} is not a checkpoint — expected the `current` \
field of an earlier audit report: {e}"
)
})
})?,
),
None => None,
};
let evidence = agentplane::audit::Evidence {
prior: prior.as_ref(),
verifier: verifier
.as_ref()
.map(|v| v as &dyn agentplane::core::Verifier),
require_signatures: audit.require_signatures,
};
let report = agentplane::audit::audit(&store, &runs, &evidence)
.await
.map_err(|e| e.to_string())?;
println!(
"{}",
serde_json::to_string_pretty(&report).map_err(|e| e.to_string())?
);
Ok(if report.is_sound() {
ExitCode::SUCCESS
} else {
ExitCode::FAILURE
})
})
}
fn drill_verb(opts: &DrillArgs) -> Result<ExitCode, String> {
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 redb = Arc::new(RedbStore::open(&opts.store).map_err(|e| e.to_string())?);
let cases: Arc<dyn agentplane::case::CaseStore> = redb;
let tenant = agentplane::core::TenantId::default();
let stores = agentplane::drill::Stores {
cases: &cases,
blobs: None,
#[cfg(feature = "keyring")]
keys: None,
tenant: &tenant,
};
let report = agentplane::drill::drill(&stores)
.await
.map_err(|e| e.to_string())?;
println!(
"{}",
serde_json::to_string_pretty(&report).map_err(|e| e.to_string())?
);
Ok(if report.is_sound() {
ExitCode::SUCCESS
} else {
ExitCode::FAILURE
})
})
}
fn forget_admissions_verb(opts: &ForgetArgs) -> Result<ExitCode, String> {
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> =
Arc::new(RedbStore::open(&opts.store).map_err(|e| e.to_string())?);
#[allow(clippy::disallowed_methods)]
let now = time::OffsetDateTime::now_utc();
let cutoff = now - std::time::Duration::from_secs(u64::from(opts.older_than_days) * 86_400);
let retired = store
.forget_admissions(cutoff)
.await
.map_err(|e| e.to_string())?;
println!(
"{}",
serde_json::json!({
"retired": retired,
"older_than_days": opts.older_than_days,
"cutoff": cutoff.unix_timestamp(),
})
);
Ok(ExitCode::SUCCESS)
})
}
fn retain_verb(opts: &RetainArgs) -> Result<ExitCode, String> {
if !opts.dry_run {
return Err(concat!(
"this binary wires no blob store and no key ring, so it cannot erase anything; ",
"it can only say what a pass would erase. Run again with --dry-run, or call ",
"`Runtime::retain` from a plane built with `.blobs(..)` and `.keyring(..)`"
)
.to_owned());
}
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 redb = RedbStore::open(&opts.store).map_err(|e| e.to_string())?;
let redb = match opts.tenant.as_deref() {
Some(name) => redb.for_tenant(
agentplane::core::TenantId::new(name).map_err(|e| format!("--tenant: {e}"))?,
),
None => redb,
};
let cases: Arc<dyn agentplane::case::CaseStore> = Arc::new(redb);
#[allow(clippy::disallowed_methods)]
let now = time::OffsetDateTime::now_utc();
let cutoff = now - std::time::Duration::from_secs(u64::from(opts.older_than_days) * 86_400);
let plan = agentplane::retention::plan(cases.as_ref(), cutoff)
.await
.map_err(|e| e.to_string())?;
println!(
"{}",
serde_json::to_string_pretty(&serde_json::json!({
"dry_run": true,
"reason": opts.reason,
"cutoff": cutoff.unix_timestamp(),
"scanned": plan.scanned,
"would_erase": plan.due,
}))
.map_err(|e| e.to_string())?
);
Ok(ExitCode::SUCCESS)
})
}
fn halt_verb(opts: &HaltArgs) -> Result<ExitCode, String> {
let scope = agentplane::quota::HaltScope::parse(&opts.scope).ok_or_else(|| {
format!(
"'{}' is not a scope: use 'tenant', 'agent:<metadata.name>', {}",
opts.scope, "or 'revision:<manifest digest>'"
)
})?;
if !opts.lift && opts.reason.is_none() {
return Err(concat!(
"--reason is required to halt: the next person to look will be somebody else, ",
"possibly at three in the morning, and why is the whole question. ",
"Use --lift to clear a halt"
)
.to_owned());
}
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 quotas = quota_store_at(&opts.store, opts.tenant.as_deref())?;
quotas
.set_halt(&scope, opts.reason.as_deref())
.await
.map_err(|e| e.to_string())?;
println!(
"{}",
serde_json::json!({
"scope": scope.key(),
"halted": !opts.lift,
"reason": opts.reason,
})
);
Ok(ExitCode::SUCCESS)
})
}
fn halts_verb(opts: &HaltsArgs) -> Result<ExitCode, String> {
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 quotas = quota_store_at(&opts.store, opts.tenant.as_deref())?;
let halts = quotas.halts().await.map_err(|e| e.to_string())?;
let rows: Vec<serde_json::Value> = halts
.iter()
.map(|h| serde_json::json!({ "scope": h.scope.key(), "reason": h.reason }))
.collect();
println!(
"{}",
serde_json::to_string_pretty(&serde_json::json!({ "halts": rows }))
.map_err(|e| e.to_string())?
);
Ok(ExitCode::SUCCESS)
})
}
fn quota_store_at(
path: &str,
tenant: Option<&str>,
) -> Result<Arc<dyn agentplane::quota::QuotaStore>, String> {
let store = RedbStore::open(path).map_err(|e| e.to_string())?;
let store = match tenant {
Some(name) => store.for_tenant(
agentplane::core::TenantId::new(name).map_err(|e| format!("--tenant: {e}"))?,
),
None => store,
};
Ok(Arc::new(store))
}
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::Schema => {
let schema = Manifest::json_schema();
println!(
"{}",
serde_json::to_string_pretty(&schema).expect("a generated schema serializes")
);
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::Audit(a) => journal_verb(&a.store, Some(&a)),
Verb::Export(a) => journal_verb(&a, None),
Verb::Drill(a) => drill_verb(&a),
Verb::ForgetAdmissions(a) => forget_admissions_verb(&a),
Verb::Retain(a) => retain_verb(&a),
Verb::Halt(a) => halt_verb(&a),
Verb::Halts(a) => halts_verb(&a),
Verb::Restore(a) => {
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 redb = Arc::new(RedbStore::open(&a.store).map_err(|e| e.to_string())?);
let store: Arc<dyn JournalStore> = redb.clone();
let cases: Arc<dyn agentplane::case::CaseStore> = redb;
let file =
std::fs::File::open(&a.file).map_err(|e| format!("reading {}: {e}", a.file))?;
let report = agentplane::export::from_jsonl(
&store,
Some(&cases),
std::io::BufReader::new(file),
)
.await
.map_err(|e| e.to_string())?;
println!(
"{}",
serde_json::to_string_pretty(&report).map_err(|e| e.to_string())?
);
Ok(if report.is_faithful() {
ExitCode::SUCCESS
} else {
ExitCode::FAILURE
})
})
}
Verb::Verify(a) => verify_verb(&a),
Verb::Run(a) => {
let manifests = manifests_at(&a.manifest)?;
execute(&manifests, &a)
}
Verb::Replay(a) => {
let manifests = manifests_at(&a.manifest)?;
replay(&manifests, &a)
}
Verb::Card(a) => card(&a),
Verb::Serve(a) => {
let manifests = manifests_at(&a.manifest)?;
serve(&manifests, &a)
}
}
}
fn read_checkpoint(path: &str) -> Result<agentplane::journal::Checkpoint, String> {
let text =
std::fs::read_to_string(path).map_err(|e| format!("reading --checkpoint {path}: {e}"))?;
if let Ok(cp) = agentplane::journal::Checkpoint::from_note(&text) {
return Ok(cp);
}
if let Ok(note) = agentplane::journal::SignedNote::parse(&text)
&& let Ok(cp) = agentplane::journal::Checkpoint::from_note(¬e.text)
{
return Ok(cp);
}
serde_json::from_str(&text).map_err(|e| {
format!(
"--checkpoint {path} is neither a tlog-checkpoint note nor the `current` \
field of an audit report: {e}"
)
})
}
fn verify_verb(opts: &VerifyArgs) -> Result<ExitCode, String> {
let verifier = verifier_from(&opts.key)?;
let expected = match &opts.checkpoint {
Some(path) => Some(read_checkpoint(path)?),
None => None,
};
let verifier = verifier
.as_ref()
.map(|v| v as &dyn agentplane::core::Verifier);
let report = if opts.file == "-" {
agentplane::export::verify(std::io::stdin().lock(), verifier, expected.as_ref())
.map_err(|e| e.to_string())
} else {
let file =
std::fs::File::open(&opts.file).map_err(|e| format!("reading {}: {e}", opts.file))?;
agentplane::export::verify(std::io::BufReader::new(file), verifier, expected.as_ref())
.map_err(|e| e.to_string())
}?;
println!(
"{}",
serde_json::to_string_pretty(&report).map_err(|e| e.to_string())?
);
Ok(if report.is_sound() {
ExitCode::SUCCESS
} else {
ExitCode::FAILURE
})
}
fn card(opts: &CardArgs) -> Result<ExitCode, String> {
let manifests = manifests_at(&opts.manifest)?;
let [manifest] = manifests.as_slice() else {
return Err(format!(
"`card` describes one agent and this file holds {}. A2A's card \
path is well-known and singular — split the file, or point this \
at the document you would serve",
manifests.len()
));
};
let card =
agentplane::peers::AgentCard::derive(manifest, &opts.url).map_err(|e| e.to_string())?;
println!(
"{}",
serde_json::to_string_pretty(&card).map_err(|e| e.to_string())?
);
Ok(ExitCode::SUCCESS)
}
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), _) if s == "-" => {
use std::io::Read as _;
let mut text = String::new();
std::io::stdin()
.read_to_string(&mut text)
.map_err(|e| format!("reading standard input: {e}"))?;
text
}
(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_on(Arc::clone(&store)),
std::slice::from_ref(manifest),
)
.await?;
for (name, client) in connect_mcp_servers(&opts.mcp, std::slice::from_ref(manifest)).await?
{
builder = builder.tool_server(name, client);
}
if let Some((registry, client)) = connect_peers(&opts.peer, std::slice::from_ref(manifest))?
{
builder = builder.peers(registry, client);
}
builder = builder
.policy(Arc::new(policy) as Arc<dyn agentplane::core::PolicyEngine>)
.agent(agentplane::runtime::Agent::new(manifest));
if !opts.push_host.is_empty() {
builder = builder.push(Arc::clone(&store) as Arc<dyn agentplane::push::PushStore>);
}
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));
spawn_drill(&runtime, opts.drill_every.unwrap_or(0));
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_drill(runtime: &Arc<Runtime>, every: u32) {
if every == 0 {
return;
}
if runtime.cases().is_none() {
eprintln!(
" drill: --drill-every was given but this plane has no case store, \
so there are no cases to walk"
);
return;
}
let plane = 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;
match plane.drill().await {
Ok(report) if !report.is_sound() => {
tracing::error!(?report, "the recovery drill found unrecoverable references");
}
Ok(report) if !report.not_checked.is_empty() => {
tracing::warn!(
?report,
"the recovery drill passed, but could not check everything"
);
}
Ok(report) => tracing::info!(cases = report.cases, "the recovery drill passed"),
Err(error) => tracing::error!(%error, "the recovery drill could not run"),
}
}
});
}
#[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(w) if w.failed > 0 => {
tracing::warn!(fired = w.fired, failed = w.failed, "timer wakes failed");
}
Ok(w) if w.fired > 0 => tracing::info!(fired = w.fired, "timers fired"),
Ok(_) => {}
Err(error) => tracing::error!(%error, "firing timers failed"),
}
match sweeper
.sweep(now, std::time::Duration::from_secs(3600))
.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(())
}
type WiredPeers = (
agentplane::peers::PeerRegistry,
Arc<dyn agentplane::peers::PeerClient>,
);
#[cfg(feature = "a2a")]
fn connect_peers(specs: &[String], manifests: &[Manifest]) -> Result<Option<WiredPeers>, String> {
use agentplane::core::Scope;
use agentplane::peers::a2a::{A2aClient, Endpoint};
use agentplane::peers::{PeerCredential, PeerGrant, PeerId, PeerRegistry, PeerRouter};
if specs.is_empty() {
return Ok(None);
}
let mut registry = PeerRegistry::new();
let mut router = PeerRouter::new();
for spec in specs {
let (name, url) = spec.split_once('=').ok_or_else(|| {
format!(
"--peer wants `<name>=<url>`, got `{spec}`. The name is the one your \
manifest's grants use: a grant `tool://reviewer/audit.check` needs \
`--peer reviewer=https://...`"
)
})?;
if name.trim().is_empty() || url.trim().is_empty() {
return Err(format!("--peer `{spec}` names no peer or no URL"));
}
let granted: Vec<String> = manifests
.iter()
.flat_map(|m| &m.spec.tools)
.filter_map(|g| agentplane::tools::ToolId::parse(&g.reference))
.filter(|id| id.server == name)
.map(|id| id.tool)
.collect();
if granted.is_empty() {
return Err(format!(
"--peer `{name}` is wired but no manifest grants a `tool://{name}/…` \
capability, so nothing could ever call it"
));
}
let peer = PeerId::new(name);
let mut grant = PeerGrant::new(Scope::of(granted.iter().cloned()));
let token_var = format!(
"AGENTPLANE_PEER_TOKEN_{}",
name.to_ascii_uppercase().replace(['.', '-'], "_")
);
if let Ok(token) = std::env::var(&token_var)
&& !token.is_empty()
{
grant = grant.with_credential(&peer, PeerCredential::for_audience(peer.clone(), token));
}
registry = registry.allow(peer.clone(), grant);
let client = A2aClient::new(Endpoint::new(url))
.map_err(|e| format!("could not build a client for peer `{name}`: {e}"))?
.allow_loopback();
router = router.peer(
peer,
Arc::new(client) as Arc<dyn agentplane::peers::PeerClient>,
);
eprintln!(" peer: {name} <- {url} ({})", granted.join(", "));
}
Ok(Some((
registry,
Arc::new(router) as Arc<dyn agentplane::peers::PeerClient>,
)))
}
#[cfg(not(feature = "a2a"))]
fn connect_peers(specs: &[String], _manifests: &[Manifest]) -> Result<Option<WiredPeers>, String> {
if specs.is_empty() {
return Ok(None);
}
Err(
"this build cannot call an A2A peer: `--peer` needs the `a2a` feature. Reinstall \
with `--features cli,a2a`, or use the `:full` container image"
.to_owned(),
)
}
#[cfg(feature = "mcp-stdio")]
async fn connect_mcp_servers(
specs: &[String],
manifests: &[Manifest],
) -> Result<Vec<(String, Arc<dyn agentplane::tools::ToolClient>)>, String> {
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 client = agentplane::tools::McpClient::connect(
name,
transport,
agentplane::tools::Destination::Local,
)
.await
.map_err(|e| format!("the MCP server `{name}` cannot be used: {e}"))?;
match client.negotiated_version() {
Some(version) => eprintln!(" mcp: {name} <- {command} (MCP {version})"),
None => eprintln!(" mcp: {name} <- {command}"),
}
match client.discover().await {
Ok(advertised) => {
for manifest in manifests {
let mut catalog = agentplane::tools::ToolCatalog::from_manifest(manifest);
for (id, adv) in &advertised {
catalog = catalog.observed(id, *adv);
}
for id in catalog.overclaiming() {
eprintln!(
" mcp: {name}: warning: `{id}` advertises more safety than \
manifest `{}` grants; the grant still rules",
manifest.metadata.name
);
}
}
}
Err(e) => {
eprintln!(" mcp: {name}: tools/list failed, advertisements not compared: {e}");
}
}
wired.push((
name.to_owned(),
Arc::new(client) as Arc<dyn agentplane::tools::ToolClient>,
));
}
Ok(wired)
}
#[cfg(not(feature = "mcp-stdio"))]
#[allow(clippy::unused_async)]
async fn connect_mcp_servers(
specs: &[String],
_manifests: &[Manifest],
) -> 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.needs_attention() => {
tracing::warn!(?report, "push delivery needs attention");
}
Ok(report) if report.deliveries > 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 require_declarative(manifests: &[Manifest]) -> Result<(), 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
));
}
}
Ok(())
}
fn conclude(outcome: &agentplane::runtime::RunOutcome) -> ExitCode {
eprintln!("run {} — {:?}", outcome.run_id, outcome.status);
if let Some(output) = &outcome.output {
println!("{}", output.peek());
}
if matches!(outcome.status, RunStatus::Succeeded) {
ExitCode::SUCCESS
} else {
ExitCode::FAILURE
}
}
fn execute(manifests: &[Manifest], opts: &RunArgs) -> Result<ExitCode, String> {
require_declarative(manifests)?;
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, manifests).await? {
builder = builder.tool_server(name, client);
}
if let Some((registry, client)) = connect_peers(&opts.peer, manifests)? {
builder = builder.peers(registry, 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 capability = entry_capability(manifests, opts.capability.as_deref())?;
let outcome = agent
.run(&capability, Tainted::trusted(opts.read_input()?))
.await
.map_err(|e| e.to_string())?;
Ok(conclude(&outcome))
})
}
fn replay(manifests: &[Manifest], opts: &ReplayArgs) -> Result<ExitCode, String> {
require_declarative(manifests)?;
let run = agentplane::core::RunId::parse(&opts.run_id)
.map_err(|e| format!("`{}` is not a run id: {e}", opts.run_id))?;
let mode = if opts.strict {
Mode::Strict
} else {
Mode::Resume
};
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::new(RedbStore::open(&opts.store).map_err(|e| e.to_string())?);
let mut builder =
with_providers(Runtime::builder_on(Arc::clone(&store)), manifests).await?;
for (name, client) in connect_mcp_servers(&opts.mcp, manifests).await? {
builder = builder.tool_server(name, client);
}
if let Some((registry, client)) = connect_peers(&opts.peer, manifests)? {
builder = builder.peers(registry, 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 = agent.replay(run, mode).await.map_err(|e| e.to_string())?;
Ok(conclude(&outcome))
})
}
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"))
}