use std::sync::{Arc, Mutex};
use agentplane::core::{
Compensation, DeadlineSpec, Effect, EffectDescriptor, EffectError, Justification, Outcome,
Recovery, Skill, SkillDescriptor, SkillError, Tainted, TaskSpec,
};
use agentplane::journal::{JournalStore, RecordKind};
use agentplane::quota::{HaltScope, TenantQuota};
use agentplane::runtime::{Mode, RunStatus, Runtime, StepCtx};
use agentplane::store::RedbStore;
use serde_json::{Value, json};
type World = Arc<Mutex<Vec<String>>>;
#[derive(Debug)]
struct Post {
world: World,
what: &'static str,
}
#[async_trait::async_trait]
impl Effect for Post {
type Output = Value;
fn descriptor(&self) -> EffectDescriptor {
EffectDescriptor::new("ledger.post", json!({ "what": self.what }))
}
fn recovery(&self) -> Recovery {
Recovery::Idempotent {
key: format!("ledger:{}", self.what),
}
}
async fn perform(&self) -> Result<Value, EffectError> {
self.world.lock().expect("world").push(self.what.to_owned());
Ok(json!({ "posted": self.what }))
}
}
#[derive(Debug)]
struct HoldThenAsk {
world: World,
}
#[async_trait::async_trait]
impl Skill for HoldThenAsk {
fn descriptor(&self) -> SkillDescriptor {
SkillDescriptor::new("dispute.hold")
}
fn compensation(&self) -> Compensation {
Compensation::Compensatable
}
async fn invoke(
&self,
cx: &mut StepCtx<'_>,
_input: Tainted<Value>,
) -> Result<Outcome, SkillError> {
cx.effect(Post {
world: Arc::clone(&self.world),
what: "hold placed",
})
.await?;
cx.deadline("review", &DeadlineSpec::days(2), None).await?;
let decision = cx
.task(
&TaskSpec::new(
"release-hold",
Justification::new("a person decides whether the hold stands", json!({})),
"review",
)
.role("dispute-officer"),
)
.await?;
Ok(Outcome::done(Tainted::trusted(
json!({ "approved": decision.approved }),
)))
}
async fn compensate(
&self,
cx: &mut StepCtx<'_>,
_out: &Tainted<Value>,
) -> Result<(), SkillError> {
cx.effect(Post {
world: Arc::clone(&self.world),
what: "hold released",
})
.await?;
Ok(())
}
}
#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
stop_one_run().await?;
halt_the_front_door().await?;
Ok(())
}
async fn stop_one_run() -> Result<(), Box<dyn std::error::Error>> {
let store = Arc::new(RedbStore::open_in_memory()?);
let world: World = Arc::default();
let rt = Runtime::builder_on(Arc::clone(&store))
.skill(HoldThenAsk {
world: Arc::clone(&world),
})
.build();
let run = rt
.run_correlated(
"dispute.hold",
Tainted::trusted(json!({ "dispute": "D-311" })),
"dispute",
&[agentplane::core::CorrelationKey::new("dispute", "D-311")],
)
.await?;
println!("1. a run places a hold, then waits for a person");
println!(" run → {}", run.status.as_str());
println!(" world → {:?}", world.lock().expect("world"));
assert!(run.status.is_suspended());
let first = rt
.request_cancel(run.run_id, "ops-carol", "counterparty withdrew the dispute")
.await?;
assert!(first, "the first request is the intervention of record");
let second = rt.request_cancel(run.run_id, "ops-bob", "me too").await?;
assert!(!second);
println!("\n ops-carol stopped it — and the stop *undid* the hold:");
println!(" world → {:?}", world.lock().expect("world"));
assert_eq!(
*world.lock().expect("world"),
["hold placed", "hold released"],
"a cancelled run must unwind what it had already done"
);
let out = rt.replay(run.run_id, Mode::Resume).await?;
println!(
" run → {} — not Failed: this was intended",
out.status.as_str()
);
assert!(out.status.is_cancelled());
let journal: Arc<dyn JournalStore> = store;
let records = journal.read(run.run_id, 1).await?;
let (actor, reason) = records
.iter()
.find_map(|r| match r.kind() {
RecordKind::RunCancelled { actor, reason } => Some((actor.clone(), reason.clone())),
_ => None,
})
.expect("the intervention is journaled");
println!(" on record → cancelled by {actor}: \"{reason}\"");
assert_eq!(actor, "ops-carol");
journal.verify(run.run_id).await?;
println!(" and the chain verifies with the intervention in it\n");
Ok(())
}
#[derive(Debug)]
struct Ack;
#[async_trait::async_trait]
impl Skill for Ack {
fn descriptor(&self) -> SkillDescriptor {
SkillDescriptor::new("desk.ack")
}
async fn invoke(
&self,
_cx: &mut StepCtx<'_>,
input: Tainted<Value>,
) -> Result<Outcome, SkillError> {
Ok(Outcome::done(input))
}
}
async fn halt_the_front_door() -> Result<(), Box<dyn std::error::Error>> {
let store = RedbStore::open_in_memory()?;
let tenant = agentplane::core::TenantId::new("acme")?;
let plane = || {
let scoped = Arc::new(store.clone().for_tenant(tenant.clone()));
Runtime::builder(scoped.clone() as Arc<dyn JournalStore>)
.tenant(tenant.clone())
.quota(
scoped as Arc<dyn agentplane::quota::QuotaStore>,
TenantQuota::default(),
)
.skill(Ack)
.build()
};
let one = plane();
let two = plane();
one.set_halt(
&HaltScope::Tenant,
Some("incident 42: ledger reconciliation is wrong"),
)
.await?;
println!("2. instance one throws the emergency stop (whole tenant)");
for (name, rt) in [("one", &one), ("two", &two)] {
match rt.run("desk.ack", Tainted::trusted(json!({}))).await {
Err(refusal) => println!(" instance {name} → refused: {refusal}"),
Ok(out) => panic!("a halted tenant admitted a run: {:?}", out.status),
}
}
one.set_halt(&HaltScope::Tenant, None).await?;
let lifted = two.run("desk.ack", Tainted::trusted(json!({}))).await?;
println!(
" lifted → instance two admits again ({})",
lifted.status.as_str()
);
assert_eq!(lifted.status, RunStatus::Succeeded);
Ok(())
}