use std::sync::Arc;
use std::sync::atomic::{AtomicUsize, Ordering};
use agentplane::prelude::*;
use agentplane::runtime::effects::Recorded;
use serde_json::{Value, json};
#[derive(Debug, Clone, Default)]
struct World(Arc<AtomicUsize>);
impl World {
fn touched(&self) -> usize {
self.0.load(Ordering::SeqCst)
}
}
#[derive(Debug)]
struct Settlement {
stages: Arc<AtomicUsize>,
crash_at: Arc<AtomicUsize>,
world: World,
}
const NO_CRASH: usize = usize::MAX;
#[async_trait::async_trait]
impl Skill for Settlement {
fn descriptor(&self) -> SkillDescriptor {
SkillDescriptor::new("settlement").provides("billing.settle")
}
async fn invoke(
&self,
cx: &mut StepCtx<'_>,
input: Tainted<Value>,
) -> Result<Outcome, SkillError> {
let n = self.stages.load(Ordering::SeqCst);
let crash_at = self.crash_at.load(Ordering::SeqCst);
for i in 0..n {
let arguments = Tainted::trusted(json!(null));
cx.sink(
Recorded::new(format!("stage-{i}")).counter(Arc::clone(&self.world.0)),
&arguments,
)
.await?;
if i == crash_at {
return Err(SkillError::Other(format!(
"simulated crash after stage {i}"
)));
}
}
let at = cx.now().await?;
Ok(Outcome::done(
input.map(|v| json!({ "settled": v, "at": at.to_string() })),
))
}
}
fn runtime(
store: &Arc<dyn JournalStore>,
stages: &Arc<AtomicUsize>,
crash_at: &Arc<AtomicUsize>,
world: &World,
) -> Arc<Runtime> {
Runtime::builder(Arc::clone(store))
.owner("example")
.skill(Settlement {
stages: Arc::clone(stages),
crash_at: Arc::clone(crash_at),
world: world.clone(),
})
.build()
}
#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
let store: Arc<dyn JournalStore> = Arc::new(RedbStore::open_in_memory()?);
let stages = Arc::new(AtomicUsize::new(3));
let crash_at = Arc::new(AtomicUsize::new(NO_CRASH));
let world = World::default();
let rt = runtime(&store, &stages, &crash_at, &world);
let first = rt
.run(
"billing.settle",
Tainted::trusted(json!({ "invoice": "INV-4711" })),
)
.await?;
println!("1. live run → {:?}", first.status);
println!(" external calls: {}", world.touched());
println!(" chain head: {}", first.chain_head);
let before = world.touched();
let replayed = rt.replay(first.run_id, Mode::Strict).await?;
println!("\n2. strict replay → {:?}", replayed.status);
println!(
" external calls: {} (unchanged: {})",
world.touched(),
world.touched() == before
);
assert_eq!(world.touched(), before, "replay must not touch the world");
assert_eq!(
first.output, replayed.output,
"replay must reproduce the output"
);
let stages = Arc::new(AtomicUsize::new(3));
let crash_at = Arc::new(AtomicUsize::new(0));
let world = World::default();
let rt = runtime(&store, &stages, &crash_at, &world);
let crashed = rt
.run(
"billing.settle",
Tainted::trusted(json!({ "invoice": "INV-4712" })),
)
.await?;
println!("\n3. run crashed → {:?}", crashed.status);
println!(" external calls: {}", world.touched());
crash_at.store(NO_CRASH, Ordering::SeqCst);
let resumed = rt.replay(crashed.run_id, Mode::Resume).await?;
println!(" resumed → {:?}", resumed.status);
println!(
" external calls: {} — stage 0 was replayed, not repeated",
world.touched()
);
assert_eq!(world.touched(), 3, "three stages total, none of them twice");
let stages = Arc::new(AtomicUsize::new(2));
let crash_at = Arc::new(AtomicUsize::new(NO_CRASH));
let world = World::default();
let rt = runtime(&store, &stages, &crash_at, &world);
let recorded = rt
.run(
"billing.settle",
Tainted::trusted(json!({ "invoice": "INV-4713" })),
)
.await?;
stages.store(3, Ordering::SeqCst);
let diverged = rt.replay(recorded.run_id, Mode::Strict).await?;
println!("\n4. changed build → {:?}", diverged.status);
assert!(
matches!(diverged.status, RunStatus::Quarantined(_)),
"divergence must be caught, not absorbed"
);
for run in [first.run_id, crashed.run_id, recorded.run_id] {
store.verify(run).await?;
}
println!("\n5. all journals verify — no record was altered after the fact");
Ok(())
}