use std::sync::Arc;
use agentplane::core::{Budget, Outcome, Skill, SkillDescriptor, SkillError, Tainted, Trust};
use agentplane::journal::JournalStore;
use agentplane::model::{ModelCall, ModelError, ModelId, ModelProvider, Usage};
use agentplane::runtime::{Mode, RunStatus, Runtime, StepCtx};
use agentplane::store::RedbStore;
use agentplane::testkit::FakeProvider;
use serde_json::{Value, json};
fn render(output: Option<&Value>) -> String {
output.map_or_else(|| "—".to_owned(), Value::to_string)
}
fn model() -> ModelId {
ModelId::new("fake", "triage-1")
}
fn schema() -> Value {
json!({
"type": "object",
"additionalProperties": false,
"required": ["severity", "summary"],
"properties": {
"severity": {"type": "string"},
"summary": {"type": "string"},
},
})
}
#[derive(Debug)]
struct Triage {
provider: Arc<FakeProvider>,
}
#[async_trait::async_trait]
impl Skill for Triage {
fn descriptor(&self) -> SkillDescriptor {
SkillDescriptor::new("triage").provides("support.triage")
}
async fn invoke(
&self,
cx: &mut StepCtx<'_>,
input: Tainted<Value>,
) -> Result<Outcome, SkillError> {
let prompt = input.map(|input| json!({ "task": "triage this ticket", "ticket": input }));
let call = ModelCall::new(
Arc::clone(&self.provider) as Arc<dyn ModelProvider>,
model(),
prompt.peek().clone(),
)
.expecting(schema());
let completion = cx.sink(call, &prompt).await?;
assert_eq!(completion.label().trust, Trust::Untrusted);
Ok(Outcome::done(completion.map(|c| {
c.structured.unwrap_or_else(|| json!({ "text": c.text }))
})))
}
}
#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
let store: Arc<dyn JournalStore> = Arc::new(RedbStore::open_in_memory()?);
let provider = FakeProvider::new();
let rt = Runtime::builder(Arc::clone(&store))
.owner("example")
.skill(Triage {
provider: Arc::clone(&provider),
})
.build();
let ticket = json!({ "id": "T-4711", "body": "checkout returns 500 for EU cards" });
let live = rt.run("support.triage", ticket.clone()).await?;
println!("1. live run → {:?}", live.status);
println!(" provider calls: {}", provider.calls());
println!(
" answer: {}",
render(live.output.as_ref().map(agentplane::Tainted::peek))
);
println!(" spend: {} tokens", live.spend.tokens);
let before = provider.calls();
let replayed = rt.replay(live.run_id, Mode::Strict).await?;
println!("\n2. strict replay → {:?}", replayed.status);
println!(
" provider calls: {} (unchanged: {})",
provider.calls(),
provider.calls() == before
);
assert_eq!(
provider.calls(),
before,
"a replay that asked the model again would pay twice and could get a \
different answer — the completion is journaled precisely so it cannot"
);
assert_eq!(
live.output, replayed.output,
"the least deterministic thing a run touches replays exactly"
);
let provider = FakeProvider::new();
provider.will_fail(ModelError::Interrupted {
model: model(),
usage: Usage {
input_tokens: 100,
output_tokens: 300,
..Usage::default()
},
detail: "connection reset mid-stream".into(),
});
let rt = Runtime::builder(Arc::clone(&store))
.owner("example")
.budget(Budget::unlimited().tokens(300))
.skill(Retries {
provider: Arc::clone(&provider),
})
.build();
let burned = rt.run("support.triage", ticket).await?;
println!("\n3. metered failure → {:?}", burned.status);
println!(
" provider calls: {} — the reworded retry never went out",
provider.calls()
);
println!(
" spend: {} tokens, on a call that returned nothing usable",
burned.spend.tokens
);
assert!(!matches!(burned.status, RunStatus::Succeeded));
assert_eq!(
provider.calls(),
1,
"400 tokens burned against a 300-token ceiling must stop the run before \
it asks again; a driver reporting the dead stream as free would let it \
retry forever against a ceiling reading zero"
);
for run in [live.run_id, burned.run_id] {
store.verify(run).await?;
}
println!("\n4. both journals verify — including the one that failed");
Ok(())
}
#[derive(Debug)]
struct Retries {
provider: Arc<FakeProvider>,
}
#[async_trait::async_trait]
impl Skill for Retries {
fn descriptor(&self) -> SkillDescriptor {
SkillDescriptor::new("triage").provides("support.triage")
}
async fn invoke(
&self,
cx: &mut StepCtx<'_>,
input: Tainted<Value>,
) -> Result<Outcome, SkillError> {
let ask = |prompt: Value| {
ModelCall::new(
Arc::clone(&self.provider) as Arc<dyn ModelProvider>,
model(),
prompt,
)
};
let first_prompt = input
.clone()
.map(|ticket| json!({ "task": "triage", "ticket": ticket }));
let _ = cx
.sink(ask(first_prompt.peek().clone()), &first_prompt)
.await;
let second_prompt =
input.map(|ticket| json!({ "task": "triage, briefly", "ticket": ticket }));
let completion = cx
.sink(ask(second_prompt.peek().clone()), &second_prompt)
.await?;
Ok(Outcome::done(completion.map(|c| json!({ "text": c.text }))))
}
}