Run one projection pass at a time per run

The startup pass called the pass directly while a signalled pass could
run for the same run, so both read one committed stream sequence and
the second insert into the stream failed on its primary key, which
stopped the restarted server. Passes now take a per-run lock, and a run
whose startup pass fails is logged and left for its next signal instead
of stopping the server. A test races four passes, a signal and the
startup pass over one run and checks the stream stays contiguous.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
This commit is contained in:
Bryan Helmkamp 2026-09-17 22:56:30 -04:00
parent e0b546d465
commit 4a20acafc7
No known key found for this signature in database
2 changed files with 58 additions and 4 deletions

View file

@ -107,6 +107,8 @@ pub struct PassReport {
pub struct StartupReport {
pub runs: usize,
pub projected: usize,
/// Runs whose pass failed and was left for the next signal.
pub failed: usize,
}
/// Why a pass could not run or commit.
@ -149,6 +151,9 @@ pub struct Projector {
store: SqliteRunStore,
platform: PlatformRecordStore,
slots: Mutex<HashMap<RunId, Slot>>,
/// One pass at a time per run: a signalled pass and the startup pass
/// over the same run never interleave their reads and writes.
passes: Mutex<HashMap<RunId, Arc<tokio::sync::Mutex<()>>>>,
/// Test-only: stop the next pass after its reads, before its view
/// transaction, as a crash there would.
fault: AtomicBool,
@ -172,6 +177,7 @@ impl Projector {
records,
pool: views,
slots: Mutex::default(),
passes: Mutex::default(),
fault: AtomicBool::new(false),
})
}
@ -260,9 +266,22 @@ impl Projector {
continue;
};
report.runs += 1;
let pass = self.project_run(run_id).await?;
if !pass.skipped {
report.projected += 1;
match self.project_run(run_id).await {
Ok(pass) => {
if !pass.skipped {
report.projected += 1;
}
}
// One run's view trailing never stops the server: the next
// signal for the run retries its pass.
Err(error) => {
warn!(
run_id = %run_id,
error = %collect_chain(&error).join(": "),
"Petri projection pass failed at startup; the next signal retries it"
);
report.failed += 1;
}
}
}
if report.projected > 0 {
@ -275,8 +294,10 @@ impl Projector {
Ok(report)
}
/// One view pass for the run.
/// One view pass for the run. Passes over one run run one at a time.
pub async fn project_run(&self, run_id: RunId) -> Result<PassReport, ProjectError> {
let pass = Arc::clone(lock(&self.passes).entry(run_id).or_default());
let _one_at_a_time = pass.lock().await;
let stored = self.load_view(&run_id).await?;
let key = RunKey::new(run_id.to_string());
let platform_head = self

View file

@ -525,6 +525,39 @@ async fn the_startup_pass_catches_up_a_view_nobody_signalled() {
assert_eq!((again.runs, again.projected), (1, 0), "nothing left to do");
}
/// Passes over one run never interleave: the startup pass and a signalled
/// pass racing over the same run commit one stream, contiguous and without
/// a duplicate.
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn concurrent_passes_over_one_run_commit_one_contiguous_stream() {
if host_plugin().is_none() {
return;
}
let scenario = parallel_scenario().await;
run_unobserved(&scenario).await;
let projector = Projector::new(scenario.pool.clone(), scenario.pool.clone());
let mut passes = Vec::new();
for _ in 0..4 {
let projector = Arc::clone(&projector);
let run_id = scenario.run_id;
passes.push(tokio::spawn(
async move { projector.project_run(run_id).await },
));
}
projector.signal(scenario.run_id);
projector
.startup_pass()
.await
.expect("the startup pass runs");
for pass in passes {
pass.await
.expect("the pass task joins")
.expect("a concurrent pass commits or skips");
}
projector.settle(scenario.run_id).await;
assert_view_equals_rebuild(&scenario.pool, scenario.run_id).await;
}
/// Every Petri record of the run, as `(log, seq, recorded_at, record_json)`.
async fn petri_rows(pool: &DbPool, run_id: RunId) -> Vec<(String, i64, i64, String)> {
sqlx::query_as(