use std::sync::Arc;
use std::sync::atomic::{AtomicUsize, Ordering};
use agentplane::core::{Budget, Outcome, Skill, SkillDescriptor, SkillError, Tainted};
use agentplane::journal::JournalStore;
use agentplane::runtime::effects::Recorded;
use agentplane::runtime::{Mode, RunStatus, Runtime, StepCtx};
use agentplane::store::RedbStore;
use serde_json::{Value, json};
#[derive(Debug)]
struct Settle {
world: Arc<AtomicUsize>,
}
#[async_trait::async_trait]
impl Skill for Settle {
fn descriptor(&self) -> SkillDescriptor {
SkillDescriptor::new("billing.settle")
}
async fn invoke(
&self,
cx: &mut StepCtx<'_>,
input: Tainted<Value>,
) -> Result<Outcome, SkillError> {
for entry in ["fees", "interest", "principal"] {
let arguments = Tainted::trusted(json!(null));
cx.sink(
Recorded::new(format!("post-{entry}")).counter(Arc::clone(&self.world)),
&arguments,
)
.await?;
}
Ok(Outcome::done(input))
}
}
fn plane(store: &Arc<dyn JournalStore>, world: &Arc<AtomicUsize>, budget: Budget) -> Arc<Runtime> {
Runtime::builder(Arc::clone(store))
.owner("billing")
.budget(budget)
.skill(Settle {
world: Arc::clone(world),
})
.build()
}
async fn count(store: &Arc<dyn JournalStore>, run: agentplane::core::RunId, kind: &str) -> usize {
store
.read(run, 1)
.await
.expect("journal readable")
.iter()
.filter(|r| r.kind().kind_str() == kind)
.count()
}
#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
let store: Arc<dyn JournalStore> = Arc::new(RedbStore::open_in_memory()?);
let world = Arc::new(AtomicUsize::new(0));
let capped = plane(&store, &world, Budget::unlimited().effects(2));
let out = capped
.run(
"billing.settle",
Tainted::trusted(json!({ "batch": "B-7" })),
)
.await?;
println!("1. two effects of budget, three postings to make");
println!(" run → {:?}", out.status);
if let RunStatus::Exhausted(why) = &out.status {
println!(" why → {why}");
}
println!(
" postings → {} — the third never started",
world.load(Ordering::SeqCst)
);
assert!(matches!(out.status, RunStatus::Exhausted(_)));
assert_eq!(world.load(Ordering::SeqCst), 2);
let again = capped.replay(out.run_id, Mode::Resume).await?;
println!("\n2. resumed under the same ceiling");
println!(" run → {}", again.status.as_str());
println!(
" postings → {} (unchanged), standing refusals → {}",
world.load(Ordering::SeqCst),
count(&store, out.run_id, "BudgetRefused").await
);
assert!(matches!(again.status, RunStatus::Exhausted(_)));
assert_eq!(world.load(Ordering::SeqCst), 2);
assert_eq!(
count(&store, out.run_id, "BudgetRefused").await,
1,
"re-concluding exhausted must consume the standing refusal, not stack another"
);
let raised = plane(&store, &world, Budget::unlimited().effects(5));
let resumed = raised.replay(out.run_id, Mode::Resume).await?;
println!("\n3. the ceiling is raised, and the run resumed");
println!(" run → {}", resumed.status.as_str());
println!(
" postings → {} — the first two replayed, the third performed once",
world.load(Ordering::SeqCst)
);
println!(
" on the record → BudgetRefused: {}, BudgetReadmitted: {}",
count(&store, out.run_id, "BudgetRefused").await,
count(&store, out.run_id, "BudgetReadmitted").await
);
assert_eq!(resumed.status, RunStatus::Succeeded);
assert_eq!(world.load(Ordering::SeqCst), 3, "each posting exactly once");
assert_eq!(count(&store, out.run_id, "BudgetReadmitted").await, 1);
let strict = raised.replay(out.run_id, Mode::Strict).await?;
assert_eq!(strict.status, RunStatus::Succeeded);
assert_eq!(
world.load(Ordering::SeqCst),
3,
"strict replay must read effects back, never perform them"
);
store.verify(out.run_id).await?;
println!(
"\n4. strict replay walks refusal, re-admission and completion as one \
record,\n and the chain verifies — a pause is history, not damage"
);
Ok(())
}