use std::sync::Barrier;
use std::sync::atomic::{AtomicI64, Ordering};
use rusqlite::Connection;
use sloop::domain::ticket::TicketState;
use sloop::domain::trigger::TriggerKind;
use sloop::outcome::Outcome;
use sloop::run_store::{Exit, RunAdmission, RunExit, RunStart, Start};
use sloop::work_state::trigger::NewTrigger;
use tempfile::TempDir;
use crate::TestStore;
use crate::model::extra_invariants;
const THREADS: usize = 8;
struct Arena {
_directory: TempDir,
db_path: std::path::PathBuf,
clock: AtomicI64,
}
impl Arena {
fn new() -> Self {
let directory = TempDir::new().expect("create tempdir");
let db_path = directory.path().join("sloop.db");
let store = TestStore::open(&db_path, 1_000);
store
.local
.insert_local_project("default", "projects/default.md", "Default", 1_000)
.expect("insert project");
Self {
_directory: directory,
db_path,
clock: AtomicI64::new(2_000),
}
}
fn now(&self) -> i64 {
self.clock.fetch_add(1, Ordering::Relaxed)
}
fn open(&self) -> TestStore {
TestStore::open(&self.db_path, self.now())
}
fn add_ticket(&self, store: &TestStore, ticket: &str) {
store
.local
.insert_local_ticket(
ticket,
"default",
&format!("tickets/{ticket}.md"),
&format!("Ticket {ticket}"),
&[],
&format!("sloop/{ticket}"),
Some("opencode"),
None,
None,
"default",
TicketState::Ready,
self.now(),
)
.expect("insert ticket");
}
fn add_trigger(&self, store: &TestStore, id: &str, ticket: &str) {
store
.local
.insert_trigger(
&NewTrigger {
id,
kind: TriggerKind::Immediate,
ticket_id: Some(ticket),
project_id: None,
eligible_at_ms: None,
interval_ms: None,
},
self.now(),
)
.expect("insert trigger");
}
fn check_invariants(&self) {
let connection = Connection::open(&self.db_path).expect("open check connection");
extra_invariants(&connection);
}
}
fn run_admission<'a>(ticket: &'a str, run_id: &'a str, trigger_id: &'a str) -> RunAdmission<'a> {
RunAdmission {
ticket_id: ticket,
run_id,
trigger_id,
flow_json: "{}",
ticket_json: "{}",
}
}
fn run_start(run_id: &str) -> RunStart<'_> {
RunStart {
run_id,
branch: "sloop/branch",
worktree_path: "/tmp/worktree",
pid: 4_242,
pid_start_time: Some(7),
process_group_id: 4_242,
worker_token: "token",
worker_socket_path: "/tmp/worker.sock",
}
}
fn run_exit(run_id: &str) -> RunExit<'_> {
RunExit {
run_id,
attempt: 1,
exit_code: Some(0),
capture_complete: true,
commits_json: "{}",
vendor_error: None,
cooldown_until_ms: None,
}
}
#[test]
fn simultaneous_claims_grant_exactly_one_winner() {
let arena = Arena::new();
let setup = arena.open();
for round in 0..20 {
let ticket = format!("T{round}");
let trigger = format!("TR{round}");
arena.add_ticket(&setup, &ticket);
arena.add_trigger(&setup, &trigger, &ticket);
let barrier = Barrier::new(THREADS);
let grants: Vec<bool> = std::thread::scope(|scope| {
let handles: Vec<_> = (0..THREADS)
.map(|thread| {
let (arena, ticket, trigger, barrier) = (&arena, &ticket, &trigger, &barrier);
scope.spawn(move || {
let store = arena.open();
let run_id = format!("{ticket}-R{thread}");
barrier.wait();
crate::claim(
&store,
&run_admission(ticket, &run_id, trigger),
60_000,
arena.now(),
)
.is_some()
})
})
.collect();
handles
.into_iter()
.map(|h| h.join().expect("join"))
.collect()
});
let winners = grants.iter().filter(|granted| **granted).count();
assert_eq!(winners, 1, "round {round}: ticket claimed {winners} times");
arena.check_invariants();
}
}
#[test]
fn simultaneous_exit_checkpoints_grant_exactly_one_owner() {
let arena = Arena::new();
let setup = arena.open();
for round in 0..20 {
let ticket = format!("T{round}");
let trigger = format!("TR{round}");
let run_id = format!("{ticket}-R0");
arena.add_ticket(&setup, &ticket);
arena.add_trigger(&setup, &trigger, &ticket);
assert!(
crate::claim(
&setup,
&run_admission(&ticket, &run_id, &trigger),
60_000,
arena.now()
)
.is_some()
);
let started = setup
.runs
.start(&run_start(&run_id), arena.now())
.expect("start");
assert_eq!(started, Start::Granted);
let barrier = Barrier::new(THREADS);
let grants: Vec<bool> = std::thread::scope(|scope| {
let handles: Vec<_> = (0..THREADS)
.map(|_| {
let (arena, run_id, barrier) = (&arena, &run_id, &barrier);
scope.spawn(move || {
let store = arena.open();
barrier.wait();
let exit = store
.runs
.record_exit(&run_exit(run_id), arena.now())
.expect("record_exit must grant or deny, never fail");
matches!(exit, Exit::Granted)
})
})
.collect();
handles
.into_iter()
.map(|h| h.join().expect("join"))
.collect()
});
let owners = grants.iter().filter(|granted| **granted).count();
assert_eq!(owners, 1, "round {round}: {owners} threads own the walk");
arena.check_invariants();
}
}
#[test]
fn simultaneous_settlements_land_exactly_once() {
const OUTCOMES: [Outcome; 4] = [
Outcome::Merged,
Outcome::Failed,
Outcome::Cancelled,
Outcome::RateLimited,
];
let arena = Arena::new();
let setup = arena.open();
for round in 0..20 {
let ticket = format!("T{round}");
let trigger = format!("TR{round}");
let run_id = format!("{ticket}-R0");
arena.add_ticket(&setup, &ticket);
arena.add_trigger(&setup, &trigger, &ticket);
crate::claim(
&setup,
&run_admission(&ticket, &run_id, &trigger),
60_000,
arena.now(),
)
.expect("claim");
setup
.runs
.start(&run_start(&run_id), arena.now())
.expect("start");
setup
.runs
.record_exit(&run_exit(&run_id), arena.now())
.expect("record exit");
let barrier = Barrier::new(THREADS);
let landed: Vec<Option<Outcome>> = std::thread::scope(|scope| {
let handles: Vec<_> = (0..THREADS)
.map(|thread| {
let (arena, run_id, barrier) = (&arena, &run_id, &barrier);
scope.spawn(move || {
let store = arena.open();
let outcome = OUTCOMES[thread % OUTCOMES.len()];
barrier.wait();
let settled = crate::settle(&store, run_id, outcome, arena.now());
settled.then_some(outcome)
})
})
.collect();
handles
.into_iter()
.map(|h| h.join().expect("join"))
.collect()
});
let winners: Vec<Outcome> = landed.into_iter().flatten().collect();
assert_eq!(winners.len(), 1, "round {round}: {winners:?} all landed");
let winner = winners[0];
let connection = Connection::open(&arena.db_path).expect("open check connection");
let run_state: String = connection
.query_row("SELECT state FROM runs WHERE id = ?1", [&run_id], |row| {
row.get(0)
})
.expect("run row");
assert_eq!(run_state, winner.as_str(), "run state is the winner's");
let ticket_state: String = connection
.query_row(
"SELECT state FROM tickets WHERE id = ?1",
[&ticket],
|row| row.get(0),
)
.expect("ticket row");
assert_eq!(
ticket_state,
TicketState::after_outcome(winner).as_str(),
"ticket state is the winner's"
);
arena.check_invariants();
}
}
#[test]
fn uncoordinated_lifecycle_hammer_preserves_invariants() {
const POOL: [&str; 4] = ["P0", "P1", "P2", "P3"];
const ITERATIONS: usize = 60;
let arena = Arena::new();
let setup = arena.open();
for ticket in POOL {
arena.add_ticket(&setup, ticket);
}
let barrier = Barrier::new(THREADS);
let completed: Vec<usize> = std::thread::scope(|scope| {
let handles: Vec<_> = (0..THREADS)
.map(|thread| {
let (arena, barrier) = (&arena, &barrier);
scope.spawn(move || {
let store = arena.open();
let mut completed = 0;
barrier.wait();
for iteration in 0..ITERATIONS {
let ticket = POOL[(thread + iteration) % POOL.len()];
let trigger = format!("{ticket}-{thread}-{iteration}");
let run_id = format!("{ticket}-{thread}-{iteration}-run");
arena.add_trigger(&store, &trigger, ticket);
if crate::claim(
&store,
&run_admission(ticket, &run_id, &trigger),
60_000,
arena.now(),
)
.is_none()
{
continue;
}
assert_eq!(
store
.runs
.start(&run_start(&run_id), arena.now())
.expect("start"),
Start::Granted
);
assert_eq!(
store
.runs
.record_exit(&run_exit(&run_id), arena.now())
.expect("record exit"),
Exit::Granted
);
let outcome = if iteration % 2 == 0 {
Outcome::Cancelled
} else {
Outcome::Merged
};
assert!(
crate::settle(&store, &run_id, outcome, arena.now()),
"the owner's settlement must land"
);
completed += 1;
}
completed
})
})
.collect();
handles
.into_iter()
.map(|h| h.join().expect("join"))
.collect()
});
let total: usize = completed.iter().sum();
assert!(total > 0, "contention must not starve every thread");
arena.check_invariants();
let connection = Connection::open(&arena.db_path).expect("open check connection");
let live_runs: i64 = connection
.query_row(
"SELECT COUNT(*) FROM runs WHERE exited_at_ms IS NULL",
[],
|row| row.get(0),
)
.expect("count live runs");
assert_eq!(live_runs, 0, "the hammer settles every run it starts");
let leases: i64 = connection
.query_row("SELECT COUNT(*) FROM leases", [], |row| row.get(0))
.expect("count leases");
assert_eq!(leases, 0, "no lease survives its settled run");
}