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(