use std::sync::{Arc, Mutex};
use agentplane::core::{Effect, EffectDescriptor, EffectError, Recovery, RetryPolicy};
use agentplane::journal::RecordKind;
use agentplane::prelude::*;
use agentplane::runtime::Invariant;
use serde_json::{Value, json};
type Ledger = Arc<Mutex<Vec<String>>>;
#[derive(Debug)]
struct Call {
kind: &'static str,
entry: String,
refuses: bool,
ledger: Ledger,
}
impl Call {
fn new(kind: &'static str, entry: impl Into<String>, ledger: &Ledger) -> Self {
Self {
kind,
entry: entry.into(),
refuses: false,
ledger: Arc::clone(ledger),
}
}
const fn refusing(mut self) -> Self {
self.refuses = true;
self
}
}
#[async_trait::async_trait]
impl Effect for Call {
type Output = Value;
fn descriptor(&self) -> EffectDescriptor {
EffectDescriptor::new(self.kind, json!({ "entry": self.entry }))
}
fn recovery(&self) -> Recovery {
Recovery::Retry
}
fn retry(&self) -> RetryPolicy {
RetryPolicy::never()
}
async fn perform(&self) -> Result<Value, EffectError> {
if self.refuses {
return Err(EffectError::Rejected(format!("{} refused", self.kind)));
}
self.ledger.lock().unwrap().push(self.entry.clone());
Ok(json!({ "reference": format!("{}-ref", self.kind) }))
}
}
#[derive(Debug)]
struct Checkout {
ledger: Ledger,
notify_fails: bool,
}
#[async_trait::async_trait]
impl Skill for Checkout {
fn descriptor(&self) -> SkillDescriptor {
SkillDescriptor::new("checkout").provides("shop.checkout")
}
async fn invoke(
&self,
cx: &mut StepCtx<'_>,
_input: Tainted<Value>,
) -> Result<Outcome, SkillError> {
let l = &self.ledger;
let mut g = cx
.group("checkout", ["inventory", "payments", "notify"])
.await
.map_err(SkillError::Step)?;
let hold = g
.reversible(
"inventory",
Call::new("stock.hold", "held ORD-42", l),
|out| {
Call::new(
"stock.release",
format!("released {}", out["reference"].as_str().unwrap_or("?")),
l,
)
},
)
.await
.map_err(SkillError::Step)?;
g.reversible(
"payments",
Call::new("card.auth", "authorised £129", l),
|out| {
Call::new(
"card.void",
format!("voided {}", out["reference"].as_str().unwrap_or("?")),
l,
)
},
)
.await
.map_err(SkillError::Step)?;
let notify = Call::new("mail.send", "emailed confirmation", l);
g.deferred(
"notify",
if self.notify_fails {
notify.refusing()
} else {
notify
},
)
.map_err(SkillError::Step)?;
g.commit(&[Invariant::new(
"the hold has a reference",
hold.peek()["reference"].is_string(),
)])
.await
.map_err(SkillError::Step)?;
Ok(Outcome::done(Tainted::trusted(json!("checked out"))))
}
}
async fn checkout(
notify_fails: bool,
) -> Result<(Ledger, agentplane::runtime::RunOutcome), Box<dyn std::error::Error>> {
let store: Arc<dyn JournalStore> = Arc::new(RedbStore::open_in_memory()?);
let ledger: Ledger = Arc::default();
let runtime = Runtime::builder(Arc::clone(&store))
.skill(Checkout {
ledger: Arc::clone(&ledger),
notify_fails,
})
.build();
let out = runtime
.run(
"shop.checkout",
Tainted::trusted(json!({ "order": "ORD-42" })),
)
.await?;
let records = store.read(out.run_id, 1).await?;
let settled = records
.iter()
.find_map(|r| match r.kind() {
RecordKind::GroupSettled { group, outcome, .. } => {
Some(format!("{group}: {}", outcome.as_str()))
}
_ => None,
})
.expect("the group was not settled");
println!(" journal says — {settled}");
store.verify(out.run_id).await?;
Ok((ledger, out))
}
#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
let (ledger, out) = checkout(false).await?;
assert_eq!(out.status, RunStatus::Succeeded);
assert_eq!(
*ledger.lock().unwrap(),
["held ORD-42", "authorised £129", "emailed confirmation"]
);
println!(
"1. committed: the gated email went out last, and only once every member had landed\n"
);
let (ledger, out) = checkout(true).await?;
assert!(matches!(out.status, RunStatus::Failed(_)));
let seen = ledger.lock().unwrap().clone();
assert_eq!(
seen,
[
"held ORD-42",
"authorised £129",
"voided card.auth-ref",
"released stock.hold-ref",
]
);
assert!(!seen.iter().any(|e| e == "emailed confirmation"));
println!("2. aborted: both members taken back in reverse, and the email was never sent");
println!(" — not sent and retracted, which is the difference deferral buys");
Ok(())
}