use std::sync::Arc;
use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
use std::time::Duration;
use agentplane::core::{Outcome, Skill, SkillDescriptor, SkillError, Tainted};
use agentplane::journal::JournalStore;
use agentplane::runtime::effects::Recorded;
use agentplane::runtime::{RunStatus, Runtime, StepCtx};
use agentplane::store::RedbStore;
use serde_json::{Value, json};
#[derive(Debug)]
struct Settlement {
stalled: Arc<AtomicBool>,
parked: Arc<AtomicBool>,
world: Arc<AtomicUsize>,
}
#[async_trait::async_trait]
impl Skill for Settlement {
fn descriptor(&self) -> SkillDescriptor {
SkillDescriptor::new("billing.settle")
}
async fn invoke(
&self,
cx: &mut StepCtx<'_>,
input: Tainted<Value>,
) -> Result<Outcome, SkillError> {
for stage in 0..3 {
let arguments = Tainted::trusted(json!(null));
cx.sink(
Recorded::new(format!("stage-{stage}")).counter(Arc::clone(&self.world)),
&arguments,
)
.await?;
if stage == 0 && self.stalled.load(Ordering::SeqCst) {
self.parked.store(true, Ordering::SeqCst);
std::future::pending::<()>().await;
}
}
Ok(Outcome::done(input))
}
}
fn instance(
store: &Arc<dyn JournalStore>,
owner: &str,
flags: &(Arc<AtomicBool>, Arc<AtomicBool>),
world: &Arc<AtomicUsize>,
) -> Arc<Runtime> {
Runtime::builder(Arc::clone(store))
.owner(owner)
.lease_ttl(Duration::from_secs(2))
.skill(Settlement {
stalled: Arc::clone(&flags.0),
parked: Arc::clone(&flags.1),
world: Arc::clone(world),
})
.build()
}
#[allow(clippy::disallowed_methods)]
fn now() -> agentplane::core::Timestamp {
agentplane::core::Timestamp::now_utc()
}
#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
let store: Arc<dyn JournalStore> = Arc::new(RedbStore::open_in_memory()?);
let flags = (
Arc::new(AtomicBool::new(true)), Arc::new(AtomicBool::new(false)), );
let world = Arc::new(AtomicUsize::new(0));
let a = instance(&store, "instance-a", &flags, &world);
let task = tokio::spawn({
let a = Arc::clone(&a);
async move {
let _ = a
.run(
"billing.settle",
Tainted::trusted(json!({ "invoice": "INV-9" })),
)
.await;
}
});
while !flags.1.load(Ordering::SeqCst) {
tokio::time::sleep(Duration::from_millis(5)).await;
}
task.abort();
drop(a);
println!("1. instance-a performed stage one, then died mid-run");
println!(" external calls → {}", world.load(Ordering::SeqCst));
assert!(store.abandoned_runs(10).await?.is_empty());
println!("\n2. for one lease TTL, the dead look exactly like the busy…");
tokio::time::sleep(Duration::from_millis(3200)).await;
let stranded = store.abandoned_runs(10).await?;
let run = *stranded.first().expect("the lease lapsed without release");
println!(" …then the lease lapses unreleased, and the store can say who:");
println!(" abandoned → {run}");
flags.0.store(false, Ordering::SeqCst); let b = instance(&store, "instance-b", &flags, &world);
let report = b.sweep(now(), Duration::from_secs(3600)).await?;
println!("\n3. instance-b sweeps");
println!(" runs recovered → {}", report.runs_recovered);
println!(
" external calls → {} — stage one was replayed, not repeated",
world.load(Ordering::SeqCst)
);
assert_eq!(report.runs_recovered, 1);
assert_eq!(world.load(Ordering::SeqCst), 3, "three stages, once each");
let outcome = b
.recorded_outcome(run)
.await?
.expect("the recovered run concluded");
assert_eq!(outcome.status, RunStatus::Succeeded);
println!(" run → {}", outcome.status.as_str());
let evidence = report.record.expect("the sweep sealed its account");
println!("\n4. the sweep journaled its takeover in its own run: {evidence}");
store.verify(run).await?;
store.verify(evidence).await?;
println!(" both chains verify — the recovery is on the record, not just done");
Ok(())
}