mod common;
use std::sync::Arc;
use std::sync::atomic::{AtomicI64, AtomicUsize, Ordering};
use async_trait::async_trait;
use common::{
ScriptedModel, TestTool, ToolBehavior, agent_builder, blocks_response, event_kinds,
fixed_clock, fixed_random, fixed_run_id, text_response, tool_result_contents, tool_use_block,
tool_use_response,
};
use salvor_core::{Effect, Event, EventEnvelope, RunId};
use salvor_runtime::{
ANSWER_TOOL, Agent, Budgets, LoopOutcome, ParkReason, RunCtx, RunOutcome, Runtime,
RuntimeError, drive_loop_structured,
};
use salvor_store::{EventStore, RunSummary, SqliteStore, StoreError};
use serde_json::{Value, json};
use wiremock::MockServer;
fn schema() -> Value {
json!({
"type": "object",
"required": ["score"],
"properties": {"score": {"type": "number"}, "note": {"type": "string"}}
})
}
fn answer_response(tool_use_id: &str, input: Value, output_tokens: u64) -> Value {
tool_use_response(tool_use_id, ANSWER_TOOL, input, 10, output_tokens)
}
async fn drive_structured(
store: Arc<dyn EventStore>,
run_id: RunId,
agent: &Agent,
input: Value,
resume_input: Option<Value>,
) -> Result<LoopOutcome, RuntimeError> {
let log = store.read_log(run_id).await.expect("log reads");
let mut ctx = RunCtx::with_hooks(store, run_id, log, fixed_clock(), fixed_random())
.expect("ctx builds over the recorded log");
if let Some(resume_input) = resume_input {
ctx.set_resume_input(resume_input);
}
let input = ctx.begin(agent.def_hash(), &input).await?;
let outcome = drive_loop_structured(&mut ctx, agent, &input, &schema()).await?;
if let LoopOutcome::Completed(output) = &outcome {
ctx.complete_run(output).await?;
}
Ok(outcome)
}
fn offered_tools(request_body: &[u8]) -> (Vec<Value>, Option<Value>) {
let body: Value = serde_json::from_slice(request_body).expect("request body is JSON");
let tools = body
.get("tools")
.and_then(Value::as_array)
.cloned()
.unwrap_or_default();
(tools, body.get("tool_choice").cloned())
}
#[tokio::test]
async fn a_valid_first_answer_is_the_loops_output_and_the_request_forces_the_call() {
let server = ScriptedModel::mount(vec![(
1,
answer_response("tu_answer", json!({"score": 0.91, "note": "tight"}), 2),
)])
.await;
let agent = agent_builder(&server.uri()).build().expect("agent builds");
let store: Arc<dyn EventStore> = Arc::new(SqliteStore::in_memory().expect("store opens"));
let run_id = fixed_run_id(60);
let outcome = drive_structured(store.clone(), run_id, &agent, json!("rate this"), None)
.await
.expect("the drive succeeds");
let LoopOutcome::Completed(output) = outcome else {
panic!("expected completion");
};
assert_eq!(output, json!({"score": 0.91, "note": "tight"}));
let requests = server.received_requests().await.expect("requests recorded");
let (tools, tool_choice) = offered_tools(&requests[0].body);
assert_eq!(tools.len(), 1, "the agent has no tools of its own");
assert_eq!(tools[0]["name"], json!(ANSWER_TOOL));
assert_eq!(tools[0]["input_schema"], schema());
assert!(
tools[0]["description"]
.as_str()
.is_some_and(|d| !d.is_empty()),
"the answer tool describes itself: {:?}",
tools[0]["description"]
);
assert_eq!(tool_choice, Some(json!({"type": "any"})));
let log = store.read_log(run_id).await.expect("log reads");
assert_eq!(
event_kinds(&log),
[
"RunStarted",
"NowObserved",
"ModelCallRequested",
"ModelCallCompleted",
"RunCompleted"
]
);
}
#[tokio::test]
async fn a_violating_answer_is_fed_back_and_the_re_ask_replays_byte_for_byte() {
let server = ScriptedModel::mount(vec![
(1, answer_response("tu_bad", json!({"score": "high"}), 2)),
(3, answer_response("tu_good", json!({"score": 0.88}), 3)),
])
.await;
let agent = agent_builder(&server.uri()).build().expect("agent builds");
let store: Arc<dyn EventStore> = Arc::new(SqliteStore::in_memory().expect("store opens"));
let run_id = fixed_run_id(61);
let outcome = drive_structured(store.clone(), run_id, &agent, json!("rate this"), None)
.await
.expect("the drive succeeds");
assert!(
matches!(outcome, LoopOutcome::Completed(ref output) if *output == json!({"score": 0.88})),
"{outcome:?}"
);
let requests = server.received_requests().await.expect("requests recorded");
let fed_back = tool_result_contents(&requests[1].body);
assert_eq!(fed_back.len(), 1, "{fed_back:?}");
assert!(
fed_back[0].contains("$.score: expected type number, got string"),
"{}",
fed_back[0]
);
assert!(
fed_back[0].contains("does not match its schema"),
"{}",
fed_back[0]
);
let log = store.read_log(run_id).await.expect("log reads");
assert_eq!(
event_kinds(&log),
[
"RunStarted",
"NowObserved",
"ModelCallRequested",
"ModelCallCompleted",
"NowObserved",
"ModelCallRequested",
"ModelCallCompleted",
"RunCompleted"
]
);
let replayed = drive_structured(store.clone(), run_id, &agent, json!("rate this"), None)
.await
.expect("the replay succeeds");
assert!(
matches!(replayed, LoopOutcome::Completed(ref output) if *output == json!({"score": 0.88}))
);
assert_eq!(store.read_log(run_id).await.expect("log reads"), log);
assert_eq!(
server.received_requests().await.expect("requests").len(),
2,
"the replayed drive reached the model zero times"
);
}
#[tokio::test]
async fn an_answer_beside_real_tool_work_is_refused_while_the_tool_still_runs() {
let server = ScriptedModel::mount(vec![
(
1,
blocks_response(
"mixed",
vec![
tool_use_block("tu_read", "lookup", json!({"q": "otters"})),
tool_use_block("tu_early", ANSWER_TOOL, json!({"score": 0.5})),
],
20,
4,
),
),
(3, answer_response("tu_answer", json!({"score": 0.75}), 3)),
])
.await;
let (tool, calls) = TestTool::new("lookup", Effect::Read, ToolBehavior::Echo);
let agent = agent_builder(&server.uri())
.tool_dyn(Box::new(tool))
.build()
.expect("agent builds");
let store: Arc<dyn EventStore> = Arc::new(SqliteStore::in_memory().expect("store opens"));
let run_id = fixed_run_id(62);
let outcome = drive_structured(store.clone(), run_id, &agent, json!("rate this"), None)
.await
.expect("the drive succeeds");
assert!(
matches!(outcome, LoopOutcome::Completed(ref output) if *output == json!({"score": 0.75})),
"{outcome:?}"
);
assert_eq!(
calls.load(Ordering::SeqCst),
1,
"the real call ran normally"
);
let requests = server.received_requests().await.expect("requests recorded");
let fed_back = tool_result_contents(&requests[1].body);
assert_eq!(fed_back.len(), 2, "{fed_back:?}");
assert!(fed_back[0].contains("otters"), "{}", fed_back[0]);
assert!(fed_back[1].contains("call it alone"), "{}", fed_back[1]);
let log = store.read_log(run_id).await.expect("log reads");
assert_eq!(
event_kinds(&log),
[
"RunStarted",
"NowObserved",
"ModelCallRequested",
"ModelCallCompleted",
"ToolCallRequested",
"ToolCallCompleted",
"NowObserved",
"ModelCallRequested",
"ModelCallCompleted",
"RunCompleted"
]
);
}
#[tokio::test]
async fn a_turn_with_no_tool_call_at_all_is_re_asked() {
let server = ScriptedModel::mount(vec![
(1, text_response("the score is about 0.9, roughly", 10, 5)),
(3, answer_response("tu_answer", json!({"score": 0.9}), 3)),
])
.await;
let agent = agent_builder(&server.uri()).build().expect("agent builds");
let store: Arc<dyn EventStore> = Arc::new(SqliteStore::in_memory().expect("store opens"));
let run_id = fixed_run_id(63);
let outcome = drive_structured(store.clone(), run_id, &agent, json!("rate this"), None)
.await
.expect("the drive succeeds");
assert!(
matches!(outcome, LoopOutcome::Completed(ref output) if *output == json!({"score": 0.9})),
"{outcome:?}"
);
let requests = server.received_requests().await.expect("requests recorded");
let body: Value = serde_json::from_slice(&requests[1].body).expect("request body is JSON");
let messages = body["messages"].as_array().expect("messages");
assert_eq!(messages.len(), 3, "input, the bare-text turn, the re-ask");
let re_ask = messages[2]["content"].as_str().expect("the re-ask is text");
assert!(re_ask.contains("That turn called no tool"), "{re_ask}");
}
#[tokio::test]
async fn repeated_identical_violations_collapse_through_the_streak_counter() {
let bad = json!({"score": "high"});
let server = ScriptedModel::mount(vec![
(1, answer_response("tu_bad_1", bad.clone(), 2)),
(3, answer_response("tu_bad_2", bad.clone(), 2)),
(5, answer_response("tu_bad_3", bad, 2)),
(7, answer_response("tu_good", json!({"score": 0.4}), 3)),
])
.await;
let agent = agent_builder(&server.uri()).build().expect("agent builds");
let store: Arc<dyn EventStore> = Arc::new(SqliteStore::in_memory().expect("store opens"));
let run_id = fixed_run_id(64);
let outcome = drive_structured(store.clone(), run_id, &agent, json!("rate this"), None)
.await
.expect("the drive succeeds");
assert!(
matches!(outcome, LoopOutcome::Completed(ref output) if *output == json!({"score": 0.4})),
"{outcome:?}"
);
let requests = server.received_requests().await.expect("requests recorded");
let first = tool_result_contents(&requests[1].body);
let second = tool_result_contents(&requests[2].body);
let third = tool_result_contents(&requests[3].body);
assert!(first[0].contains("expected type number"), "{}", first[0]);
assert!(
second
.last()
.expect("a result")
.contains("2 consecutive times"),
"{second:?}"
);
assert!(
third
.last()
.expect("a result")
.contains("3 consecutive times"),
"{third:?}"
);
}
#[tokio::test]
async fn a_steps_budget_crossing_mid_re_ask_parks_and_a_resume_continues() {
let bad = json!({"score": "high"});
let server = ScriptedModel::mount(vec![
(1, answer_response("tu_bad_1", bad.clone(), 2)),
(3, answer_response("tu_bad_2", bad, 2)),
(5, answer_response("tu_good", json!({"score": 0.6}), 3)),
])
.await;
let agent = agent_builder(&server.uri())
.budgets(Budgets {
max_steps: Some(2),
..Budgets::default()
})
.build()
.expect("agent builds");
let store: Arc<dyn EventStore> = Arc::new(SqliteStore::in_memory().expect("store opens"));
let run_id = fixed_run_id(65);
let outcome = drive_structured(store.clone(), run_id, &agent, json!("rate this"), None)
.await
.expect("the drive itself succeeds");
assert!(
matches!(
outcome,
LoopOutcome::Parked(ParkReason::BudgetExceeded { budget, observed })
if budget.limit == 2.0 && observed == 2.0
),
"{outcome:?}"
);
let log = store.read_log(run_id).await.expect("log reads");
assert_eq!(
event_kinds(&log),
[
"RunStarted",
"NowObserved",
"ModelCallRequested",
"ModelCallCompleted",
"NowObserved",
"ModelCallRequested",
"ModelCallCompleted",
"NowObserved",
"BudgetExceeded"
]
);
let outcome = drive_structured(
store.clone(),
run_id,
&agent,
json!("rate this"),
Some(json!({"extend": {"steps": 2}})),
)
.await
.expect("the resumed drive succeeds");
assert!(
matches!(outcome, LoopOutcome::Completed(ref output) if *output == json!({"score": 0.6})),
"{outcome:?}"
);
}
#[tokio::test]
async fn a_real_tool_named_salvor_answer_refuses_before_any_model_call() {
let server = ScriptedModel::mount(vec![(
1,
answer_response("tu_answer", json!({"score": 0.1}), 2),
)])
.await;
let (tool, calls) = TestTool::new(ANSWER_TOOL, Effect::Read, ToolBehavior::Echo);
let agent = agent_builder(&server.uri())
.tool_dyn(Box::new(tool))
.build()
.expect("agent builds");
let store: Arc<dyn EventStore> = Arc::new(SqliteStore::in_memory().expect("store opens"));
let run_id = fixed_run_id(66);
let error = drive_structured(store.clone(), run_id, &agent, json!("rate this"), None)
.await
.expect_err("a colliding tool name refuses the drive");
assert!(
matches!(error, RuntimeError::AnswerToolNameTaken),
"{error}"
);
let log = store.read_log(run_id).await.expect("log reads");
assert_eq!(event_kinds(&log), ["RunStarted"]);
assert!(
server
.received_requests()
.await
.expect("requests recorded")
.is_empty()
);
assert_eq!(calls.load(Ordering::SeqCst), 0);
}
struct KillStore {
inner: Arc<dyn EventStore>,
remaining: AtomicI64,
}
impl KillStore {
fn new(inner: Arc<dyn EventStore>, allow: i64) -> Self {
Self {
inner,
remaining: AtomicI64::new(allow),
}
}
}
#[async_trait]
impl EventStore for KillStore {
async fn append(&self, envelope: &EventEnvelope) -> Result<(), StoreError> {
if self.remaining.fetch_sub(1, Ordering::SeqCst) <= 0 {
return Err(StoreError::Backend("simulated crash".to_owned()));
}
self.inner.append(envelope).await
}
async fn read_log(&self, run_id: RunId) -> Result<Vec<EventEnvelope>, StoreError> {
self.inner.read_log(run_id).await
}
async fn list_runs(&self) -> Result<Vec<RunSummary>, StoreError> {
self.inner.list_runs().await
}
async fn claim_call(
&self,
claimant: salvor_store::CallClaimant<'_>,
) -> Result<salvor_store::CallClaim, StoreError> {
self.inner.claim_call(claimant).await
}
async fn lookup_call(
&self,
tool: &str,
idempotency_key: &str,
) -> Result<Option<salvor_store::CallCommitment>, StoreError> {
self.inner.lookup_call(tool, idempotency_key).await
}
async fn append_settling_call(
&self,
envelope: &EventEnvelope,
claimant: salvor_store::CallClaimant<'_>,
) -> Result<(), StoreError> {
if self.remaining.fetch_sub(1, Ordering::SeqCst) <= 0 {
return Err(StoreError::Backend("simulated crash".to_owned()));
}
self.inner.append_settling_call(envelope, claimant).await
}
}
async fn sweep_server() -> MockServer {
ScriptedModel::mount(vec![
(
1,
tool_use_response("tu_read", "lookup", json!({"q": "otters"}), 20, 4),
),
(3, answer_response("tu_bad", json!({"score": "high"}), 2)),
(5, answer_response("tu_good", json!({"score": 0.93}), 3)),
])
.await
}
fn sweep_agent(server_uri: &str) -> (Agent, Arc<AtomicUsize>) {
let (tool, calls) = TestTool::new("lookup", Effect::Read, ToolBehavior::Echo);
let agent = agent_builder(server_uri)
.tool_dyn(Box::new(tool))
.build()
.expect("agent builds");
(agent, calls)
}
#[tokio::test]
async fn a_structured_run_recovers_identically_from_every_kill_boundary() {
let server = sweep_server().await;
let (agent, _calls) = sweep_agent(&server.uri());
let store: Arc<dyn EventStore> = Arc::new(SqliteStore::in_memory().expect("store opens"));
let control_id = fixed_run_id(70);
let outcome = drive_structured(store.clone(), control_id, &agent, json!("rate this"), None)
.await
.expect("the control run completes");
assert!(
matches!(outcome, LoopOutcome::Completed(ref output) if *output == json!({"score": 0.93}))
);
let control = store.read_log(control_id).await.expect("log reads");
assert_eq!(
control.len(),
13,
"the reference structured run records 13 events"
);
for allow in 1..control.len() {
let server = sweep_server().await;
let (agent, calls) = sweep_agent(&server.uri());
let store: Arc<dyn EventStore> = Arc::new(SqliteStore::in_memory().expect("store opens"));
let run_id = fixed_run_id(71);
let killer: Arc<dyn EventStore> = Arc::new(KillStore::new(
store.clone(),
i64::try_from(allow).expect("a small budget"),
));
let error = drive_structured(killer, run_id, &agent, json!("rate this"), None)
.await
.expect_err("the kill store aborts the drive");
assert!(matches!(error, RuntimeError::Store(_)), "{error}");
assert_eq!(
store.read_log(run_id).await.expect("log reads").len(),
allow,
"exactly the budgeted number of events persisted"
);
let outcome = drive_structured(store.clone(), run_id, &agent, json!("rate this"), None)
.await
.unwrap_or_else(|error| panic!("recovery after {allow} appends failed: {error}"));
assert!(
matches!(
outcome,
LoopOutcome::Completed(ref output) if *output == json!({"score": 0.93})
),
"recovery after {allow} appends: {outcome:?}"
);
let recovered = store.read_log(run_id).await.expect("log reads");
assert_eq!(
recovered.len(),
control.len(),
"recovery after {allow} appends produced a different log length"
);
for (position, (left, right)) in recovered.iter().zip(control.iter()).enumerate() {
assert_eq!(
left.event, right.event,
"recovery after {allow} appends diverged at seq {position}"
);
}
let executions = calls.load(Ordering::SeqCst);
assert!(
(1..=2).contains(&executions),
"the read tool ran {executions} times for a kill after {allow} appends"
);
}
}
#[tokio::test]
async fn an_agent_that_declares_a_schema_drives_the_structured_loop_through_the_runtime() {
let server = ScriptedModel::mount(vec![(
1,
answer_response("tu_runtime", json!({"score": 0.42, "note": "thin"}), 2),
)])
.await;
let agent = agent_builder(&server.uri())
.output_schema(schema())
.build()
.expect("agent builds");
let store = Arc::new(SqliteStore::in_memory().expect("store opens"));
let runtime = Runtime::with_hooks(store.clone(), fixed_clock(), fixed_random());
let run_id = fixed_run_id(72);
let outcome = runtime
.start_with_id(&agent, run_id, json!("rate this"))
.await
.expect("the run completes");
let RunOutcome::Completed { output, .. } = outcome else {
panic!("expected completion, got {outcome:?}");
};
assert_eq!(output, json!({"score": 0.42, "note": "thin"}));
let log = store.read_log(run_id).await.expect("log reads");
assert_eq!(
event_kinds(&log),
[
"RunStarted",
"NowObserved",
"ModelCallRequested",
"ModelCallCompleted",
"RunCompleted",
]
);
let Event::RunCompleted { output } = &log.last().expect("a terminal").event else {
panic!("the last event is the terminal");
};
assert_eq!(*output, json!({"score": 0.42, "note": "thin"}));
let requests = server.received_requests().await.expect("requests recorded");
let (tools, tool_choice) = offered_tools(&requests[0].body);
assert_eq!(tools.len(), 1);
assert_eq!(tools[0]["name"], json!(ANSWER_TOOL));
assert_eq!(tools[0]["input_schema"], schema());
assert_eq!(tool_choice, Some(json!({"type": "any"})));
}
#[tokio::test]
async fn an_agent_without_a_schema_still_drives_the_plain_loop_through_the_runtime() {
let server = ScriptedModel::mount(vec![(1, text_response("about 0.42", 10, 2))]).await;
let agent = agent_builder(&server.uri()).build().expect("agent builds");
let store = Arc::new(SqliteStore::in_memory().expect("store opens"));
let runtime = Runtime::with_hooks(store.clone(), fixed_clock(), fixed_random());
let run_id = fixed_run_id(73);
let outcome = runtime
.start_with_id(&agent, run_id, json!("rate this"))
.await
.expect("the run completes");
let RunOutcome::Completed { output, .. } = outcome else {
panic!("expected completion, got {outcome:?}");
};
assert_eq!(output, json!("about 0.42"));
let requests = server.received_requests().await.expect("requests recorded");
let body: Value = serde_json::from_slice(&requests[0].body).expect("request body is JSON");
assert!(body.get("tools").is_none(), "{body}");
assert!(body.get("tool_choice").is_none(), "{body}");
}