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(ValidateArgs),
Schema,
Digest(FileArgs),
Audit(AuditArgs),
Export(StoreArgs),
Verify(VerifyArgs),
Restore(RestoreArgs),
Drill(DrillArgs),
ForgetAdmissions(ForgetArgs),
Retain(RetainArgs),
Halt(HaltArgs),
Halts(HaltsArgs),
Hold(HoldArgs),
}
#[derive(clap::Args, Debug)]
struct RetainArgs {
#[command(flatten)]
at: StoreRef,
#[arg(long)]
older_than_days: u32,
#[arg(long)]
reason: String,
#[arg(long)]
dry_run: bool,
}
#[derive(clap::Args, Debug)]
struct HoldArgs {
#[command(flatten)]
at: StoreRef,
#[arg(long)]
case: Option<String>,
#[arg(long)]
reason: Option<String>,
#[arg(long)]
actor: Option<String>,
#[arg(long, conflicts_with_all = ["reason", "actor"])]
lift: bool,
}
#[derive(clap::Args, Debug)]
struct HaltArgs {
#[command(flatten)]
at: StoreRef,
#[arg(long, default_value = "tenant")]
scope: String,
#[arg(long)]
reason: Option<String>,
#[arg(long)]
actor: Option<String>,
#[arg(long, conflicts_with_all = ["reason", "actor"])]
lift: bool,
}
#[derive(clap::Args, Debug)]
struct HaltsArgs {
#[command(flatten)]
at: StoreRef,
}
#[derive(clap::Args, Debug)]
struct ForgetArgs {
#[command(flatten)]
at: StoreRef,
#[arg(long)]
older_than_days: u32,
}
#[derive(clap::Args, Debug)]
struct DrillArgs {
#[command(flatten)]
at: StoreRef,
}
#[derive(clap::Args, Debug)]
struct RestoreArgs {
file: String,
#[command(flatten)]
at: StoreRef,
}
#[derive(clap::Args, Debug)]
struct VerifyArgs {
file: String,
#[arg(long)]
key: Vec<String>,
#[arg(long)]
checkpoint: Option<String>,
#[arg(long = "witness")]
witness: Vec<String>,
#[arg(long = "witness-key")]
witness_key: Vec<String>,
#[arg(long)]
origin: 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,
#[arg(long = "witness")]
witness: Vec<String>,
#[arg(long = "witness-key")]
witness_key: Vec<String>,
#[arg(long)]
origin: Option<String>,
}
#[derive(clap::Args, Debug)]
struct StoreArgs {
#[command(flatten)]
at: StoreRef,
#[arg(long)]
outcome: Vec<String>,
#[arg(long, default_value_t = 1000)]
limit: usize,
}
#[derive(clap::Args, Debug, Clone)]
struct StoreRef {
#[arg(long, env = "AGENTPLANE_STORE")]
store: String,
#[arg(long, env = "AGENTPLANE_TENANT")]
tenant: Option<String>,
}
impl StoreRef {
async fn open(&self) -> Result<Backend, String> {
Backend::open(&self.store, self.tenant.as_deref()).await
}
}
#[derive(clap::Args, Debug, Clone)]
struct MaybeStoreRef {
#[arg(long, env = "AGENTPLANE_STORE")]
store: Option<String>,
#[arg(long, env = "AGENTPLANE_TENANT")]
tenant: Option<String>,
}
impl MaybeStoreRef {
async fn open(&self) -> Result<Option<Backend>, String> {
match &self.store {
Some(spec) => Backend::open(spec, self.tenant.as_deref()).await.map(Some),
None => Ok(None),
}
}
}
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 ValidateArgs {
manifest: String,
#[arg(long = "require-annotation", value_name = "KEY")]
require_annotation: Vec<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>,
#[command(flatten)]
at: MaybeStoreRef,
#[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,
#[command(flatten)]
at: StoreRef,
#[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>,
#[command(flatten)]
at: MaybeStoreRef,
#[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>,
#[arg(
long,
value_name = "SECS",
env = "AGENTPLANE_DRAIN_SECS",
default_value_t = 25
)]
drain_secs: u64,
}
#[derive(Debug, Default, serde::Serialize)]
struct Anchor {
#[serde(skip)]
checkpoint: Option<agentplane::journal::Checkpoint>,
#[serde(skip_serializing_if = "Option::is_none")]
obtained_from: Option<String>,
#[serde(skip_serializing_if = "Vec::is_empty")]
cosigned_by: Vec<String>,
#[serde(skip_serializing_if = "Vec::is_empty")]
unreached: Vec<String>,
#[serde(skip_serializing_if = "Vec::is_empty")]
split_view: Vec<String>,
}
#[derive(serde::Serialize)]
struct AuditDocument<'a> {
anchor: &'a Anchor,
#[serde(flatten)]
report: &'a agentplane::audit::AuditReport,
}
#[derive(serde::Serialize)]
struct VerifyDocument<'a> {
anchor: &'a Anchor,
#[serde(flatten)]
report: &'a agentplane::export::VerifyReport,
}
fn witness_keys(keys: &[String]) -> Result<Vec<agentplane::journal::TrustedWitness>, String> {
let mut out = Vec::new();
for entry in keys {
let Some((name, encoded)) = entry.split_once('=') else {
return Err(format!(
"--witness-key takes <name>=<base64 Ed25519 public key>, got '{entry}' — \
the name is the one the witness signs its lines with, and the key is \
what makes a cosignature checkable rather than a string"
));
};
let bytes = base64::Engine::decode(&base64::engine::general_purpose::STANDARD, encoded)
.map_err(|e| format!("--witness-key {name}: not base64: {e}"))?;
let key: [u8; 32] = bytes.try_into().map_err(|b: Vec<u8>| {
format!(
"--witness-key {name}: an Ed25519 public key is 32 bytes, this is {}",
b.len()
)
})?;
out.push(agentplane::journal::TrustedWitness::ed25519(name, key));
}
Ok(out)
}
async fn anchor_from_witnesses(
prefixes: &[String],
keys: &[String],
origin: &str,
) -> Result<Anchor, String> {
let mut anchor = Anchor::default();
if prefixes.is_empty() {
if !keys.is_empty() {
return Err("--witness-key was given with no --witness to use it against".to_owned());
}
return Ok(anchor);
}
let trusted = witness_keys(keys)?;
if trusted.is_empty() {
return Err(
"--witness needs at least one --witness-key: a checkpoint fetched from a URL \
nobody's signature covers is a stranger's claim presented as an independent \
anchor"
.to_owned(),
);
}
let mut held: Vec<(String, agentplane::journal::Checkpoint)> = Vec::new();
for prefix in prefixes {
let reader = agentplane::journal::WitnessReader::new(prefix, trusted.clone())
.map_err(|e| format!("--witness {prefix}: {e}"))?;
match reader.latest(origin).await {
Ok(Some(cosigned)) => {
eprintln!(
"witness {prefix}: log '{}' at size {} with root {}, {} cosignature(s)",
cosigned.checkpoint.origin,
cosigned.checkpoint.size,
cosigned.checkpoint.root.to_hex(),
cosigned.cosignatures.len(),
);
if anchor
.checkpoint
.as_ref()
.is_none_or(|b| b.size < cosigned.checkpoint.size)
{
anchor.checkpoint = Some(cosigned.checkpoint.clone());
anchor.obtained_from = Some(format!("witness {prefix}"));
anchor.cosigned_by = cosigned
.cosignatures
.iter()
.map(|c| c.key_id.clone())
.collect();
}
held.push((prefix.clone(), cosigned.checkpoint));
}
Ok(None) => {
let said = format!(
"witness {prefix}: has never cosigned log '{origin}' — no anchor from \
this one, and nothing here is evidence about deletion"
);
eprintln!("{said}");
anchor.unreached.push(said);
}
Err(e) => {
let said = format!("witness {prefix}: {e}");
eprintln!("{said}");
anchor.unreached.push(said);
}
}
}
for split in agentplane::journal::split_views(&held) {
eprintln!("{split}");
anchor.split_view.push(split.to_string());
}
Ok(anchor)
}
async fn audit_report(
store: &Arc<dyn JournalStore>,
runs: &[agentplane::RunId],
audit: &AuditArgs,
) -> Result<ExitCode, String> {
let verifier = verifier_from(&audit.key)?;
let origin = match &audit.origin {
Some(o) => o.clone(),
None => store.checkpoint().await.map_err(|e| e.to_string())?.origin,
};
let fetched = anchor_from_witnesses(&audit.witness, &audit.witness_key, &origin).await?;
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 mut anchor = fetched;
match (prior, anchor.checkpoint.as_ref()) {
(Some(saved), Some(found)) if found.size >= saved.size => {}
(Some(saved), _) => {
anchor.obtained_from = Some(match &audit.prior {
Some(path) => format!("file {path}"),
None => "file".to_owned(),
});
anchor.cosigned_by.clear();
anchor.checkpoint = Some(saved);
}
(None, _) => {}
}
let evidence = agentplane::audit::Evidence {
prior: anchor.checkpoint.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(&AuditDocument {
anchor: &anchor,
report: &report,
})
.map_err(|e| e.to_string())?
);
Ok(if report.is_sound() && anchor.split_view.is_empty() {
ExitCode::SUCCESS
} else {
ExitCode::FAILURE
})
}
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 backend = opts.at.open().await?;
let store = backend.journal();
let cases = backend.cases();
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 opts.outcome.is_empty() {
let flight = agentplane::export::runs_in_flight(&store, opts.limit)
.await
.map_err(|e| e.to_string())?;
if flight.truncated {
truncated.push("in-flight runs".to_owned());
}
for (run, why) in &flight.unreadable {
eprintln!("warning: in-flight run {run} could not be read: {why}");
}
if !flight.runs.is_empty() {
eprintln!(
"including {} run(s) still in flight — sleeping, awaiting a message, \
or waiting on a person. The Merkle log commits to sealed runs only, \
so these are carried and the checkpoint does not cover them",
flight.runs.len()
);
}
runs.extend(flight.runs);
}
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);
};
audit_report(&store, &runs, audit).await
})
}
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 backend = opts.at.open().await?;
let cases = backend.cases();
let tenant = opts
.at
.tenant
.as_deref()
.map(|name| agentplane::core::TenantId::new(name).map_err(|e| format!("--tenant: {e}")))
.transpose()?
.unwrap_or_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 cutoff_before(
now: agentplane::core::Timestamp,
days: u32,
) -> Result<agentplane::core::Timestamp, String> {
now.checked_sub(time::Duration::days(i64::from(days)))
.ok_or_else(|| {
format!(
"--older-than-days {days} reaches back past the first instant this \
runtime can name"
)
})
}
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 = opts.at.open().await?.journal();
#[allow(clippy::disallowed_methods)]
let now = time::OffsetDateTime::now_utc();
let cutoff = cutoff_before(now, opts.older_than_days)?;
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 cases = opts.at.open().await?.cases();
#[allow(clippy::disallowed_methods)]
let now = time::OffsetDateTime::now_utc();
let cutoff = cutoff_before(now, opts.older_than_days)?;
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 hold_verb(opts: &HoldArgs) -> Result<ExitCode, String> {
if opts.case.is_none() && (opts.lift || opts.reason.is_some()) {
return Err(
"--lift and --reason act on one matter: name it with --case, or pass \
neither to list every hold standing"
.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 cases = opts.at.open().await?.cases();
let Some(case) = opts.case.as_deref() else {
let mut standing = Vec::new();
let mut after = None;
loop {
let page = cases.holds(after, 256).await.map_err(|e| e.to_string())?;
if page.is_empty() {
break;
}
after = page.last().map(|(c, _)| *c);
for (id, hold) in page {
standing.push(serde_json::json!({
"case": id.to_string(),
"placed_at": hold.placed_at.unix_timestamp(),
"reason": hold.reason,
"by": hold.by.actor(),
"basis": hold.by.basis().as_str(),
}));
}
}
println!(
"{}",
serde_json::to_string_pretty(&serde_json::json!({ "holds": standing }))
.map_err(|e| e.to_string())?
);
return Ok(ExitCode::SUCCESS);
};
let case = agentplane::core::CaseId::parse(case).map_err(|e| format!("--case: {e}"))?;
if opts.lift {
let lifted = cases.release_hold(case).await.map_err(|e| e.to_string())?;
println!(
"{}",
serde_json::to_string_pretty(&serde_json::json!({
"case": case.to_string(),
"lifted": lifted,
}))
.map_err(|e| e.to_string())?
);
return Ok(ExitCode::SUCCESS);
}
let Some(reason) = opts.reason.as_deref() else {
return Err(concat!(
"--reason is required to place a hold: a preservation nobody can account for ",
"is indistinguishable from a sweep that quietly stopped working. ",
"Use --lift to release one"
)
.to_owned());
};
let Some(actor) = opts.actor.as_deref() else {
return Err(concat!(
"--actor is required to place a hold: a preservation order the runtime ",
"cannot check is worth the name beside it, and the row records that this ",
"one was asserted rather than authenticated"
)
.to_owned());
};
let by = agentplane::core::Operator::asserted(actor).map_err(|e| e.to_string())?;
#[allow(clippy::disallowed_methods)]
let now = time::OffsetDateTime::now_utc();
let placed = cases
.place_hold(
case,
&agentplane::core::LegalHold {
placed_at: now,
reason: reason.to_owned(),
by,
},
)
.await
.map_err(|e| e.to_string())?;
println!(
"{}",
serde_json::to_string_pretty(&serde_json::json!({
"case": case.to_string(),
"placed": placed,
"in_force": cases.hold(case).await.map_err(|e| e.to_string())?
.map(|h| serde_json::json!({
"placed_at": h.placed_at.unix_timestamp(),
"reason": h.reason,
"by": h.by.actor(),
"basis": h.by.basis().as_str(),
})),
}))
.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 {}",
opts.scope,
agentplane::quota::HaltScope::forms()
)
})?;
let thrown = if opts.lift {
None
} else {
let reason = opts.reason.as_deref().ok_or(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"
))?;
let actor = opts.actor.as_deref().ok_or(concat!(
"--actor is required to halt: the runtime cannot check an emergency stop, ",
"so the name beside it is the whole of its evidence. It is recorded as ",
"asserted — nothing here verified it"
))?;
let by = agentplane::core::Operator::asserted(actor).map_err(|e| e.to_string())?;
Some((by, reason))
};
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 = opts.at.open().await?.quotas();
let printed = if let Some((by, reason)) = &thrown {
{
#[allow(clippy::disallowed_methods)]
let now = time::OffsetDateTime::now_utc();
quotas
.set_halt(&scope, by, now, reason)
.await
.map_err(|e| e.to_string())?;
serde_json::json!({
"scope": scope.key(),
"halted": true,
"reason": reason,
"by": by.actor(),
"basis": by.basis().as_str(),
})
}
} else {
{
let was_standing = quotas.lift_halt(&scope).await.map_err(|e| e.to_string())?;
serde_json::json!({
"scope": scope.key(),
"halted": false,
"was_standing": was_standing,
})
}
};
println!("{printed}");
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 = opts.at.open().await?.quotas();
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,
"by": h.by.actor(),
"basis": h.by.basis().as_str(),
"thrown_at": h.at.unix_timestamp(),
})
})
.collect();
println!(
"{}",
serde_json::to_string_pretty(&serde_json::json!({ "halts": rows }))
.map_err(|e| e.to_string())?
);
Ok(ExitCode::SUCCESS)
})
}
enum Backend {
Embedded(Arc<RedbStore>, agentplane::core::TenantId),
#[cfg(feature = "postgres")]
Shared(
Arc<agentplane::store::PostgresStore>,
agentplane::core::TenantId,
),
}
impl Backend {
#[allow(clippy::unused_async, clippy::unused_async_trait_impl)]
async fn open(spec: &str, tenant: Option<&str>) -> Result<Self, String> {
let tenant = tenant
.map(|name| agentplane::core::TenantId::new(name).map_err(|e| format!("--tenant: {e}")))
.transpose()?;
if is_connection_string(spec) {
#[cfg(not(feature = "postgres"))]
return Err(
"--store names a PostgreSQL database and this build cannot open one. \
Reinstall with `--features cli,postgres`, or use the `:full` container \
image, which is built with it"
.to_owned(),
);
#[cfg(feature = "postgres")]
return Self::shared(spec, tenant).await;
}
let store = RedbStore::open(spec).map_err(|e| held_by_a_plane(&e.to_string()))?;
let tenant = tenant.unwrap_or_default();
Ok(Self::Embedded(
Arc::new(store.for_tenant(tenant.clone())),
tenant,
))
}
#[cfg(feature = "postgres")]
async fn shared(url: &str, tenant: Option<agentplane::core::TenantId>) -> Result<Self, String> {
let store = agentplane::store::PostgresStore::connect(url)
.await
.map_err(|e| e.to_string())?;
let tenant = tenant.unwrap_or_default();
Ok(Self::Shared(
Arc::new(store.for_tenant(tenant.clone())),
tenant,
))
}
fn stores(&self) -> agentplane::runtime::Stores {
match self {
Self::Embedded(s, _) => agentplane::runtime::Stores::on(Arc::clone(s)),
#[cfg(feature = "postgres")]
Self::Shared(s, _) => agentplane::runtime::Stores::on(Arc::clone(s)),
}
}
fn journal(&self) -> Arc<dyn JournalStore> {
match self {
Self::Embedded(s, _) => Arc::clone(s) as _,
#[cfg(feature = "postgres")]
Self::Shared(s, _) => Arc::clone(s) as _,
}
}
fn cases(&self) -> Arc<dyn agentplane::case::CaseStore> {
match self {
Self::Embedded(s, _) => Arc::clone(s) as _,
#[cfg(feature = "postgres")]
Self::Shared(s, _) => Arc::clone(s) as _,
}
}
fn quotas(&self) -> Arc<dyn agentplane::quota::QuotaStore> {
match self {
Self::Embedded(s, _) => Arc::clone(s) as _,
#[cfg(feature = "postgres")]
Self::Shared(s, _) => Arc::clone(s) as _,
}
}
#[cfg(feature = "push")]
fn push(&self) -> Arc<dyn agentplane::push::PushStore> {
match self {
Self::Embedded(s, _) => Arc::clone(s) as _,
#[cfg(feature = "postgres")]
Self::Shared(s, _) => Arc::clone(s) as _,
}
}
fn tenant(&self) -> agentplane::core::TenantId {
match self {
Self::Embedded(_, t) => t.clone(),
#[cfg(feature = "postgres")]
Self::Shared(_, t) => t.clone(),
}
}
#[cfg(all(feature = "a2a-server", feature = "cedar"))]
const fn describe(&self) -> &'static str {
match self {
Self::Embedded(..) => "embedded redb file",
#[cfg(feature = "postgres")]
Self::Shared(..) => "shared PostgreSQL store",
}
}
}
fn held_by_a_plane(detail: &str) -> String {
if !detail.contains("already open") {
return detail.to_owned();
}
format!(
"{detail}\n\nAn embedded redb store admits one writer process, and something \
else is holding this one — most likely `agentplane serve`. Either stop that \
process and run this again, or put the plane on a shared store \
(`--store postgres://…`), where an operator verb and a serving plane \
coexist. To act on a *running* embedded plane meanwhile, the operator API \
is the surface that reaches it."
)
}
fn is_connection_string(spec: &str) -> bool {
spec.starts_with("postgres://") || spec.starts_with("postgresql://")
}
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) => validate(&a),
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::Hold(a) => hold_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 backend = a.at.open().await?;
let store = backend.journal();
let cases = backend.cases();
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 saved = match &opts.checkpoint {
Some(path) => Some(read_checkpoint(path)?),
None => None,
};
let fetched = match (&opts.origin, opts.witness.is_empty()) {
(Some(origin), false) => {
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(anchor_from_witnesses(
&opts.witness,
&opts.witness_key,
origin,
))?
}
(None, false) => {
return Err(
"--witness needs --origin: the log's name cannot come from the file being \
checked, because that header is written by whoever wrote the file"
.to_owned(),
);
}
(_, true) => {
if !opts.witness_key.is_empty() {
return Err(
"--witness-key was given with no --witness to use it against".to_owned(),
);
}
Anchor::default()
}
};
let mut anchor = fetched;
match (saved, anchor.checkpoint.as_ref()) {
(Some(saved), Some(found)) if found.size >= saved.size => {}
(Some(saved), _) => {
anchor.obtained_from = Some(match &opts.checkpoint {
Some(path) => format!("file {path}"),
None => "file".to_owned(),
});
anchor.cosigned_by.clear();
anchor.checkpoint = Some(saved);
}
(None, _) => {}
}
let expected = anchor.checkpoint.clone();
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(&VerifyDocument {
anchor: &anchor,
report: &report,
})
.map_err(|e| e.to_string())?
);
Ok(if report.is_sound() && anchor.split_view.is_empty() {
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 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 backend = opts.at.open().await?.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 mut builder = with_providers(
Runtime::builder_with(backend.stores()).tenant(backend.tenant()),
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(backend.push());
}
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, &backend)?;
serve_until_stopped(
&runtime,
server,
operator_auth,
opts,
manifest,
url,
&backend,
)
.await?;
Ok(ExitCode::SUCCESS)
})
}
#[cfg(all(feature = "a2a-server", feature = "cedar"))]
async fn serve_until_stopped(
runtime: &Arc<Runtime>,
server: agentplane::api::a2a::A2aServer,
operator_auth: Arc<dyn agentplane::api::Authenticator>,
opts: &ServeArgs,
manifest: &Manifest,
url: &str,
backend: &Backend,
) -> Result<(), String> {
let addr = opts.addr.as_str();
let (stop_tx, stop_rx) = tokio::sync::watch::channel(false);
let mut background: Vec<Task> = Vec::new();
if let Some(worker) = server.push_worker() {
background.extend(spawn_push_worker(
worker,
opts.sweep_every.unwrap_or(DEFAULT_SWEEP_SECONDS),
stop_rx.clone(),
));
}
background.extend(spawn_sweeper(
runtime,
opts.sweep_every.unwrap_or(DEFAULT_SWEEP_SECONDS),
stop_rx.clone(),
));
background.extend(spawn_drill(
runtime,
opts.drill_every.unwrap_or(0),
stop_rx.clone(),
));
if let Some(operator_addr) = opts.operator_addr.as_deref() {
background.push(
spawn_operator_surface(runtime, operator_auth, operator_addr, stop_rx.clone()).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");
eprintln!(" store: {}", backend.describe());
eprintln!(" stop: SIGTERM drains for up to {}s", opts.drain_secs);
let mut peer = tokio::spawn(async move {
axum::serve(listener, server.router())
.with_graceful_shutdown(stopping(stop_rx))
.await
});
let mut ended = tokio::select! {
() = stop_requested() => None,
joined = &mut peer => Some(joined),
};
let still_serving = ended.is_none();
let _ = stop_tx.send(true);
let grace = std::time::Duration::from_secs(opts.drain_secs);
let deadline = tokio::time::Instant::now() + grace;
let (report, closed) = tokio::join!(
runtime.drain(grace),
stop_serving(&mut peer, background, still_serving, deadline),
);
ended = ended.or(closed);
report_drain(&report);
match ended {
Some(Ok(listening)) => listening.map_err(|e| format!("the server stopped: {e}")),
Some(Err(e)) => Err(format!("the server task failed: {e}")),
None => Ok(()),
}
}
#[cfg(all(feature = "a2a-server", feature = "cedar"))]
async fn stop_serving(
peer: &mut tokio::task::JoinHandle<std::io::Result<()>>,
background: Vec<Task>,
still_serving: bool,
deadline: tokio::time::Instant,
) -> Option<Result<std::io::Result<()>, tokio::task::JoinError>> {
let mut ended = None;
if still_serving {
match tokio::time::timeout_at(deadline, peer).await {
Ok(joined) => ended = Some(joined),
Err(_) => eprintln!(" stop: connections were still open at the grace period"),
}
}
for task in background {
if tokio::time::timeout_at(deadline, task).await.is_err() {
break;
}
}
ended
}
#[cfg(all(feature = "a2a-server", feature = "cedar"))]
fn report_drain(report: &agentplane::runtime::DrainReport) {
if report.is_complete() {
eprintln!(
"stopped; {} background runs finished first",
report.settled()
);
return;
}
let unfinished: Vec<String> = report.unfinished.iter().map(ToString::to_string).collect();
tracing::warn!(
settled = report.settled(),
unfinished = ?unfinished,
"the grace period ended with runs still executing; they are left for the recovery sweep"
);
eprintln!(
"stopped; {} background runs finished, {} left to the recovery sweep: {}",
report.settled(),
unfinished.len(),
unfinished.join(" ")
);
}
#[cfg(all(feature = "a2a-server", feature = "cedar"))]
type Stop = tokio::sync::watch::Receiver<bool>;
#[cfg(all(feature = "a2a-server", feature = "cedar"))]
type Task = tokio::task::JoinHandle<()>;
#[cfg(all(feature = "a2a-server", feature = "cedar"))]
async fn next_tick(tick: &mut tokio::time::Interval, stop: &mut Stop) -> bool {
tokio::select! {
_ = tick.tick() => true,
_ = stop.changed() => false,
}
}
#[cfg(all(feature = "a2a-server", feature = "cedar"))]
async fn stopping(mut stop: Stop) {
let _ = stop.changed().await;
}
#[cfg(all(feature = "a2a-server", feature = "cedar"))]
async fn stop_requested() {
#[cfg(unix)]
{
use tokio::signal::unix::{SignalKind, signal};
let term = match signal(SignalKind::terminate()) {
Ok(s) => Some(s),
Err(error) => {
tracing::error!(%error, "could not listen for SIGTERM; this process will not drain");
None
}
};
let terminated = async move {
match term {
Some(mut term) => {
term.recv().await;
}
None => std::future::pending().await,
}
};
tokio::select! {
() = terminated => {}
r = tokio::signal::ctrl_c() => {
if let Err(error) = r {
tracing::error!(%error, "could not listen for SIGINT");
}
}
}
}
#[cfg(not(unix))]
{
if let Err(error) = tokio::signal::ctrl_c().await {
tracing::error!(%error, "could not listen for an interrupt");
}
}
}
#[cfg(all(feature = "a2a-server", feature = "cedar"))]
fn spawn_drill(runtime: &Arc<Runtime>, every: u32, stop: Stop) -> Option<Task> {
if every == 0 {
return None;
}
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 None;
}
let plane = Arc::clone(runtime);
Some(tokio::spawn(async move {
let mut tick = tokio::time::interval(std::time::Duration::from_secs(u64::from(every)));
let mut stop = stop;
loop {
if !next_tick(&mut tick, &mut stop).await {
break;
}
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, stop: Stop) -> Option<Task> {
if every == 0 {
return None;
}
let sweeper = Arc::clone(runtime);
Some(tokio::spawn(async move {
let mut tick = tokio::time::interval(std::time::Duration::from_secs(u64::from(every)));
let mut stop = stop;
loop {
if !next_tick(&mut tick, &mut stop).await {
break;
}
#[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,
stop: Stop,
) -> Result<Task, 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");
Ok(tokio::spawn(async move {
let served = axum::serve(listener, api.router())
.with_graceful_shutdown(stopping(stop))
.await;
if let Err(error) = served {
tracing::error!(%error, "the operator surface stopped");
}
}))
}
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);
if !url.starts_with("https://") {
return Err(format!(
"--peer `{name}` is `{url}`, and a peer is reached over HTTPS. A card or a \
peer answer steers the calls that follow it, so a plaintext hop would let \
the network choose them. For a peer on this machine, build with \
`--features cli,a2a,testkit`."
));
}
let client = A2aClient::new(Endpoint::new(url))
.map_err(|e| format!("could not build a client for peer `{name}`: {e}"))?;
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>,
)))
}
fn validate(a: &ValidateArgs) -> Result<ExitCode, String> {
let mut missing = Vec::new();
for m in &manifests_at(&a.manifest)? {
let absent: Vec<&String> = a
.require_annotation
.iter()
.filter(|key| !m.metadata.annotations.contains_key(*key))
.collect();
if absent.is_empty() {
println!("ok: {} {}", m.metadata.name, m.metadata.version);
} else {
for key in absent {
println!("MISSING: {} — annotation '{key}'", m.metadata.name);
missing.push(format!("{}: {key}", m.metadata.name));
}
}
}
if missing.is_empty() {
return Ok(ExitCode::SUCCESS);
}
Err(format!(
"{} required annotation(s) absent: {}",
missing.len(),
missing.join(", ")
))
}
#[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],
backend: &Backend,
) -> 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(
backend.push(),
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,
stop: Stop,
) -> Option<Task> {
if every == 0 {
return None;
}
Some(tokio::spawn(async move {
let mut tick = tokio::time::interval(std::time::Duration::from_secs(u64::from(every)));
let mut stop = stop;
loop {
if !next_tick(&mut tick, &mut stop).await {
break;
}
#[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 (stores, tenant) = if let Some(backend) = opts.at.open().await? {
(backend.stores(), backend.tenant())
} else {
eprintln!("note: journaling to memory; this run will not survive the process");
(
agentplane::runtime::Stores::on(Arc::new(
RedbStore::open_in_memory().map_err(|e| e.to_string())?,
)),
agentplane::core::TenantId::default(),
)
};
let mut builder =
with_providers(Runtime::builder_with(stores).tenant(tenant), 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 backend = opts.at.open().await?;
let mut builder = with_providers(
Runtime::builder_with(backend.stores()).tenant(backend.tenant()),
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 = "fake-model")]
"fake" => Ok(agentplane::model::fake::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 = "fake-model")]
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"))
}
#[cfg(test)]
mod tests {
use super::cutoff_before;
#[test]
fn a_retention_window_past_the_calendar_is_a_message_not_an_abort() {
let now = time::macros::datetime!(2026-09-13 12:00:00 UTC);
assert!(cutoff_before(now, 90).is_ok());
let err = cutoff_before(now, u32::MAX).expect_err("a window of 11m years");
assert!(err.contains("--older-than-days"), "{err}");
}
}