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(metrics: bool, verifying: bool, format: LogFormat) {
use tracing_subscriber::{EnvFilter, fmt};
let mut default = if metrics {
"warn,agentplane=info".to_owned()
} else {
"warn,agentplane=info,agentplane.metric=off".to_owned()
};
if verifying {
for target in [
agentplane::runtime::telemetry::NONDETERMINISM,
agentplane::runtime::telemetry::QUARANTINED,
] {
default.push(',');
default.push_str(target);
default.push_str("=off");
}
}
let filter = EnvFilter::try_from_default_env().unwrap_or_else(|_| EnvFilter::new(default));
let builder = fmt().with_env_filter(filter).with_writer(std::io::stderr);
let _ = match format {
LogFormat::Text => builder.try_init(),
LogFormat::Json => builder.json().try_init(),
};
}
#[derive(clap::ValueEnum, Clone, Copy, Debug, Default, PartialEq, Eq)]
enum LogFormat {
#[default]
Text,
Json,
}
mod exit {
pub const OK: u8 = 0;
pub const FINDING: u8 = 1;
pub const USAGE: u8 = 2;
pub const SUSPENDED: u8 = 3;
pub const OPERATIONAL: u8 = 4;
pub const PARTIAL: u8 = 5;
}
const EXIT_STATUS_HELP: &str = "Exit status:
0 ok
1 a finding or a negative answer (a failed run, an audit finding, needs attention)
2 usage: the command as typed cannot be carried out
3 a run is suspended, waiting for a person, a timer or an event
4 operational: a store, witness, network or file could not be used
5 partial: --limit truncated the answer, or a strict replay could not replay a run";
#[derive(Debug)]
enum Fault {
Usage(String),
Operational(String),
}
impl Fault {
const fn status(&self) -> u8 {
match self {
Self::Usage(_) => exit::USAGE,
Self::Operational(_) => exit::OPERATIONAL,
}
}
}
impl std::fmt::Display for Fault {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
Self::Usage(m) | Self::Operational(m) => f.write_str(m),
}
}
}
impl From<String> for Fault {
fn from(message: String) -> Self {
Self::Operational(message)
}
}
fn usage(message: impl Into<String>) -> Fault {
Fault::Usage(message.into())
}
fn build_description() -> &'static str {
static TEXT: std::sync::OnceLock<String> = std::sync::OnceLock::new();
TEXT.get_or_init(|| {
format!(
"{}\nfeatures: {}",
env!("CARGO_PKG_VERSION"),
compiled_features().join(", ")
)
})
}
fn compiled_features() -> Vec<&'static str> {
let all = [
("a2a", cfg!(feature = "a2a")),
("a2a-server", cfg!(feature = "a2a-server")),
("acp", cfg!(feature = "acp")),
("bedrock", cfg!(feature = "bedrock")),
("cedar", cfg!(feature = "cedar")),
("cli", cfg!(feature = "cli")),
("fake-model", cfg!(feature = "fake-model")),
("http", cfg!(feature = "http")),
("keyring", cfg!(feature = "keyring")),
("keyring-vault", cfg!(feature = "keyring-vault")),
("manifest", cfg!(feature = "manifest")),
("mcp", cfg!(feature = "mcp")),
("mcp-http", cfg!(feature = "mcp-http")),
("mcp-server", cfg!(feature = "mcp-server")),
("mcp-stdio", cfg!(feature = "mcp-stdio")),
("media", cfg!(feature = "media")),
("opendal", cfg!(feature = "opendal")),
("postgres", cfg!(feature = "postgres")),
("providers", cfg!(feature = "providers")),
("push", cfg!(feature = "push")),
("redb", cfg!(feature = "redb")),
("signing", cfg!(feature = "signing")),
("testkit", cfg!(feature = "testkit")),
("witness-http", cfg!(feature = "witness-http")),
];
all.iter()
.filter(|(_, on)| *on)
.map(|(name, _)| *name)
.collect()
}
#[derive(clap::Parser, Debug)]
#[command(
name = "agentplane",
version,
about = "Run an agent that is only a file",
long_version = build_description(),
after_help = EXIT_STATUS_HELP,
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 {
Init(InitArgs),
Run(RunArgs),
Replay(ReplayArgs),
Card(CardArgs),
Serve(Box<ServeArgs>),
Validate(ValidateArgs),
Schema,
Digest(DigestArgs),
Audit(AuditArgs),
Export(ExportArgs),
Verify(VerifyArgs),
Policy(PolicyArgs),
Restore(RestoreArgs),
Drill(DrillArgs),
ForgetAdmissions(ForgetArgs),
Retention(RetentionArgs),
Halt(HaltArgs),
Waiting(WaitingArgs),
Attention(WaitingArgs),
Hold(HoldArgs),
Reconcile(ReconcileArgs),
Quarantine(QuarantineArgs),
Tasks(TasksArgs),
Decide(DecideArgs),
Acknowledge(AcknowledgeArgs),
#[cfg(feature = "push")]
Rearm(RearmArgs),
Cancel(CancelArgs),
}
#[derive(clap::Args, Debug)]
struct ActingAs {
#[arg(long)]
actor: String,
}
impl ActingAs {
fn operator(&self) -> Result<agentplane::core::Operator, Fault> {
agentplane::core::Operator::asserted(&self.actor).map_err(|e| usage(e.to_string()))
}
}
#[derive(clap::Args, Debug)]
struct ReconcileArgs {
run_id: String,
#[command(flatten)]
at: StoreRef,
#[arg(long)]
effect: String,
#[arg(long)]
outcome: String,
#[arg(long)]
output: Option<String>,
#[arg(long)]
note: String,
#[command(flatten)]
who: ActingAs,
}
#[derive(clap::Args, Debug)]
struct QuarantineArgs {
run_id: String,
#[command(flatten)]
at: StoreRef,
#[arg(long)]
decision: String,
#[arg(long)]
reason: String,
#[command(flatten)]
who: ActingAs,
}
#[derive(clap::ValueEnum, Clone, Copy, Debug, PartialEq, Eq)]
enum Verdict {
Approve,
Reject,
}
#[derive(clap::Args, Debug)]
struct TasksArgs {
#[command(flatten)]
at: StoreRef,
#[arg(long, value_name = "TASK")]
show: Option<String>,
#[arg(long = "role")]
roles: Vec<String>,
#[arg(long, default_value_t = 100)]
limit: usize,
}
#[derive(clap::Args, Debug)]
struct DecideArgs {
task_id: String,
#[arg(value_enum)]
verdict: Verdict,
#[command(flatten)]
at: StoreRef,
#[arg(long)]
reason: String,
#[arg(long = "role")]
roles: Vec<String>,
#[command(flatten)]
who: ActingAs,
}
#[derive(clap::Args, Debug)]
struct AcknowledgeArgs {
case_id: String,
#[command(flatten)]
at: StoreRef,
#[arg(long)]
obligation: String,
#[arg(long)]
note: String,
#[command(flatten)]
who: ActingAs,
}
#[derive(clap::Args, Debug)]
struct CancelArgs {
run_id: String,
#[command(flatten)]
at: StoreRef,
#[arg(long)]
reason: String,
#[command(flatten)]
who: ActingAs,
}
#[cfg(feature = "push")]
#[derive(clap::Args, Debug)]
struct RearmArgs {
run_id: String,
#[command(flatten)]
at: StoreRef,
#[arg(long)]
id: String,
}
#[derive(clap::Args, Debug)]
struct RetentionArgs {
#[command(subcommand)]
act: RetentionAct,
}
#[derive(clap::Subcommand, Debug)]
enum RetentionAct {
Plan(RetentionPlanArgs),
}
#[derive(clap::Args, Debug)]
struct RetentionPlanArgs {
#[command(flatten)]
at: StoreRef,
#[arg(long)]
older_than_days: u32,
}
#[derive(clap::Subcommand, Debug)]
enum Listing {
List(ListArgs),
}
#[derive(clap::Args, Debug)]
struct ListArgs {
#[command(flatten)]
at: StoreRef,
}
#[derive(clap::Args, Debug)]
#[command(args_conflicts_with_subcommands = true, subcommand_negates_reqs = true)]
struct HoldArgs {
#[command(subcommand)]
list: Option<Listing>,
#[command(flatten)]
at: Option<StoreRef>,
#[arg(long, required = true)]
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)]
#[command(args_conflicts_with_subcommands = true, subcommand_negates_reqs = true)]
struct HaltArgs {
#[command(subcommand)]
list: Option<Listing>,
#[command(flatten)]
at: Option<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,
#[arg(long, global = true)]
json: bool,
}
#[derive(clap::Args, Debug)]
struct WaitingArgs {
#[command(flatten)]
at: StoreRef,
#[arg(long, default_value_t = 100)]
limit: usize,
}
#[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,
#[arg(long)]
last: bool,
}
#[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)]
struct ExportArgs {
#[command(flatten)]
store: StoreArgs,
#[arg(long)]
allow_partial: bool,
}
#[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, Fault> {
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>, Fault> {
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 InitArgs {
#[arg(default_value = "agent.yaml")]
path: String,
#[arg(long)]
tools: bool,
#[arg(long, default_value = "my-agent")]
name: String,
#[command(flatten)]
out: JsonFlag,
}
#[derive(clap::Args, Debug, Clone, Copy)]
struct JsonFlag {
#[arg(long)]
json: bool,
}
#[derive(clap::Args, Debug)]
struct DigestArgs {
manifest: String,
#[command(flatten)]
out: JsonFlag,
}
#[derive(clap::Args, Debug)]
struct ValidateArgs {
manifest: String,
#[arg(long = "require-annotation", value_name = "KEY")]
require_annotation: Vec<String>,
#[arg(long)]
json: bool,
}
#[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", requires = "acting_as")]
peer: Vec<String>,
#[arg(long, value_name = "SUBJECT")]
acting_as: Option<String>,
#[arg(long, value_name = "NAMESPACE=VALUE")]
correlate: Vec<String>,
}
#[derive(clap::Args, Debug)]
#[command(after_help = REPLAY_EXIT_HELP)]
struct ReplayArgs {
run_id: Option<String>,
#[command(flatten)]
at: MaybeStoreRef,
#[arg(long, value_name = "EXPORT", requires = "strict")]
from: Vec<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>,
}
const REPLAY_EXIT_HELP: &str = "Exit status of `--strict`:
0 every run replayed was verified
1 a run diverged from its record
4 a store, file or journal could not be used
5 partial: a run could not be replayed (named in the report), and none diverged";
#[derive(clap::Args, Debug)]
struct CardArgs {
manifest: String,
#[arg(long, env = "AGENTPLANE_URL")]
url: String,
}
#[derive(clap::Args, Debug)]
struct PolicyArgs {
#[command(subcommand)]
act: PolicyAct,
}
#[derive(clap::Subcommand, Debug)]
enum PolicyAct {
Check(PolicyCheckArgs),
}
#[derive(clap::Args, Debug)]
#[cfg_attr(not(feature = "cedar"), allow(dead_code))]
struct PolicyCheckArgs {
#[arg(long)]
bundle: String,
#[arg(long)]
from: String,
#[arg(long)]
candidate: Option<String>,
#[arg(long, env = "AGENTPLANE_TENANT")]
tenant: Option<String>,
#[command(flatten)]
out: JsonFlag,
}
#[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,
#[arg(
long,
value_enum,
env = "AGENTPLANE_LOG_FORMAT",
default_value_t = LogFormat::Text
)]
log_format: LogFormat,
}
#[derive(Debug, Default, serde::Serialize)]
struct Anchor {
#[serde(skip)]
checkpoints: Vec<agentplane::audit::Anchor>,
#[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,
truncated: &'a Truncation,
#[serde(flatten)]
report: &'a agentplane::audit::AuditReport,
}
#[derive(Debug, Default, serde::Serialize)]
struct Truncation {
limit: usize,
reached: Vec<String>,
}
impl Truncation {
const fn is_partial(&self) -> bool {
!self.reached.is_empty()
}
}
#[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(),
);
anchor.checkpoints.push(agentplane::audit::Anchor::new(
cosigned.checkpoint.clone(),
format!("witness {prefix}"),
));
anchor.cosigned_by.extend(
cosigned
.cosignatures
.iter()
.map(|c| format!("{prefix}:{}", c.key_id)),
);
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,
truncated: &Truncation,
) -> Result<ExitCode, Fault> {
let verifier = verifier_from(&audit.key).map_err(usage)?;
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
.map_err(usage)?;
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;
if let Some(saved) = prior {
anchor.checkpoints.push(agentplane::audit::Anchor::new(
saved,
match &audit.prior {
Some(path) => format!("file {path}"),
None => "file".to_owned(),
},
));
}
let evidence = agentplane::audit::Evidence {
anchors: &anchor.checkpoints,
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,
truncated,
report: &report,
})
.map_err(|e| e.to_string())?
);
Ok(ExitCode::from(audit_status(
report.is_sound() && anchor.split_view.is_empty(),
truncated,
)))
}
const fn refuses_partial_export(truncated: &Truncation, allow_partial: bool) -> bool {
truncated.is_partial() && !allow_partial
}
const fn audit_status(sound: bool, truncated: &Truncation) -> u8 {
if !sound {
exit::FINDING
} else if truncated.is_partial() {
exit::PARTIAL
} else {
exit::OK
}
}
async fn export_runs(
store: &Arc<dyn JournalStore>,
cases: &Arc<dyn agentplane::case::CaseStore>,
runs: &[agentplane::RunId],
truncation: &Truncation,
allow_partial: bool,
) -> Result<ExitCode, Fault> {
if refuses_partial_export(truncation, allow_partial) {
eprintln!(
"agentplane: refusing a partial export; raise --limit, narrow \
--outcome, or pass --allow-partial to write it anyway"
);
return Ok(ExitCode::from(exit::PARTIAL));
}
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::from(exit::FINDING));
}
Ok(ExitCode::from(if truncation.is_partial() {
exit::PARTIAL
} else {
exit::OK
}))
}
fn journal_verb(
opts: &StoreArgs,
audit: Option<&AuditArgs>,
allow_partial: bool,
) -> Result<ExitCode, Fault> {
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 truncation = Truncation {
limit: opts.limit,
reached: truncated,
};
return export_runs(&store, &cases, &runs, &truncation, allow_partial).await;
};
audit_report(
&store,
&runs,
audit,
&Truncation {
limit: opts.limit,
reached: truncated,
},
)
.await
})
}
fn drill_verb(opts: &DrillArgs) -> Result<ExitCode, Fault> {
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();
if opts.last {
let record = cases.last_drill().await.map_err(|e| e.to_string())?;
let Some(record) = record else {
println!("{}", serde_json::json!({ "drilled": false }));
return Ok(ExitCode::from(exit::FINDING));
};
let sound = record.sound;
println!(
"{}",
serde_json::to_string_pretty(&serde_json::json!({
"drilled": true,
"at": record.at.to_string(),
"sound": sound,
"cases": record.cases,
"findings": record.findings,
"not_checked": record.not_checked,
"origin": record.origin,
"log_size": record.size,
}))
.map_err(|e| e.to_string())?
);
return Ok(if sound {
ExitCode::SUCCESS
} else {
ExitCode::from(exit::FINDING)
});
}
let tenant = opts
.at
.tenant
.as_deref()
.map(|name| {
agentplane::core::TenantId::new(name).map_err(|e| usage(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())?;
#[allow(clippy::disallowed_methods)]
let at = agentplane::core::Timestamp::now_utc();
let checkpoint = backend
.journal()
.checkpoint()
.await
.map_err(|e| e.to_string())?;
cases
.record_drill(&agentplane::case::DrillRecord {
at,
sound: report.is_sound(),
cases: report.cases as u64,
findings: report.findings.len() as u64,
not_checked: report.not_checked.len() as u64,
origin: checkpoint.origin.clone(),
size: checkpoint.size,
})
.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::from(exit::FINDING)
})
})
}
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, Fault> {
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).map_err(usage)?;
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 retention_plan_verb(opts: &RetentionPlanArgs) -> Result<ExitCode, Fault> {
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).map_err(usage)?;
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!({
"cutoff": cutoff.unix_timestamp(),
"scanned": plan.scanned,
"would_erase": plan.due,
}))
.map_err(|e| e.to_string())?
);
Ok(ExitCode::SUCCESS)
})
}
fn missing_store() -> String {
"--store is required (or set AGENTPLANE_STORE): the plane's store, a redb file or a \
`postgres://` connection string"
.to_owned()
}
fn holds_verb(at: &StoreRef) -> Result<ExitCode, Fault> {
blocking(async {
let cases = at.open().await?.cases();
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())?
);
Ok(ExitCode::SUCCESS)
})
}
fn hold_verb(opts: &HoldArgs) -> Result<ExitCode, Fault> {
let (at, case) = match (&opts.list, &opts.at, &opts.case) {
(Some(Listing::List(list)), ..) => return holds_verb(&list.at),
(None, Some(at), Some(case)) => (at, case.as_str()),
(None, None, _) => return Err(usage(missing_store())),
(None, Some(_), None) => {
return Err(usage(
"name the matter with --case, or run `hold list` 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 = at.open().await?.cases();
let case =
agentplane::core::CaseId::parse(case).map_err(|e| usage(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(lift_status(lifted));
}
let Some(reason) = opts.reason.as_deref() else {
return Err(usage(
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(usage(
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| usage(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, Fault> {
let at = match (&opts.list, &opts.at) {
(Some(Listing::List(list)), _) => return halts_verb(&list.at, opts.json),
(None, Some(at)) => at,
(None, None) => return Err(usage(missing_store())),
};
let scope = agentplane::quota::HaltScope::parse(&opts.scope).ok_or_else(|| {
usage(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_else(|| {
usage(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_else(|| {
usage(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| usage(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 = 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())?;
if !opts.json {
println!(
"halted {}: {reason} (by {}, {})",
scope.key(),
by.actor(),
by.basis().as_str()
);
return Ok(ExitCode::SUCCESS);
}
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())?;
if opts.json {
println!(
"{}",
serde_json::json!({
"scope": scope.key(),
"halted": false,
"was_standing": was_standing,
})
);
} else if was_standing {
println!("lifted {}", scope.key());
} else {
println!("no halt was standing on {}; nothing lifted", scope.key());
}
return Ok(lift_status(was_standing));
}
};
println!("{printed}");
Ok(ExitCode::SUCCESS)
})
}
fn lift_status(was_standing: bool) -> ExitCode {
ExitCode::from(if was_standing {
exit::OK
} else {
exit::FINDING
})
}
fn halts_verb(at: &StoreRef, json: bool) -> Result<ExitCode, Fault> {
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 = 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();
if json {
println!(
"{}",
serde_json::to_string_pretty(&serde_json::json!({ "halts": rows }))
.map_err(|e| e.to_string())?
);
} else if halts.is_empty() {
println!("no halts standing");
} else {
for h in &halts {
println!(
"{}\tthrown {} by {} ({}): {}",
h.scope.key(),
h.at.unix_timestamp(),
h.by.actor(),
h.by.basis().as_str(),
h.reason
);
}
}
Ok(ExitCode::SUCCESS)
})
}
fn waiting_verb(opts: &WaitingArgs) -> Result<ExitCode, Fault> {
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();
let waiting = store
.waiting_runs(opts.limit)
.await
.map_err(|e| e.to_string())?;
let rows: Vec<serde_json::Value> = waiting
.iter()
.map(|w| {
serde_json::json!({
"run": w.run.to_string(),
"waiting_for": w.reason,
"until": w.reason.until().to_string(),
})
})
.collect();
println!(
"{}",
serde_json::to_string_pretty(&serde_json::json!({ "waiting": rows }))
.map_err(|e| e.to_string())?
);
Ok(ExitCode::SUCCESS)
})
}
fn blocking<T, E: From<String>>(
f: impl std::future::Future<Output = Result<T, E>>,
) -> Result<T, E> {
tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.map_err(|e| E::from(format!("could not start the async runtime: {e}")))?
.block_on(f)
}
fn reconcile_verb(opts: &ReconcileArgs) -> Result<ExitCode, Fault> {
let assertion = match opts.outcome.as_str() {
"landed" => {
let raw = opts.output.as_deref().unwrap_or("null");
agentplane::core::Assertion::Landed(
serde_json::from_str(raw)
.map_err(|e| usage(format!("--output is not JSON: {e}")))?,
)
}
"did-not-happen" => {
if opts.output.is_some() {
return Err(usage(
"--output belongs to `--outcome landed`: an effect that did not happen \
produced no result to read back"
.to_owned(),
));
}
agentplane::core::Assertion::DidNotHappen
}
other => {
return Err(usage(format!(
"'{other}' is not an outcome: use `landed` or `did-not-happen`"
)));
}
};
let by = opts.who.operator()?;
let run = agentplane::core::RunId::parse(&opts.run_id).map_err(|e| usage(e.to_string()))?;
let effect =
agentplane::core::EffectKey::from_hex(&opts.effect).map_err(|e| usage(e.to_string()))?;
blocking(async move {
let backend = opts.at.open().await?;
let plane = Runtime::builder_with(backend.stores())
.tenant(backend.tenant())
.build();
plane
.reconcile_effect(run, effect, assertion, &by, &opts.note)
.await
.map_err(|e| e.to_string())?;
println!(
"{}",
serde_json::json!({
"run": run.to_string(),
"effect": opts.effect,
"outcome": opts.outcome,
"by": by.actor(),
"basis": by.basis().as_str(),
})
);
Ok(ExitCode::SUCCESS)
})
}
fn quarantine_verb(opts: &QuarantineArgs) -> Result<ExitCode, Fault> {
use agentplane::core::QuarantineDecision;
let decision = match opts.decision.as_str() {
"reopen" => QuarantineDecision::Reopen,
"abandon" => QuarantineDecision::Abandon,
other => {
return Err(usage(format!(
"'{other}' is not a decision: use `reopen` or `abandon`"
)));
}
};
let by = opts.who.operator()?;
let run = agentplane::core::RunId::parse(&opts.run_id).map_err(|e| usage(e.to_string()))?;
blocking(async move {
let backend = opts.at.open().await?;
let plane = Runtime::builder_with(backend.stores())
.tenant(backend.tenant())
.build();
plane
.record_quarantine_decision(run, &by, &opts.reason, decision)
.await
.map_err(|e| e.to_string())?;
println!(
"{}",
serde_json::json!({
"run": run.to_string(),
"decision": opts.decision,
"by": by.actor(),
"basis": by.basis().as_str(),
"applied": false,
"next": "recorded; the next resume of this run applies it",
})
);
Ok(ExitCode::SUCCESS)
})
}
fn task_json(task: &agentplane::core::Task, whole: bool) -> serde_json::Value {
let j = &task.justification;
let shown = task.rendering();
let mut out = serde_json::json!({
"task": task.id.to_string(),
"run": task.run.to_string(),
"case": task.case.map(|c| c.to_string()),
"kind": task.kind,
"state": task.state.as_str(),
"priority": task.priority.as_str(),
"summary": shown.summary,
"proposed_action": shown.proposed_action,
"withheld": shown.withheld.map(|w| format!("{} — {w}", w.as_str())),
"candidate_roles": task.candidate_roles,
"assignee": task.assignee,
"due_at": task.due_at.map(|d| d.to_string()),
"has_untrusted_prose": j.has_untrusted_prose(),
"escaped": shown.escaped,
"mixed_script": shown.mixed_script,
});
if whole {
out["confidence"] = serde_json::json!(j.confidence);
out["cost"] = serde_json::json!(shown.cost);
out["evidence"] = serde_json::json!(shown.evidence);
out["excluded_actors"] = serde_json::json!(task.excluded_actors);
out["on_expiry"] = serde_json::json!(task.on_expiry.as_str());
out["escalate_to"] = serde_json::json!(task.escalate_to);
}
out
}
fn tasks_verb(opts: &TasksArgs) -> Result<ExitCode, Fault> {
blocking(async move {
let tasks = opts.at.open().await?.tasks();
if let Some(show) = &opts.show {
let id = agentplane::core::TaskId::parse(show)
.map_err(|e| usage(format!("`{show}` is not a task id: {e}")))?;
let task = tasks
.task(id)
.await
.map_err(|e| e.to_string())?
.ok_or_else(|| format!("no task {id} on this plane"))?;
println!(
"{}",
serde_json::to_string_pretty(&task_json(&task, true)).map_err(|e| e.to_string())?
);
return Ok(ExitCode::SUCCESS);
}
let mut queued = tasks
.queue(&opts.roles, opts.limit.saturating_add(1))
.await
.map_err(|e| e.to_string())?;
let truncated = queued.len() > opts.limit;
queued.truncate(opts.limit);
println!(
"{}",
serde_json::to_string_pretty(&serde_json::json!({
"tasks": queued.iter().map(|t| task_json(t, false)).collect::<Vec<_>>(),
"truncated": truncated,
}))
.map_err(|e| e.to_string())?
);
Ok(ExitCode::SUCCESS)
})
}
fn decide_verb(opts: &DecideArgs) -> Result<ExitCode, Fault> {
let by = opts.who.operator()?;
let id = agentplane::core::TaskId::parse(&opts.task_id)
.map_err(|e| usage(format!("`{}` is not a task id: {e}", opts.task_id)))?;
let approved = opts.verdict == Verdict::Approve;
let decision = if approved {
agentplane::core::Decision::approve(by.clone(), opts.reason.clone())
} else {
agentplane::core::Decision::reject(by.clone(), opts.reason.clone())
};
blocking(async move {
let backend = opts.at.open().await?;
let plane = Runtime::builder_with(backend.stores())
.tenant(backend.tenant())
.lease_ttl(std::time::Duration::from_secs(2))
.build();
let delivery = match plane.decide_task(id, &decision, &opts.roles).await {
Ok(delivery) => delivery,
Err(agentplane::core::RuntimeError::ProposalWithheld { reason, .. }) => {
eprintln!("{}", withheld_refusal(id, reason));
return Ok(ExitCode::from(exit::FINDING));
}
Err(e) => return Err(e.to_string().into()),
};
println!(
"{}",
serde_json::json!({
"task": id.to_string(),
"approved": approved,
"by": by.actor(),
"basis": by.basis().as_str(),
"delivery": format!("{delivery:?}"),
})
);
if delivery.resumed_run().is_none()
&& let Some(task) = backend.tasks().task(id).await.map_err(|e| e.to_string())?
{
eprintln!(
"recorded. Run {} reads it when it resumes:\n agentplane replay {} \
--manifest <file>{}",
task.run,
task.run,
where_flags(Some(&opts.at.store), opts.at.tenant.as_deref())
);
}
Ok(ExitCode::SUCCESS)
})
}
fn withheld_refusal(id: agentplane::core::TaskId, reason: agentplane::core::Withheld) -> String {
use agentplane::core::Withheld;
let why = match reason {
Withheld::Sealed => {
"its proposal is sealed at rest and this terminal holds no key ring, so it \
cannot show you what you would be approving. Approve it on the plane that \
holds the key ring, or reject it here — a rejection needs no proposal"
}
Withheld::Erased => {
"its proposal was erased — the key it was sealed under is destroyed — so \
nobody can be shown what an approval would approve. Reject it"
}
Withheld::Undecodable => {
"its sealed proposal opened and does not decode, so what it holds cannot be \
shown. Reject it, or restore the row from a backup and decide it then"
}
};
format!("task {id} was not approved: {why}. Nothing was recorded; the task is still open.")
}
fn acknowledge_verb(opts: &AcknowledgeArgs) -> Result<ExitCode, Fault> {
let by = opts.who.operator()?;
let case = agentplane::core::CaseId::parse(&opts.case_id).map_err(|e| usage(e.to_string()))?;
blocking(async move {
let backend = opts.at.open().await?;
#[allow(clippy::disallowed_methods)]
let at = agentplane::core::Timestamp::now_utc();
let note = agentplane::core::BreachNote {
by: by.clone(),
note: opts.note.clone(),
at,
};
let recorded = backend
.cases()
.acknowledge_breach(case, &opts.obligation, ¬e)
.await
.map_err(|e| e.to_string())?;
println!(
"{}",
serde_json::json!({
"case": opts.case_id,
"obligation": opts.obligation,
"recorded": recorded,
"by": by.actor(),
})
);
Ok(ExitCode::SUCCESS)
})
}
fn cancel_verb(opts: &CancelArgs) -> Result<ExitCode, Fault> {
let by = opts.who.operator()?;
let run = agentplane::core::RunId::parse(&opts.run_id).map_err(|e| usage(e.to_string()))?;
blocking(async move {
let backend = opts.at.open().await?;
let plane = Runtime::builder_with(backend.stores())
.tenant(backend.tenant())
.build();
let first = plane
.request_cancel(run, &by, &opts.reason)
.await
.map_err(|e| e.to_string())?;
println!(
"{}",
serde_json::json!({
"run": run.to_string(),
"by": by.actor(),
"basis": by.basis().as_str(),
"requested": first,
})
);
Ok(ExitCode::SUCCESS)
})
}
#[cfg(feature = "push")]
fn rearm_verb(opts: &RearmArgs) -> Result<ExitCode, Fault> {
let run = agentplane::core::RunId::parse(&opts.run_id).map_err(|e| usage(e.to_string()))?;
blocking(async move {
let backend = opts.at.open().await?;
let plane = Runtime::builder_with(backend.stores())
.tenant(backend.tenant())
.push(backend.push())
.build();
let rearmed = plane
.rearm_push(run, &opts.id)
.await
.map_err(|e| e.to_string())?;
println!(
"{}",
serde_json::json!({
"run": run.to_string(),
"id": opts.id,
"rearmed": rearmed,
})
);
Ok(lift_status(rearmed))
})
}
fn attention_verb(opts: &WaitingArgs) -> Result<ExitCode, Fault> {
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 plane = Runtime::builder_with(backend.stores())
.tenant(backend.tenant())
.quota(backend.quotas(), agentplane::quota::TenantQuota::default())
.build();
#[allow(clippy::disallowed_methods)]
let now = agentplane::core::Timestamp::now_utc();
let found = plane
.attention(now, opts.limit)
.await
.map_err(|e| e.to_string())?;
println!(
"{}",
serde_json::to_string_pretty(&serde_json::json!({
"needs_attention": found.any(),
"conditions": found
.conditions
.iter()
.map(|c| serde_json::json!({
"condition": c.kind,
"found": c.found,
"at_least": c.at_least,
"subjects": c.subjects,
"unlisted": c.unlisted,
"remedy": c.remedy.cli,
}))
.collect::<Vec<_>>(),
"not_checked": found.not_checked,
}))
.map_err(|e| e.to_string())?
);
Ok(if found.any() {
ExitCode::from(exit::FINDING)
} else {
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, Fault> {
let tenant = tenant
.map(|name| {
agentplane::core::TenantId::new(name).map_err(|e| usage(format!("--tenant: {e}")))
})
.transpose()?;
if is_connection_string(spec) {
#[cfg(not(feature = "postgres"))]
return Err(usage(
"--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",
));
#[cfg(feature = "postgres")]
return Self::shared(spec, tenant).await;
}
let store =
RedbStore::open(spec).map_err(|e| Fault::from(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, Fault> {
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 tasks(&self) -> Arc<dyn agentplane::case::TaskStore> {
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 digest_verb(a: &DigestArgs) -> Result<ExitCode, Fault> {
let manifests = manifests_at(&a.manifest).map_err(usage)?;
if a.out.json {
let rows = manifests
.iter()
.map(|m| {
Ok(serde_json::json!({
"name": m.metadata.name,
"version": m.metadata.version,
"digest": m.digest().map_err(|e| e.to_string())?.to_hex(),
}))
})
.collect::<Result<Vec<_>, String>>()?;
println!("{}", serde_json::json!({ "digests": rows }));
} else 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)
}
fn dispatch(cli: Cli) -> Result<ExitCode, Fault> {
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) => digest_verb(&a),
Verb::Audit(a) => journal_verb(&a.store, Some(&a), false),
Verb::Export(a) => journal_verb(&a.store, None, a.allow_partial),
Verb::Drill(a) => drill_verb(&a),
Verb::ForgetAdmissions(a) => forget_admissions_verb(&a),
Verb::Retention(RetentionArgs {
act: RetentionAct::Plan(a),
}) => retention_plan_verb(&a),
Verb::Halt(a) => halt_verb(&a),
Verb::Hold(a) => hold_verb(&a),
Verb::Reconcile(a) => reconcile_verb(&a),
Verb::Quarantine(a) => quarantine_verb(&a),
Verb::Tasks(a) => tasks_verb(&a),
Verb::Decide(a) => decide_verb(&a),
Verb::Init(a) => init_verb(&a),
Verb::Acknowledge(a) => acknowledge_verb(&a),
#[cfg(feature = "push")]
Verb::Rearm(a) => rearm_verb(&a),
Verb::Cancel(a) => cancel_verb(&a),
Verb::Waiting(a) => waiting_verb(&a),
Verb::Attention(a) => attention_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::from(exit::FINDING)
})
})
}
Verb::Verify(a) => verify_verb(&a),
Verb::Policy(PolicyArgs {
act: PolicyAct::Check(a),
}) => policy_check_verb(&a),
Verb::Run(a) => {
let manifests = manifests_at(&a.manifest).map_err(usage)?;
execute(&manifests, &a)
}
Verb::Replay(a) => {
let manifests = manifests_at(&a.manifest).map_err(usage)?;
replay(&manifests, &a)
}
Verb::Card(a) => card(&a),
Verb::Serve(a) => {
let manifests = manifests_at(&a.manifest).map_err(usage)?;
serve(&manifests, &a)
}
}
}
const STARTER: &str = r#"# yaml-language-server: $schema=https://hupe1980.github.io/agentplane/agent.schema.json
apiVersion: agentplane.hupe1980.github.io/v1alpha1
kind: Agent
metadata: { name: NAME, version: "0.1.0" }
spec:
execution: { kind: completion }
identity:
role: "Summarise the input in one sentence"
constraints: "No speculation."
capabilities: { provides: [NAME] }
models:
# No key yet? `provider: fake` answers with a value of the output schema.
privileged: { provider: anthropic, model: claude-sonnet-5 }
output:
schema:
type: object
additionalProperties: false
required: [summary]
properties: { summary: { type: string } }
budgets: { max_tokens: 20000, max_steps: 4 }
"#;
const STARTER_TOOLS: &str = r#"# yaml-language-server: $schema=https://hupe1980.github.io/agentplane/agent.schema.json
#
# agentplane run FILE --input '{"ticket": "T-1"}' \
# --mcp "tickets=python3 examples/mcp-server.py"
apiVersion: agentplane.hupe1980.github.io/v1alpha1
kind: Agent
metadata: { name: NAME, version: "0.1.0" }
spec:
execution: { kind: tool-calling, max_turns: 4 }
identity:
role: "Answer a support question using the ticket tool"
constraints: "One sentence. Cite the ticket id."
capabilities: { provides: [NAME] }
models:
# No key yet? `provider: fake` answers with a value of the output schema.
privileged: { provider: anthropic, model: claude-sonnet-5 }
tools:
- ref: "tool://tickets/read"
mutates: false
description: "Read a ticket by id"
arguments:
type: object
additionalProperties: false
required: [id]
properties: { id: { type: string } }
budgets: { max_tokens: 20000, max_steps: 8 }
"#;
fn init_verb(opts: &InitArgs) -> Result<ExitCode, Fault> {
let template = if opts.tools { STARTER_TOOLS } else { STARTER };
let text = template
.replace("NAME", &opts.name)
.replace("FILE", &opts.path);
let parsed =
Manifest::parse_all(&text).map_err(|e| usage(format!("--name {}: {e}", opts.name)))?;
let mut file = std::fs::OpenOptions::new()
.write(true)
.create_new(true)
.open(&opts.path)
.map_err(|e| format!("writing {}: {e}", opts.path))?;
std::io::Write::write_all(&mut file, text.as_bytes())
.map_err(|e| format!("writing {}: {e}", opts.path))?;
let digest = parsed
.first()
.map(|m| m.digest().map(agentplane::Digest::to_hex))
.transpose()
.map_err(|e| e.to_string())?
.unwrap_or_default();
if opts.out.json {
println!(
"{}",
serde_json::json!({ "wrote": opts.path, "digest": digest })
);
} else {
println!("wrote {} ({digest})", opts.path);
}
eprintln!(
"next:\n agentplane validate {path}\n agentplane run {path} --input '{{}}'{mcp}",
path = opts.path,
mcp = if opts.tools {
" --mcp \"tickets=<command>\""
} else {
""
}
);
Ok(ExitCode::SUCCESS)
}
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}"
)
})
}
#[cfg(feature = "cedar")]
const BUNDLE_FILES: [&str; 3] = ["policy.cedar", "schema.json", "entities.json"];
#[cfg(feature = "cedar")]
fn load_policy_bundle(path: &str) -> Result<agentplane::policy::CedarEngine, Fault> {
let at = std::path::Path::new(path);
let meta =
std::fs::metadata(at).map_err(|e| format!("reading the policy bundle {path}: {e}"))?;
let read = |file: &std::path::Path| {
std::fs::read_to_string(file).map_err(|e| format!("reading {}: {e}", file.display()))
};
let (rules, schema, entities) = if meta.is_dir() {
let entries = std::fs::read_dir(at).map_err(|e| format!("reading {path}: {e}"))?;
for entry in entries {
let name = entry
.map_err(|e| format!("reading {path}: {e}"))?
.file_name()
.to_string_lossy()
.into_owned();
if !name.starts_with('.') && !BUNDLE_FILES.contains(&name.as_str()) {
return Err(usage(format!(
"the policy bundle {path} holds `{name}`, which no bundle reads: a bundle \
directory holds {} and nothing else, so a rule cannot sit in a file \
the bundle's identity does not cover",
BUNDLE_FILES.join(", ")
)));
}
}
let rules = at.join(BUNDLE_FILES[0]);
if !rules.is_file() {
return Err(usage(format!(
"the policy bundle {path} has no {}",
BUNDLE_FILES[0]
)));
}
let optional = |name: &str| {
let file = at.join(name);
if file.is_file() {
read(&file).map(Some)
} else {
Ok(None)
}
};
(
read(&rules)?,
optional(BUNDLE_FILES[1])?,
optional(BUNDLE_FILES[2])?,
)
} else {
(read(at)?, None, None)
};
agentplane::policy::CedarEngine::from_bundle(&rules, schema.as_deref(), entities.as_deref())
.map_err(|e| usage(format!("the policy bundle {path} was refused: {e}")))
}
#[cfg(feature = "cedar")]
fn policy_check_verb(opts: &PolicyCheckArgs) -> Result<ExitCode, Fault> {
use agentplane::policy::check::{Check, CheckError, TenantSource, Verdict};
let bundle = load_policy_bundle(&opts.bundle)?;
let candidate = opts
.candidate
.as_deref()
.map(load_policy_bundle)
.transpose()?;
let (tenant, source) = match &opts.tenant {
Some(name) => (
agentplane::core::TenantId::new(name).map_err(|e| usage(format!("--tenant: {e}")))?,
TenantSource::Supplied,
),
None => (agentplane::core::TenantId::default(), TenantSource::Default),
};
let check = Check::new(tenant.as_str(), source);
let candidate = candidate
.as_ref()
.map(|c| c as &dyn agentplane::core::PolicyEngine);
let rt = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.map_err(|e| format!("could not start the async runtime: {e}"))?;
let report = if opts.from == "-" {
rt.block_on(check.run(std::io::stdin().lock(), &bundle, candidate))
} else {
let file =
std::fs::File::open(&opts.from).map_err(|e| format!("reading {}: {e}", opts.from))?;
rt.block_on(check.run(std::io::BufReader::new(file), &bundle, candidate))
}
.map_err(|e| match e {
CheckError::NotAnExport(_) => usage(format!("{}: {e}", opts.from)),
CheckError::Io(_) | CheckError::Keys(_) => Fault::Operational(e.to_string()),
})?;
if opts.out.json {
println!(
"{}",
serde_json::to_string_pretty(&report).map_err(|e| e.to_string())?
);
} else {
print!("{}", policy_report_text(&report));
}
Ok(match report.verdict() {
Verdict::Clean => ExitCode::SUCCESS,
Verdict::Findings => ExitCode::from(exit::FINDING),
Verdict::Partial => ExitCode::from(exit::PARTIAL),
})
}
#[cfg(feature = "cedar")]
fn policy_report_text(report: &agentplane::policy::check::Report) -> String {
use agentplane::policy::check::{Finding, Mode};
use std::fmt::Write as _;
let line = |f: &Finding| {
let mut at = String::new();
if let Some(step) = f.step {
let _ = write!(at, "step {step} ");
}
if let Some(key) = f.effect_key {
let _ = write!(at, "effect {} ", key.to_hex());
}
format!("{at}{} on {}: {}", f.action, f.resource, f.reason)
};
let mut out = String::new();
let _ = write!(out, "bundle {}", report.bundle.to_hex());
if let Some(candidate) = report.candidate {
let _ = write!(out, ", candidate {}", candidate.to_hex());
}
let _ = writeln!(
out,
"; tenant {} ({})",
report.tenant.value,
match report.tenant.source {
agentplane::policy::check::TenantSource::Supplied => "supplied",
agentplane::policy::check::TenantSource::Default => "assumed: none was supplied",
}
);
for run in &report.runs {
match run.mode {
Mode::Recorded => {
let _ = writeln!(
out,
"run {}: {} evaluated, {} finding(s)",
run.run,
run.evaluated,
run.findings.len()
);
}
Mode::Mismatch => {
let _ = writeln!(
out,
"run {}: recorded bundle {} is not the one supplied — not evaluated",
run.run,
run.recorded_bundle
.map(agentplane::core::Digest::to_hex)
.unwrap_or_default()
);
}
Mode::Ungoverned => {
let _ = writeln!(
out,
"run {}: ungoverned — no bundle is on its record, so no gate ran",
run.run
);
}
}
for f in &run.findings {
let _ = writeln!(out, " finding: {}", line(f));
}
if let Some(diff) = &run.diff {
for f in &diff.newly_denied {
let _ = writeln!(out, " newly denied: {}", line(f));
}
for f in &diff.malformed_under_candidate {
let _ = writeln!(out, " malformed under candidate: {}", line(f));
}
}
let mut reasons = std::collections::BTreeMap::new();
for n in &run.not_evaluable {
*reasons.entry(n.reason).or_insert(0usize) += 1;
}
for (reason, count) in reasons {
let _ = writeln!(
out,
" not evaluable: {count} {}",
serde_json::to_value(reason)
.ok()
.and_then(|v| v.as_str().map(ToOwned::to_owned))
.unwrap_or_default()
);
}
}
for run in &report.unreadable {
let _ = writeln!(out, "run {run}: the export could not read it");
}
let _ = writeln!(out, "not in any export: {}", report.outside_export);
out
}
fn verify_verb(opts: &VerifyArgs) -> Result<ExitCode, Fault> {
let verifier = verifier_from(&opts.key).map_err(usage)?;
let saved = match &opts.checkpoint {
Some(path) => Some(read_checkpoint(path).map_err(usage)?),
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,
))
.map_err(usage)?
}
(None, false) => {
return Err(usage(
"--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(usage(
"--witness-key was given with no --witness to use it against".to_owned(),
));
}
Anchor::default()
}
};
let mut anchor = fetched;
if let Some(saved) = saved {
anchor.checkpoints.push(agentplane::audit::Anchor::new(
saved,
match &opts.checkpoint {
Some(path) => format!("file {path}"),
None => "file".to_owned(),
},
));
}
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, &anchor.checkpoints)
.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, &anchor.checkpoints)
.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::from(exit::FINDING)
})
}
fn card(opts: &CardArgs) -> Result<ExitCode, Fault> {
let manifests = manifests_at(&opts.manifest).map_err(usage)?;
let [manifest] = manifests.as_slice() else {
return Err(usage(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 {
let cli = <Cli as clap::Parser>::parse();
let (metrics, format) = match &cli.verb {
Verb::Serve(serve) => (true, serve.log_format),
_ => (false, LogFormat::Text),
};
let verifying = matches!(&cli.verb, Verb::Replay(replay) if replay.strict);
install_tracing(metrics, verifying, format);
match dispatch(cli) {
Ok(code) => code,
Err(fault) => {
eprintln!("agentplane: {fault}");
ExitCode::from(fault.status())
}
}
}
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, Fault> {
use agentplane::api::a2a::A2aServer;
use agentplane::api::tokens::TokenAuthenticator;
let [manifest] = manifests else {
return Err(usage(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_else(|| {
usage(
"`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_else(|| {
usage(
"`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_else(|| {
usage(
"`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 = load_policy_bundle(policy_path)?;
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_else(|| {
usage(
"`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)).map_err(usage)?
{
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;
}
#[allow(clippy::disallowed_methods)]
let at = agentplane::core::Timestamp::now_utc();
match plane.drill(at).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);
}
refuse_ambiguous_peers(specs)?;
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 grants: Vec<(String, bool)> = manifests
.iter()
.flat_map(|m| &m.spec.tools)
.filter_map(|g| {
agentplane::tools::ToolId::parse(&g.reference).map(|id| (id, g.mutates))
})
.filter(|(id, _)| id.server == name)
.map(|(id, mutates)| (id.tool, mutates))
.collect();
let granted: Vec<String> = grants.iter().map(|(tool, _)| tool.clone()).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()));
if grants.iter().all(|(_, mutates)| !mutates) {
grant = grant.read_only();
}
let token_var = peer_token_var(name);
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 local = cfg!(feature = "testkit")
&& ["http://localhost:", "http://127.0.0.1:", "http://[::1]:"]
.iter()
.any(|prefix| url.starts_with(prefix));
if !url.starts_with("https://") && !local {
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}"))?;
#[cfg(feature = "testkit")]
let client = if local {
client.allow_loopback()
} else {
client
};
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_attr(not(feature = "a2a"), allow(dead_code))]
fn peer_token_var(name: &str) -> String {
format!(
"AGENTPLANE_PEER_TOKEN_{}",
name.to_ascii_uppercase().replace(['.', '-'], "_")
)
}
#[cfg_attr(not(feature = "a2a"), allow(dead_code))]
fn refuse_ambiguous_peers(specs: &[String]) -> Result<(), String> {
let mut seen: std::collections::HashMap<String, &str> = std::collections::HashMap::new();
for spec in specs {
let name = spec.split_once('=').map_or(spec.as_str(), |(name, _)| name);
let var = peer_token_var(name);
if let Some(earlier) = seen.insert(var.clone(), name) {
return Err(format!(
"--peer `{earlier}` and --peer `{name}` both read their token from {var}, \
so one would present the other's credential — rename one of them"
));
}
}
Ok(())
}
fn validate(a: &ValidateArgs) -> Result<ExitCode, Fault> {
let text = std::fs::read_to_string(&a.manifest)
.map_err(|e| usage(format!("reading {}: {e}", a.manifest)))?;
let manifests = match Manifest::parse_all(&text) {
Ok(manifests) => manifests,
Err(e) => {
if a.json {
println!(
"{}",
serde_json::json!({ "valid": false, "error": e.to_string() })
);
} else {
eprintln!("invalid: {e}");
}
return Ok(ExitCode::from(exit::FINDING));
}
};
let mut documents = Vec::new();
let mut missing = Vec::new();
for m in &manifests {
let absent: Vec<&String> = a
.require_annotation
.iter()
.filter(|key| !m.metadata.annotations.contains_key(*key))
.collect();
let bound = declared_bound(m);
if !a.json {
if absent.is_empty() {
println!("ok: {} {}", m.metadata.name, m.metadata.version);
}
for key in &absent {
println!("MISSING: {} — annotation '{key}'", m.metadata.name);
}
for line in &bound.lines {
println!(" {line}");
}
}
missing.extend(
absent
.iter()
.map(|key| format!("{}: {key}", m.metadata.name)),
);
documents.push(serde_json::json!({
"name": m.metadata.name,
"version": m.metadata.version,
"missing_annotations": absent,
"bound": bound.json,
}));
}
if a.json {
println!(
"{}",
serde_json::json!({ "valid": missing.is_empty(), "documents": documents })
);
} else if !missing.is_empty() {
eprintln!(
"{} required annotation(s) absent: {}",
missing.len(),
missing.join(", ")
);
}
Ok(ExitCode::from(if missing.is_empty() {
exit::OK
} else {
exit::FINDING
}))
}
struct DeclaredBound {
lines: Vec<String>,
json: serde_json::Value,
}
#[allow(clippy::too_many_lines)]
fn declared_bound(m: &Manifest) -> DeclaredBound {
let budget = m.budget();
let width = budget.max_parallel_steps.unwrap_or(1) as u64;
let mut lines = Vec::new();
let roles: Vec<(&str, &agentplane::manifest::ModelRef)> = m
.spec
.models
.as_ref()
.map(|models| {
[
("privileged", models.privileged.as_ref()),
("quarantined", models.quarantined.as_ref()),
]
.into_iter()
.filter_map(|(name, r)| r.map(|r| (name, r)))
.collect()
})
.unwrap_or_default();
let mut role_json = Vec::new();
for (name, r) in &roles {
let line = match r.call_bound() {
Some(call) => format!(
"one call, {name} {}/{}: {} tokens, {} minor units",
r.provider, r.model, call.tokens, call.minor_units
),
None => format!(
"one call, {name} {}/{}: unbounded — spec.models.{name}.max_input_tokens is \
not declared",
r.provider, r.model
),
};
lines.push(line);
role_json.push(serde_json::json!({
"role": name,
"call": r.call_bound().map(|c| serde_json::json!({
"tokens": c.tokens,
"minor_units": c.minor_units,
})),
}));
}
let call = m.call_bound();
let call_gaps: Vec<String> = if roles.is_empty() {
vec!["spec.models (no model role states what one call can cost)".to_owned()]
} else {
roles
.iter()
.filter(|(_, r)| r.call_bound().is_none())
.map(|(name, _)| format!("spec.models.{name}.max_input_tokens"))
.collect()
};
let mut units = serde_json::Map::new();
for (unit, ceiling, field, per_call) in [
(
"tokens",
budget.max_tokens,
"spec.budgets.max_tokens",
call.map(|c| c.tokens),
),
(
"minor units",
budget.max_minor_units,
"spec.budgets.max_minor_units",
call.map(|c| c.minor_units),
),
] {
let mut unbounded: Vec<String> = Vec::new();
if ceiling.is_none() {
unbounded.push(field.to_owned());
}
if per_call.is_none() {
unbounded.extend(call_gaps.iter().cloned());
}
let total = match (ceiling, per_call) {
(Some(c), Some(p)) => Some(c.saturating_add(width.saturating_mul(p))),
_ => None,
};
let shown = |v: Option<u64>| v.map_or_else(|| "?".to_owned(), |v| v.to_string());
lines.push(match total {
Some(total) => format!(
"worst case, {unit}: {total} = ceiling {} + width {width} × one call {}",
shown(ceiling),
shown(per_call)
),
None => format!(
"worst case, {unit}: no total — unbounded: {}",
unbounded.join(", ")
),
});
units.insert(
unit.replace(' ', "_"),
serde_json::json!({
"ceiling": ceiling,
"per_call": per_call,
"width": width,
"total": total,
"unbounded": unbounded,
}),
);
}
let turns = m
.spec
.execution
.as_ref()
.filter(|e| e.kind == agentplane::manifest::ExecutionKind::ToolCalling)
.map(|e| e.max_turns);
if let Some(turns) = turns {
lines.push(match call {
Some(c) => format!(
"model calls: at most {turns} turns × one call = {} tokens, {} minor units",
u64::from(turns).saturating_mul(c.tokens),
u64::from(turns).saturating_mul(c.minor_units)
),
None => format!(
"model calls: at most {turns} turns, no total — unbounded: {}",
call_gaps.join(", ")
),
});
}
DeclaredBound {
lines,
json: serde_json::json!({
"roles": role_json,
"units": units,
"turns": turns,
}),
}
}
#[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, Fault> {
Err(usage(
"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",
))
}
#[cfg(not(feature = "cedar"))]
#[allow(clippy::unnecessary_wraps)]
fn policy_check_verb(_opts: &PolicyCheckArgs) -> Result<ExitCode, Fault> {
Err(usage(
"this build cannot check policy: `policy check` evaluates a Cedar bundle and needs \
the `cedar` feature. Reinstall with `--features cli,cedar`, or use the `:full` \
container image",
))
}
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(())
}
struct Resume<'a> {
manifest: &'a str,
store: Option<&'a str>,
tenant: Option<&'a str>,
}
fn shell_quote(word: &str) -> String {
let safe = !word.is_empty()
&& word
.chars()
.all(|c| c.is_ascii_alphanumeric() || "_-./:@%+=,".contains(c));
if safe {
word.to_owned()
} else {
format!("'{}'", word.replace('\'', r"'\''"))
}
}
fn without_password(spec: &str) -> String {
let Some((scheme, rest)) = spec.split_once("://") else {
return spec.to_owned();
};
let authority_end = rest.find(['/', '?']).unwrap_or(rest.len());
let (authority, tail) = rest.split_at(authority_end);
let authority = match authority.rsplit_once('@') {
Some((userinfo, host)) => {
let user = userinfo.split_once(':').map_or(userinfo, |(user, _)| user);
format!("{user}@{host}")
}
None => authority.to_owned(),
};
let tail = match tail.split_once('?') {
Some((path, query)) => {
let kept: Vec<&str> = query
.split('&')
.filter(|pair| !pair.to_ascii_lowercase().starts_with("password="))
.collect();
if kept.is_empty() {
path.to_owned()
} else {
format!("{path}?{}", kept.join("&"))
}
}
None => tail.to_owned(),
};
format!("{scheme}://{authority}{tail}")
}
fn where_flags(store: Option<&str>, tenant: Option<&str>) -> String {
let from_env = std::env::var("AGENTPLANE_STORE").ok();
let store = match store {
None => " --store <file>".to_owned(),
Some(s) if from_env.as_deref() == Some(s) => r#" --store "$AGENTPLANE_STORE""#.to_owned(),
Some(s) if is_connection_string(s) => {
format!(" --store {}", shell_quote(&without_password(s)))
}
Some(s) => format!(" --store {}", shell_quote(s)),
};
match tenant {
Some(t) => format!("{store} --tenant {}", shell_quote(t)),
None => store,
}
}
fn conclude(outcome: &agentplane::runtime::RunOutcome, resume: &Resume<'_>) -> ExitCode {
if let RunStatus::Suspended(reason) = &outcome.status {
suspended(outcome.run_id, reason, resume);
return ExitCode::from(exit::SUSPENDED);
}
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::from(exit::FINDING)
}
}
fn suspended(
run: agentplane::core::RunId,
reason: &agentplane::core::SuspendReason,
at: &Resume<'_>,
) {
use agentplane::core::SuspendReason;
let store = where_flags(at.store, at.tenant);
let replay = format!(
"agentplane replay {run} --manifest {}{store}",
shell_quote(at.manifest)
);
match reason {
SuspendReason::AwaitingEvent {
kind, correlation, ..
} => {
let task = correlation
.iter()
.find(|k| k.namespace == "task")
.and_then(|k| agentplane::core::TaskId::parse(&k.value).ok());
if let Some(task) = task {
eprintln!("run {run} is waiting for a person to decide {task}");
eprintln!("next:");
eprintln!(" agentplane tasks --show {task}{store}");
eprintln!(" agentplane decide {task} approve --reason '…' --actor <you>{store}");
eprintln!(" {replay}");
} else {
eprintln!("run {run} is waiting for a `{kind}` event ({reason})");
eprintln!("next: deliver it, then\n {replay}");
}
}
SuspendReason::AwaitingTime { until } => {
eprintln!("run {run} is waiting until {until}");
eprintln!("next, once that instant has passed:\n {replay}");
}
other => {
eprintln!("run {run} is waiting: {other}");
eprintln!("next:\n {replay}");
}
}
if at.store.is_none() {
eprintln!(
"note: this run journaled to memory and ended with the process — run it \
again with --store to keep it"
);
}
}
fn correlation(flags: &[String]) -> Result<Vec<agentplane::core::CorrelationKey>, String> {
if flags.is_empty() {
return Ok(vec![agentplane::core::CorrelationKey::new(
"invocation",
agentplane::core::RunId::generate().to_string(),
)]);
}
flags
.iter()
.map(|flag| {
let (namespace, value) = flag
.split_once('=')
.filter(|(n, v)| !n.trim().is_empty() && !v.trim().is_empty())
.ok_or_else(|| {
format!(
"--correlate wants `<namespace>=<value>`, got `{flag}` — a \
`$correlation/customer` subject needs `--correlate customer=C-7`"
)
})?;
Ok(agentplane::core::CorrelationKey::new(namespace, value))
})
.collect()
}
fn chain_for(
subject: &str,
manifests: &[Manifest],
peers: &[String],
) -> Result<agentplane::core::Delegation, String> {
let peer_names: Vec<&str> = peers
.iter()
.filter_map(|p| p.split_once('=').map(|(n, _)| n))
.collect();
let mut scope: Vec<String> = manifests
.iter()
.flat_map(|m| m.spec.capabilities.provides.iter().cloned())
.collect();
scope.extend(
manifests
.iter()
.flat_map(|m| &m.spec.tools)
.filter_map(|g| agentplane::tools::ToolId::parse(&g.reference))
.filter(|id| peer_names.contains(&id.server.as_str()))
.map(|id| id.tool),
);
if subject.trim().is_empty() {
return Err("--acting-as names nobody".to_owned());
}
Ok(agentplane::core::Delegation::root(
agentplane::core::Principal::new(subject, agentplane::core::Scope::of(scope)),
))
}
fn in_cli_terms(error: &agentplane::runtime::BuildError, manifests: &[Manifest]) -> String {
use agentplane::runtime::BuildError;
match error {
BuildError::DeclarativeToolsUnreachable { agent, kind, .. } => {
let servers: Vec<String> = manifests
.iter()
.filter(|m| &m.metadata.name == agent)
.flat_map(|m| &m.spec.tools)
.filter_map(|g| agentplane::tools::ToolId::parse(&g.reference))
.map(|id| id.server)
.filter(|s| s != "agent")
.fold(Vec::new(), |mut acc, s| {
if !acc.contains(&s) {
acc.push(s);
}
acc
});
format!(
"agent '{agent}' declares `execution.kind: {kind}` and grants tools on {}, \
but nothing reaches them. Name the process that serves each one: \
`--mcp {}=<command>` for an MCP server, or `--peer <name>=<url>` for an \
A2A peer",
servers.join(", "),
servers.first().map_or("<server>", String::as_str),
)
}
BuildError::PeerIsAlsoAToolServer { server } => format!(
"'{server}' is named by both --mcp and --peer; a grant `tool://{server}/…` \
cannot mean both a tool call and a delegating hop — drop one of the two flags"
),
BuildError::UnknownProvider { agent, provider } => format!(
"agent '{agent}' names provider '{provider}', and this binary ships {}",
shipped_providers().join(", ")
),
other => other.to_string(),
}
}
fn execute(manifests: &[Manifest], opts: &RunArgs) -> Result<ExitCode, Fault> {
require_declarative(manifests).map_err(usage)?;
let chain = opts
.acting_as
.as_deref()
.map(|subject| chain_for(subject, manifests, &opts.peer))
.transpose()
.map_err(usage)?;
let keys = correlation(&opts.correlate).map_err(usage)?;
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).map_err(usage)? {
builder = builder.peers(registry, client);
}
if let Some(chain) = chain {
builder = builder.acting_as(chain);
}
for manifest in manifests {
builder = builder.agent(agentplane::runtime::Agent::new(manifest));
}
let agent = builder
.try_build()
.map_err(|e| in_cli_terms(&e, manifests))?;
let capability = entry_capability(manifests, opts.capability.as_deref()).map_err(usage)?;
let outcome = agent
.run_correlated(
&capability,
Tainted::trusted(opts.read_input().map_err(usage)?),
&capability,
&keys,
)
.await
.map_err(|e| e.to_string())?;
Ok(conclude(
&outcome,
&Resume {
manifest: &opts.manifest,
store: opts.at.store.as_deref(),
tenant: opts.at.tenant.as_deref(),
},
))
})
}
fn replay(manifests: &[Manifest], opts: &ReplayArgs) -> Result<ExitCode, Fault> {
require_declarative(manifests).map_err(usage)?;
if opts.strict {
return verify_replay(manifests, opts);
}
let run_id = opts.run_id.as_deref().ok_or_else(|| {
usage("a resume needs the run to resume: `agentplane replay <run> --manifest …`")
})?;
let store = opts.at.store.as_deref().ok_or_else(|| {
usage("a resume needs the store the run lives in: --store <file|postgres://…>")
})?;
let run = agentplane::core::RunId::parse(run_id)
.map_err(|e| usage(format!("`{run_id}` is not a run id: {e}")))?;
let mode = 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 = Backend::open(store, opts.at.tenant.as_deref()).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).map_err(usage)? {
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| in_cli_terms(&e, manifests))?;
let outcome = match agent.replay(run, mode).await {
Err(agentplane::core::RuntimeError::LeaseHeld { remaining_secs, .. })
if remaining_secs <= 5 =>
{
tokio::time::sleep(std::time::Duration::from_secs(remaining_secs + 1)).await;
agent.replay(run, mode).await
}
other => other,
}
.map_err(|e| e.to_string())?;
Ok(conclude(
&outcome,
&Resume {
manifest: &opts.manifest,
store: Some(store),
tenant: opts.at.tenant.as_deref(),
},
))
})
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord)]
enum Replayed {
Verified,
CannotReplay,
Diverged,
Unreadable,
}
impl Replayed {
fn of(verdict: &agentplane::runtime::Verdict) -> Self {
use agentplane::runtime::Finding;
match verdict.finding {
Finding::Verified { .. } => Self::Verified,
Finding::CannotReplay(_) => Self::CannotReplay,
_ => Self::Diverged,
}
}
}
fn replay_exit(runs: &[Replayed]) -> u8 {
match runs.iter().max() {
None | Some(Replayed::Verified) => exit::OK,
Some(Replayed::CannotReplay) => exit::PARTIAL,
Some(Replayed::Diverged) => exit::FINDING,
Some(Replayed::Unreadable) => exit::OPERATIONAL,
}
}
fn verify_replay(manifests: &[Manifest], opts: &ReplayArgs) -> Result<ExitCode, Fault> {
if !opts.mcp.is_empty() || !opts.peer.is_empty() {
return Err(usage(
"--strict dispatches nothing, so --mcp could only start a server and --peer only \
dial one; drop them — every tool result and peer reply is read from the record",
));
}
if !opts.from.is_empty() && opts.at.store.is_some() {
return Err(usage(
"--from and --store name two sources; replay one of them (AGENTPLANE_STORE counts \
as --store)",
));
}
let wanted = opts
.run_id
.as_deref()
.map(|id| {
agentplane::core::RunId::parse(id)
.map_err(|e| usage(format!("`{id}` is not a run id: {e}")))
})
.transpose()?;
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 mut sources: Vec<(
agentplane::runtime::Stores,
agentplane::core::TenantId,
Vec<agentplane::core::RunId>,
)> = Vec::new();
if opts.from.is_empty() {
let store = opts.at.store.as_deref().ok_or_else(|| {
usage("--strict needs a source: --store <file|postgres://…> or --from <export>")
})?;
let run = wanted.ok_or_else(|| {
usage("a strict replay of a store names its run; replay every run of an export with --from")
})?;
let backend = Backend::open(store, opts.at.tenant.as_deref()).await?;
sources.push((backend.stores(), backend.tenant(), vec![run]));
} else {
for path in &opts.from {
let file = std::fs::File::open(path)
.map_err(|e| format!("could not read the export {path}: {e}"))?;
let source =
agentplane::export::open_for_replay(std::io::BufReader::new(file))
.await
.map_err(|e| format!("{path}: {e}"))?;
let runs = match wanted {
Some(run) if source.runs.contains(&run) => vec![run],
Some(_) => continue,
None => source.runs.clone(),
};
sources.push((
agentplane::runtime::Stores::on(source.store),
agentplane::core::TenantId::default(),
runs,
));
}
if let Some(run) = wanted
&& sources.is_empty()
{
return Err(format!("no export named holds run {run}").into());
}
}
let mut results = Vec::new();
for (stores, tenant, runs) in sources {
for run in runs {
results.push(verify_one(manifests, &stores, &tenant, run).await?);
}
}
Ok(ExitCode::from(replay_exit(&results)))
})
}
async fn verify_one(
manifests: &[Manifest],
stores: &agentplane::runtime::Stores,
tenant: &agentplane::core::TenantId,
run: agentplane::core::RunId,
) -> Result<Replayed, Fault> {
let history = match stores.journal.read(run, 1).await {
Ok(history) => history,
Err(e) => {
eprintln!("run {run} — cannot be read: {e}");
return Ok(Replayed::Unreadable);
}
};
let mut builder = agentplane::runtime::replay_only::wire(
Runtime::builder_with(stores.clone()).tenant(tenant.clone()),
manifests,
&history,
);
for manifest in manifests {
builder = builder.agent(agentplane::runtime::Agent::new(manifest));
}
let plane = builder
.try_build()
.map_err(|e| usage(in_cli_terms(&e, manifests)))?;
match plane.verify(run).await {
Ok(verdict) => {
eprint!("{verdict}");
Ok(Replayed::of(&verdict))
}
Err(e) => {
eprintln!("run {run} — cannot be read: {e}");
Ok(Replayed::Unreadable)
}
}
}
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::{
EXIT_STATUS_HELP, Fault, Truncation, audit_status, cutoff_before, declared_bound, exit,
lift_status, refuse_ambiguous_peers, refuses_partial_export, shell_quote, where_flags,
without_password,
};
#[cfg(feature = "cedar")]
fn scratch(name: &str) -> std::path::PathBuf {
let dir = std::env::temp_dir().join(format!(
"agentplane-bundle-{name}-{}",
agentplane::core::RunId::generate()
));
std::fs::create_dir_all(&dir).expect("scratch dir");
dir
}
#[cfg(feature = "cedar")]
#[test]
fn a_bundle_directory_and_its_rules_file_have_one_identity() {
use agentplane::core::PolicyEngine as _;
let rules = "permit(principal, action, resource);";
let dir = scratch("same");
std::fs::write(dir.join("policy.cedar"), rules).unwrap();
let single = scratch("single").join("served.cedar");
std::fs::write(&single, rules).unwrap();
let from_dir = super::load_policy_bundle(dir.to_str().unwrap()).expect("dir");
let from_file = super::load_policy_bundle(single.to_str().unwrap()).expect("file");
assert_eq!(from_dir.bundle(), from_file.bundle());
assert_eq!(
from_dir.bundle(),
agentplane::policy::CedarEngine::new(rules)
.unwrap()
.bundle(),
"the loader is not the engine `serve` used to build"
);
std::fs::write(dir.join("entities.json"), "[]").unwrap();
let with_entities = super::load_policy_bundle(dir.to_str().unwrap()).expect("dir");
assert_ne!(
with_entities.bundle(),
from_file.bundle(),
"the entities beside the rules did not reach the bundle identity"
);
}
#[cfg(feature = "cedar")]
#[test]
fn a_bundle_directory_holding_another_file_is_refused() {
let dir = scratch("stray");
std::fs::write(
dir.join("policy.cedar"),
"permit(principal, action, resource);",
)
.unwrap();
std::fs::write(
dir.join("extra.cedar"),
"forbid(principal, action, resource);",
)
.unwrap();
match super::load_policy_bundle(dir.to_str().unwrap()) {
Err(Fault::Usage(why)) => assert!(why.contains("extra.cedar"), "{why}"),
Err(other) => panic!("refused as the wrong kind of fault: {other}"),
Ok(_) => panic!("a rule in a file the bundle does not read was accepted"),
}
}
fn agent(models: &str, budgets: &str) -> agentplane::manifest::Manifest {
agentplane::manifest::Manifest::parse(&format!(
"apiVersion: agentplane.hupe1980.github.io/v1alpha1
kind: Agent
metadata: {{ name: bounded, version: \"1.0.0\" }}
spec:
execution: {{ kind: tool-calling, max_turns: 4 }}
capabilities: {{ provides: [bounded.answer] }}
models:
{models}
budgets: {budgets}
"
))
.expect("the fixture parses")
}
const PRICE: &str =
"pricing: { input: 1000000, output: 2000000, cache_read: 100000, cache_write: 1250000 }";
#[test]
fn validate_derives_the_worst_case_from_the_declaration() {
let m = agent(
&format!(
" privileged: {{ provider: fake, model: m, max_tokens: 100, max_input_tokens: 900, {PRICE} }}"
),
"{ max_tokens: 10000, max_minor_units: 50000, max_parallel_steps: 2 }",
);
let bound = declared_bound(&m);
let tokens = &bound.json["units"]["tokens"];
assert_eq!(tokens["per_call"], 1000);
assert_eq!(tokens["total"], 12_000);
let money = &bound.json["units"]["minor_units"];
assert_eq!(money["per_call"], 1325);
assert_eq!(money["total"], 50_000 + 2 * 1325);
assert_eq!(bound.json["turns"], 4);
assert!(
bound
.lines
.iter()
.any(|l| l.contains("12000 = ceiling 10000 + width 2 × one call 1000")),
"{:?}",
bound.lines
);
}
#[test]
fn validate_names_every_unbounded_term() {
let m = agent(
&format!(" privileged: {{ provider: fake, model: m, max_tokens: 100, {PRICE} }}"),
"{ max_tokens: 10000 }",
);
let bound = declared_bound(&m);
for unit in ["tokens", "minor_units"] {
assert!(
bound.json["units"][unit]["total"].is_null(),
"a total was printed for {unit} although a term is unbounded: {}",
bound.json
);
}
let tokens: Vec<String> =
serde_json::from_value(bound.json["units"]["tokens"]["unbounded"].clone())
.expect("names");
assert_eq!(tokens, vec!["spec.models.privileged.max_input_tokens"]);
let money: Vec<String> =
serde_json::from_value(bound.json["units"]["minor_units"]["unbounded"].clone())
.expect("names");
assert_eq!(
money,
vec![
"spec.budgets.max_minor_units",
"spec.models.privileged.max_input_tokens"
]
);
assert!(
bound.lines.iter().any(|l| l.contains("no total")),
"{:?}",
bound.lines
);
}
use std::process::ExitCode;
fn truncation(reached: &[&str]) -> Truncation {
Truncation {
limit: 10,
reached: reached.iter().map(|s| (*s).to_owned()).collect(),
}
}
#[test]
fn the_exit_statuses_are_one_table() {
use clap::CommandFactory as _;
assert_eq!(usage_status(), exit::USAGE);
assert_eq!(
Fault::from("store down".to_owned()).status(),
exit::OPERATIONAL
);
assert_ne!(exit::FINDING, exit::OPERATIONAL);
assert_eq!(audit_status(true, &truncation(&[])), exit::OK);
assert_eq!(
audit_status(true, &truncation(&["succeeded"])),
exit::PARTIAL
);
assert_eq!(
audit_status(false, &truncation(&["succeeded"])),
exit::FINDING
);
assert_eq!(lift_status(true), ExitCode::SUCCESS);
assert_eq!(lift_status(false), ExitCode::from(exit::FINDING));
for (code, word) in [
(exit::OK, "ok"),
(exit::FINDING, "finding"),
(exit::USAGE, "usage"),
(exit::SUSPENDED, "suspended"),
(exit::OPERATIONAL, "operational"),
(exit::PARTIAL, "partial"),
] {
assert!(
EXIT_STATUS_HELP
.lines()
.any(|l| l.trim_start().starts_with(&format!("{code} ")) && l.contains(word)),
"--help does not say that {code} means {word}: {EXIT_STATUS_HELP}"
);
}
let help = super::Cli::command()
.get_after_help()
.map(ToString::to_string)
.unwrap_or_default();
assert_eq!(help, EXIT_STATUS_HELP, "--help prints another table");
}
fn usage_status() -> u8 {
super::usage("bad flag").status()
}
#[test]
fn a_truncated_export_is_refused_without_allow_partial() {
assert!(refuses_partial_export(
&truncation(&["in-flight runs"]),
false
));
assert!(!refuses_partial_export(
&truncation(&["in-flight runs"]),
true
));
assert!(!refuses_partial_export(&truncation(&[]), false));
}
#[test]
fn a_printed_next_step_carries_the_tenant_and_no_password() {
let flags = where_flags(
Some("postgres://ada:s3cret@db.internal:5432/plane?password=hunter2&sslmode=require"),
Some("acme corp"),
);
assert!(
!flags.contains("s3cret"),
"the URL's password was printed: {flags}"
);
assert!(
!flags.contains("hunter2"),
"the password parameter was printed: {flags}"
);
assert!(flags.contains("ada@db.internal:5432/plane"), "{flags}");
assert!(flags.contains("sslmode=require"), "{flags}");
assert!(
flags.ends_with(" --tenant 'acme corp'"),
"the tenant is missing or unquoted: {flags}"
);
assert_eq!(
without_password("postgres://db/plane"),
"postgres://db/plane",
"a URL with no password is printed as it is"
);
assert_eq!(shell_quote("it's"), r"'it'\''s'");
assert_eq!(where_flags(Some("runs.redb"), None), " --store runs.redb");
}
#[test]
fn ambiguous_peer_names_are_refused_at_boot() {
let specs = |names: &[&str]| -> Vec<String> {
names
.iter()
.map(|n| format!("{n}=https://peer.example"))
.collect()
};
let err = refuse_ambiguous_peers(&specs(&["billing.eu", "billing-eu"]))
.expect_err("both read AGENTPLANE_PEER_TOKEN_BILLING_EU");
assert!(
err.contains("billing.eu") && err.contains("billing-eu"),
"{err}"
);
assert!(refuse_ambiguous_peers(&specs(&["billing", "reviewer"])).is_ok());
}
#[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}");
}
#[test]
fn a_corpus_exits_with_its_worst_verdict() {
use super::{Replayed, replay_exit};
assert_eq!(replay_exit(&[]), exit::OK);
assert_eq!(replay_exit(&[Replayed::Verified]), exit::OK);
assert_eq!(
replay_exit(&[Replayed::Verified, Replayed::CannotReplay]),
exit::PARTIAL
);
assert_eq!(
replay_exit(&[
Replayed::CannotReplay,
Replayed::Diverged,
Replayed::Verified
]),
exit::FINDING
);
assert_eq!(
replay_exit(&[Replayed::Diverged, Replayed::Unreadable]),
exit::OPERATIONAL
);
}
#[test]
fn strict_replay_wires_only_replay_only_drivers() {
use agentplane::journal::JournalStore;
const ACME: &str = r#"
apiVersion: agentplane.hupe1980.github.io/v1alpha1
kind: Agent
metadata: { name: summariser, version: "1.0.0" }
spec:
execution: { kind: completion }
identity: { role: "Summarise a support ticket" }
capabilities: { provides: [support.summarise] }
models:
privileged: { provider: acme, model: sum-1 }
budgets: { max_tokens: 10000 }
"#;
let manifests = vec![agentplane::manifest::Manifest::parse(ACME).expect("parses")];
tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.expect("runtime")
.block_on(async {
let store = std::sync::Arc::new(
agentplane::store::RedbStore::open_in_memory().expect("store"),
);
let recorded = agentplane::runtime::Runtime::builder(
std::sync::Arc::clone(&store) as std::sync::Arc<dyn JournalStore>
)
.provider("acme", agentplane::model::fake::FakeProvider::new())
.agent(agentplane::runtime::Agent::new(&manifests[0]))
.build()
.run(
"support.summarise",
agentplane::core::Tainted::trusted(serde_json::json!({ "t": 1 })),
)
.await
.expect("recorded");
let replayed = super::verify_one(
&manifests,
&agentplane::runtime::Stores::on(store),
&agentplane::core::TenantId::default(),
recorded.run_id,
)
.await
.expect("a strict replay needs no driver this binary can build");
assert_eq!(replayed, super::Replayed::Verified);
});
}
fn listed(
summary: &str,
proposed: serde_json::Value,
withheld: Option<agentplane::core::Withheld>,
) -> agentplane::core::Task {
use agentplane::core::{
EffectDescriptor, EffectKey, Justification, OnExpiry, Phase, Priority, RunId, StepId,
Tainted, Task, TaskId, TaskState,
};
let run = RunId::generate();
Task {
id: TaskId::derive(
run,
EffectKey::for_effect(
StepId(0),
Phase::Forward,
0,
1,
&EffectDescriptor::new("approval", serde_json::json!({})),
),
),
run,
case: None,
kind: "approval".into(),
justification: Justification::new(Tainted::trusted(summary.to_owned()), proposed),
candidate_roles: Vec::new(),
escalate_to: Vec::new(),
assignee: None,
priority: Priority::Normal,
state: TaskState::Open,
on_expiry: OnExpiry::Deny,
excluded_actors: Vec::new(),
created_at: time::OffsetDateTime::UNIX_EPOCH,
due_at: None,
withheld,
}
}
#[test]
fn the_terminal_shows_a_withheld_proposal_as_withheld() {
let task = listed(
"Refund the invoice",
serde_json::json!({ "$sealed": "AAECAwQFBgc=" }),
Some(agentplane::core::Withheld::Sealed),
);
let shown = super::task_json(&task, true);
assert_eq!(shown["proposed_action"], serde_json::Value::Null, "{shown}");
assert!(
shown["withheld"]
.as_str()
.is_some_and(|w| w.starts_with("sealed") && w.contains("no key ring")),
"{shown}"
);
assert!(!shown.to_string().contains("$sealed"), "{shown}");
}
#[test]
fn the_terminal_escapes_what_a_reviewer_cannot_see() {
let task = listed(
"Pay\u{200B} the vendor",
serde_json::json!({ "amount": "\u{202E}0001" }),
None,
);
let shown = super::task_json(&task, false);
assert_eq!(shown["escaped"], true, "{shown}");
assert_eq!(
shown["proposed_action"]["amount"], "\\u{202E}0001",
"{shown}"
);
let printed = shown.to_string();
assert!(
!printed.contains('\u{202E}') && !printed.contains('\u{200B}'),
"{printed}"
);
}
}