mod common;
use std::pin::Pin;
use std::sync::Arc;
use std::sync::atomic::{AtomicUsize, Ordering};
use std::time::Duration;
use common::{
CountBehavior, ScriptedModel, SseReader, TestServer, agent_factory, app_state, counter,
memory_store, open_store, register_agent, sample_toml,
};
use salvor_core::{Effect, Event, EventEnvelope, Performer, RunId, SequenceNumber};
use salvor_runtime::due_runs;
use salvor_server::{AgentFactory, AppState};
use salvor_store::EventStore;
use serde_json::json;
use tempfile::tempdir;
use time::OffsetDateTime;
use time::macros::datetime;
const NOW: OffsetDateTime = datetime!(2026-07-10 12:00:00 UTC);
async fn seed_sleeping(
store: &Arc<dyn EventStore>,
agent_hash: &str,
wake_at: OffsetDateTime,
) -> RunId {
let run_id = RunId::new();
for (seq, event) in [
Event::RunStarted {
agent_def_hash: agent_hash.to_owned(),
input: json!("go"),
labels: None,
driven_by: None,
},
Event::SleepStarted { wake_at },
]
.into_iter()
.enumerate()
{
let envelope = EventEnvelope::new(run_id, SequenceNumber::new(seq as u64), NOW, event);
store.append(&envelope).await.expect("append");
}
run_id
}
async fn seed_sleeping_client_driven(
store: &Arc<dyn EventStore>,
wake_at: OffsetDateTime,
) -> RunId {
let run_id = RunId::new();
for (seq, event) in [
Event::RunStarted {
agent_def_hash: "sha256:any".to_owned(),
input: json!("go"),
labels: None,
driven_by: Some(Performer::Client),
},
Event::SleepStarted { wake_at },
]
.into_iter()
.enumerate()
{
let envelope = EventEnvelope::new(run_id, SequenceNumber::new(seq as u64), NOW, event);
store.append(&envelope).await.expect("append");
}
run_id
}
fn counting_factory(inner: AgentFactory, builds: Arc<AtomicUsize>) -> AgentFactory {
Arc::new(move |definition| {
builds.fetch_add(1, Ordering::SeqCst);
let inner = inner.clone();
Box::pin(async move { inner(definition).await })
as Pin<Box<dyn std::future::Future<Output = _> + Send>>
})
}
async fn state_with_counter(
store: Arc<dyn EventStore>,
) -> (AppState, Arc<AtomicUsize>, wiremock::MockServer) {
let model = ScriptedModel::mount(vec![]).await;
let factory = agent_factory(
model.uri(),
"record",
Effect::Read,
CountBehavior::Record,
counter(),
);
let builds = Arc::new(AtomicUsize::new(0));
let state = app_state(store, counting_factory(factory, builds.clone()));
(state, builds, model)
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn a_run_past_its_deadline_is_woken_by_the_server_alone() {
let store = memory_store();
let (state, builds, _model) = state_with_counter(store.clone()).await;
let state = state.with_wake_interval(Duration::from_millis(20));
let server = TestServer::spawn(state).await;
let client = reqwest::Client::new();
let agent = register_agent(&client, &server.base, sample_toml(), None).await;
let builds_after_register = builds.load(Ordering::SeqCst);
let run_id = seed_sleeping(&store, &agent, NOW - time::Duration::hours(1)).await;
for _ in 0..300 {
if builds.load(Ordering::SeqCst) > builds_after_register {
return;
}
tokio::time::sleep(Duration::from_millis(10)).await;
}
panic!(
"the sweeper never re-drove run {}; it was due an hour before the injected clock",
run_id.as_uuid()
);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn a_pass_drives_the_due_run_and_leaves_the_rest_alone() {
let store = memory_store();
let (state, _builds, _model) = state_with_counter(store.clone()).await;
let server = TestServer::spawn(state.clone()).await;
let client = reqwest::Client::new();
let agent = register_agent(&client, &server.base, sample_toml(), None).await;
let due = seed_sleeping(&store, &agent, NOW - time::Duration::days(1)).await;
let not_due = seed_sleeping(&store, &agent, NOW + time::Duration::days(1)).await;
let driven = salvor_server::sweep(&state).await;
assert_eq!(driven, vec![due], "exactly the run past its deadline");
let untouched = store.read_log(not_due).await.expect("log reads");
assert_eq!(untouched.len(), 2, "the run not yet due is not appended to");
assert_eq!(
due_runs(store.as_ref(), NOW)
.await
.expect("selection reads")
.len(),
1,
"and it is still not due"
);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn a_run_already_being_driven_here_is_skipped() {
let store = memory_store();
let (state, builds, _model) = state_with_counter(store.clone()).await;
let server = TestServer::spawn(state.clone()).await;
let client = reqwest::Client::new();
let agent = register_agent(&client, &server.base, sample_toml(), None).await;
let run_id = seed_sleeping(&store, &agent, NOW - time::Duration::hours(6)).await;
state.begin_run(run_id);
let before = builds.load(Ordering::SeqCst);
assert!(
salvor_server::sweep(&state).await.is_empty(),
"a run with a driver on it is not driven again"
);
assert_eq!(
builds.load(Ordering::SeqCst),
before,
"and the skip happens before anything is rebuilt"
);
state.end_run(run_id);
assert_eq!(
salvor_server::sweep(&state).await,
vec![run_id],
"nothing about the run changed; only the claim did"
);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn a_client_driven_runs_due_timer_is_left_alone() {
let store = memory_store();
let (state, builds, _model) = state_with_counter(store.clone()).await;
let run_id = seed_sleeping(&store, "sha256:any", NOW - time::Duration::hours(1)).await;
state.lease_client_run(run_id, false);
let before = builds.load(Ordering::SeqCst);
assert!(
salvor_server::sweep(&state).await.is_empty(),
"a client-driven run's due timer is not driven from here"
);
assert_eq!(
builds.load(Ordering::SeqCst),
before,
"and nothing was rebuilt to try"
);
assert_eq!(
store.read_log(run_id).await.expect("log reads").len(),
2,
"the log is untouched"
);
assert_eq!(
due_runs(store.as_ref(), NOW).await.expect("reads").len(),
1,
"still due, waiting on the client"
);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn a_client_driven_run_known_only_from_its_log_is_left_alone() {
let store = memory_store();
let (state, builds, _model) = state_with_counter(store.clone()).await;
let run_id = seed_sleeping_client_driven(&store, NOW - time::Duration::hours(1)).await;
assert!(
!state.is_client_run(run_id),
"no lease here: the log is the only evidence the sweeper has"
);
let before = builds.load(Ordering::SeqCst);
assert!(
salvor_server::sweep(&state).await.is_empty(),
"a client-driven run's due timer is not driven from here, lease or no lease"
);
assert_eq!(
builds.load(Ordering::SeqCst),
before,
"and nothing was rebuilt to try"
);
assert_eq!(
store.read_log(run_id).await.expect("log reads").len(),
2,
"the log is untouched"
);
assert_eq!(
due_runs(store.as_ref(), NOW).await.expect("reads").len(),
1,
"still due, waiting on the client"
);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn an_unregistered_agent_leaves_the_run_asleep() {
let store = memory_store();
let (state, _builds, _model) = state_with_counter(store.clone()).await;
let run_id = seed_sleeping(
&store,
"sha256:never-registered-here",
NOW - time::Duration::hours(1),
)
.await;
assert!(
salvor_server::sweep(&state).await.is_empty(),
"the server cannot wake a run it has no definition for"
);
let log = store.read_log(run_id).await.expect("log reads");
assert_eq!(log.len(), 2, "the run's log is untouched");
assert_eq!(
due_runs(store.as_ref(), NOW).await.expect("reads").len(),
1,
"and it is still due, so a later pass gets another chance"
);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn an_unwakeable_run_is_warned_about_once_not_every_pass() {
let store = memory_store();
let (state, _builds, _model) = state_with_counter(store.clone()).await;
let run_id = seed_sleeping(
&store,
"sha256:never-registered-here",
NOW - time::Duration::hours(1),
)
.await;
assert!(
!state.unwakeable_warned(run_id),
"nothing has swept yet, so there is no record"
);
assert!(salvor_server::sweep(&state).await.is_empty());
assert!(
state.unwakeable_warned(run_id),
"the first pass that cannot wake it records the sighting"
);
assert!(salvor_server::sweep(&state).await.is_empty());
assert!(
state.unwakeable_warned(run_id),
"the record still stands after a second pass"
);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn the_warned_record_clears_once_the_run_wakes() {
let store = memory_store();
let (state, _builds, _model) = state_with_counter(store.clone()).await;
let hash = "sha256:not-registered-yet".to_owned();
let run_id = seed_sleeping(&store, &hash, NOW - time::Duration::hours(1)).await;
assert!(salvor_server::sweep(&state).await.is_empty());
assert!(
state.unwakeable_warned(run_id),
"the first pass records the sighting, hash unregistered"
);
state.register_agent(salvor_server::RegisteredAgent {
definition: salvor_server::AgentDefinition {
format: salvor_server::DefFormat::Toml,
body: sample_toml().as_bytes().to_vec(),
},
agent_hash: hash,
name: None,
});
assert_eq!(
salvor_server::sweep(&state).await,
vec![run_id],
"now buildable, the same due run wakes"
);
assert!(
!state.unwakeable_warned(run_id),
"a woken run's record does not carry forward"
);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn a_run_that_will_not_drive_does_not_stop_the_pass() {
let store = memory_store();
let (state, _builds, _model) = state_with_counter(store.clone()).await;
let server = TestServer::spawn(state.clone()).await;
let client = reqwest::Client::new();
let agent = register_agent(&client, &server.base, sample_toml(), None).await;
seed_sleeping(
&store,
"sha256:never-registered-here",
NOW - time::Duration::days(5),
)
.await;
let wakeable = seed_sleeping(&store, &agent, NOW - time::Duration::hours(1)).await;
assert_eq!(
salvor_server::sweep(&state).await,
vec![wakeable],
"the run after the failure is still driven"
);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn a_zero_interval_turns_the_sweeper_off() {
let store = memory_store();
let (state, builds, _model) = state_with_counter(store.clone()).await;
let state = state.with_wake_interval(Duration::ZERO);
let server = TestServer::spawn(state).await;
let client = reqwest::Client::new();
let agent = register_agent(&client, &server.base, sample_toml(), None).await;
let builds_after_register = builds.load(Ordering::SeqCst);
let run_id = seed_sleeping(&store, &agent, NOW - time::Duration::days(30)).await;
tokio::time::sleep(Duration::from_millis(300)).await;
assert_eq!(
builds.load(Ordering::SeqCst),
builds_after_register,
"nothing rebuilt the agent, so nothing swept"
);
assert_eq!(
store.read_log(run_id).await.expect("log reads").len(),
2,
"and the run is exactly as it was seeded"
);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn a_sleeping_runs_stream_ends_and_reports_its_deadline() {
let dir = tempdir().expect("tempdir");
let store = open_store(&dir.path().join("salvor.db"));
let (state, _builds, _model) = state_with_counter(store.clone()).await;
let state = state.with_wake_interval(Duration::ZERO);
let server = TestServer::spawn(state).await;
let wake_at = NOW + time::Duration::days(7);
let run_id = seed_sleeping(&store, "sha256:any", wake_at).await;
let client = reqwest::Client::new();
let mut stream = SseReader::open(
&client,
&server.base,
&run_id.as_uuid().to_string(),
None,
None,
None,
)
.await;
let frames = stream.read_to_end().await;
let end = frames.last().expect("the stream sends frames");
assert!(end.is_end(), "the stream closes on a resting run");
let status = &end.json()["status"];
assert_eq!(status["state"], "sleeping", "the state it rested at");
assert_eq!(
status["wake_at"], "2026-07-17T12:00:00Z",
"and the instant it may continue at"
);
assert_eq!(
frames.len(),
3,
"both recorded events, then the end frame: {frames:?}"
);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn an_early_resume_is_refused_and_a_due_one_drives() {
let store = memory_store();
let (state, _builds, _model) = state_with_counter(store.clone()).await;
let state = state.with_wake_interval(Duration::ZERO);
let server = TestServer::spawn(state).await;
let client = reqwest::Client::new();
let agent = register_agent(&client, &server.base, sample_toml(), None).await;
let early = seed_sleeping(&store, &agent, NOW + time::Duration::minutes(29)).await;
let refusal = client
.post(format!(
"{}/v1/runs/{}/resume",
server.base,
early.as_uuid()
))
.json(&json!({}))
.send()
.await
.expect("the request is made");
assert_eq!(refusal.status(), 409, "a state conflict, not a bad request");
let body: serde_json::Value = refusal.json().await.expect("a json body");
assert_eq!(body["error"]["code"], "still_sleeping");
assert_eq!(
body["error"]["details"]["wake_at"], "2026-07-10T12:29:00Z",
"the deadline it is waiting for"
);
assert_eq!(
body["error"]["details"]["remaining_seconds"], 1740,
"and how long is left of the wait"
);
assert_eq!(
store.read_log(early).await.expect("log reads").len(),
2,
"a refused resume records nothing"
);
let due = seed_sleeping(&store, &agent, NOW - time::Duration::minutes(1)).await;
let accepted = client
.post(format!("{}/v1/runs/{}/resume", server.base, due.as_uuid()))
.json(&json!({}))
.send()
.await
.expect("the request is made");
assert_eq!(
accepted.status(),
202,
"a due run is re-driven, which is the whole of waking"
);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn get_reports_a_sleeping_run_as_overdue_past_its_deadline() {
let store = memory_store();
let (state, _builds, _model) = state_with_counter(store.clone()).await;
let state = state.with_wake_interval(Duration::ZERO);
let server = TestServer::spawn(state).await;
let client = reqwest::Client::new();
let not_due = seed_sleeping(&store, "sha256:any", NOW + time::Duration::minutes(5)).await;
let overdue = seed_sleeping(&store, "sha256:any", NOW - time::Duration::minutes(2)).await;
let not_due_body: serde_json::Value = client
.get(format!("{}/v1/runs/{}", server.base, not_due.as_uuid()))
.send()
.await
.expect("the request is made")
.json()
.await
.expect("a json body");
assert_eq!(not_due_body["status"]["state"], "sleeping");
assert!(
not_due_body["status"].get("overdue").is_none(),
"not due yet, so no overdue key at all: {not_due_body}"
);
assert!(not_due_body["status"].get("overdue_seconds").is_none());
let overdue_body: serde_json::Value = client
.get(format!("{}/v1/runs/{}", server.base, overdue.as_uuid()))
.send()
.await
.expect("the request is made")
.json()
.await
.expect("a json body");
assert_eq!(overdue_body["status"]["state"], "sleeping");
assert_eq!(overdue_body["status"]["overdue"], true);
assert_eq!(
overdue_body["status"]["overdue_seconds"], 120,
"whole seconds since the deadline passed"
);
}