use std::sync::Arc;
use std::sync::atomic::{AtomicUsize, Ordering};
use agentplane::core::Tainted;
use agentplane::journal::JournalStore;
use agentplane::manifest::Manifest;
use agentplane::model::ModelProvider;
use agentplane::runtime::{Agent, Mode, RunStatus, Runtime};
use agentplane::store::RedbStore;
use agentplane::testkit::FakeProvider;
use agentplane::tools::{Tool, ToolBox, ToolFailure};
use serde_json::{Value, json};
const TELLER: &str = r#"
apiVersion: agentplane.hupe1980.github.io/v1alpha1
kind: Agent
metadata: { name: teller, version: "1.0.0" }
spec:
identity:
role: "Answer questions about a ledger account."
constraints: "Use the tools. Do not guess a balance."
capabilities:
provides: [ledger.ask]
models:
privileged: { provider: fake, model: teller-1 }
security:
max_sensitivity_egress: internal
tools:
# Exactly what this agent may reach. A tool absent from here cannot be
# called however the model spells it — the grant is the operator's decision
# and the model's suggestion is only a suggestion.
- ref: tool://ledger/read
mutates: false
max_sensitivity: internal
description: Read a ledger account's balance.
- ref: tool://ledger/post
# Mutating, and it names the one argument that carries authority. The
# amount is ordinary content a model may choose; **which account** is a
# selector, and a model does not get to pick it.
#
# Declaring the field is also what makes the grant *reachable* at all: a
# `mutates: true` grant with no `protected_fields` on a tool-calling agent
# is refused at parse, because the whole argument bundle carries the
# completion's untrusted label and the taint gate would refuse every call
# — a grant that reads as a capability and fires never.
mutates: true
max_sensitivity: internal
description: Post an amount to an account.
protected_fields:
- path: /account
require_trusted: true
execution: { kind: tool-calling, max_turns: 4 }
budgets:
max_tokens: 100000
"#;
static READS: AtomicUsize = AtomicUsize::new(0);
static POSTS: AtomicUsize = AtomicUsize::new(0);
fn reset() {
READS.store(0, Ordering::Relaxed);
POSTS.store(0, Ordering::Relaxed);
}
#[derive(Debug, serde::Deserialize, schemars::JsonSchema)]
struct ReadBalance {
account: String,
}
#[async_trait::async_trait]
impl Tool for ReadBalance {
const SERVER: &'static str = "ledger";
const NAME: &'static str = "read";
fn mutates() -> bool {
false
}
async fn call(self) -> Result<Value, ToolFailure> {
READS.fetch_add(1, Ordering::Relaxed);
println!(" → read {}", self.account);
Ok(json!({ "account": self.account, "balance": 42 }))
}
}
#[derive(Debug, serde::Deserialize, schemars::JsonSchema)]
struct PostEntry {
#[allow(dead_code)]
account: String,
amount: i64,
}
#[async_trait::async_trait]
impl Tool for PostEntry {
const SERVER: &'static str = "ledger";
const NAME: &'static str = "post";
async fn call(self) -> Result<Value, ToolFailure> {
POSTS.fetch_add(1, Ordering::Relaxed);
Ok(json!({ "posted": self.amount }))
}
}
fn plane(provider: &Arc<FakeProvider>) -> (Arc<Runtime>, Arc<dyn JournalStore>) {
let store: Arc<dyn JournalStore> = Arc::new(RedbStore::open_in_memory().expect("store"));
let manifest = Manifest::parse(TELLER).expect("the agent parses");
let tools = ToolBox::new().with::<ReadBalance>().with::<PostEntry>();
let driver: Arc<dyn ModelProvider> = provider.clone();
let rt = Runtime::builder(Arc::clone(&store))
.provider("fake", driver)
.agent(Agent::new(&manifest))
.toolbox(tools)
.build();
(rt, store)
}
fn ask(question: &str) -> Tainted<Value> {
Tainted::trusted(json!({ "question": question }))
}
#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
let provider = FakeProvider::new();
provider.will_call_tool("call_1", "ledger__read", json!({ "account": "AC-1" }));
provider.will_say("AC-1 holds 42.");
let (rt, store) = plane(&provider);
println!("1. the model chooses a tool");
let first = rt.run("ledger.ask", ask("what is in AC-1?")).await?;
let out = &first;
assert_eq!(out.status, RunStatus::Succeeded);
assert_eq!(READS.load(Ordering::Relaxed), 1);
let asked = provider.asked();
let read = asked[0]
.tools
.iter()
.find(|tool| tool.name == "ledger__read")
.expect("the typed read tool was offered");
assert_eq!(
read.parameters["properties"]["account"]["type"], "string",
"the model was not shown the schema derived from ReadBalance"
);
println!(" answered: {}", out.output.as_ref().unwrap().peek());
let before = (provider.calls(), READS.load(Ordering::Relaxed));
let replayed = rt.replay(out.run_id, Mode::Strict).await?;
assert_eq!(replayed.output, out.output);
assert_eq!(
(provider.calls(), READS.load(Ordering::Relaxed)),
before,
"strict replay called the model or the tool again"
);
println!(" strict replay reassembled it with zero calls");
reset();
let provider = FakeProvider::new();
provider.will_call_tool("call_1", "ledger__write", json!({ "account": "AC-1" }));
provider.will_say("I could not do that, so here is the balance instead.");
let (rt, _) = plane(&provider);
let out = rt.run("ledger.ask", ask("empty AC-1")).await?;
println!("\n2. the model asks for a tool nobody granted");
assert_eq!(
READS.load(Ordering::Relaxed),
0,
"an ungranted tool was called"
);
assert_eq!(out.status, RunStatus::Succeeded);
println!(" → refused, reported to the model, and the run continued");
reset();
let provider = FakeProvider::new();
provider.will_call_tool(
"call_1",
"ledger__post",
json!({ "account": "AC-1", "amount": 1_000_000 }),
);
provider.will_say("I was not able to post that.");
let (rt, _) = plane(&provider);
let out = rt.run("ledger.ask", ask("put a million in AC-1")).await?;
println!("\n3. the model asks to *change* something");
assert_eq!(
POSTS.load(Ordering::Relaxed),
0,
"a model chose the arguments of a mutating call"
);
assert_eq!(out.status, RunStatus::Succeeded);
println!(" → refused: a model-chosen value may not select the protected `/account`");
println!(" (the amount beside it is ordinary content, and would have been fine)");
reset();
let provider = FakeProvider::new();
for i in 0..8 {
provider.will_call_tool(
format!("call_{i}"),
"ledger__read",
json!({ "account": "AC-1" }),
);
}
let (rt, _) = plane(&provider);
println!("\n4. the model never stops asking");
let out = rt.run("ledger.ask", ask("loop")).await?;
match &out.status {
RunStatus::Failed(why) => println!(" → {why}"),
other => panic!("an unbounded loop was allowed to finish: {other:?}"),
}
let records = store.read(first.run_id, 1).await?;
let kinds: Vec<&str> = records.iter().map(|r| r.kind().kind_str()).collect();
let effects = kinds.iter().filter(|k| **k == "EffectStarted").count();
println!(
"\n5. the first run left {} records, {effects} of them effects — one per \
model turn and one per tool call",
records.len()
);
assert!(
effects >= 2,
"a model turn and a tool call should both be effects: {kinds:?}"
);
store.verify(first.run_id).await?;
println!(" and its chain verifies end to end");
Ok(())
}