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, fixed_clock, fixed_random, fixed_run_id,
text_response, tool_use_response,
};
use salvor_core::{Effect, Event, EventEnvelope, RunId, SequenceNumber};
use salvor_runtime::{Agent, RunOutcome, Runtime, RuntimeError};
use salvor_store::{EventStore, RunSummary, SqliteStore, StoreError};
use serde_json::json;
use time::macros::datetime;
use wiremock::MockServer;
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 scripted_server() -> MockServer {
ScriptedModel::mount(vec![
(
1,
tool_use_response("tu_write", "publish", json!({"doc": "otters"}), 100, 10),
),
(3, text_response("published", 120, 12)),
])
.await
}
fn reference_agent(server_uri: &str) -> (Agent, Arc<AtomicUsize>) {
let (write_tool, write_calls) = TestTool::new("publish", Effect::Write, ToolBehavior::Echo);
let agent = agent_builder(server_uri)
.tool_dyn(Box::new(write_tool))
.build()
.expect("agent builds");
(agent, write_calls)
}
async fn dangling_write() -> (Arc<dyn EventStore>, Agent, RunId, Arc<AtomicUsize>) {
let server = scripted_server().await;
let (agent, write_calls) = reference_agent(&server.uri());
Box::leak(Box::new(server));
let store: Arc<dyn EventStore> = Arc::new(SqliteStore::in_memory().expect("store opens"));
let run_id = fixed_run_id(40);
let killed = Runtime::with_hooks(
Arc::new(KillStore::new(store.clone(), 5)),
fixed_clock(),
fixed_random(),
);
killed
.start_with_id(&agent, run_id, json!("publish otters"))
.await
.expect_err("the kill store aborts the drive");
let log = store.read_log(run_id).await.expect("log reads");
assert_eq!(log.len(), 5, "the log ends at the write intent");
assert!(
matches!(
log[4].event,
Event::ToolCallRequested {
effect: Effect::Write,
..
}
),
"the last event is the dangling write intent",
);
assert_eq!(write_calls.load(Ordering::SeqCst), 1, "the write ran once");
(store, agent, run_id, write_calls)
}
#[tokio::test]
async fn resolve_records_one_completion_then_recovers() {
let (store, agent, run_id, write_calls) = dangling_write().await;
let recovering = Runtime::with_hooks(store.clone(), fixed_clock(), fixed_random());
recovering
.recover(&agent, run_id)
.await
.expect_err("a dangling write refuses automatic recovery");
assert_eq!(
write_calls.load(Ordering::SeqCst),
1,
"recovery ran nothing"
);
let resolver = Runtime::with_hooks(store.clone(), fixed_clock(), fixed_random());
resolver
.resolve(run_id, json!({"echo": {"doc": "otters"}}))
.await
.expect("resolve records the completion");
let log = store.read_log(run_id).await.expect("log reads");
assert_eq!(log.len(), 6, "resolve appends exactly one event");
match &log[5].event {
Event::ToolCallCompleted { seq, output } => {
assert_eq!(*seq, SequenceNumber::new(4), "correlates to the intent seq");
assert_eq!(*output, json!({"echo": {"doc": "otters"}}));
}
other => panic!("expected a tool completion, got {other:?}"),
}
assert_eq!(
write_calls.load(Ordering::SeqCst),
1,
"resolve executed nothing"
);
let again = resolver.resolve(run_id, json!({"echo": 2})).await;
assert!(
matches!(again, Err(RuntimeError::NotReconcilable { .. })),
"second resolve refuses, got {again:?}"
);
let outcome = recovering
.recover(&agent, run_id)
.await
.expect("recovery completes after resolution");
assert!(matches!(
outcome,
RunOutcome::Completed { ref output, .. } if *output == json!("published")
));
assert_eq!(
write_calls.load(Ordering::SeqCst),
1,
"the write is never re-executed"
);
}
async fn store_with_log(events: Vec<Event>) -> (Arc<dyn EventStore>, RunId) {
let store: Arc<dyn EventStore> = Arc::new(SqliteStore::in_memory().expect("store opens"));
let run_id = fixed_run_id(41);
for (i, event) in events.into_iter().enumerate() {
let envelope = EventEnvelope::new(
run_id,
SequenceNumber::new(i as u64),
datetime!(2026-07-09 12:00:00 UTC),
event,
);
store.append(&envelope).await.expect("append seeds the log");
}
(store, run_id)
}
fn started() -> Event {
Event::RunStarted {
agent_def_hash: "sha256:agent".into(),
input: json!("go"),
labels: None,
}
}
#[tokio::test]
async fn resolve_refuses_every_non_reconciliation_state() {
let (store, run_id) =
store_with_log(vec![started(), Event::RunCompleted { output: json!(1) }]).await;
let runtime = Runtime::with_hooks(store, fixed_clock(), fixed_random());
assert!(matches!(
runtime.resolve(run_id, json!(1)).await,
Err(RuntimeError::NotReconcilable { .. })
));
let (store, run_id) = store_with_log(vec![
started(),
Event::Suspended {
reason: "awaiting approval".into(),
input_schema: json!({"type": "object"}),
},
])
.await;
let runtime = Runtime::with_hooks(store, fixed_clock(), fixed_random());
assert!(matches!(
runtime.resolve(run_id, json!(1)).await,
Err(RuntimeError::NotReconcilable { .. })
));
let (store, run_id) = store_with_log(vec![
started(),
Event::ToolCallRequested {
seq: SequenceNumber::new(1),
tool: "search".into(),
input: json!({"q": "otters"}),
effect: Effect::Read,
idempotency_key: None,
},
])
.await;
let runtime = Runtime::with_hooks(store, fixed_clock(), fixed_random());
assert!(matches!(
runtime.resolve(run_id, json!(1)).await,
Err(RuntimeError::NotReconcilable { .. })
));
let store: Arc<dyn EventStore> = Arc::new(SqliteStore::in_memory().expect("store opens"));
let runtime = Runtime::with_hooks(store, fixed_clock(), fixed_random());
assert!(matches!(
runtime.resolve(fixed_run_id(99), json!(1)).await,
Err(RuntimeError::UnknownRun { .. })
));
}