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,
}
#[cfg(feature = "dev")]
#[path = "agentplane/dev.rs"]
mod dev;
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;
pub const UNVERIFIABLE: u8 = 6;
}
const EXIT_STATUS_HELP: &str = "Exit status:
0 ok
1 a finding or a negative answer (a failed run, an audit or drill finding, an
unused grant, 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, a strict replay could not replay a run,
policy check evaluated nothing, grants met an incomplete export or unreadable
calls, or subject a cut scan or an unreadable run
6 unverifiable: verify or restore met an export under a canon this build does not
implement";
#[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 admission_fault(e: agentplane::core::RuntimeError) -> Fault {
match e {
e @ agentplane::core::RuntimeError::SubjectUnbound { .. } => usage(e.to_string()),
e => e.to_string().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")),
("dev", cfg!(feature = "dev")),
("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,
Openapi,
Digest(DigestArgs),
Audit(AuditArgs),
Export(ExportArgs),
Disclosures(DisclosuresArgs),
Verify(VerifyArgs),
Bind(BindArgs),
Policy(PolicyArgs),
Content(ContentArgs),
Grants(GrantsArgs),
Subject(SubjectArgs),
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),
History(HistoryArgs),
#[cfg(feature = "dev")]
Dev(DevArgs),
}
#[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 HistoryArgs {
run: String,
#[command(flatten)]
at: StoreRef,
#[arg(long, value_name = "SEQ")]
from: Option<u64>,
#[arg(long)]
json: bool,
}
#[cfg(feature = "dev")]
#[derive(clap::Args, Debug)]
struct DevArgs {
manifest: String,
#[arg(long, default_value_t = 0)]
port: u16,
#[arg(long, value_name = "DIR")]
scratch: Option<String>,
#[arg(long, env = "AGENTPLANE_TENANT")]
tenant: Option<String>,
#[arg(long, value_name = "NAME=COMMAND")]
mcp: Vec<String>,
#[arg(long, value_name = "NAME=URL", requires = "acting_as")]
peer: Vec<String>,
#[arg(long)]
allow_live: bool,
#[arg(long, value_name = "SUBJECT")]
acting_as: Option<String>,
}
#[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>,
#[arg(long, value_name = "HEX")]
digest: Option<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 HoldListing {
List(HoldListArgs),
}
#[derive(clap::Args, Debug)]
struct HoldListArgs {
#[command(flatten)]
at: StoreRef,
#[arg(long)]
released: bool,
#[arg(long, default_value_t = 100)]
limit: usize,
}
#[derive(clap::Subcommand, Debug)]
enum HaltListing {
List(HaltListArgs),
}
#[derive(clap::Args, Debug)]
struct HaltListArgs {
#[command(flatten)]
at: StoreRef,
#[arg(long)]
lifted: bool,
#[arg(long, default_value_t = 100)]
limit: usize,
}
#[derive(clap::Args, Debug)]
#[command(args_conflicts_with_subcommands = true, subcommand_negates_reqs = true)]
struct HoldArgs {
#[command(subcommand)]
list: Option<HoldListing>,
#[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 = "reason")]
lift: bool,
}
#[derive(clap::Args, Debug)]
#[command(args_conflicts_with_subcommands = true, subcommand_negates_reqs = true)]
struct HaltArgs {
#[command(subcommand)]
list: Option<HaltListing>,
#[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 = "reason")]
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>,
#[arg(long = "grader-verdict", value_name = "FILE")]
grader_verdict: Vec<String>,
#[arg(long = "grader-key", value_name = "KEY_ID=HEX")]
grader_key: Vec<String>,
}
#[derive(clap::Args, Debug)]
struct BindArgs {
export: String,
#[arg(long)]
run: String,
#[arg(long = "last-seq")]
last_seq: Option<u64>,
#[arg(long)]
content: String,
#[arg(long)]
out: 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>,
#[arg(long = "max-checkpoint-age", value_name = "SECS")]
max_checkpoint_age: Option<u64>,
}
#[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,
#[arg(long = "case", conflicts_with_all = ["outcome", "allow_partial"])]
cases: Vec<String>,
#[arg(long = "run", conflicts_with_all = ["outcome", "allow_partial"])]
runs: Vec<String>,
#[arg(long)]
to: Option<String>,
#[arg(long)]
actor: Option<String>,
#[arg(long)]
output: Option<String>,
}
#[derive(clap::Args, Debug)]
struct DisclosuresArgs {
#[command(flatten)]
at: StoreRef,
#[arg(long = "case")]
cases: Vec<String>,
#[arg(long = "run")]
runs: Vec<String>,
}
#[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,
#[arg(long, value_name = "DIR", conflicts_with_all = ["tools", "name", "path"])]
serve: Option<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>,
#[arg(long, value_name = "DIGEST")]
expect_digest: Option<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 ContentArgs {
#[command(subcommand)]
act: ContentAct,
}
#[derive(clap::Subcommand, Debug)]
enum ContentAct {
Check(ContentCheckArgs),
}
#[derive(clap::Args, Debug)]
struct ContentCheckArgs {
manifest: String,
#[arg(long)]
at: String,
#[arg(long)]
value: Option<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 GrantsArgs {
#[arg(long)]
from: String,
#[arg(long, required = true)]
manifest: Vec<String>,
#[arg(long)]
propose: Option<String>,
#[command(flatten)]
out: JsonFlag,
}
#[derive(clap::Args, Debug)]
struct SubjectArgs {
subject: String,
#[command(flatten)]
at: StoreRef,
#[arg(long, default_value_t = 1000)]
limit: usize,
#[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, env = "AGENTPLANE_MCP_ADDR")]
mcp_addr: Option<String>,
#[arg(long, value_name = "NAME")]
mcp_agent: Vec<String>,
#[arg(long, value_name = "HOST")]
mcp_allowed_host: Vec<String>,
#[arg(long, value_name = "ORIGIN")]
mcp_allowed_origin: Vec<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 = "witness-submit", value_name = "URL")]
witness_submit: Vec<String>,
#[arg(long = "witness-key", value_name = "NAME=BASE64")]
witness_key: Vec<String>,
#[arg(long = "witness-quorum", value_name = "N")]
witness_quorum: Option<usize>,
#[arg(long = "log-key", value_name = "NAME=PATH")]
log_key: Option<String>,
#[arg(long = "witness-interval", value_name = "SECS")]
witness_interval: Option<u64>,
#[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,
#[serde(skip_serializing_if = "<[_]>::is_empty")]
grader_verdicts: &'a [agentplane::grader_verdict::SidecarReport],
}
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)
}
#[cfg(all(feature = "a2a-server", feature = "cedar"))]
fn with_submission_witnesses(
builder: agentplane::runtime::RuntimeBuilder,
opts: &ServeArgs,
) -> Result<agentplane::runtime::RuntimeBuilder, String> {
if opts.witness_submit.is_empty() {
if !opts.witness_key.is_empty()
|| opts.witness_quorum.is_some()
|| opts.log_key.is_some()
|| opts.witness_interval.is_some()
{
return Err(
"--witness-key, --witness-quorum, --log-key and --witness-interval need \
--witness-submit: there is no witness to submit to"
.to_owned(),
);
}
return Ok(builder);
}
let trusted = witness_keys(&opts.witness_key)?;
if trusted.is_empty() {
return Err("--witness-submit needs at least one --witness-key".to_owned());
}
let Some((name, path)) = opts.log_key.as_deref().and_then(|k| k.split_once('=')) else {
return Err(
"--witness-submit needs --log-key <name>=<path to a 64-hex-character seed>: \
a witness recognises a log by its signed note"
.to_owned(),
);
};
let seed_hex = std::fs::read_to_string(path).map_err(|e| format!("--log-key {path}: {e}"))?;
let seed: [u8; 32] = hex::decode(seed_hex.trim())
.ok()
.and_then(|b| b.try_into().ok())
.ok_or_else(|| format!("--log-key {path}: not a 32-byte seed as 64 hex characters"))?;
let signer = agentplane::policy::Ed25519Signer::new(name, &seed);
let public = signer.verifying_key();
let log = agentplane::journal::LogKey::ed25519(name, public, Arc::new(signer))
.map_err(|e| format!("--log-key {name}: {e}"))?;
let mut witnesses: Vec<Arc<dyn agentplane::journal::Witness>> = Vec::new();
for prefix in &opts.witness_submit {
let witness = agentplane::journal::HttpWitness::new(prefix, log.clone(), trusted.clone())
.map_err(|e| format!("--witness-submit {prefix}: {e}"))?;
witnesses.push(Arc::new(witness));
}
let required = opts.witness_quorum.unwrap_or(witnesses.len());
if required > witnesses.len() {
return Err(format!(
"--witness-quorum {required} is above the {} witness(es) given",
witnesses.len()
));
}
let quorum = agentplane::journal::WitnessQuorum::of(required)
.map_err(|e| format!("--witness-quorum {required}: {e}"))?;
let mut builder = builder.witnesses(witnesses, quorum);
if let Some(secs) = opts.witness_interval {
let sweep = opts.sweep_every.unwrap_or(DEFAULT_SWEEP_SECONDS);
if sweep != 0 && secs < u64::from(sweep) {
return Err(format!(
"--witness-interval {secs} is shorter than --sweep-every {sweep}: the sweep \
re-submits, so an interval shorter than it cannot be kept"
));
}
if sweep == 0 {
eprintln!(
"--witness-interval {secs} with --sweep-every 0: the interval is kept only as \
often as your own scheduler sweeps"
);
}
builder = builder.witness_interval(std::time::Duration::from_secs(secs));
}
Ok(builder)
}
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::from_cosigned(
&cosigned,
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(),
},
));
}
#[allow(clippy::disallowed_methods)]
let now = agentplane::core::Timestamp::now_utc();
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,
freshness: audit
.max_checkpoint_age
.map(|secs| agentplane::audit::Freshness {
now,
max_age: std::time::Duration::from_secs(secs),
}),
};
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, 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 grants_verb(opts: &GrantsArgs) -> Result<ExitCode, Fault> {
let manifests = opts
.manifest
.iter()
.map(|path| {
let text =
std::fs::read_to_string(path).map_err(|e| usage(format!("reading {path}: {e}")))?;
Manifest::parse(&text).map_err(|e| usage(format!("{path}: {e}")))
})
.collect::<Result<Vec<_>, Fault>>()?;
let grants = agentplane::grants::Grants::new(&manifests);
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(grants.run(std::io::stdin().lock()))
} else {
let file =
std::fs::File::open(&opts.from).map_err(|e| format!("reading {}: {e}", opts.from))?;
rt.block_on(grants.run(std::io::BufReader::new(file)))
}
.map_err(|e| match e {
agentplane::grants::GrantsError::NotAnExport(_)
| agentplane::grants::GrantsError::Manifest(_) => usage(format!("{}: {e}", opts.from)),
_ => Fault::Operational(e.to_string()),
})?;
if let Some(dir) = &opts.propose {
std::fs::create_dir_all(dir).map_err(|e| format!("creating {dir}: {e}"))?;
for row in &report.digests {
let (Some(digest), Some(proposal)) = (&row.digest, &row.proposal) else {
continue;
};
let yaml = agentplane::manifest::registry::to_yaml(&proposal.manifest)
.map_err(|e| e.to_string())?;
let header = format!(
"# Proposed by `agentplane grants` from {}: unsigned, not published.\n\
# Equal to the input after parsing except for the removed grants: {:?}.\n\
# Comments and key order are not kept. Export {}.\n",
opts.from,
proposal.removed,
if report.window.complete() {
"complete"
} else {
"INCOMPLETE"
},
);
let path = std::path::Path::new(dir).join(format!("{digest}.yaml"));
std::fs::write(&path, header + &yaml)
.map_err(|e| format!("writing {}: {e}", path.display()))?;
}
}
if opts.out.json {
println!(
"{}",
serde_json::to_string_pretty(&report).map_err(|e| e.to_string())?
);
} else {
print!("{report}");
}
Ok(if report.has_unused() {
ExitCode::from(exit::FINDING)
} else if report.partial() {
ExitCode::from(exit::PARTIAL)
} else {
ExitCode::SUCCESS
})
}
fn subject_verb(opts: &SubjectArgs) -> 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 tenant = backend.tenant();
let report = agentplane::subject::Trace::new(tenant.as_str())
.report(
&backend.journal(),
backend.memory().as_ref(),
&opts.subject,
opts.limit,
)
.await
.map_err(|e| Fault::Operational(e.to_string()))?;
if opts.out.json {
println!(
"{}",
serde_json::to_string_pretty(&report).map_err(|e| e.to_string())?
);
} else {
print!("{report}");
}
Ok(if report.partial() {
ExitCode::from(exit::PARTIAL)
} else {
ExitCode::SUCCESS
})
})
}
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 found =
agentplane::export::runs_to_read(&store, &wanted, opts.outcome.is_empty(), opts.limit)
.await
.map_err(|e| e.to_string())?;
let runs = found.runs;
let truncated = found.reached;
for (run, why) in &found.unreadable {
eprintln!("warning: in-flight run {run} could not be read: {why}");
}
if found.in_flight > 0 {
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",
found.in_flight
);
}
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, released: Option<usize>) -> Result<ExitCode, Fault> {
blocking(async {
let backend = at.open().await?;
if let Some(limit) = released {
let plane = backend.plane().build();
let mut records = plane
.released_holds(limit.saturating_add(1))
.await
.map_err(|e| e.to_string())?;
let truncated = records.len() > limit;
records.truncate(limit);
let rows: Vec<serde_json::Value> = records.iter().filter_map(release_row).collect();
println!(
"{}",
serde_json::to_string_pretty(&serde_json::json!({
"released": rows,
"truncated": truncated,
}))
.map_err(|e| e.to_string())?
);
return Ok(listed_status(truncated));
}
let cases = backend.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(HoldListing::List(list)), ..) => {
return holds_verb(&list.at, list.released.then_some(list.limit));
}
(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}"))?;
if opts.lift {
let by = lifter(opts.actor.as_deref(), "release a hold")?;
let case =
agentplane::core::CaseId::parse(case).map_err(|e| usage(format!("--case: {e}")))?;
return rt.block_on(async {
let plane = at.open().await?.plane().build();
#[allow(clippy::disallowed_methods)]
let now = time::OffsetDateTime::now_utc();
let record = plane
.release_hold(case, &by, now)
.await
.map_err(|e| e.to_string())?;
println!(
"{}",
serde_json::to_string_pretty(&serde_json::json!({
"case": case.to_string(),
"lifted": record.is_some(),
"removed": record.is_some_and(|r| r.removed),
"by": by.actor(),
"basis": by.basis().as_str(),
"record": record.map(|r| r.record.to_string()),
}))
.map_err(|e| e.to_string())?
);
Ok(lift_status(record.is_some()))
});
}
rt.block_on(async {
let cases = at.open().await?.cases();
let case =
agentplane::core::CaseId::parse(case).map_err(|e| usage(format!("--case: {e}")))?;
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(HaltListing::List(list)), _) => {
return halts_verb(&list.at, list.lifted.then_some(list.limit), 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()
))
})?;
if opts.lift {
let by = lifter(opts.actor.as_deref(), "lift a halt")?;
return lift_one_halt(at, &scope, &by, opts.json);
}
let thrown = {
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()))?;
(by, reason)
};
let rt = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.map_err(|e| format!("could not start the async runtime: {e}"))?;
let (by, reason) = &thrown;
rt.block_on(async {
let quotas = at.open().await?.quotas();
#[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!(
"{}",
serde_json::json!({
"scope": scope.key(),
"halted": true,
"reason": reason,
"by": by.actor(),
"basis": by.basis().as_str(),
})
);
} else {
println!(
"halted {}: {reason} (by {}, {})",
scope.key(),
by.actor(),
by.basis().as_str()
);
}
Ok(ExitCode::SUCCESS)
})
}
fn lift_one_halt(
at: &StoreRef,
scope: &agentplane::quota::HaltScope,
by: &agentplane::core::Operator,
json: bool,
) -> Result<ExitCode, Fault> {
blocking(async {
let plane = at.open().await?.plane().build();
#[allow(clippy::disallowed_methods)]
let now = time::OffsetDateTime::now_utc();
let record = plane
.lift_halt(scope, by, now)
.await
.map_err(|e| e.to_string())?;
if json {
println!(
"{}",
serde_json::json!({
"scope": scope.key(),
"halted": false,
"was_standing": record.is_some(),
"removed": record.is_some_and(|r| r.removed),
"by": by.actor(),
"basis": by.basis().as_str(),
"record": record.map(|r| r.record.to_string()),
})
);
} else if let Some(lifted) = record {
let removal = if lifted.removed {
""
} else {
"; another lift removed the row first"
};
println!(
"lifted {} (by {}, {}; recorded in run {}{removal})",
scope.key(),
by.actor(),
by.basis().as_str(),
lifted.record
);
} else {
println!("no halt was standing on {}; nothing lifted", scope.key());
}
Ok(lift_status(record.is_some()))
})
}
fn lifter(actor: Option<&str>, act: &str) -> Result<agentplane::core::Operator, Fault> {
let Some(actor) = actor else {
return Err(usage(format!(
"--actor is required to {act}: the lift is recorded under that name, and the \
record says it was asserted rather than authenticated"
)));
};
agentplane::core::Operator::asserted(actor).map_err(|e| usage(e.to_string()))
}
fn lift_status(was_standing: bool) -> ExitCode {
ExitCode::from(if was_standing {
exit::OK
} else {
exit::FINDING
})
}
fn release_row(record: &agentplane::journal::Record) -> Option<serde_json::Value> {
let agentplane::journal::RecordKind::HoldReleased {
by,
at,
placed_by,
placed_at,
} = record.kind()
else {
return None;
};
Some(serde_json::json!({
"case": record.body.case.map(|c| c.to_string()),
"by": by.actor(),
"basis": by.basis().as_str(),
"released_at": at.unix_timestamp(),
"placed_by": placed_by.actor(),
"placed_basis": placed_by.basis().as_str(),
"placed_at": placed_at.unix_timestamp(),
"record": record.body.run.to_string(),
}))
}
fn lift_row(record: &agentplane::journal::Record) -> Option<serde_json::Value> {
let agentplane::journal::RecordKind::HaltLifted {
scope,
by,
at,
reason,
thrown_by,
thrown_at,
} = record.kind()
else {
return None;
};
Some(serde_json::json!({
"scope": scope,
"by": by.actor(),
"basis": by.basis().as_str(),
"lifted_at": at.unix_timestamp(),
"reason": reason,
"thrown_by": thrown_by.actor(),
"thrown_basis": thrown_by.basis().as_str(),
"thrown_at": thrown_at.unix_timestamp(),
"record": record.body.run.to_string(),
}))
}
fn halts_verb(at: &StoreRef, lifted: Option<usize>, 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 backend = at.open().await?;
if let Some(limit) = lifted {
let plane = backend.plane().build();
let mut records = plane
.lifted_halts(limit.saturating_add(1))
.await
.map_err(|e| e.to_string())?;
let truncated = records.len() > limit;
records.truncate(limit);
let rows: Vec<serde_json::Value> = records.iter().filter_map(lift_row).collect();
if json {
println!(
"{}",
serde_json::to_string_pretty(&serde_json::json!({
"lifted": rows,
"truncated": truncated,
}))
.map_err(|e| e.to_string())?
);
} else if rows.is_empty() {
println!("no lifts recorded");
} else {
for row in &rows {
println!(
"{}\tlifted {} by {} ({}); thrown {} by {} ({}): {}\trun {}",
row["scope"].as_str().unwrap_or_default(),
row["lifted_at"],
row["by"].as_str().unwrap_or_default(),
row["basis"].as_str().unwrap_or_default(),
row["thrown_at"],
row["thrown_by"].as_str().unwrap_or_default(),
row["thrown_basis"].as_str().unwrap_or_default(),
row["reason"].as_str().unwrap_or_default(),
row["record"].as_str().unwrap_or_default(),
);
}
if truncated {
println!("truncated: --limit {limit} cut the listing short");
}
}
return Ok(listed_status(truncated));
}
let quotas = backend.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 mut waiting = store
.waiting_runs(opts.limit.saturating_add(1))
.await
.map_err(|e| e.to_string())?;
let truncated = waiting.len() > opts.limit;
waiting.truncate(opts.limit);
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,
"truncated": truncated,
}))
.map_err(|e| e.to_string())?
);
Ok(listed_status(truncated))
})
}
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 = backend.plane().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 = backend.plane().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,
"digest": j.digest().to_hex(),
});
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(listed_status(truncated))
})
}
fn listed_status(truncated: bool) -> ExitCode {
ExitCode::from(if truncated { exit::PARTIAL } else { exit::OK })
}
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 expected = opts
.digest
.as_deref()
.map(|hex| {
agentplane::core::Digest::from_hex(hex)
.map_err(|e| usage(format!("--digest `{hex}` is not a digest: {e}")))
})
.transpose()?;
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 = backend
.plane()
.lease_ttl(std::time::Duration::from_secs(2))
.build();
let delivery = match plane
.decide_task_at(id, &decision, &opts.roles, expected)
.await
{
Ok(delivery) => delivery,
Err(e @ agentplane::core::RuntimeError::TaskChanged { .. }) => {
eprintln!(
"{e}\n agentplane tasks --show {id}{}",
where_flags(Some(&opts.at.store), opts.at.tenant.as_deref())
);
return Ok(ExitCode::from(exit::FINDING));
}
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 = backend.plane().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 = backend.plane().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 = backend.plane().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
})
})
}
const HISTORY_PAGE: usize = 500;
fn history_verb(opts: &HistoryArgs) -> Result<ExitCode, Fault> {
let run = agentplane::core::RunId::parse(&opts.run)
.map_err(|e| usage(format!("`{}` is not a run id: {e}", opts.run)))?;
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 lines = history_lines(&backend.journal(), run, opts.from, opts.json).await?;
let Some(lines) = lines else {
eprintln!(
"no run {run} in {}{}",
without_password(&opts.at.store),
opts.at
.tenant
.as_deref()
.map_or_else(String::new, |t| format!(" (tenant {t})"))
);
return Ok(ExitCode::from(exit::FINDING));
};
let mut out = std::io::stdout().lock();
for line in lines {
std::io::Write::write_all(&mut out, line.as_bytes())
.and_then(|()| std::io::Write::write_all(&mut out, b"\n"))
.map_err(|e| format!("could not write the timeline: {e}"))?;
}
Ok(ExitCode::SUCCESS)
})
}
async fn history_lines(
journal: &Arc<dyn JournalStore>,
run: agentplane::core::RunId,
from: Option<u64>,
json: bool,
) -> Result<Option<Vec<String>>, Fault> {
let start = from.unwrap_or(1).max(1);
let mut next = start;
let mut lines = Vec::new();
loop {
let page = journal
.read_page(run, next, HISTORY_PAGE)
.await
.map_err(|e| e.to_string())?;
for record in &page {
let view = agentplane::journal::view::record_view(record);
lines.push(if json {
serde_json::to_string(&view).map_err(|e| e.to_string())?
} else {
timeline_line(&view)
});
next = record.seq() + 1;
}
if page.len() < HISTORY_PAGE {
break;
}
}
if lines.is_empty()
&& (start == 1
|| journal
.read_page(run, 1, 1)
.await
.map_err(|e| e.to_string())?
.is_empty())
{
return Ok(None);
}
Ok(Some(lines))
}
fn timeline_line(view: &agentplane::journal::view::RecordView) -> String {
use std::fmt::Write as _;
let mut payload = view.record.clone();
if let Some(fields) = payload.as_object_mut() {
fields.remove("kind");
}
let mut line = format!("{:>5} {}", view.seq, view.kind);
if let Some(step) = &view.step {
let _ = write!(line, " step {step}");
}
if view.phase != "forward" {
let _ = write!(line, " {}", view.phase);
}
line.push_str(" ");
line.push_str(&serde_json::to_string(&payload).unwrap_or_default());
agentplane::core::visible::escape(&line).0
}
#[cfg(feature = "dev")]
const DEV_MARKER: &str = ".agentplane-dev";
#[cfg(feature = "dev")]
const DEV_MARKER_TEXT: &str = "created by `agentplane dev`; not a deployment's store\n";
#[cfg(feature = "dev")]
const DEV_STORE: &str = "dev.redb";
#[cfg(feature = "dev")]
fn dev_store(scratch: Option<&str>, tenant: Option<&str>) -> Result<Backend, Fault> {
if let Some(other) = tenant.filter(|t| *t != agentplane::api::dev::TENANT) {
return Err(usage(format!(
"`agentplane dev` runs as tenant `{}` only, and `{other}` was named (by --tenant or \
AGENTPLANE_TENANT) — a deployment's tenant is not a scratch plane",
agentplane::api::dev::TENANT
)));
}
let tenant =
agentplane::core::TenantId::new(agentplane::api::dev::TENANT).map_err(|e| e.to_string())?;
let Some(dir) = scratch else {
let store = RedbStore::open_in_memory().map_err(|e| e.to_string())?;
return Ok(Backend::Embedded(
Arc::new(store.for_tenant(tenant.clone())),
tenant,
));
};
if is_connection_string(dir) {
return Err(usage(
"--scratch names a database; `agentplane dev` opens only a directory it created",
));
}
let dir = std::path::Path::new(dir);
let kind = |path: &std::path::Path| std::fs::symlink_metadata(path).map(|m| m.file_type());
if let Ok(found) = kind(dir) {
if found.is_symlink() {
return Err(usage(format!(
"--scratch {} is a symbolic link; `agentplane dev` opens only a directory it \
created, and a link can name any other",
dir.display()
)));
}
if !found.is_dir() {
return Err(usage(format!(
"--scratch {} is a file; `agentplane dev` opens only a directory it created, so \
no deployment's store is one it writes to",
dir.display()
)));
}
}
std::fs::create_dir_all(dir).map_err(|e| format!("could not create {}: {e}", dir.display()))?;
let held: Vec<String> = std::fs::read_dir(dir)
.map_err(|e| format!("could not read {}: {e}", dir.display()))?
.filter_map(Result::ok)
.map(|entry| entry.file_name().to_string_lossy().into_owned())
.collect();
let marker = dir.join(DEV_MARKER);
let ours = kind(&marker).is_ok_and(|k| k.is_file())
&& std::fs::read(&marker).is_ok_and(|bytes| bytes == DEV_MARKER_TEXT.as_bytes());
if held.is_empty() {
std::fs::write(&marker, DEV_MARKER_TEXT)
.map_err(|e| format!("could not mark {}: {e}", dir.display()))?;
} else if !ours {
return Err(usage(format!(
"--scratch {} holds files and no dev marker; `agentplane dev` opens only an empty \
directory or one it marked",
dir.display()
)));
} else if let Some(other) = held.iter().find(|n| *n != DEV_MARKER && *n != DEV_STORE) {
return Err(usage(format!(
"--scratch {} holds `{other}` beside its dev store; `agentplane dev` opens only a \
directory holding nothing else",
dir.display()
)));
} else if kind(&dir.join(DEV_STORE)).is_ok_and(|k| !k.is_file()) {
return Err(usage(format!(
"--scratch {}/{DEV_STORE} is not a regular file; `agentplane dev` opens only the \
store it created there",
dir.display()
)));
}
let path = dir.join(DEV_STORE);
let store = RedbStore::open(&path).map_err(|e| Fault::from(held_by_a_plane(&e.to_string())))?;
Ok(Backend::Embedded(
Arc::new(store.for_tenant(tenant.clone())),
tenant,
))
}
#[cfg(feature = "dev")]
fn refuse_live_without_consent(opts: &DevArgs) -> Result<(), Fault> {
if !opts.allow_live && (!opts.mcp.is_empty() || !opts.peer.is_empty()) {
return Err(usage(
"--mcp and --peer reach real systems, and approving a task on the page performs \
a real effect through them; add --allow-live to say they are yours to act on",
));
}
Ok(())
}
#[cfg(feature = "dev")]
async fn bind_dev(port: u16) -> Result<tokio::net::TcpListener, Fault> {
tokio::net::TcpListener::bind((std::net::Ipv4Addr::LOCALHOST, port))
.await
.map_err(|e| usage(format!("could not listen on 127.0.0.1:{port}: {e}")))
}
#[cfg(feature = "dev")]
fn dev_actor() -> String {
let user = std::env::var("USER")
.or_else(|_| std::env::var("USERNAME"))
.ok()
.filter(|u| {
!u.trim().is_empty()
&& u.chars()
.all(|c| c.is_ascii_alphanumeric() || "-_.".contains(c))
})
.unwrap_or_else(|| "author".to_owned());
format!("dev:{user}")
}
#[cfg(feature = "dev")]
fn dev_verb(opts: &DevArgs) -> Result<ExitCode, Fault> {
refuse_live_without_consent(opts)?;
let manifests = manifests_at(&opts.manifest).map_err(usage)?;
require_declarative(&manifests).map_err(usage)?;
let token = fresh_token()?;
let actor = dev_actor();
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 {
let backend = dev_store(opts.scratch.as_deref(), opts.tenant.as_deref())?;
let bench = dev::Bench::start(
backend,
dev::Wiring {
file: opts.manifest.clone(),
mcp: opts.mcp.clone(),
peer: opts.peer.clone(),
acting_as: opts.acting_as.clone(),
policy: Arc::new(agentplane::api::dev::DevPolicy::new(&actor)),
},
manifests,
)
.await?;
let auth = agentplane::api::tokens::TokenAuthenticator::new(vec![
agentplane::api::tokens::TokenEntry {
token: token.clone(),
actor,
roles: Vec::new(),
tenant: Some(agentplane::api::dev::TENANT.to_owned()),
scope: None,
not_after: None,
},
])
.map_err(|e| e.to_string())?;
let listener = bind_dev(opts.port).await?;
let port = listener
.local_addr()
.map_err(|e| format!("the listener has no address: {e}"))?
.port();
println!("http://127.0.0.1:{port}/#t={token}");
let streams = agentplane::api::dev::Workbench::streams(&bench);
axum::serve(
listener,
agentplane::api::dev::router(Arc::new(bench), Arc::new(auth), port),
)
.with_graceful_shutdown(async move {
let _ = tokio::signal::ctrl_c().await;
streams.close();
})
.await
.map_err(|e| format!("the dev page stopped: {e}"))?;
Ok(ExitCode::SUCCESS)
})
}
enum Backend {
Embedded(Arc<RedbStore>, agentplane::core::TenantId),
#[cfg(feature = "postgres")]
Shared(
Arc<agentplane::store::PostgresStore>,
agentplane::core::TenantId,
),
}
impl Backend {
#[allow(clippy::unused_async, clippy::unused_async_trait_impl)]
async fn open(spec: &str, tenant: Option<&str>) -> Result<Self, 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 plane(&self) -> RuntimeBuilder {
Runtime::builder_with(self.stores()).tenant(self.tenant())
}
fn in_memory() -> Result<Self, Fault> {
let tenant = agentplane::core::TenantId::default();
let store = RedbStore::open_in_memory().map_err(|e| e.to_string())?;
Ok(Self::Embedded(
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 disclosures(&self) -> Arc<dyn agentplane::disclosure::DisclosureRegister> {
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 memory(&self) -> Arc<dyn agentplane::memory::MemoryStore> {
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 content_check_verb(a: &ContentCheckArgs) -> Result<ExitCode, Fault> {
use agentplane::content::At;
use std::io::Read as _;
let at = match a.at.split_once(':') {
None if a.at == "admission" => At::Admission,
Some(("source", kind)) => At::Source(kind),
Some(("sink", kind)) => At::Sink(kind),
_ => {
return Err(usage(format!(
"--at '{}': use admission, source:<kind> or sink:<kind>",
a.at
)));
}
};
let manifests = manifests_at(&a.manifest).map_err(usage)?;
let [manifest] = manifests.as_slice() else {
return Err(usage(format!(
"{} holds {} agents; content check reads one",
a.manifest,
manifests.len()
)));
};
let text = if let Some(path) = &a.value {
std::fs::read_to_string(path).map_err(|e| usage(format!("reading {path}: {e}")))?
} else {
let mut text = String::new();
std::io::stdin()
.read_to_string(&mut text)
.map_err(|e| usage(format!("reading standard input: {e}")))?;
text
};
let value: serde_json::Value =
serde_json::from_str(&text).map_err(|e| usage(format!("the value is not JSON: {e}")))?;
let Some(content) = manifest.spec.security.content.as_ref() else {
println!("{}", serde_json::json!({ "evaluated": [] }));
return Ok(ExitCode::SUCCESS);
};
let outcome = content.rules().map_err(usage)?.at(at, &value);
let hits = |hits: &[agentplane::content::Hit]| {
hits.iter()
.map(|h| serde_json::json!({ "rule": h.rule, "pointer": h.pointer }))
.collect::<Vec<_>>()
};
let not_evaluable: Vec<&str> = content.checks_at(at).map(|c| c.id.as_str()).collect();
println!(
"{}",
serde_json::to_string_pretty(&serde_json::json!({
"evaluated": outcome.evaluated,
"refused": hits(&outcome.refused),
"redacted": hits(&outcome.redactions),
"classified": outcome.classified,
"sensitivity": outcome.sensitivity,
"checks_not_evaluable": not_evaluable,
}))
.map_err(|e| e.to_string())?
);
Ok(if outcome.refused.is_empty() {
ExitCode::SUCCESS
} else {
ExitCode::from(exit::FINDING)
})
}
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::Openapi => openapi_verb(),
Verb::Digest(a) => digest_verb(&a),
Verb::Audit(a) => journal_verb(&a.store, Some(&a), false),
Verb::Export(a) if !a.cases.is_empty() || !a.runs.is_empty() => disclose_verb(&a),
Verb::Export(a) => journal_verb(&a.store, None, a.allow_partial),
Verb::Disclosures(a) => disclosures_verb(&a),
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::History(a) => history_verb(&a),
#[cfg(feature = "dev")]
Verb::Dev(a) => dev_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 open = || {
std::fs::File::open(&a.file)
.map(std::io::BufReader::new)
.map_err(|e| format!("reading {}: {e}", a.file))
};
if let Some(why) = restore_unverifiable(open()?) {
eprintln!("restore: {why}");
return Ok(ExitCode::from(exit::UNVERIFIABLE));
}
let backend = a.at.open().await?;
let store = backend.journal();
let cases = backend.cases();
let file = open()?;
let report = agentplane::export::from_jsonl(&store, Some(&cases), 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::Bind(a) => bind_verb(&a),
Verb::Policy(PolicyArgs {
act: PolicyAct::Check(a),
}) => policy_check_verb(&a),
Verb::Content(ContentArgs {
act: ContentAct::Check(a),
}) => content_check_verb(&a),
Verb::Grants(a) => grants_verb(&a),
Verb::Subject(a) => subject_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> {
if let Some(dir) = &opts.serve {
return init_serve_verb(dir, opts.out);
}
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)
}
const SERVED_FILES: [&str; 7] = [
"agent.yaml",
"policy.cedar",
"tokens.yaml",
"framework.token",
"postgres.password",
"store.env",
"compose.yaml",
];
#[cfg(all(feature = "a2a-server", feature = "cedar", feature = "mcp-server-http"))]
const SERVED_MANIFEST: &str = include_str!("../../examples/served-starter.yaml");
#[cfg(all(feature = "a2a-server", feature = "cedar", feature = "mcp-server-http"))]
const SERVED_POLICY: &str = include_str!("../../examples/serve-policy.cedar");
#[cfg(all(feature = "a2a-server", feature = "cedar", feature = "mcp-server-http"))]
const SERVED_COMPOSE: &str = include_str!("../../examples/compose.yaml");
#[cfg(all(feature = "a2a-server", feature = "cedar", feature = "mcp-server-http"))]
#[derive(Debug)]
struct Served {
paths: Vec<String>,
digest: String,
}
#[cfg(any(
feature = "dev",
all(feature = "a2a-server", feature = "cedar", feature = "mcp-server-http")
))]
fn fresh_token() -> Result<String, Fault> {
use rand::TryRng as _;
let mut bytes = [0_u8; agentplane::api::tokens::MIN_TOKEN_BYTES];
rand::rngs::SysRng
.try_fill_bytes(&mut bytes)
.map_err(|e| format!("the operating system's random source failed: {e}"))?;
Ok(bytes.iter().fold(String::new(), |mut hex, b| {
let _ = std::fmt::Write::write_fmt(&mut hex, format_args!("{b:02x}"));
hex
}))
}
#[cfg(all(feature = "a2a-server", feature = "cedar", feature = "mcp-server-http"))]
#[derive(Default)]
struct Written {
paths: Vec<std::path::PathBuf>,
kept: bool,
}
#[cfg(all(feature = "a2a-server", feature = "cedar", feature = "mcp-server-http"))]
impl Written {
fn create(
&mut self,
dir: &std::path::Path,
name: &str,
text: &str,
secret: bool,
) -> Result<(), Fault> {
let path = dir.join(name);
let mut open = std::fs::OpenOptions::new();
open.write(true).create_new(true);
#[cfg(unix)]
{
use std::os::unix::fs::OpenOptionsExt as _;
open.mode(if secret { 0o600 } else { 0o644 });
}
#[cfg(not(unix))]
let _ = secret;
let mut file = open
.open(&path)
.map_err(|e| format!("writing {}: {e}", path.display()))?;
self.paths.push(path.clone());
std::io::Write::write_all(&mut file, text.as_bytes())
.map_err(|e| format!("writing {}: {e}", path.display()).into())
}
fn keep(mut self) -> Vec<String> {
self.kept = true;
self.paths.iter().map(|p| p.display().to_string()).collect()
}
}
#[cfg(all(feature = "a2a-server", feature = "cedar", feature = "mcp-server-http"))]
impl Drop for Written {
fn drop(&mut self) {
if !self.kept {
for path in &self.paths {
let _ = std::fs::remove_file(path);
}
}
}
}
#[cfg(all(feature = "a2a-server", feature = "cedar", feature = "mcp-server-http"))]
fn plane_user(uid: u32, gid: u32) -> Result<String, Fault> {
if uid == 0 {
return Err(usage(
"init --serve ran as root, so the token file is root's and the plane would run \
as root to read it; run it as the user the plane should run as (in a container, \
`docker run --user \"$(id -u):$(id -g)\"`), so nothing was written",
));
}
Ok(format!("{uid}:{gid}"))
}
#[cfg(all(feature = "a2a-server", feature = "cedar", feature = "mcp-server-http"))]
fn init_serve(dir: &std::path::Path) -> Result<Served, Fault> {
std::fs::create_dir_all(dir).map_err(|e| format!("creating {}: {e}", dir.display()))?;
if let Some(existing) = SERVED_FILES
.iter()
.find(|name| dir.join(name).symlink_metadata().is_ok())
{
return Err(usage(format!(
"{} exists; init --serve writes nothing over a file, so nothing was written",
dir.join(existing).display()
)));
}
let [peer, framework, operator] = [fresh_token()?, fresh_token()?, fresh_token()?];
let tokens = format!(
"# The callers this plane accepts, generated by `agentplane init --serve`.\n\
# Keep this file secret; `serve` reads it from a mounted file, never from\n\
# the environment. Roles are inputs to policy.cedar, not grants.\n\
- token: \"{peer}\"\n actor: peer-1\n roles: [peer]\n\n\
- token: \"{framework}\"\n actor: framework-1\n roles: [framework]\n\n\
- token: \"{operator}\"\n actor: ops-1\n roles: [operator]\n"
);
let password = fresh_token()?;
let store = format!(
"# The plane's journal, read by `serve` as AGENTPLANE_STORE. Keep this file\n\
# secret: it holds the Postgres password.\n\
AGENTPLANE_STORE=postgres://agentplane:{password}@postgres:5432/agentplane?sslmode=disable\n"
);
let parsed = Manifest::parse_all(SERVED_MANIFEST)
.map_err(|e| format!("the served starter does not parse: {e}"))?;
agentplane::policy::CedarEngine::from_bundle(SERVED_POLICY, None, None)
.map_err(|e| format!("the shipped policy was refused: {e}"))?;
agentplane::api::tokens::TokenAuthenticator::from_yaml(&tokens)
.map_err(|e| format!("the generated tokens were refused: {e}"))?;
let digest = parsed
.first()
.map(|m| m.digest().map(agentplane::Digest::to_hex))
.transpose()
.map_err(|e| e.to_string())?
.unwrap_or_default();
let mut written = Written::default();
written.create(dir, SERVED_FILES[0], SERVED_MANIFEST, false)?;
written.create(dir, SERVED_FILES[1], SERVED_POLICY, false)?;
written.create(dir, SERVED_FILES[2], &tokens, true)?;
written.create(dir, SERVED_FILES[3], &framework, true)?;
written.create(dir, SERVED_FILES[4], &format!("{password}\n"), true)?;
written.create(dir, SERVED_FILES[5], &store, true)?;
#[cfg(unix)]
let owner = {
use std::os::unix::fs::MetadataExt as _;
let token_file = dir.join(SERVED_FILES[2]);
let meta = std::fs::metadata(&token_file)
.map_err(|e| format!("reading {}: {e}", token_file.display()))?;
plane_user(meta.uid(), meta.gid())?
};
#[cfg(not(unix))]
let owner = "65532:65532".to_owned();
let compose = SERVED_COMPOSE
.replace("AGENTPLANE_VERSION", env!("CARGO_PKG_VERSION"))
.replace("PLANE_USER", &owner);
written.create(dir, SERVED_FILES[6], &compose, false)?;
Ok(Served {
paths: written.keep(),
digest,
})
}
#[cfg(all(feature = "a2a-server", feature = "cedar", feature = "mcp-server-http"))]
fn served_report(served: &Served, json: bool) -> String {
if json {
return serde_json::json!({ "wrote": served.paths, "digest": served.digest }).to_string();
}
let mut out = String::new();
for (index, path) in served.paths.iter().enumerate() {
out.push_str("wrote ");
out.push_str(path);
if index == 0 {
let _ = std::fmt::Write::write_fmt(&mut out, format_args!(" ({})", served.digest));
}
out.push('\n');
}
out
}
#[cfg(all(feature = "a2a-server", feature = "cedar", feature = "mcp-server-http"))]
fn init_serve_verb(dir: &str, out: JsonFlag) -> Result<ExitCode, Fault> {
let served = init_serve(std::path::Path::new(dir))?;
print!("{}", served_report(&served, out.json));
if out.json {
println!();
}
eprintln!(
"next:\n docker compose -f {dir}/compose.yaml up --wait\n \
then a framework quickstart: \
https://hupe1980.github.io/agentplane/docs/getting-started/#zero-to-governed"
);
Ok(ExitCode::SUCCESS)
}
#[cfg(all(feature = "a2a-server", feature = "cedar"))]
fn a2a_endpoint(url: &str) -> Result<&str, Fault> {
if url.ends_with("/a2a") {
return Ok(url);
}
let base = url.trim_end_matches('/');
let fix = if base.ends_with("/a2a") {
base.to_owned()
} else {
format!("{base}/a2a")
};
Err(usage(format!(
"--url {url} does not end in /a2a: the Agent Card publishes it verbatim and this \
plane serves A2A at /a2a only, so a client following the card would be answered \
404. Pass --url {fix}"
)))
}
#[cfg(not(all(feature = "a2a-server", feature = "cedar", feature = "mcp-server-http")))]
#[allow(clippy::unnecessary_wraps)]
fn init_serve_verb(_dir: &str, _out: JsonFlag) -> Result<ExitCode, Fault> {
Err(usage(format!(
"this build cannot write a served plane: `init --serve` writes for `serve --mcp-addr`, \
which needs the `a2a-server`, `cedar` and `mcp-server-http` features. Use the \
`:full` container image, or reinstall with `--features \
cli,a2a-server,cedar,mcp-server-http`; it would have written {}",
SERVED_FILES.join(", ")
)))
}
fn read_checkpoint(
path: &str,
) -> Result<
(
agentplane::journal::Checkpoint,
Option<agentplane::journal::SignedNote>,
),
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, None));
}
if let Ok(note) = agentplane::journal::SignedNote::parse(&text)
&& let Ok(cp) = agentplane::journal::Checkpoint::from_note(¬e.text)
{
return Ok((cp, Some(note)));
}
serde_json::from_str(&text)
.map(|cp| (cp, None))
.map_err(|e| {
format!(
"--checkpoint {path} is neither a tlog-checkpoint note nor the `current` \
field of an audit report: {e}"
)
})
}
struct FileAnchor {
anchor: agentplane::audit::Anchor,
cosigned_by: Vec<String>,
signed_note: bool,
}
fn checkpoint_anchor(path: &str, keys: &[String]) -> Result<FileAnchor, String> {
let (checkpoint, note) = read_checkpoint(path)?;
let from = format!("file {path}");
let signed_note = note.is_some();
let Some(note) = note.filter(|_| !keys.is_empty()) else {
return Ok(FileAnchor {
anchor: agentplane::audit::Anchor::new(checkpoint, from),
cosigned_by: Vec::new(),
signed_note,
});
};
let trusted = witness_keys(keys)?;
let cosignatures = agentplane::journal::cosignatures_in(¬e, &trusted);
let cosigned_by = cosignatures
.iter()
.map(|c| format!("{path}:{}", c.key_id))
.collect();
let anchor = if cosignatures.is_empty() {
agentplane::audit::Anchor::new(checkpoint, from)
} else {
agentplane::audit::Anchor::from_cosigned(
&agentplane::journal::CosignedCheckpoint {
checkpoint,
cosignatures,
},
from,
)
};
Ok(FileAnchor {
anchor,
cosigned_by,
signed_note,
})
}
#[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(checkpoint_anchor(path, &opts.witness_key).map_err(usage)?),
None => None,
};
let note_checked = saved.as_ref().is_some_and(|file| file.signed_note);
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() && !note_checked {
return Err(usage(
"--witness-key was given with no --witness and no signed-note \
--checkpoint to use it against"
.to_owned(),
));
}
Anchor::default()
}
};
let mut anchor = fetched;
if let Some(saved) = saved {
anchor.checkpoints.push(saved.anchor);
anchor.cosigned_by.extend(saved.cosigned_by);
}
let verifier = verifier
.as_ref()
.map(|v| v as &dyn agentplane::core::Verifier);
if opts.grader_verdict.is_empty() && !opts.grader_key.is_empty() {
return Err(usage(
"--grader-key was given with no --grader-verdict to use it against".to_owned(),
));
}
let graders =
verifier_from(&opts.grader_key).map_err(|e| usage(e.replace("--key", "--grader-key")))?;
let sidecars = opts
.grader_verdict
.iter()
.map(|path| std::fs::read(path).map_err(|e| format!("reading {path}: {e}")))
.collect::<Result<Vec<_>, _>>()?;
let input: Box<dyn std::io::BufRead> = if opts.file == "-" {
Box::new(std::io::stdin().lock())
} else {
let file =
std::fs::File::open(&opts.file).map_err(|e| format!("reading {}: {e}", opts.file))?;
Box::new(std::io::BufReader::new(file))
};
let checked = agentplane::grader_verdict::check(
input,
verifier,
&anchor.checkpoints,
&sidecars,
graders
.as_ref()
.map(|v| v as &dyn agentplane::core::Verifier),
)
.map_err(|e| e.to_string())?;
let report = &checked.export;
println!(
"{}",
serde_json::to_string_pretty(&VerifyDocument {
anchor: &anchor,
report,
grader_verdicts: &checked.sidecars,
})
.map_err(|e| e.to_string())?
);
if checked.any_refused() {
return Ok(ExitCode::from(exit::FINDING));
}
Ok(verify_status(report, anchor.split_view.is_empty()))
}
fn selection_of(cases: &[String], runs: &[String]) -> Result<agentplane::export::Selection, Fault> {
Ok(agentplane::export::Selection {
cases: cases
.iter()
.map(|c| {
agentplane::core::CaseId::parse(c).map_err(|e| usage(format!("--case {c}: {e}")))
})
.collect::<Result<_, _>>()?,
runs: runs
.iter()
.map(|r| {
agentplane::core::RunId::parse(r).map_err(|e| usage(format!("--run {r}: {e}")))
})
.collect::<Result<_, _>>()?,
})
}
const REGISTER_RUNG: &str = "read from the operator's disclosure register — an unchained row \
whoever administers the store can edit or delete, so an empty list does not show that \
nothing left the plane";
fn disclose_verb(opts: &ExportArgs) -> Result<ExitCode, Fault> {
let selection = selection_of(&opts.cases, &opts.runs)?;
let Some(to) = opts.to.as_deref().filter(|t| !t.trim().is_empty()) else {
return Err(usage(
"--to is required to disclose: an erasure names a copy by who received it".to_owned(),
));
};
let Some(actor) = opts.actor.as_deref() else {
return Err(usage(
"--actor is required to disclose: the register records who made the copy".to_owned(),
));
};
let by = agentplane::core::Operator::asserted(actor).map_err(|e| usage(e.to_string()))?;
let rt = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.map_err(|e| format!("could not start the async runtime: {e}"))?;
rt.block_on(async {
let backend = opts.store.at.open().await?;
#[allow(clippy::disallowed_methods)]
let at = time::OffsetDateTime::now_utc();
let request = agentplane::disclosure::Request {
selection,
recipient: to.to_owned(),
by,
at,
};
let staging = match &opts.output {
Some(_) => None,
None => Some(private_dir()?),
};
let destination = staging.as_ref().map_or_else(
|| std::path::PathBuf::from(opts.output.as_deref().unwrap_or_default()),
|dir| dir.join("package.jsonl"),
);
let disclosed = agentplane::disclosure::disclose(
&backend.journal(),
&backend.cases(),
backend.disclosures().as_ref(),
&request,
&destination,
)
.await;
let delivered = match disclosed {
Ok(act) => {
if staging.is_some() {
let bytes = std::fs::read(&destination).map_err(|e| e.to_string());
let written = bytes.and_then(|b| {
use std::io::Write as _;
std::io::stdout().lock().write_all(&b).map_err(|e| {
format!(
"disclosure {} is recorded and was not delivered: {e}",
act.id
)
})
});
written.map(|()| act).map_err(Fault::from)
} else {
Ok(act)
}
}
Err(agentplane::disclosure::DiscloseError::Write(e))
if matches!(
e.kind(),
std::io::ErrorKind::NotFound | std::io::ErrorKind::InvalidInput
) =>
{
Err(usage(e.to_string()))
}
Err(e) => Err(Fault::from(e.to_string())),
};
if let Some(dir) = &staging {
let _ = std::fs::remove_dir_all(dir);
}
let act = delivered?;
eprintln!(
"{}",
serde_json::to_string_pretty(&serde_json::json!({
"disclosure": act,
"register": REGISTER_RUNG,
}))
.map_err(|e| e.to_string())?
);
Ok(ExitCode::SUCCESS)
})
}
fn private_dir() -> Result<std::path::PathBuf, Fault> {
let dir = std::env::temp_dir().join(format!(
"agentplane-disclosure-{}-{}",
std::process::id(),
agentplane::core::RunId::generate()
));
let mut builder = std::fs::DirBuilder::new();
#[cfg(unix)]
std::os::unix::fs::DirBuilderExt::mode(&mut builder, 0o700);
builder
.create(&dir)
.map_err(|e| format!("staging {}: {e}", dir.display()))?;
Ok(dir)
}
fn disclosures_verb(opts: &DisclosuresArgs) -> Result<ExitCode, Fault> {
if opts.cases.is_empty() && opts.runs.is_empty() {
return Err(usage("name a matter with --case or --run".to_owned()));
}
let selection = selection_of(&opts.cases, &opts.runs)?;
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 acts = opts
.at
.open()
.await?
.disclosures()
.disclosures(&selection.cases, &selection.runs)
.await
.map_err(|e| e.to_string())?;
println!(
"{}",
serde_json::to_string_pretty(&serde_json::json!({
"disclosures": acts,
"register": REGISTER_RUNG,
}))
.map_err(|e| e.to_string())?
);
Ok(ExitCode::SUCCESS)
})
}
fn bind_verb(opts: &BindArgs) -> Result<ExitCode, Fault> {
let run = agentplane::core::RunId::parse(&opts.run)
.map_err(|e| usage(format!("--run {}: {e}", opts.run)))?;
let content =
std::fs::read(&opts.content).map_err(|e| format!("reading {}: {e}", opts.content))?;
let input: Box<dyn std::io::BufRead> = if opts.export == "-" {
Box::new(std::io::stdin().lock())
} else {
let file = std::fs::File::open(&opts.export)
.map_err(|e| format!("reading {}: {e}", opts.export))?;
Box::new(std::io::BufReader::new(file))
};
let sidecar = agentplane::grader_verdict::bind(input, run, opts.last_seq, content)
.map_err(|e| e.to_string())?;
let json = serde_json::to_string_pretty(&sidecar).map_err(|e| e.to_string())?;
std::fs::write(&opts.out, json + "\n").map_err(|e| format!("writing {}: {e}", opts.out))?;
println!("{}", sidecar.signing_digest().to_hex());
Ok(ExitCode::SUCCESS)
}
fn restore_unverifiable(input: impl std::io::BufRead) -> Option<String> {
let header = input.lines().next()?.ok()?;
agentplane::export::foreign_canon(&header)
}
fn verify_status(report: &agentplane::export::VerifyReport, one_view: bool) -> ExitCode {
ExitCode::from(if !report.findings.is_empty() || !one_view {
exit::FINDING
} else if report.unverifiable.is_some() {
exit::UNVERIFIABLE
} else if report.is_sound() {
exit::OK
} else {
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 = a2a_agent(manifests)?;
refuse_mcp_listener(opts)?;
let url = opts.url.as_deref().ok_or_else(|| {
usage(
"`serve` needs --url: the A2A endpoint callers reach this plane at, `/a2a` under \
its public address. 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 url = a2a_endpoint(url)?;
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(backend.plane(), 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);
}
builder = builder.policy(Arc::new(policy) as Arc<dyn agentplane::core::PolicyEngine>);
for m in manifests {
builder = builder.agent(agentplane::runtime::Agent::new(m));
}
if !opts.push_host.is_empty() {
builder = builder.push(backend.push());
}
let builder = with_submission_witnesses(builder, opts).map_err(usage)?;
let runtime = builder.try_build().map_err(|e| e.to_string())?;
let mcp = mcp_surface(&runtime, Arc::clone(&auth), manifests, opts)?;
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,
mcp,
operator_auth,
opts,
manifest,
url,
&backend,
)
.await?;
Ok(ExitCode::SUCCESS)
})
}
#[cfg(all(feature = "a2a-server", feature = "cedar"))]
#[allow(clippy::too_many_arguments)]
async fn serve_until_stopped(
runtime: &Arc<Runtime>,
server: agentplane::api::a2a::A2aServer,
mcp: Option<McpSurface>,
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?,
);
}
if let Some(mcp) = mcp {
background.push(spawn_mcp_surface(mcp, 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");
}
}))
}
#[cfg(all(feature = "a2a-server", feature = "cedar"))]
fn a2a_agent(manifests: &[Manifest]) -> Result<&Manifest, Fault> {
if let [only] = manifests {
return Ok(only);
}
let orchestrators: Vec<&Manifest> = manifests
.iter()
.filter(|m| {
m.spec
.topology
.as_ref()
.is_some_and(|t| t.role == agentplane::manifest::Role::Orchestrator)
})
.collect();
match orchestrators.as_slice() {
[desk] => Ok(desk),
found => Err(usage(format!(
"`serve` hosts a room on A2A through its one orchestrator, and this file holds \
{} agents of which {} declare `topology.role: orchestrator`. A2A's card path \
is well-known and singular — declare exactly one orchestrator, or split the file",
manifests.len(),
found.len()
))),
}
}
#[cfg(all(feature = "a2a-server", feature = "cedar"))]
fn refuse_mcp_listener(opts: &ServeArgs) -> Result<(), Fault> {
let Some(addr) = opts.mcp_addr.as_deref() else {
if opts.mcp_agent.is_empty() {
return Ok(());
}
return Err(usage(
"--mcp-agent names an agent to serve on the MCP listener, and no \
--mcp-addr opens one",
));
};
if !cfg!(feature = "mcp-server-http") {
return Err(usage(
"this build cannot serve MCP over HTTP: `--mcp-addr` needs the \
`mcp-server-http` feature. Reinstall with \
`--features cli,a2a-server,cedar,mcp-server-http`, or use the `:full` \
container image, which is built with it",
));
}
if addr == opts.addr || opts.operator_addr.as_deref() == Some(addr) {
return Err(usage(format!(
"--mcp-addr {addr} is another listener's address. Each surface has its own \
socket and its own action vocabulary — give MCP a port of its own"
)));
}
if !loopback_bind(addr) && opts.mcp_allowed_host.is_empty() {
return Err(usage(format!(
"--mcp-addr {addr} is not a loopback address and no --mcp-allowed-host names \
the host callers reach it by, so every request would be refused at its \
`Host` header. Name it, as `host` or `host:port`"
)));
}
Ok(())
}
#[cfg(all(feature = "a2a-server", feature = "cedar"))]
fn loopback_bind(addr: &str) -> bool {
addr.parse::<std::net::SocketAddr>().map_or_else(
|_| {
addr.rsplit_once(':')
.is_some_and(|(host, _)| host == "localhost")
},
|a| a.ip().is_loopback(),
)
}
#[cfg(all(feature = "a2a-server", feature = "cedar", feature = "mcp-server-http"))]
type McpSurface = (agentplane::tools::serve_http::McpHttp, String);
#[cfg(all(
feature = "a2a-server",
feature = "cedar",
not(feature = "mcp-server-http")
))]
type McpSurface = std::convert::Infallible;
#[cfg(all(feature = "a2a-server", feature = "cedar", feature = "mcp-server-http"))]
fn mcp_surface(
runtime: &Arc<Runtime>,
auth: Arc<dyn agentplane::api::Authenticator>,
manifests: &[Manifest],
opts: &ServeArgs,
) -> Result<Option<McpSurface>, String> {
use agentplane::tools::serve::McpServer;
use agentplane::tools::serve_http::{HttpConfig, McpHttp};
let Some(addr) = opts.mcp_addr.clone() else {
return Ok(None);
};
let chosen = mcp_served(manifests, &opts.mcp_agent)?;
let server = McpServer::new(Arc::clone(runtime), &chosen).map_err(|e| mcp_refusal(&e))?;
let mut config = HttpConfig::new();
for host in &opts.mcp_allowed_host {
config = config.allow_host(host.clone());
}
for origin in &opts.mcp_allowed_origin {
config = config.allow_origin(origin.clone());
}
let http = McpHttp::new(server, auth, &config).map_err(|e| format!("--mcp-addr: {e}"))?;
Ok(Some((http, addr)))
}
#[cfg(all(
feature = "a2a-server",
feature = "cedar",
not(feature = "mcp-server-http")
))]
#[allow(clippy::unnecessary_wraps)]
fn mcp_surface(
_runtime: &Arc<Runtime>,
_auth: Arc<dyn agentplane::api::Authenticator>,
_manifests: &[Manifest],
_opts: &ServeArgs,
) -> Result<Option<McpSurface>, String> {
Ok(None)
}
#[cfg(all(feature = "a2a-server", feature = "cedar", feature = "mcp-server-http"))]
fn mcp_served(manifests: &[Manifest], named: &[String]) -> Result<Vec<Manifest>, String> {
if named.is_empty() {
return Ok(manifests.to_vec());
}
if let Some(unknown) = named
.iter()
.find(|name| !manifests.iter().any(|m| &m.metadata.name == *name))
{
return Err(format!(
"--mcp-agent {unknown}: no agent in the file is named that"
));
}
Ok(manifests
.iter()
.filter(|m| named.contains(&m.metadata.name))
.cloned()
.collect())
}
#[cfg(all(feature = "a2a-server", feature = "cedar", feature = "mcp-server-http"))]
fn mcp_refusal(error: &agentplane::tools::serve::ServeError) -> String {
match error {
agentplane::tools::serve::ServeError::NoInputSchema { .. } => {
format!("--mcp-addr: {error} — `--mcp-agent NAME` serves only the agents it names")
}
_ => format!("--mcp-addr: {error}"),
}
}
#[cfg(all(feature = "a2a-server", feature = "cedar", feature = "mcp-server-http"))]
async fn spawn_mcp_surface((http, addr): McpSurface, stop: Stop) -> Result<Task, String> {
let listener = tokio::net::TcpListener::bind(&addr)
.await
.map_err(|e| format!("could not bind the MCP surface {addr}: {e}"))?;
eprintln!(
" mcp: http://{addr}{}",
agentplane::tools::serve_http::MCP_PATH
);
let router = http.router();
Ok(tokio::spawn(async move {
let served = axum::serve(listener, router)
.with_graceful_shutdown(async move {
stopping(stop).await;
http.close();
})
.await;
if let Err(error) = served {
tracing::error!(%error, "the MCP surface stopped");
}
}))
}
#[cfg(all(
feature = "a2a-server",
feature = "cedar",
not(feature = "mcp-server-http")
))]
#[allow(clippy::unused_async)]
async fn spawn_mcp_surface(surface: McpSurface, _stop: Stop) -> Result<Task, String> {
match surface {}
}
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(feature = "http")]
fn openapi_text() -> String {
serde_json::to_string_pretty(&agentplane::api::openapi::document())
.expect("a generated document serializes")
}
#[cfg(feature = "http")]
#[allow(clippy::unnecessary_wraps)]
fn openapi_verb() -> Result<ExitCode, Fault> {
println!("{}", openapi_text());
Ok(ExitCode::SUCCESS)
}
#[cfg(not(feature = "http"))]
#[allow(clippy::unnecessary_wraps)]
fn openapi_verb() -> Result<ExitCode, Fault> {
Err(usage(
"this build cannot describe the operator API: `openapi` needs the `http` feature. \
Reinstall with `--features cli,http`, or read the published document at \
https://hupe1980.github.io/agentplane/openapi.json",
))
}
#[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(),
}
}
async fn build_plane(
backend: &Backend,
manifests: &[Manifest],
mcp: &[String],
peers: &[String],
chain: Option<agentplane::core::Delegation>,
policy: Option<Arc<dyn agentplane::core::PolicyEngine>>,
streams: Option<Arc<dyn agentplane::runtime::RunStreamObserver>>,
) -> Result<Arc<Runtime>, Fault> {
let mut builder = with_providers(backend.plane(), manifests).await?;
if let Some(streams) = streams {
builder = builder.observe_model_streams(streams);
}
for (name, client) in connect_mcp_servers(mcp, manifests).await? {
builder = builder.tool_server(name, client);
}
if let Some((registry, client)) = connect_peers(peers, manifests).map_err(usage)? {
builder = builder.peers(registry, client);
}
if let Some(chain) = chain {
builder = builder.acting_as(chain);
}
if let Some(policy) = policy {
builder = builder.policy(policy);
}
for manifest in manifests {
builder = builder.agent(agentplane::runtime::Agent::new(manifest));
}
builder
.try_build()
.map_err(|e| Fault::from(in_cli_terms(&e, manifests)))
}
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 pinned = opts
.expect_digest
.as_deref()
.map(|hex| {
agentplane::core::Digest::from_hex(hex)
.map_err(|e| usage(format!("--expect-digest `{hex}` is not a digest: {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 backend = if let Some(backend) = opts.at.open().await? {
backend
} else {
eprintln!("note: journaling to memory; this run will not survive the process");
Backend::in_memory()?
};
let agent = build_plane(
&backend, manifests, &opts.mcp, &opts.peer, chain, None, None,
)
.await?;
let capability = entry_capability(manifests, opts.capability.as_deref()).map_err(usage)?;
let mut terms = agentplane::runtime::RunTerms::default().correlated(&capability, &keys);
if let Some(digest) = pinned {
terms = terms.expect_declaration(digest);
}
let agentplane::runtime::Admission::Fresh(outcome) = agent
.run_under(
&capability,
Tainted::trusted(opts.read_input().map_err(usage)?),
terms,
)
.await
.map_err(admission_fault)?
else {
return Err("an unkeyed run was answered as a keyed one"
.to_owned()
.into());
};
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(backend.plane(), 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<(Backend, 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, 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((
Backend::Embedded(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 (backend, runs) in sources {
for run in runs {
results.push(verify_one(manifests, &backend, run).await?);
}
}
Ok(ExitCode::from(replay_exit(&results)))
})
}
async fn verify_one(
manifests: &[Manifest],
backend: &Backend,
run: agentplane::core::RunId,
) -> Result<Replayed, Fault> {
match strict_verdict(manifests, backend, run).await? {
Ok(verdict) => {
eprint!("{verdict}");
Ok(Replayed::of(&verdict))
}
Err(e) => {
eprintln!("run {run} — cannot be read: {e}");
Ok(Replayed::Unreadable)
}
}
}
async fn strict_verdict(
manifests: &[Manifest],
backend: &Backend,
run: agentplane::core::RunId,
) -> Result<Result<agentplane::runtime::Verdict, String>, Fault> {
let history = match backend.journal().read(run, 1).await {
Ok(history) => history,
Err(e) => return Ok(Err(e.to_string())),
};
let mut builder = agentplane::runtime::replay_only::wire(backend.plane(), 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)))?;
Ok(plane.verify(run).await.map_err(|e| e.to_string()))
}
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" => {
let fake = agentplane::model::fake::FakeProvider::new();
fake.streaming();
Ok(fake)
}
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 {
#[test]
fn an_unbound_subject_exits_as_usage() {
let unbound = agentplane::core::RuntimeError::SubjectUnbound {
binding: "$input/customer/id".to_owned(),
reason: "it selects nothing in the run's input".to_owned(),
};
assert_eq!(super::admission_fault(unbound).status(), super::exit::USAGE);
assert_eq!(
super::admission_fault(agentplane::core::RuntimeError::Draining).status(),
super::exit::OPERATIONAL
);
}
#[test]
fn a_cosigned_note_file_names_its_witness_in_verify() {
let golden = std::path::Path::new(env!("CARGO_MANIFEST_DIR")).join("tests/golden");
let note = golden
.join("checkpoint.cosigned.note")
.display()
.to_string();
let keys = std::fs::read_to_string(golden.join("keys.txt")).expect("keys.txt");
let witness = keys
.lines()
.find_map(|l| l.strip_prefix("--witness-key "))
.expect("a witness key")
.to_owned();
let file = super::checkpoint_anchor(¬e, std::slice::from_ref(&witness))
.expect("the note reads");
assert!(file.signed_note);
assert_eq!(file.cosigned_by, vec![format!("{note}:golden-witness")]);
assert_eq!(
file.anchor
.witnessed
.iter()
.map(|t| (t.key_id.as_str(), t.timestamp))
.collect::<Vec<_>>(),
vec![("golden-witness", 1_700_000_600)],
"the witness's signed time travels with the anchor, for `audit`'s freshness rule"
);
let unkeyed = super::checkpoint_anchor(¬e, &[]).expect("the note reads");
assert!(unkeyed.cosigned_by.is_empty() && unkeyed.anchor.witnessed.is_empty());
let stranger = format!(
"golden-witness={}",
base64::Engine::encode(
&base64::engine::general_purpose::STANDARD,
agentplane::policy::Ed25519Signer::new("x", &[3u8; 32]).verifying_key()
)
);
let wrong = super::checkpoint_anchor(¬e, &[stranger]).expect("the note reads");
assert!(
wrong.cosigned_by.is_empty(),
"a line is a cosignature only under the key that made it"
);
}
use super::{
EXIT_STATUS_HELP, Fault, Truncation, audit_status, cutoff_before, declared_bound, exit,
lift_status, refuse_ambiguous_peers, refuses_partial_export, restore_unverifiable,
shell_quote, verify_status, where_flags, without_password,
};
use std::sync::Arc;
#[cfg(feature = "http")]
#[test]
fn the_openapi_verb_prints_the_published_document() {
let published = std::fs::read_to_string(concat!(
env!("CARGO_MANIFEST_DIR"),
"/site/static/openapi.json"
))
.expect("site/static/openapi.json exists");
assert_eq!(
super::openapi_text().trim_end(),
published.trim_end(),
"`agentplane openapi` and site/static/openapi.json disagree"
);
}
#[cfg(not(feature = "http"))]
#[test]
fn the_openapi_verb_names_the_feature_it_needs() {
match super::openapi_verb() {
Err(fault @ Fault::Usage(_)) => {
assert_eq!(fault.status(), super::exit::USAGE);
assert!(fault.to_string().contains("`http`"), "{fault}");
}
other => panic!("a build without `http` described the operator API: {other:?}"),
}
}
#[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));
let mut report = agentplane::export::VerifyReport {
checkpoint: agentplane::journal::Checkpoint {
origin: String::new(),
size: 0,
root: agentplane::core::Digest::ZERO,
},
sound: Vec::new(),
findings: Vec::new(),
not_checked: Vec::new(),
records: 0,
cases: 0,
complete: true,
unverifiable: None,
selection: None,
};
assert_eq!(verify_status(&report, true), ExitCode::SUCCESS);
report.unverifiable = Some("unknown canon".to_owned());
assert_eq!(
verify_status(&report, true),
ExitCode::from(exit::UNVERIFIABLE)
);
let version = agentplane::export::FORMAT_VERSION;
let foreign =
format!("{{\"kind\":\"agentplane.export\",\"version\":{version},\"canon\":999}}\n");
assert!(restore_unverifiable(foreign.as_bytes()).is_some());
let native = format!(
"{{\"kind\":\"agentplane.export\",\"version\":{version},\"canon\":{}}}\n",
agentplane::core::canon::VERSION
);
assert!(restore_unverifiable(native.as_bytes()).is_none());
let not_an_export = "{\"kind\":\"RunAdmitted\",\"canon\":999}\n";
assert!(restore_unverifiable(not_an_export.as_bytes()).is_none());
let other_version = format!(
"{{\"kind\":\"agentplane.export\",\"version\":{},\"canon\":999}}\n",
version + 1
);
assert!(restore_unverifiable(other_version.as_bytes()).is_none());
report.findings.push("damage".to_owned());
assert_eq!(verify_status(&report, true), 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"),
(exit::UNVERIFIABLE, "unverifiable"),
] {
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,
&super::Backend::Embedded(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);
});
}
#[cfg(all(feature = "a2a-server", feature = "cedar"))]
fn serve_witness_args(extra: &[&str]) -> super::ServeArgs {
let dir = std::env::temp_dir().join(format!(
"agentplane-logkey-{}",
agentplane::core::RunId::generate()
));
std::fs::create_dir_all(&dir).expect("scratch dir");
let seed = dir.join("log.seed");
std::fs::write(&seed, "07".repeat(32)).expect("seed");
let witness_key = format!(
"w={}",
base64::Engine::encode(
&base64::engine::general_purpose::STANDARD,
agentplane::policy::Ed25519Signer::new("w", &[9u8; 32]).verifying_key()
)
);
let log_key = format!("log.example/plane={}", seed.display());
let mut args = vec![
"agentplane",
"serve",
"agent.yaml",
"--witness-submit",
"http://127.0.0.1:9",
"--witness-key",
&witness_key,
"--log-key",
&log_key,
];
args.extend_from_slice(extra);
match <super::Cli as clap::Parser>::try_parse_from(args)
.expect("parses")
.verb
{
super::Verb::Serve(opts) => *opts,
_ => unreachable!("a serve line"),
}
}
#[test]
#[cfg(all(feature = "a2a-server", feature = "cedar"))]
fn serve_refuses_an_interval_shorter_than_its_sweep() {
let builder = || {
agentplane::runtime::Runtime::builder(std::sync::Arc::new(
agentplane::store::RedbStore::open_in_memory().expect("store"),
)
as std::sync::Arc<dyn agentplane::journal::JournalStore>)
};
let short = serve_witness_args(&["--sweep-every", "60", "--witness-interval", "30"]);
let refused = super::with_submission_witnesses(builder(), &short)
.expect_err("an interval the sweep cannot keep");
assert!(
refused.contains("30") && refused.contains("60"),
"the refusal does not name both values: {refused}"
);
let external = serve_witness_args(&["--sweep-every", "0", "--witness-interval", "30"]);
assert!(super::with_submission_witnesses(builder(), &external).is_ok());
}
#[test]
#[cfg(all(feature = "a2a-server", feature = "cedar"))]
fn serve_wires_its_submission_witnesses_into_the_sweep() {
use agentplane::journal::JournalStore;
const AGENT: &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 manifest = agentplane::manifest::Manifest::parse(AGENT).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"),
);
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(&manifest))
.build()
.run(
"support.summarise",
agentplane::core::Tainted::trusted(serde_json::json!({ "t": 1 })),
)
.await
.expect("a sealed run");
let plane = super::with_submission_witnesses(
agentplane::runtime::Runtime::builder(
std::sync::Arc::clone(&store) as std::sync::Arc<dyn JournalStore>
),
&serve_witness_args(&[]),
)
.expect("wired")
.build();
#[allow(clippy::disallowed_methods)]
let now = agentplane::core::Timestamp::now_utc();
let report = plane
.sweep(now, std::time::Duration::from_secs(60))
.await
.expect("a sweep");
assert_eq!(
report.witness_shortfall, 1,
"serve's witnesses never reached the sweep: {report:?}"
);
});
}
fn cli(args: &[&str]) -> Result<std::process::ExitCode, Fault> {
super::dispatch(<super::Cli as clap::Parser>::try_parse_from(args).expect("parses"))
}
#[test]
fn a_halt_thrown_at_the_terminal_refuses_a_run_on_the_same_store() {
let dir = std::env::temp_dir().join(format!(
"agentplane-halt-{}",
agentplane::core::RunId::generate()
));
std::fs::create_dir_all(&dir).expect("scratch dir");
let store = dir.join("plane.redb");
let store = store.to_str().expect("utf-8 path");
let manifest = concat!(env!("CARGO_MANIFEST_DIR"), "/examples/summariser.yaml");
let input = r#"{"ticket":"printer on fire"}"#;
cli(&[
"agentplane",
"run",
manifest,
"--input",
input,
"--store",
store,
])
.expect("an unhalted plane runs the example");
cli(&[
"agentplane",
"halt",
"--store",
store,
"--reason",
"incident 7",
"--actor",
"ops",
])
.expect("the halt is thrown");
let refused = cli(&[
"agentplane",
"run",
manifest,
"--input",
input,
"--store",
store,
])
.expect_err("a run under a standing tenant halt must be refused at admission");
assert!(
refused.to_string().contains("halt"),
"the refusal must name the halt: {refused}"
);
let _ = std::fs::remove_dir_all(&dir);
}
fn store_with_a_case(tag: &str) -> (std::path::PathBuf, String, String) {
use agentplane::case::CaseStore;
let dir = std::env::temp_dir().join(format!(
"agentplane-{tag}-{}",
agentplane::core::RunId::generate()
));
std::fs::create_dir_all(&dir).expect("scratch dir");
let path = dir.join("plane.redb");
let store = path.to_str().expect("utf-8 path").to_owned();
let case = {
let cases = agentplane::store::RedbStore::open(&store).expect("store");
tokio::runtime::Builder::new_current_thread()
.build()
.expect("runtime")
.block_on(
cases.correlate_or_open(
"matter",
&[agentplane::core::CorrelationKey::new("matter", "M-1")],
agentplane::core::Timestamp::from_unix_timestamp(1_700_000_000)
.expect("an instant"),
),
)
.expect("a case")
.case_id()
.to_string()
};
(dir, store, case)
}
fn lift_records(store: &str, outcome: &str) -> Vec<agentplane::journal::RecordKind> {
use agentplane::journal::JournalStore;
let journal = agentplane::store::RedbStore::open(store).expect("store");
tokio::runtime::Builder::new_current_thread()
.build()
.expect("runtime")
.block_on(async {
let mut kinds = Vec::new();
for run in journal.runs_by_outcome(outcome, 100).await.expect("index") {
let page = journal.read_page(run, 1, 1).await.expect("records");
kinds.extend(page.into_iter().map(|r| r.kind().clone()));
}
kinds
})
}
#[test]
fn a_terminal_lift_is_recorded_as_asserted() {
let (dir, store, case) = store_with_a_case("lift");
let store = store.as_str();
cli(&[
"agentplane",
"halt",
"--store",
store,
"--reason",
"incident 7",
"--actor",
"ops",
])
.expect("the halt is thrown");
cli(&[
"agentplane",
"hold",
"--store",
store,
"--case",
&case,
"--reason",
"order 9",
"--actor",
"ops",
])
.expect("the hold is placed");
assert_eq!(
cli(&[
"agentplane",
"halt",
"--store",
store,
"--lift",
"--actor",
"ops-erin"
])
.expect("the halt is lifted"),
std::process::ExitCode::SUCCESS
);
assert_eq!(
cli(&[
"agentplane",
"hold",
"--store",
store,
"--case",
&case,
"--lift",
"--actor",
"ops-erin",
])
.expect("the hold is released"),
std::process::ExitCode::SUCCESS
);
let lifts = lift_records(store, "halt-lifted");
assert_eq!(lifts.len(), 1, "{lifts:?}");
let agentplane::journal::RecordKind::HaltLifted { by, thrown_by, .. } = &lifts[0] else {
panic!("a lift run holds a lift record: {lifts:?}");
};
assert_eq!((by.actor(), by.basis().as_str()), ("ops-erin", "asserted"));
assert_eq!(thrown_by.actor(), "ops");
let releases = lift_records(store, "hold-released");
assert_eq!(releases.len(), 1, "{releases:?}");
let agentplane::journal::RecordKind::HoldReleased { by, .. } = &releases[0] else {
panic!("a release run holds a release record: {releases:?}");
};
assert_eq!((by.actor(), by.basis().as_str()), ("ops-erin", "asserted"));
the_histories_list(store);
let _ = std::fs::remove_dir_all(&dir);
}
fn the_histories_list(store: &str) {
for args in [
&["agentplane", "halt", "list", "--store", store, "--lifted"][..],
&[
"agentplane",
"halt",
"list",
"--store",
store,
"--lifted",
"--json",
][..],
&["agentplane", "hold", "list", "--store", store, "--released"][..],
] {
assert_eq!(
cli(args).expect("the history lists"),
std::process::ExitCode::SUCCESS,
"{args:?}"
);
}
for args in [
&[
"agentplane",
"halt",
"list",
"--store",
store,
"--lifted",
"--limit",
"0",
][..],
&[
"agentplane",
"hold",
"list",
"--store",
store,
"--released",
"--limit",
"0",
][..],
] {
assert_eq!(
cli(args).expect("the history lists"),
std::process::ExitCode::from(exit::PARTIAL),
"{args:?}"
);
}
}
#[test]
fn a_terminal_lift_without_an_actor_is_refused() {
use agentplane::case::CaseStore;
use agentplane::quota::QuotaStore;
let (dir, store, case) = store_with_a_case("unattributed");
let store = store.as_str();
cli(&[
"agentplane",
"halt",
"--store",
store,
"--reason",
"incident 7",
"--actor",
"ops",
])
.expect("the halt is thrown");
cli(&[
"agentplane",
"hold",
"--store",
store,
"--case",
&case,
"--reason",
"order 9",
"--actor",
"ops",
])
.expect("the hold is placed");
for args in [
&["agentplane", "halt", "--store", store, "--lift"][..],
&[
"agentplane",
"hold",
"--store",
store,
"--case",
&case,
"--lift",
][..],
] {
let refused = cli(args).expect_err("a lift nobody is named for");
assert!(
matches!(refused, Fault::Usage(_)) && refused.to_string().contains("--actor"),
"{args:?}: {refused}"
);
}
let plane = agentplane::store::RedbStore::open(store).expect("store");
tokio::runtime::Builder::new_current_thread()
.build()
.expect("runtime")
.block_on(async {
assert_eq!(
plane.halts().await.expect("halts").len(),
1,
"the halt stands"
);
let id = agentplane::core::CaseId::parse(&case).expect("case id");
assert!(
plane.hold(id).await.expect("hold").is_some(),
"the hold stands"
);
});
drop(plane);
assert_eq!(lift_records(store, "halt-lifted"), []);
assert_eq!(lift_records(store, "hold-released"), []);
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn content_check_exits_on_what_a_rule_would_refuse() {
let dir = std::env::temp_dir().join(format!(
"agentplane-content-{}",
agentplane::core::RunId::generate()
));
std::fs::create_dir_all(&dir).expect("scratch dir");
let manifest = dir.join("agent.yaml");
std::fs::write(
&manifest,
"apiVersion: agentplane.hupe1980.github.io/v1alpha1
kind: Agent
metadata: { name: ruled, version: \"1.0.0\" }
spec:
capabilities: { provides: [work.do] }
budgets: {}
security:
content:
rules:
- id: codename
match: {contains: [falcon], case: fold}
at: {sinks: [model.complete]}
then: refuse
",
)
.expect("manifest");
let check = |value: &str, at: &str| {
let path = dir.join("value.json");
std::fs::write(&path, value).expect("value");
cli(&[
"agentplane",
"content",
"check",
manifest.to_str().expect("utf-8"),
"--at",
at,
"--value",
path.to_str().expect("utf-8"),
])
};
assert_eq!(
check(r#"{"q": "Falcon"}"#, "sink:model.complete").expect("judged"),
std::process::ExitCode::from(exit::FINDING)
);
assert_eq!(
check(r#"{"q": "Falcon"}"#, "sink:tool.call").expect("judged"),
std::process::ExitCode::SUCCESS,
"a rule applies only where it is declared"
);
assert!(matches!(
check(r#"{"q": "x"}"#, "outbound"),
Err(Fault::Usage(_))
));
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn a_listing_cut_short_by_its_limit_exits_partial() {
use agentplane::case::TaskStore;
let dir = std::env::temp_dir().join(format!(
"agentplane-listing-{}",
agentplane::core::RunId::generate()
));
std::fs::create_dir_all(&dir).expect("scratch dir");
let path = dir.join("plane.redb");
let store = path.to_str().expect("utf-8 path");
{
let tasks = agentplane::store::RedbStore::open(store).expect("store");
tokio::runtime::Builder::new_current_thread()
.build()
.expect("runtime")
.block_on(async {
for summary in ["first", "second"] {
tasks
.open(&listed(summary, serde_json::json!({}), None))
.await
.expect("opened");
}
});
}
let status = |limit: &str| {
cli(&["agentplane", "tasks", "--store", store, "--limit", limit]).expect("listed")
};
assert_eq!(
status("1"),
std::process::ExitCode::from(exit::PARTIAL),
"a page of one over two tasks is a partial answer"
);
assert_eq!(status("2"), std::process::ExitCode::SUCCESS);
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn a_stale_digest_at_the_terminal_is_refused() {
use agentplane::case::TaskStore;
let dir = std::env::temp_dir().join(format!(
"agentplane-stale-{}",
agentplane::core::RunId::generate()
));
std::fs::create_dir_all(&dir).expect("scratch dir");
let path = dir.join("plane.redb");
let store = path.to_str().expect("utf-8 path");
let task = listed("Refund the invoice", serde_json::json!({}), None);
let rt = tokio::runtime::Builder::new_current_thread()
.build()
.expect("runtime");
{
let tasks = agentplane::store::RedbStore::open(store).expect("store");
rt.block_on(tasks.open(&task)).expect("opened");
}
let id = task.id.to_string();
let decide = |digest: &str| {
cli(&[
"agentplane",
"decide",
&id,
"reject",
"--reason",
"no",
"--actor",
"ops",
"--store",
store,
"--digest",
digest,
])
};
let stale = agentplane::core::Digest::of(b"another version").to_hex();
assert_eq!(
decide(&stale).expect("a refusal, not a fault"),
std::process::ExitCode::from(exit::FINDING)
);
let tasks = agentplane::store::RedbStore::open(store).expect("store");
let after = rt.block_on(tasks.task(task.id)).unwrap().unwrap();
assert!(
after.state.is_pending() && after.assignee.is_none(),
"{after:?}"
);
drop(tasks);
assert!(matches!(decide("not-hex"), Err(Fault::Usage(_))));
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn every_plane_this_binary_builds_comes_through_backend_plane() {
let source = include_str!("agentplane.rs");
let body = source
.split("#[cfg(test)]\nmod tests {")
.next()
.expect("the file has a body");
assert_eq!(
body.matches("Runtime::builder_with(").count(),
1,
"only Backend::plane may start a runtime builder"
);
for other in ["Runtime::builder(", "Runtime::builder_on("] {
assert_eq!(
body.matches(other).count(),
0,
"only Backend::plane may start a runtime builder, and {other} does"
);
}
}
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_eq!(
shown["digest"],
task.justification.digest().to_hex(),
"{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}"
);
}
#[cfg(all(feature = "a2a-server", feature = "cedar", feature = "mcp-server-http"))]
#[test]
fn mcp_agent_serves_only_the_agents_it_names() {
let agent = |name: &str| {
agentplane::manifest::Manifest::parse(&format!(
"apiVersion: agentplane.hupe1980.github.io/v1alpha1\n\
kind: Agent\n\
metadata: {{ name: {name}, version: \"1.0.0\" }}\n\
spec:\n identity: {{ role: r }}\n capabilities: {{ provides: [{name}.do] }}\n budgets: {{}}\n"
))
.expect("manifest")
};
let file = [agent("triage"), agent("specialist")];
let all = super::mcp_served(&file, &[]).expect("every agent");
assert_eq!(all.len(), 2, "no --mcp-agent must serve every agent");
let one = super::mcp_served(&file, &["triage".to_owned()]).expect("one agent");
assert_eq!(
one.iter()
.map(|m| m.metadata.name.as_str())
.collect::<Vec<_>>(),
["triage"],
"--mcp-agent served an agent it did not name"
);
let unknown = super::mcp_served(&file, &["nobody".to_owned()]);
assert!(
unknown.is_err_and(|e| e.contains("nobody")),
"a name matching no agent was not refused"
);
let refusal = super::mcp_refusal(&agentplane::tools::serve::ServeError::NoInputSchema {
agent: "specialist".into(),
capability: "specialist.do".into(),
});
assert!(
refusal.contains("--mcp-agent"),
"the refusal does not name what this command can do: {refusal}"
);
}
#[cfg(all(feature = "a2a-server", feature = "cedar", feature = "mcp-server-http"))]
fn served_dir(name: &str) -> std::path::PathBuf {
std::env::temp_dir().join(format!(
"agentplane-served-{name}-{}",
agentplane::core::RunId::generate()
))
}
#[cfg(all(feature = "a2a-server", feature = "cedar", feature = "mcp-server-http"))]
fn tokens_in(file: &str) -> Vec<(String, String)> {
let entries: Vec<serde_json::Value> = serde_yaml_ng::from_str(file).expect("yaml");
entries
.iter()
.map(|e| {
(
e["actor"].as_str().expect("actor").to_owned(),
e["token"].as_str().expect("token").to_owned(),
)
})
.collect()
}
#[cfg(all(feature = "a2a-server", feature = "cedar", feature = "mcp-server-http"))]
#[test]
fn init_serve_writes_three_distinct_accepted_tokens() {
let dir = served_dir("tokens");
let served = super::init_serve(&dir).expect("an empty directory");
let again = super::init_serve(&served_dir("again")).expect("a second plane");
for name in super::SERVED_FILES {
assert!(dir.join(name).is_file(), "{name} was not written");
}
let file = std::fs::read_to_string(dir.join("tokens.yaml")).unwrap();
agentplane::api::tokens::TokenAuthenticator::from_yaml(&file).expect("serve accepts it");
let tokens = tokens_in(&file);
let actors: Vec<&str> = tokens.iter().map(|(a, _)| a.as_str()).collect();
assert_eq!(actors, ["peer-1", "framework-1", "ops-1"]);
assert!(tokens.iter().all(|(_, t)| t.len() >= 64), "{tokens:?}");
let other = std::fs::read_to_string(
again
.paths
.iter()
.find(|p| p.ends_with("tokens.yaml"))
.unwrap(),
)
.unwrap();
let distinct: std::collections::BTreeSet<String> = tokens
.iter()
.chain(&tokens_in(&other))
.map(|(_, t)| t.clone())
.collect();
assert_eq!(distinct.len(), 6, "a token repeats across two planes");
assert_eq!(
std::fs::read_to_string(dir.join("framework.token")).unwrap(),
tokens[1].1,
"framework.token is not the framework caller's token"
);
#[cfg(unix)]
for name in [
"tokens.yaml",
"framework.token",
"postgres.password",
"store.env",
] {
use std::os::unix::fs::PermissionsExt as _;
let mode = std::fs::metadata(dir.join(name))
.unwrap()
.permissions()
.mode();
assert_eq!(mode & 0o777, 0o600, "{name} is mode {mode:o}");
}
for json in [false, true] {
let printed = super::served_report(&served, json);
assert!(printed.contains(&served.digest) && printed.contains("compose.yaml"));
for (_, token) in &tokens {
assert!(!printed.contains(token.as_str()), "a token was printed");
}
}
}
#[cfg(all(feature = "a2a-server", feature = "cedar", feature = "mcp-server-http"))]
#[test]
fn init_serve_refuses_and_writes_nothing_when_a_file_exists() {
for name in super::SERVED_FILES {
let dir = served_dir("refused");
std::fs::create_dir_all(&dir).unwrap();
std::fs::write(dir.join(name), "reviewed\n").unwrap();
let err = super::init_serve(&dir).expect_err("an existing file");
assert!(err.to_string().contains(name), "{err}");
assert!(err.to_string().contains("so nothing was written"), "{err}");
assert_eq!(
std::fs::read_to_string(dir.join(name)).unwrap(),
"reviewed\n"
);
let present: Vec<String> = std::fs::read_dir(&dir)
.unwrap()
.map(|e| e.unwrap().file_name().to_string_lossy().into_owned())
.collect();
assert_eq!(present, [name], "a refused init --serve wrote {present:?}");
}
}
#[cfg(all(feature = "a2a-server", feature = "cedar", feature = "mcp-server-http"))]
#[test]
fn init_serve_writes_the_shipped_policy() {
let dir = served_dir("policy");
super::init_serve(&dir).expect("an empty directory");
let root = std::path::Path::new(env!("CARGO_MANIFEST_DIR"));
assert_eq!(
std::fs::read(dir.join("policy.cedar")).unwrap(),
std::fs::read(root.join("examples/serve-policy.cedar")).unwrap(),
"init --serve wrote a policy that is not the shipped bundle"
);
let manifest = std::fs::read_to_string(dir.join("agent.yaml")).unwrap();
let parsed = agentplane::manifest::Manifest::parse_all(&manifest).expect("it parses");
super::mcp_served(&parsed, &[]).expect("the MCP listener serves it");
let compose = std::fs::read_to_string(dir.join("compose.yaml")).unwrap();
assert!(
compose.contains(&format!("agentplane:{}-full", env!("CARGO_PKG_VERSION"))),
"the compose file does not run this version"
);
assert!(!compose.contains("PLANE_USER") && !compose.contains("AGENTPLANE_VERSION"));
}
#[cfg(all(feature = "a2a-server", feature = "cedar", feature = "mcp-server-http"))]
#[test]
fn init_serve_gives_postgres_a_password_off_the_command_line() {
let dir = served_dir("password");
super::init_serve(&dir).expect("an empty directory");
let password = std::fs::read_to_string(dir.join("postgres.password")).unwrap();
let password = password.trim();
assert!(password.len() >= 64, "{password:?}");
let store = std::fs::read_to_string(dir.join("store.env")).unwrap();
assert!(
store.contains(&format!(
"AGENTPLANE_STORE=postgres://agentplane:{password}@postgres:5432/"
)),
"store.env does not connect with the generated password: {store}"
);
let compose = std::fs::read_to_string(dir.join("compose.yaml")).unwrap();
let yaml: serde_json::Value = serde_yaml_ng::from_str(&compose).expect("yaml");
let postgres = &yaml["services"]["postgres"];
assert_eq!(
postgres["environment"]["POSTGRES_PASSWORD_FILE"], "/run/secrets/postgres_password",
"{postgres}"
);
assert!(
postgres["environment"]["POSTGRES_HOST_AUTH_METHOD"].is_null(),
"{postgres}"
);
assert_eq!(
yaml["secrets"]["postgres_password"]["file"], "./postgres.password",
"{yaml}"
);
let plane = &yaml["services"]["plane"];
assert_eq!(plane["env_file"][0], "./store.env", "{plane}");
let command = plane["command"].to_string();
assert!(
!command.contains("--store") && !command.contains("postgres://"),
"{command}"
);
assert!(
!compose.contains(password),
"the compose file holds the password"
);
for port in plane["ports"].as_array().expect("published ports") {
let port = port.as_str().expect("a short-syntax port");
assert!(
port.starts_with("127.0.0.1:"),
"{port} is published beyond loopback"
);
}
}
#[cfg(all(feature = "a2a-server", feature = "cedar", feature = "mcp-server-http"))]
#[test]
fn init_serve_runs_the_plane_as_the_token_files_owner_and_never_root() {
let err = super::plane_user(0, 0).expect_err("root");
assert!(err.to_string().contains("--user"), "{err}");
assert_eq!(super::plane_user(1000, 1000).unwrap(), "1000:1000");
#[cfg(unix)]
{
use std::os::unix::fs::MetadataExt as _;
let dir = served_dir("owner");
super::init_serve(&dir).expect("an empty directory");
let meta = std::fs::metadata(dir.join("tokens.yaml")).unwrap();
let compose = std::fs::read_to_string(dir.join("compose.yaml")).unwrap();
let user = format!("user: \"{}:{}\"", meta.uid(), meta.gid());
assert!(
compose.contains(&user),
"the compose file does not hold {user}"
);
}
}
#[cfg(all(feature = "a2a-server", feature = "cedar", feature = "mcp-server-http"))]
#[test]
fn a_failed_init_serve_removes_what_it_wrote() {
let dir = served_dir("partial");
std::fs::create_dir_all(&dir).unwrap();
{
let mut written = super::Written::default();
written
.create(&dir, "agent.yaml", "a\n", false)
.expect("the first file");
written
.create(&dir, "tokens.yaml", "t\n", true)
.expect("the second file");
written
.create(&dir, "no-such-dir/compose.yaml", "c\n", false)
.expect_err("a file under a missing directory");
}
let present: Vec<_> = std::fs::read_dir(&dir).unwrap().collect();
assert!(present.is_empty(), "a failed init --serve left {present:?}");
super::init_serve(&dir).expect("a re-run after the failure");
}
#[cfg(all(feature = "a2a-server", feature = "cedar"))]
#[test]
fn serve_refuses_a_url_that_is_not_the_a2a_endpoint() {
for (bare, fix) in [
("https://agent.example.com", "https://agent.example.com/a2a"),
(
"https://agent.example.com/",
"https://agent.example.com/a2a",
),
("http://h:8080/a2a/", "http://h:8080/a2a"),
] {
let err = super::a2a_endpoint(bare).expect_err(bare);
assert!(err.to_string().ends_with(&format!("--url {fix}")), "{err}");
}
assert_eq!(
super::a2a_endpoint("http://localhost:8080/a2a").unwrap(),
"http://localhost:8080/a2a"
);
}
fn temp_dir(name: &str) -> std::path::PathBuf {
let dir = std::env::temp_dir().join(format!(
"agentplane-{name}-{}",
agentplane::core::RunId::generate()
));
std::fs::create_dir_all(&dir).expect("scratch dir");
dir
}
fn succeeded_in(store: &str) -> Vec<agentplane::core::RunId> {
let journal = agentplane::store::RedbStore::open(store).expect("store");
tokio::runtime::Builder::new_current_thread()
.build()
.expect("runtime")
.block_on(agentplane::journal::JournalStore::runs_by_outcome(
&journal,
"succeeded",
10,
))
.expect("listed")
}
#[test]
fn history_prints_a_hostile_record_escaped() {
let dir = temp_dir("history");
let path = dir.join("plane.redb");
let store = path.to_str().expect("utf-8 path");
let manifest = concat!(env!("CARGO_MANIFEST_DIR"), "/examples/summariser.yaml");
let hostile =
"\u{1b}]52;c;ZXZpbA==\u{7}\u{9b}31m<img src=x onerror=alert(1)>\u{202E}gnp.exe";
let input = serde_json::json!({ "ticket": hostile }).to_string();
cli(&[
"agentplane",
"run",
manifest,
"--input",
&input,
"--store",
store,
])
.expect("the example runs");
let run = succeeded_in(store)[0];
let journal: Arc<dyn agentplane::journal::JournalStore> =
Arc::new(agentplane::store::RedbStore::open(store).expect("store"));
let rt = tokio::runtime::Builder::new_current_thread()
.build()
.expect("runtime");
let lines = rt
.block_on(super::history_lines(&journal, run, None, false))
.expect("read")
.expect("the run exists");
assert!(lines.len() >= 2, "{lines:?}");
let seqs: Vec<u64> = lines
.iter()
.map(|l| l.split_whitespace().next().unwrap().parse().unwrap())
.collect();
assert!(seqs.windows(2).all(|w| w[1] == w[0] + 1), "{seqs:?}");
let printed = lines.join("\n");
assert!(
printed.contains("\\u{202E}gnp.exe"),
"the bidi override was not shown escaped: {printed}"
);
for raw in ['\u{1b}', '\u{7}', '\u{9b}', '\u{202E}'] {
assert!(
!printed.contains(raw),
"U+{:04X} reached the terminal raw",
raw as u32
);
}
let later = rt
.block_on(super::history_lines(&journal, run, Some(2), false))
.expect("read")
.expect("the run exists");
assert!(later[0].trim_start().starts_with("2 "), "{later:?}");
let unknown = agentplane::core::RunId::generate();
for from in [None, Some(2)] {
assert!(
rt.block_on(super::history_lines(&journal, unknown, from, false))
.expect("read")
.is_none(),
"an unknown run read from {from:?} was answered as one that exists"
);
}
assert!(
rt.block_on(super::history_lines(&journal, run, Some(1_000), false))
.expect("read")
.is_some_and(|lines| lines.is_empty()),
"a reader past the end of a run that exists was told it does not"
);
drop(journal);
assert_eq!(
cli(&[
"agentplane",
"history",
&unknown.to_string(),
"--store",
store
])
.expect("an answer"),
std::process::ExitCode::from(exit::FINDING)
);
let _ = std::fs::remove_dir_all(&dir);
}
#[cfg(feature = "dev")]
#[test]
fn dev_refuses_a_store_it_did_not_create() {
let refused = |scratch: Option<&str>, tenant: Option<&str>| {
matches!(super::dev_store(scratch, tenant), Err(Fault::Usage(_)))
};
assert!(refused(Some("postgres://ops@db/plane"), None));
assert!(refused(None, Some("acme")));
let deployed = temp_dir("dev-deployed");
let file = deployed.join("dev.redb");
drop(agentplane::store::RedbStore::open(&file).expect("a deployment's store"));
assert!(refused(file.to_str(), None), "a redb file was opened");
assert!(
refused(deployed.to_str(), None),
"an unmarked directory was opened"
);
assert_eq!(
std::fs::read_dir(&deployed).unwrap().count(),
1,
"a refused directory was written to"
);
let scratch = temp_dir("dev-scratch");
let scratch = scratch.to_str().expect("utf-8 path");
drop(super::dev_store(Some(scratch), Some("dev")).expect("an empty directory opens"));
drop(super::dev_store(Some(scratch), None).expect("a marked directory reopens"));
#[cfg(unix)]
{
let linked = temp_dir("dev-linked");
let link = linked.join("scratch");
std::os::unix::fs::symlink(scratch, &link).expect("a link");
assert!(
refused(link.to_str(), None),
"a symbolic link to a marked directory was followed"
);
let store = std::path::Path::new(scratch).join("dev.redb");
std::fs::remove_file(&store).expect("the dev store");
std::os::unix::fs::symlink(&file, &store).expect("a link");
assert!(
refused(Some(scratch), None),
"a dev.redb linking to a deployment's store was opened"
);
std::fs::remove_file(&store).expect("the link");
let _ = std::fs::remove_dir_all(&linked);
}
let marker = std::path::Path::new(scratch).join(".agentplane-dev");
std::fs::write(&marker, "").expect("an empty marker");
assert!(
refused(Some(scratch), None),
"a marker without the exact bytes dev writes was accepted"
);
std::fs::write(&marker, super::DEV_MARKER_TEXT).expect("the marker");
drop(super::dev_store(Some(scratch), None).expect("a marked directory reopens"));
std::fs::write(std::path::Path::new(scratch).join("other"), "x").unwrap();
assert!(
refused(Some(scratch), None),
"a marked directory holding more was opened"
);
let _ = std::fs::remove_dir_all(&deployed);
let _ = std::fs::remove_dir_all(scratch);
}
#[cfg(feature = "dev")]
#[test]
fn dev_refuses_live_transports_without_consent() {
let args = |extra: &[&str]| {
let mut line = vec!["agentplane", "dev", "agent.yaml"];
line.extend_from_slice(extra);
match <super::Cli as clap::Parser>::try_parse_from(line)
.expect("parses")
.verb
{
super::Verb::Dev(args) => args,
other => panic!("parsed as {other:?}"),
}
};
for extra in [
&["--mcp", "files=mcp-files"][..],
&[
"--peer",
"billing=https://billing.example/a2a",
"--acting-as",
"ada",
][..],
] {
let refused = super::refuse_live_without_consent(&args(extra))
.expect_err("a live transport without consent");
assert!(refused.to_string().contains("--allow-live"), "{refused}");
let mut consented = extra.to_vec();
consented.push("--allow-live");
super::refuse_live_without_consent(&args(&consented)).expect("consented");
}
super::refuse_live_without_consent(&args(&[])).expect("nothing live");
}
#[cfg(feature = "dev")]
#[test]
fn the_dev_listener_binds_loopback_only() {
let rt = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.expect("runtime");
let listener = rt.block_on(super::bind_dev(0)).expect("bound");
let addr = listener.local_addr().expect("an address");
assert!(addr.ip().is_loopback(), "the dev page listens on {addr}");
assert_ne!(addr.port(), 0);
}
#[cfg(feature = "dev")]
const DEV_PORT: u16 = 47_312;
#[cfg(feature = "dev")]
const DEV_TOKEN: &str = "fedcba9876543210fedcba9876543210fedcba9876543210fedcba9876543210";
#[cfg(feature = "dev")]
async fn dev_session(
manifest: &str,
) -> (axum::Router, Arc<dyn agentplane::journal::JournalStore>) {
dev_session_as(manifest, None).await
}
#[cfg(feature = "dev")]
async fn dev_session_as(
manifest: &str,
acting_as: Option<&str>,
) -> (axum::Router, Arc<dyn agentplane::journal::JournalStore>) {
let backend = super::dev_store(None, None).expect("memory");
let journal = backend.journal();
let manifests = super::manifests_at(manifest).expect("the manifest");
let bench = super::dev::Bench::start(
backend,
super::dev::Wiring {
file: manifest.to_owned(),
mcp: Vec::new(),
peer: Vec::new(),
acting_as: acting_as.map(str::to_owned),
policy: Arc::new(agentplane::api::dev::DevPolicy::new("dev:author")),
},
manifests,
)
.await
.expect("the plane builds");
let auth = agentplane::api::tokens::TokenAuthenticator::new(vec![
agentplane::api::tokens::TokenEntry {
token: DEV_TOKEN.to_owned(),
actor: "dev:author".to_owned(),
roles: Vec::new(),
tenant: Some("dev".to_owned()),
scope: None,
not_after: None,
},
])
.expect("tokens");
(
agentplane::api::dev::router(Arc::new(bench), Arc::new(auth), DEV_PORT),
journal,
)
}
#[cfg(feature = "dev")]
async fn page(
router: &axum::Router,
method: &str,
path: &str,
body: Option<serde_json::Value>,
) -> (u16, serde_json::Value) {
use tower::ServiceExt as _;
let request = axum::http::Request::builder()
.method(method)
.uri(path)
.header("host", format!("127.0.0.1:{DEV_PORT}"))
.header("origin", format!("http://127.0.0.1:{DEV_PORT}"))
.header("authorization", format!("Bearer {DEV_TOKEN}"))
.header("content-type", "application/json")
.body(axum::body::Body::from(
body.map(|b| b.to_string()).unwrap_or_default(),
))
.expect("a request");
let response = router.clone().oneshot(request).await.expect("an answer");
let status = response.status().as_u16();
let bytes = axum::body::to_bytes(response.into_body(), usize::MAX)
.await
.expect("a body");
(
status,
serde_json::from_slice(&bytes).unwrap_or(serde_json::Value::Null),
)
}
#[cfg(feature = "dev")]
async fn served_history(router: &axum::Router, run: &str) -> Vec<serde_json::Value> {
let mut records = Vec::new();
let mut from = 1;
loop {
let (status, page_) = page(
router,
"GET",
&format!("/api/runs/{run}/history?from={from}"),
None,
)
.await;
assert_eq!(status, 200, "{page_}");
for record in page_["records"].as_array().expect("records") {
records.push(record.clone());
}
match page_["next_from"].as_u64() {
Some(next) => from = next,
None => return records,
}
}
}
#[cfg(feature = "dev")]
#[test]
fn history_json_is_the_history_routes_records() {
let manifest = concat!(env!("CARGO_MANIFEST_DIR"), "/examples/summariser.yaml");
tokio::runtime::Builder::new_multi_thread()
.enable_all()
.build()
.expect("runtime")
.block_on(async {
let (router, journal) = dev_session(manifest).await;
let (status, started) = page(
&router,
"POST",
"/dev/runs",
Some(serde_json::json!({ "input": { "ticket": "printer\u{202E} on fire" } })),
)
.await;
assert_eq!(status, 200, "{started}");
let run_text = started["run"].as_str().expect("a run").to_owned();
let run = agentplane::core::RunId::parse(&run_text).expect("a run id");
let served = served_history(&router, &run_text).await;
let printed: Vec<serde_json::Value> =
super::history_lines(&journal, run, None, true)
.await
.expect("read")
.expect("the run exists")
.iter()
.map(|l| serde_json::from_str(l).expect("one JSON object per line"))
.collect();
assert_ne!(printed, Vec::<serde_json::Value>::new());
assert_eq!(printed, served);
});
}
#[cfg(feature = "dev")]
#[test]
fn the_dev_timeline_escapes_hidden_characters() {
let manifest = concat!(env!("CARGO_MANIFEST_DIR"), "/examples/summariser.yaml");
tokio::runtime::Builder::new_multi_thread()
.enable_all()
.build()
.expect("runtime")
.block_on(async {
let (router, _) = dev_session(manifest).await;
let (_, started) = page(
&router,
"POST",
"/dev/runs",
Some(serde_json::json!({ "input": { "printer\u{202E}": "on\u{200B} fire" } })),
)
.await;
let run = started["run"].as_str().expect("a run").to_owned();
let (status, shown) =
page(&router, "GET", &format!("/dev/runs/{run}/history"), None).await;
assert_eq!(status, 200, "{shown}");
let text = shown["records"].to_string();
assert!(
!text.contains('\u{202E}') && !text.contains('\u{200B}'),
"{text}"
);
assert!(
text.contains("printer\\\\u{202E}"),
"a key was not escaped: {text}"
);
assert!(text.contains("on\\\\u{200B} fire"), "{text}");
assert_eq!(shown["escaped"], true);
let served = served_history(&router, &run).await;
assert_eq!(
shown["records"].as_array().expect("records").len(),
served.len()
);
let (status, _) = page(&router, "GET", "/dev/runs/not-a-run/history", None).await;
assert_eq!(status, 400);
});
}
#[cfg(feature = "dev")]
async fn governed_by(
journal: &Arc<dyn agentplane::journal::JournalStore>,
run: agentplane::core::RunId,
) -> serde_json::Value {
super::history_lines(journal, run, None, true)
.await
.expect("read")
.expect("the run exists")
.iter()
.map(|l| serde_json::from_str::<serde_json::Value>(l).expect("JSON"))
.find(|r| r["kind"] == "RunAdmitted")
.expect("an admission")["record"]["governed_by"]
.clone()
}
#[cfg(feature = "dev")]
#[test]
fn dev_and_run_admit_the_same_declaration() {
let manifest = concat!(env!("CARGO_MANIFEST_DIR"), "/examples/summariser.yaml");
let dir = temp_dir("dev-same");
let path = dir.join("plane.redb");
let store = path.to_str().expect("utf-8 path");
cli(&[
"agentplane",
"run",
manifest,
"--input",
"{\"ticket\":\"t\"}",
"--store",
store,
])
.expect("run");
let ran = succeeded_in(store)[0];
let rt = tokio::runtime::Builder::new_multi_thread()
.enable_all()
.build()
.expect("runtime");
let by_run = rt.block_on(async {
let journal: Arc<dyn agentplane::journal::JournalStore> =
Arc::new(agentplane::store::RedbStore::open(store).expect("store"));
governed_by(&journal, ran).await
});
let by_dev = rt.block_on(async {
let (router, journal) = dev_session(manifest).await;
let (_, started) = page(
&router,
"POST",
"/dev/runs",
Some(serde_json::json!({ "input": { "ticket": "t" } })),
)
.await;
let run = agentplane::core::RunId::parse(started["run"].as_str().expect("a run"))
.expect("a run id");
governed_by(&journal, run).await
});
assert!(by_run.is_object(), "{by_run}");
assert_eq!(by_run, by_dev);
let _ = std::fs::remove_dir_all(&dir);
}
#[cfg(feature = "dev")]
#[test]
fn a_rebuild_scopes_the_acting_as_chain_to_the_saved_file() {
let dir = temp_dir("dev-chain");
let file = dir.join("agent.yaml");
let original = std::fs::read_to_string(concat!(
env!("CARGO_MANIFEST_DIR"),
"/examples/summariser.yaml"
))
.expect("the example");
std::fs::write(&file, &original).expect("written");
let manifest = file.to_str().expect("utf-8 path").to_owned();
tokio::runtime::Builder::new_multi_thread()
.enable_all()
.build()
.expect("runtime")
.block_on(async {
let (router, _) = dev_session_as(&manifest, Some("ada")).await;
let start = |capability: &str| {
serde_json::json!({
"input": { "ticket": "t" },
"capability": capability,
})
};
let (status, ran) = page(
&router,
"POST",
"/dev/runs",
Some(start("support.summarise")),
)
.await;
assert_eq!(status, 200, "{ran}");
assert_eq!(ran["status"], "succeeded", "{ran}");
let triager = original
.replace("name: summariser", "name: triager")
.replace("support.summarise", "support.triage");
std::fs::write(&file, format!("{original}---\n{triager}")).expect("edited");
let (status, declared) = page(&router, "GET", "/dev/manifest", None).await;
assert_eq!(status, 200, "{declared}");
assert!(declared["refused"].is_null(), "{declared}");
let (status, ran) =
page(&router, "POST", "/dev/runs", Some(start("support.triage"))).await;
assert_eq!(status, 200, "{ran}");
assert_eq!(ran["status"], "succeeded", "{ran}");
});
let _ = std::fs::remove_dir_all(&dir);
}
#[cfg(feature = "dev")]
#[test]
#[allow(clippy::too_many_lines)]
fn a_run_started_on_the_dev_page_is_decided_replayed_and_exported() {
let dir = temp_dir("dev-loop");
let file = dir.join("agent.yaml");
let original = std::fs::read_to_string(concat!(
env!("CARGO_MANIFEST_DIR"),
"/examples/approval.yaml"
))
.expect("the example");
std::fs::write(&file, &original).expect("written");
let manifest = file.to_str().expect("utf-8 path").to_owned();
tokio::runtime::Builder::new_multi_thread()
.enable_all()
.build()
.expect("runtime")
.block_on(async {
let (router, journal) = dev_session(&manifest).await;
let (status, declared) = page(&router, "GET", "/dev/manifest", None).await;
assert_eq!(status, 200);
assert_eq!(
declared["agents"][0]["name"], "approved-summary",
"{declared}"
);
let start = |ticket: &str| {
serde_json::json!({
"input": { "ticket": ticket },
"correlate": ["customer=C-7"],
})
};
let (status, started) =
page(&router, "POST", "/dev/runs", Some(start("T-1"))).await;
assert_eq!(status, 200, "{started}");
assert_eq!(started["status"], "suspended", "{started}");
let run = started["run"].as_str().expect("a run").to_owned();
let run_id = agentplane::core::RunId::parse(&run).expect("a run id");
let served = served_history(&router, &run).await;
let held = journal.read(run_id, 1).await.expect("the journal");
assert_eq!(served.len(), held.len());
assert_eq!(
served.last().expect("records")["seq"],
held.last().expect("records").seq()
);
let (_, second) = page(&router, "POST", "/dev/runs", Some(start("T-2"))).await;
let second_history =
served_history(&router, second["run"].as_str().expect("a run")).await;
assert_eq!(served[0]["case"], second_history[0]["case"]);
assert!(served[0]["case"].is_string());
let (status, worklist) = page(&router, "GET", "/api/tasks", None).await;
assert_eq!(status, 200, "{worklist}");
let task = worklist["tasks"]
.as_array()
.expect("tasks")
.iter()
.find(|t| t["run"] == run.as_str())
.expect("the run's task")
.clone();
assert!(task["rendering"]["summary"].is_string(), "{task}");
let decide = format!("/api/tasks/{}/decide", task["id"].as_str().expect("an id"));
let stale = agentplane::core::Digest::of(b"another version").to_hex();
let (status, _) = page(
&router,
"POST",
&decide,
Some(serde_json::json!({ "approved": true, "reason": "ok", "digest": stale })),
)
.await;
assert_eq!(status, 412, "a stale digest was not refused");
let (status, decided) = page(
&router,
"POST",
&decide,
Some(serde_json::json!({
"approved": true,
"reason": "checked",
"digest": task["digest"],
})),
)
.await;
assert_eq!(status, 200, "{decided}");
assert_eq!(decided["decided_by"], "dev:author");
let mut concluded = false;
for _ in 0..50 {
let (_, view) = page(&router, "GET", &format!("/api/runs/{run}"), None).await;
if view["status"] == "succeeded" {
concluded = true;
break;
}
tokio::time::sleep(std::time::Duration::from_millis(100)).await;
}
assert!(
concluded,
"the decided run did not resume to its conclusion"
);
let (status, verdicts) = page(
&router,
"POST",
"/dev/replay",
Some(serde_json::json!({ "run": run })),
)
.await;
assert_eq!(status, 200, "{verdicts}");
assert_eq!(verdicts[0]["verdict"], "reproduced", "{verdicts}");
let (status, export) = page(&router, "GET", "/dev/export", None).await;
assert_eq!(status, 200, "{export}");
let report = &export["report"];
assert_eq!(report["findings"], serde_json::json!([]), "{report}");
let reverified = agentplane::export::verify(
export["export"].as_str().expect("the export").as_bytes(),
None,
&[],
)
.expect("readable");
assert_eq!(
serde_json::to_value(&reverified).expect("a report"),
*report,
"the bytes offered are not the bytes verified"
);
std::fs::write(&file, original.replace("One sentence.", "Two sentences."))
.expect("edited");
let (status, verdicts) = page(
&router,
"POST",
"/dev/replay",
Some(serde_json::json!({ "run": run })),
)
.await;
assert_eq!(status, 200, "{verdicts}");
assert_eq!(verdicts[0]["verdict"], "diverged", "{verdicts}");
});
let _ = std::fs::remove_dir_all(&dir);
}
}