use std::sync::Arc;
use std::sync::atomic::{AtomicUsize, Ordering};
use agentplane::core::{Outcome, Skill, SkillDescriptor, SkillError, Tainted, Trust};
use agentplane::journal::JournalStore;
use agentplane::model::openai::OpenAi;
use agentplane::model::{Completion, ModelCall, ModelError, ModelId, ModelProvider, Request};
use agentplane::runtime::{Mode, RunStatus, Runtime, StepCtx};
use agentplane::store::RedbStore;
use serde_json::{Value, json};
const MODEL: &str = "gpt-4o-mini";
#[derive(Debug)]
struct Counted {
inner: Arc<dyn ModelProvider>,
calls: Arc<AtomicUsize>,
}
#[async_trait::async_trait]
impl ModelProvider for Counted {
async fn complete(&self, request: Request<'_>) -> Result<Completion, ModelError> {
self.calls.fetch_add(1, Ordering::SeqCst);
self.inner.complete(request).await
}
}
fn schema() -> Value {
json!({
"type": "object",
"additionalProperties": false,
"required": ["severity", "summary"],
"properties": {
"severity": { "type": "string", "enum": ["low", "high"] },
"summary": { "type": "string" }
}
})
}
#[derive(Debug)]
struct Triage {
provider: Arc<dyn ModelProvider>,
}
#[async_trait::async_trait]
impl Skill for Triage {
fn descriptor(&self) -> SkillDescriptor {
SkillDescriptor::new("triage").provides("triage")
}
async fn invoke(
&self,
cx: &mut StepCtx<'_>,
input: Tainted<Value>,
) -> Result<Outcome, SkillError> {
let prompt = input.map(|ticket| {
json!({
"instruction": "Classify this support ticket. Answer only in the given schema.",
"ticket": ticket
})
});
let call = ModelCall::new(
Arc::clone(&self.provider),
ModelId::new("openai", 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 Ok(key) = std::env::var("OPENAI_API_KEY") else {
eprintln!("OPENAI_API_KEY is not set — this example calls a real model.");
eprintln!(
"For a version that needs no key: cargo run --example model_run --features redb,testkit"
);
return Ok(());
};
let calls = Arc::new(AtomicUsize::new(0));
let provider = Arc::new(Counted {
inner: Arc::new(OpenAi::new(key)?),
calls: Arc::clone(&calls),
}) as Arc<dyn ModelProvider>;
let store: Arc<dyn JournalStore> = Arc::new(RedbStore::open_in_memory()?);
let rt = Runtime::builder(Arc::clone(&store))
.owner("openai-live-example")
.skill(Triage {
provider: Arc::clone(&provider),
})
.build();
let out = rt
.run(
"triage",
json!({ "text": "The printer is on fire and the office is being evacuated." }),
)
.await?;
println!("run {} → {:?}", out.run_id, out.status);
println!(
"answer {}",
out.output
.as_ref()
.map_or(&Value::Null, agentplane::Tainted::peek)
);
println!(
"spend {} tokens, {} calls to OpenAI",
out.spend.tokens,
calls.load(Ordering::SeqCst)
);
assert_eq!(out.status, RunStatus::Succeeded);
let before = calls.load(Ordering::SeqCst);
let replayed = rt.replay(out.run_id, Mode::Strict).await?;
let after = calls.load(Ordering::SeqCst);
println!("replay {:?}", replayed.status);
println!("calls {before} before, {after} after — the model was not asked again");
assert_eq!(
before, after,
"strict replay called the provider, so replay costs money and can \
differ from the history it claims to reproduce"
);
assert_eq!(replayed.output, out.output);
println!("\nthe answer replayed from the journal, byte for byte, without a network call");
Ok(())
}