#![allow(clippy::disallowed_methods)]
#![allow(clippy::cast_precision_loss)]
use std::sync::Arc;
use std::time::{Duration, Instant};
use agentplane::core::{Effect, EffectDescriptor, EffectError, ProtectedField, Recovery};
use agentplane::prelude::*;
use serde_json::{Value, json};
const RULES: &str = r#"
@id("permit-effects")
permit (principal, action == Action::"effect:perform", resource);
@id("permit-admission")
permit (principal, action == Action::"run:admit", resource);
@id("permit-release")
permit (principal, action == Action::"information_flow.release", resource);
"#;
const RULE_COUNT: usize = 3;
#[derive(Debug)]
struct Burst(usize);
#[async_trait::async_trait]
impl Skill for Burst {
fn descriptor(&self) -> SkillDescriptor {
SkillDescriptor::new("burst").provides("perf.burst")
}
async fn invoke(
&self,
cx: &mut StepCtx<'_>,
input: Tainted<Value>,
) -> Result<Outcome, SkillError> {
for _ in 0..self.0 {
let _ = cx.now().await?;
}
Ok(Outcome::done(input))
}
}
#[derive(Debug)]
struct Pay {
args: Value,
protected: Vec<ProtectedField>,
}
#[async_trait::async_trait]
impl Effect for Pay {
type Output = Value;
fn descriptor(&self) -> EffectDescriptor {
EffectDescriptor::new("ledger.pay", self.args.clone())
}
fn mutates(&self) -> bool {
true
}
fn recovery(&self) -> Recovery {
Recovery::Retry
}
fn sink_arguments(&self) -> Option<&Value> {
Some(&self.args)
}
fn protected_fields(&self) -> &[ProtectedField] {
&self.protected
}
async fn perform(&self) -> Result<Value, EffectError> {
Ok(json!({ "ok": true }))
}
}
#[derive(Debug)]
struct Sinking(usize);
#[async_trait::async_trait]
impl Skill for Sinking {
fn descriptor(&self) -> SkillDescriptor {
SkillDescriptor::new("sinking").provides("perf.sink")
}
async fn invoke(
&self,
cx: &mut StepCtx<'_>,
input: Tainted<Value>,
) -> Result<Outcome, SkillError> {
let protected = vec![ProtectedField::trusted("/account")];
for i in 0..self.0 {
let args = Tainted::trusted(json!({ "account": "DE-1", "seq": i }));
let effect = Pay {
args: args.peek().clone(),
protected: protected.clone(),
};
let _ = cx.sink(effect, &args).await?;
}
Ok(Outcome::done(input))
}
}
fn store(on_disk: bool, tag: &str) -> Result<Arc<dyn JournalStore>, Box<dyn std::error::Error>> {
if on_disk {
let path = std::env::temp_dir().join(format!(
"agentplane-bench-{tag}-{}.redb",
std::process::id()
));
let _ = std::fs::remove_file(&path);
Ok(Arc::new(RedbStore::open(&path)?))
} else {
Ok(Arc::new(RedbStore::open_in_memory()?))
}
}
async fn time(
runtime: &Arc<agentplane::runtime::Runtime>,
capability: &str,
expected: usize,
) -> Result<(Duration, Duration), Box<dyn std::error::Error>> {
let started = Instant::now();
let outcome = runtime.run(capability, Tainted::trusted(json!({}))).await?;
let live = started.elapsed();
let started = Instant::now();
runtime.replay(outcome.run_id, Mode::Strict).await?;
let replay = started.elapsed();
if outcome.status != RunStatus::Succeeded {
return Err(format!("{capability} did not succeed: {:?}", outcome.status).into());
}
let journaled = runtime.journal().read(outcome.run_id, 0).await?.len();
if journaled != expected {
return Err(format!(
"{capability} journaled {journaled} records, not {expected} — the axis is \
not performing the work it is being timed for"
)
.into());
}
Ok((live, replay))
}
#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
let on_disk = std::env::var("DISK").is_ok();
let n: usize = std::env::var("N")
.ok()
.and_then(|v| v.parse().ok())
.unwrap_or(if on_disk { 400 } else { 2000 });
let repeats: usize = std::env::var("REPEATS")
.ok()
.and_then(|v| v.parse().ok())
.unwrap_or(3);
let records = n * 2 + 5;
let mut baseline = Vec::new();
let mut replay = Duration::MAX;
for i in 0..repeats {
let rt = Runtime::builder(store(on_disk, &format!("journal-{i}"))?)
.skill(Burst(n))
.build();
let (live, read) = time(&rt, "perf.burst", records).await?;
baseline.push(live);
replay = replay.min(read);
}
let journal = *baseline.iter().min().expect("at least one run");
let spread = baseline
.iter()
.max()
.expect("at least one run")
.saturating_sub(journal);
let mut policy = Duration::MAX;
for i in 0..repeats {
let engine = agentplane::policy::CedarEngine::new(RULES)?;
let rt = Runtime::builder(store(on_disk, &format!("policy-{i}"))?)
.skill(Burst(n))
.policy(Arc::new(engine))
.build();
policy = policy.min(time(&rt, "perf.burst", records).await?.0);
}
let mut sink = Duration::MAX;
for i in 0..repeats {
let rt = Runtime::builder(store(on_disk, &format!("sink-{i}"))?)
.skill(Sinking(n))
.build();
sink = sink.min(time(&rt, "perf.sink", records).await?.0);
}
let per = |d: Duration| d.as_secs_f64() / n as f64 * 1000.0;
let resolves = per(spread);
let delta = |d: Duration| {
let d = per(d) - per(journal);
if d.abs() <= resolves {
" under the spread".to_owned()
} else {
format!("{d:+8.3} ms vs journal")
}
};
println!(
"{n} effects × {repeats} runs, redb {}, {RULE_COUNT} policy rules\n \
journal {:>8.3} ms/effect {:>8.0} effects/sec canonicalize, chain, commit\n \
policy {:>8.3} ms/effect {} one authorize per effect\n \
sink {:>8.3} ms/effect {} label gate, one protected field\n \
replay {:>8.3} ms/effect {:>8.0} effects/sec the read path, nothing performed\n \
spread {resolves:>8.3} ms/effect across {repeats} baseline runs, what this resolves",
if on_disk { "on disk" } else { "in memory" },
per(journal),
n as f64 / journal.as_secs_f64(),
per(policy),
delta(policy),
per(sink),
delta(sink),
per(replay),
n as f64 / replay.as_secs_f64(),
);
Ok(())
}