use std::sync::{Arc, Mutex};
use agentplane::batch::{BatchItem, BatchStatus, BatchStore, ItemOutcome, ItemSource, SourceError};
use agentplane::core::{
ArgSource, BatchId, Effect, EffectDescriptor, EffectError, PlanIR, PlanNode, Recovery,
RetryPolicy, Spend,
};
use agentplane::prelude::*;
use agentplane::runtime::BatchSpec;
use serde_json::{Value, json};
type Ledger = Arc<Mutex<Vec<String>>>;
#[derive(Debug)]
struct Settle {
meter: String,
ledger: Ledger,
}
#[async_trait::async_trait]
impl Effect for Settle {
type Output = Value;
fn descriptor(&self) -> EffectDescriptor {
EffectDescriptor::new("meter.settle", json!({ "meter": self.meter }))
}
fn mutates(&self) -> bool {
true
}
fn recovery(&self) -> Recovery {
Recovery::Retry
}
fn retry(&self) -> RetryPolicy {
RetryPolicy::never()
}
fn spend(&self, _out: &Value) -> Spend {
Spend {
tokens: 0,
minor_units: 250,
}
}
async fn perform(&self) -> Result<Value, EffectError> {
self.ledger.lock().expect("ledger").push(self.meter.clone());
Ok(json!({ "settled": self.meter }))
}
}
#[derive(Debug)]
struct Settler {
ledger: Ledger,
}
#[async_trait::async_trait]
impl Skill for Settler {
fn descriptor(&self) -> SkillDescriptor {
SkillDescriptor::new("settle").provides("settle")
}
async fn invoke(
&self,
cx: &mut StepCtx<'_>,
input: Tainted<Value>,
) -> Result<Outcome, SkillError> {
let meter = input.peek()["meter"]
.as_str()
.unwrap_or_default()
.to_owned();
if meter == "M-004" {
return Ok(Outcome::fail(format!("{meter} has no valid tariff")));
}
cx.effect(Settle {
meter: meter.clone(),
ledger: Arc::clone(&self.ledger),
})
.await?;
Ok(Outcome::done(Tainted::trusted(json!({ "settled": meter }))))
}
}
#[derive(Debug)]
struct Meters(Vec<String>);
impl Meters {
fn upto(n: usize) -> Self {
Self((1..=n).map(|i| format!("M-{i:03}")).collect())
}
}
#[async_trait::async_trait]
impl ItemSource for Meters {
async fn next(&self, after: Option<&str>, limit: usize) -> Result<Vec<BatchItem>, SourceError> {
Ok(self
.0
.iter()
.filter(|k| after.is_none_or(|a| k.as_str() > a))
.take(limit)
.map(|k| BatchItem::new(k, json!({ "meter": k })))
.collect())
}
}
fn plan() -> PlanIR {
PlanIR::new(vec![
PlanNode::new(0, "settle")
.arg("input", ArgSource::run_input())
.terminal(),
])
}
fn plane(store: &Arc<RedbStore>, ledger: &Ledger) -> Arc<Runtime> {
Runtime::builder(Arc::clone(store) as Arc<dyn JournalStore>)
.owner("settlement")
.batches(Arc::clone(store) as Arc<dyn BatchStore>)
.skill(Settler {
ledger: Arc::clone(ledger),
})
.build()
}
#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
let store = Arc::new(RedbStore::open_in_memory()?);
let batches = Arc::clone(&store) as Arc<dyn BatchStore>;
let id = BatchId::generate();
let first: Ledger = Arc::default();
let windowed = BatchSpec::new(plan(), Arc::new(Meters::upto(6)))
.page(1)
.max_items(2);
let report = plane(&store, &first).run_batch(id, &windowed).await?;
println!("1. first pass, a window of two");
println!(" status → {:?}", report.status);
println!(" settled → {:?}", first.lock().expect("ledger"));
println!(" cursor → {:?}", report.cursor.as_deref());
assert_eq!(report.status, BatchStatus::Running);
assert_eq!(report.cursor.as_deref(), Some("M-002"));
let second: Ledger = Arc::default();
let whole = BatchSpec::new(plan(), Arc::new(Meters::upto(6)));
let report = plane(&store, &second).run_batch(id, &whole).await?;
println!("\n2. resumed");
println!(" settled now → {:?}", second.lock().expect("ledger"));
println!(" — M-001 and M-002 are absent: they were not settled again");
assert_eq!(
*second.lock().expect("ledger"),
vec!["M-003", "M-005", "M-006"],
"a resume that re-settles is the re-issued invoice this design refuses"
);
println!("\n3. the batch is finished");
println!(" status → {:?}", report.status);
assert_eq!(
report.status,
BatchStatus::Completed {
succeeded: 5,
failed: 1,
quarantined: 0,
}
);
assert!(
report.needs_attention(),
"a batch with a failed item needs somebody"
);
let backlog = batches.items_needing_attention(id, 100).await?;
println!("\n4. and this is which one");
for item in &backlog {
println!(
" {} → {:?} (run {})",
item.key,
item.outcome
.as_ref()
.map_or("in flight", ItemOutcome::as_str),
item.run
);
}
assert_eq!(backlog.len(), 1);
assert_eq!(backlog[0].key, "M-004");
println!("\n5. cost");
println!(
" batch → {} minor units, summed from the five items that \
settled",
report.spend.minor_units
);
println!(" M-004 → billed for nothing: it posted nothing");
assert_eq!(report.spend.minor_units, 5 * 250);
let third: Ledger = Arc::default();
let again = plane(&store, &third).run_batch(id, &whole).await?;
println!("\n6. the whole batch re-run from item one");
println!(
" settled → {:?} — nothing",
third.lock().expect("ledger")
);
println!(
" status → {:?} — and it still says what happened",
again.status
);
assert!(
third.lock().expect("ledger").is_empty(),
"losing the cursor must be slow, never wrong"
);
assert_eq!(again.status, report.status);
println!(
"\nOne act, six items, five settlements — each performed exactly once,\n\
and the one that did not is named rather than counted."
);
Ok(())
}