From 4a20acafc78c9432aafd05ce569faf2a67e922af Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Thu, 17 Sep 2026 22:56:30 -0400 Subject: [PATCH] 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 --- lib/components/fabro-petri/src/projector.rs | 29 +++++++++++++--- .../fabro-petri/tests/projection.rs | 33 +++++++++++++++++++ 2 files changed, 58 insertions(+), 4 deletions(-) diff --git a/lib/components/fabro-petri/src/projector.rs b/lib/components/fabro-petri/src/projector.rs index 6421180a3..8c5e7e96f 100644 --- a/lib/components/fabro-petri/src/projector.rs +++ b/lib/components/fabro-petri/src/projector.rs @@ -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>, + /// 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>>>, /// 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 { + 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 diff --git a/lib/components/fabro-petri/tests/projection.rs b/lib/components/fabro-petri/tests/projection.rs index e231515f7..69057e9c5 100644 --- a/lib/components/fabro-petri/tests/projection.rs +++ b/lib/components/fabro-petri/tests/projection.rs @@ -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(