use std::sync::Arc;
use agentplane::case::{CaseStore, EventStore, TaskStore};
use agentplane::core::{
AwaitSpec, Calendar, CalendarError, CaseStatus, CorrelationKey, DeadlineSpec, DeadlineState,
Decision, Delivery, Digest, InboundEvent, Justification, OnExpiry, Outcome, Priority, Skill,
SkillDescriptor, SkillError, Tainted, TaskSpec, Timestamp,
};
use agentplane::journal::{JournalStore, RecordKind};
use agentplane::runtime::{RunStatus, Runtime, StepCtx};
use agentplane::store::RedbStore;
use serde_json::{Value, json};
#[derive(Debug)]
struct WorkingDays;
impl Calendar for WorkingDays {
fn resolve(&self, from: Timestamp, spec: &DeadlineSpec) -> Result<Timestamp, CalendarError> {
if spec.kind != "working-days" {
return Err(CalendarError::UnknownKind(spec.kind.clone()));
}
let n = spec
.params
.get("n")
.and_then(Value::as_i64)
.ok_or_else(|| CalendarError::BadParams {
kind: spec.kind.clone(),
detail: "expected `n`".into(),
})?;
let mut at = from;
let mut left = n;
while left > 0 {
at += time::Duration::days(1);
if !matches!(
at.weekday(),
time::Weekday::Saturday | time::Weekday::Sunday
) {
left -= 1;
}
}
Ok(at)
}
fn digest(&self) -> Digest {
Digest::of(b"example.calendar.working-days.v1")
}
}
#[derive(Debug)]
struct SendRequest;
#[async_trait::async_trait]
impl Skill for SendRequest {
fn descriptor(&self) -> SkillDescriptor {
SkillDescriptor::new("send-request").provides("switch.request")
}
async fn invoke(
&self,
cx: &mut StepCtx<'_>,
input: Tainted<Value>,
) -> Result<Outcome, SkillError> {
let document = input
.peek()
.get("document")
.and_then(Value::as_str)
.ok_or_else(|| SkillError::Input("input needs a `document` field".into()))?
.to_owned();
cx.note(format!("dispatching switch request for {document}"))
.await?;
let due = cx
.deadline(
"acknowledgement",
&DeadlineSpec::new("working-days", json!({ "n": 5 })),
Some(time::Duration::days(1)),
)
.await?;
let (_, at) = cx.case_state().await?;
let at = cx
.put_case_state(at, json!({ "stage": "awaiting-acknowledgement" }))
.await?;
cx.set_case_status(CaseStatus::AwaitingExternal).await?;
cx.note(format!("acknowledgement owed by {}", due.resolved_at))
.await?;
cx.deadline(
"decision",
&DeadlineSpec::new("working-days", json!({ "n": 2 })),
None,
)
.await?;
let ack = cx
.await_event(
&AwaitSpec::new("acknowledgement.received", "acknowledgement")
.correlate(CorrelationKey::new("document", &document)),
)
.await?;
cx.meet_deadline("acknowledgement").await?;
let rejected = ack.peek().get("status").and_then(Value::as_str) == Some("rejected");
if rejected {
let at = cx
.put_case_state(at, json!({ "stage": "awaiting-decision" }))
.await?;
cx.set_case_status(CaseStatus::AwaitingHuman).await?;
let decision = cx
.task(
&TaskSpec::new(
"rejection-handling",
Justification::new(
"counterparty rejected the switch request",
json!({ "action": "resubmit-with-corrected-meter" }),
)
.confidence(0.55)
.cost("one further exchange, ~5 working days")
.evidence(format!("rejection payload: {}", ack.peek())),
"decision",
)
.role("mako-operator")
.priority(Priority::High)
.excluding("agent:switch-bot")
.on_expiry(OnExpiry::Escalate),
)
.await?;
cx.put_case_state(
at,
json!({
"stage": "decided",
"approved": decision.approved,
"by": decision.actor,
}),
)
.await?;
cx.set_case_status(CaseStatus::Open).await?;
return Ok(Outcome::done(Tainted::trusted(json!({
"outcome": "human-decided",
"approved": decision.approved,
"by": decision.actor,
}))));
}
cx.put_case_state(at, json!({ "stage": "acknowledged" }))
.await?;
cx.set_case_status(CaseStatus::Open).await?;
Ok(Outcome::done(input.zip(ack).map(
|(sent, reply)| json!({ "sent": sent, "reply": reply }),
)))
}
}
#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
let store = Arc::new(RedbStore::open_in_memory()?);
let rt = Runtime::builder(store.clone() as Arc<dyn JournalStore>)
.cases(store.clone() as Arc<dyn CaseStore>)
.events(store.clone() as Arc<dyn EventStore>)
.tasks(store.clone() as Arc<dyn TaskStore>)
.calendar(Arc::new(WorkingDays))
.skill(SendRequest)
.build();
let keys = [CorrelationKey::new("document", "DOC-4711")];
let sent = rt
.run_in_case(
"switch.request",
json!({ "document": "DOC-4711", "meter": "51238696781" }),
"supplier-switch",
&keys,
)
.await?;
println!("day 0 request → {}", sent.status.as_str());
if let RunStatus::Suspended(reason) = &sent.status {
println!(" waiting → {reason}");
}
let case_id = store.correlate(&keys).await?.expect("case was opened");
println!(" case → {case_id}");
println!(
" cost → {} waiting run(s): a row, not a thread",
store.waiting(10).await?.len()
);
match store.close(case_id).await {
Err(e) => println!(" close → refused: {e}"),
Ok(()) => panic!("an open obligation must block closing"),
}
let ack = InboundEvent::new(
"urn:clearing:counterparty-a",
"MSG-88219",
"acknowledgement.received",
json!({ "status": "rejected", "code": "E_0624", "detail": "meter unknown" }),
)
.correlate(CorrelationKey::new("document", "DOC-4711"));
let delivery = rt.deliver(&ack).await?;
println!("\nday 1 ack → rejected (E_0624)");
println!(" delivery → {delivery:?}");
assert_eq!(delivery, Delivery::Resumed { run: sent.run_id });
assert_eq!(rt.deliver(&ack).await?, Delivery::Duplicate);
println!(" retry → Duplicate (deduplicated by message id)");
handle_rejection(&rt, &store).await?;
let case = store.case(case_id).await?.unwrap();
println!(" state → {}", case.state);
finish_and_close(&rt, &store, case_id, sent.run_id).await?;
demonstrate_early_arrival(&rt).await?;
Ok(())
}
async fn handle_rejection(
rt: &Runtime,
store: &Arc<RedbStore>,
) -> Result<(), Box<dyn std::error::Error>> {
let task = store
.queue(&["mako-operator".to_owned()], 10)
.await?
.pop()
.expect("a rejection must reach a human");
println!(
"\n escalated → task {} ({:?})",
task.kind, task.priority
);
println!(" proposal → {}", task.justification.proposed_action);
println!(" confidence→ {:?}", task.justification.confidence);
println!(" cost → {:?}", task.justification.cost);
let self_approval = rt
.decide_task(
task.id,
&Decision::approve("agent:switch-bot", "I am sure"),
&["mako-operator".to_owned()],
)
.await;
println!(
" self-appr → refused ({})",
self_approval.unwrap_err()
);
rt.decide_task(
task.id,
&Decision::approve("frank", "meter id corrected in the master data"),
&["mako-operator".to_owned()],
)
.await?;
println!(" decided → approved by frank");
Ok(())
}
async fn finish_and_close(
_rt: &Runtime,
store: &Arc<RedbStore>,
case_id: agentplane::core::CaseId,
run_id: agentplane::core::RunId,
) -> Result<(), Box<dyn std::error::Error>> {
for d in store.deadlines(case_id).await? {
if d.state == DeadlineState::Pending {
store
.set_deadline_state(case_id, &d.name, DeadlineState::Met)
.await?;
}
}
assert!(
store
.deadlines(case_id)
.await?
.iter()
.all(|d| d.state == DeadlineState::Met)
);
store.close(case_id).await?;
println!(
"\nclosed → {:?}",
store.case(case_id).await?.unwrap().status
);
let keys = [CorrelationKey::new("document", "DOC-4711")];
assert!(store.correlate(&keys).await?.is_none());
println!("keys released → a new message opens a new case");
let records = store.read(run_id, 1).await?;
let obligations_journaled = records
.iter()
.filter(|r| {
matches!(
r.kind(),
RecordKind::DeadlineRegistered { .. } | RecordKind::DeadlineTransition { .. }
)
})
.count();
println!(
"\naudit → {} records, {} of them obligation events",
records.len(),
obligations_journaled
);
println!(" every one carries the case id");
assert!(records.iter().all(|r| r.body.case == Some(case_id)));
store.verify(run_id).await?;
println!(" the chain verifies across the suspension");
Ok(())
}
async fn demonstrate_early_arrival(rt: &Runtime) -> Result<(), Box<dyn std::error::Error>> {
let early_keys = [CorrelationKey::new("document", "DOC-9999")];
let early = InboundEvent::new(
"urn:clearing:counterparty-a",
"MSG-EARLY",
"acknowledgement.received",
json!({ "status": "very prompt" }),
)
.correlate(CorrelationKey::new("document", "DOC-9999"));
println!("\nearly reply → {:?}", rt.deliver(&early).await?);
let racy = rt
.run_in_case(
"switch.request",
json!({ "document": "DOC-9999", "meter": "51238696782" }),
"supplier-switch",
&early_keys,
)
.await?;
println!("then the request → {}", racy.status.as_str());
println!(" the buffered reply satisfied the wait immediately");
assert_eq!(racy.status, RunStatus::Succeeded);
Ok(())
}