From 1ed30db1d6606e77f32a9e854cbc511087751ae3 Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Fri, 18 Sep 2026 19:58:03 -0400 Subject: [PATCH] Keep a live run's replay and view between projector passes Each projector pass replayed the run whole through `replay_since` to rebuild the engine and invocation state the derivation needs, so a pass cost the run's length. The projector now keeps, per live run and behind the run's pass lock, Petri's `RunReplay` and the view as the last committed pass left it (`projector/cache.rs`), and a pass advances the replay over the records past the ones it consumed: it reads and folds only the new records. The SQLite store answers `read_from` with `seq >= ?`, and the signalling store forwards it. The rules hold as before. Records first: the cache moves only after the view transaction commits, and a pass that commits nothing (a platform record landed under it, a fault before the transaction) keeps the events it derived as pending for the next pass. The cache is never checkpointed and never a source of facts: it is dropped when the run records its finish, after ten idle minutes, when the stored view moves under it, when the run is deleted, and with the process; the first pass after that rebuilds it by a full replay, filtered to the held positions. A torn tail fails the advance, which leaves the replay where it stood, so the view holds and the pass is retried. `PassReport::replayed_records` says how many records a pass fed through the derivation. Two projection tests: the records of a parallel run land in batches and every pass replays at most its batch, with the finished run's cache dropped; a restart and the idle period drop the cache and the next pass replays the run so far once, then only its new records, and the view equals the rebuild throughout. `test_support` exposes whether a cache is held and the idle sweep. Co-Authored-By: Claude Fable 5.1 --- lib/components/fabro-petri/src/projector.rs | 244 +++++++++++++----- .../fabro-petri/src/projector/cache.rs | 132 ++++++++++ lib/components/fabro-petri/src/run_store.rs | 8 +- .../fabro-petri/src/test_support.rs | 15 ++ .../fabro-petri/tests/projection.rs | 181 ++++++++++++- 5 files changed, 515 insertions(+), 65 deletions(-) create mode 100644 lib/components/fabro-petri/src/projector/cache.rs diff --git a/lib/components/fabro-petri/src/projector.rs b/lib/components/fabro-petri/src/projector.rs index 39da40bfe..2dea3dcf0 100644 --- a/lib/components/fabro-petri/src/projector.rs +++ b/lib/components/fabro-petri/src/projector.rs @@ -12,14 +12,29 @@ //! Petri log, the last platform record consumed, and the delivery sequence //! (`stream_seq`) it assigned to each item. The view therefore trails a //! committed record and never leads one. No projection state of Petri's is -//! checkpointed: each pass replays the run through `replay_since`, which -//! rebuilds the engine and invocation state the derivation needs and -//! delivers only the events past the held positions. +//! checkpointed: a pass derives the events past the held positions from +//! the records alone, and every stored view equals a full replay +//! (`replay_run`) of the records it holds. //! //! A pass that finds new platform records committed between its read and //! its write leaves the view alone and runs again, so the `runs` row never //! moves backwards behind a concurrent lifecycle write. //! +//! # The live run's cache +//! +//! A pass keeps in memory, per live run, Petri's replay of the run (a +//! `RunReplay`: the coordinator state, each execution's engine state, the +//! projection) and the view as the pass last committed it, so the next +//! pass reads and folds only the records past the ones the view holds and +//! costs the new records, not the run's length. The cache is never a +//! source of facts and never checkpointed: it is dropped when the run +//! records its finish, after ten idle minutes, when the stored view moves +//! under it, when the run is deleted, and with the process, and the first +//! pass after that rebuilds it by a full replay. A pass that commits +//! nothing (a platform record landed under it, or it failed before its +//! view transaction) keeps the events it derived for the next pass, so +//! nothing is derived twice or lost. +//! //! # Where it runs //! //! In the server. [`Projector::signal`] schedules a pass for a run: the @@ -38,6 +53,8 @@ //! as incomplete with the replay's error; `inspect_run` decides //! completeness once the run has recorded its finish. +mod cache; + use std::collections::{BTreeMap, BTreeSet, HashMap}; use std::sync::atomic::{AtomicBool, Ordering}; use std::sync::{Arc, Mutex, MutexGuard, PoisonError}; @@ -53,10 +70,11 @@ use petri_execution::{Access, CoordinatorEvent, RunKey, RunStore as _, inspect}; use petri_runtime::engine::Event; use petri_store::StoreError; use serde::{Deserialize, Serialize}; -use tokio::sync::{Mutex as AsyncMutex, broadcast}; +use tokio::sync::broadcast; use tokio::time; use tracing::{debug, info, warn}; +use self::cache::{Caches, IDLE, RunCache}; use crate::SqliteRunStore; use crate::projection::{self, FoldState, Item, RecordHealth, RunView}; @@ -98,6 +116,11 @@ pub struct PassReport { pub contended: bool, pub petri_events: usize, pub platform_records: usize, + /// How many of the run's records the pass fed through Petri's + /// derivation, before the held positions trimmed their events: the + /// pass's cost. The records past the cache for a live run, the whole + /// run for a pass that rebuilt it. + pub replayed_records: usize, /// The last delivery sequence the view holds. pub stream_seq: u64, pub positions: Positions, @@ -129,6 +152,7 @@ pub enum ProjectError { } /// The stored view of a run, as the projection tables hold it. +#[derive(Clone)] struct StoredView { view: RunView, positions: Positions, @@ -148,21 +172,22 @@ struct Slot { /// the server both are the one database; a test may hand it the run /// summary store's own pool for the views. pub struct Projector { - records: DbPool, - pool: DbPool, - 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>>>, + records: DbPool, + pool: DbPool, + 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), and the + /// cache each live run's passes continue from. + pub(crate) caches: Caches, /// Test-only: stop the next pass after its reads, before its view /// transaction, as a crash there would. - fault: AtomicBool, + fault: AtomicBool, /// Sent after each committed pass that wrote stream rows: the run whose /// stream grew. A wake-up for the stream's readers, never a source of /// facts; a reader that lags re-reads from its cursor. - committed: broadcast::Sender, + committed: broadcast::Sender, } impl std::fmt::Debug for Projector { @@ -183,7 +208,7 @@ impl Projector { records, pool: views, slots: Mutex::default(), - passes: Mutex::default(), + caches: Caches::default(), fault: AtomicBool::new(false), committed: broadcast::channel(COMMIT_SIGNAL_CAPACITY).0, }) @@ -214,6 +239,11 @@ impl Projector { /// and its stream. The caller has ended the run's worker, so no writer /// holds the lease. pub async fn delete_run(&self, run_id: RunId) -> Result<(), ProjectError> { + // Under the run's pass lock: no pass reads the rows being deleted, + // and no cache outlives them. + let pass = self.caches.pass_of(run_id); + let mut cache = pass.lock().await; + *cache = None; let id = run_id.to_string(); let mut views = self.pool.begin().await.map_err(ProjectError::Database)?; for delete in [ @@ -370,9 +400,34 @@ impl Projector { /// 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?; + self.caches.sweep(IDLE); + let pass = self.caches.pass_of(run_id); + let mut slot = pass.lock().await; + // The view tables are the source of truth: a cache that no longer + // describes them (another projector committed a pass) is dropped. + let (positions, stream_seq) = stored_positions(&self.pool, run_id) + .await? + .unwrap_or_default(); + let mut run = match slot.take() { + Some(cache) if cache.matches(&positions, stream_seq) => cache, + Some(_) => { + debug!(run_id = %run_id, "the stored view moved under the run's cache; rebuilding it"); + RunCache::over(self.load_view(&run_id).await?) + } + None => RunCache::over(self.load_view(&run_id).await?), + }; + let report = self.pass(run_id, &mut run).await; + // A finished run's records are complete: its cache is dropped, and + // the passes its late platform records take rebuild the view whole. + if !run.view.view.state.finished_run() { + *slot = Some(run); + } + report + } + + /// The pass over the run's cache: read what is committed past the + /// positions the cache's view holds, fold it, and write the view. + async fn pass(&self, run_id: RunId, run: &mut RunCache) -> Result { let key = RunKey::new(run_id.to_string()); let platform_head = self .platform @@ -381,6 +436,7 @@ impl Projector { .map_err(ProjectError::Store)? .unwrap_or(0); let petri_heads = self.petri_heads(&run_id).await?; + let stored = &run.view; let at_head = platform_head == stored.positions.platform_seq && petri_heads.iter().all(|(log, head)| { stored @@ -396,25 +452,35 @@ impl Projector { contended: false, petri_events: 0, platform_records: 0, + replayed_records: 0, stream_seq: stored.stream_seq, - positions: stored.positions, - health: stored.view.state.health, + positions: stored.positions.clone(), + health: stored.view.state.health.clone(), }); } - let StoredView { - mut view, - mut positions, - mut stream_seq, - } = stored; let platform_records = self .platform - .read_after(&run_id, positions.platform_seq) + .read_after(&run_id, stored.positions.platform_seq) .await .map_err(ProjectError::Store)?; + let mut replayed_records = 0; let (events, replay_failure) = match self.store.open(&key, Access::Read).await { - Ok(logs) => match events::replay_since(&*logs, &positions.held()).await { - Ok(events) => (events, None), + Ok(logs) => match run.replay.advance(&*logs).await { + Ok(new) => { + replayed_records = new.iter().filter(|event| event.id.index == 0).count(); + // A rebuilt replay derives the run whole: only the events + // past the view's positions are new to it. + let held = run.view.positions.held(); + let mut events = std::mem::take(&mut run.pending); + events.extend(new.into_iter().filter(|event| { + held.get(&event.id.source) + .is_none_or(|last| event.id > *last) + })); + (events, None) + } + // The replay stood still and is retried by the next pass; + // what it derived before stays pending. Err(error) => { let chain = collect_chain(&error).join(": "); warn!(run_id = %run_id, error = %chain, "Petri run does not replay; the view holds"); @@ -425,6 +491,9 @@ impl Projector { Err(error) => return Err(ProjectError::Open(error)), }; + let mut view = run.view.view.clone(); + let mut positions = run.view.positions.clone(); + let mut stream_seq = run.view.stream_seq; let run_finished = view.state.finished_run() || events.iter().any(|event| { matches!( @@ -475,12 +544,85 @@ impl Projector { }; rows.push(row); } + drop(items); view.state.health = self.health(&key, &view.state, replay_failure).await?; if self.fault.swap(false, Ordering::SeqCst) { + run.pending = events; return Err(ProjectError::Injected); } + let written = self + .write_view( + run_id, + &view, + &positions, + stream_seq, + &rows, + platform_head_seen, + ) + .await; + match written { + Ok(true) => {} + Ok(false) => { + debug!(run_id = %run_id, "platform records landed during the pass; running it again"); + run.pending = events; + return Ok(PassReport { + run_id, + skipped: false, + contended: true, + petri_events: 0, + platform_records: 0, + replayed_records, + stream_seq: 0, + positions: Positions::default(), + health: RecordHealth::default(), + }); + } + Err(error) => { + run.pending = events; + return Err(error); + } + } + debug!( + run_id = %run_id, + petri_events = events.len(), + platform_records = platform_records.len(), + replayed_records, + stream_seq, + "Petri projection pass committed" + ); + let petri_events = events.len(); + let health = view.state.health.clone(); + run.committed(view, positions.clone(), stream_seq); + if !rows.is_empty() { + // No receiver is not an error: nobody follows the stream. + let _ = self.committed.send(run_id); + } + Ok(PassReport { + run_id, + skipped: false, + contended: false, + petri_events, + platform_records: platform_records.len(), + replayed_records, + stream_seq, + positions, + health, + }) + } + /// The view transaction: the projection row, the stream rows and the + /// `runs` row, committed together, unless a platform record landed + /// since the pass read them (`false`: the view is left alone). + async fn write_view( + &self, + run_id: RunId, + view: &RunView, + positions: &Positions, + stream_seq: u64, + rows: &[StreamRow], + platform_head_seen: u64, + ) -> Result { let mut tx = self .pool .begin_with("BEGIN IMMEDIATE") @@ -494,23 +636,13 @@ impl Projector { .await .map_err(ProjectError::Database)?; if u64::try_from(head_now).unwrap_or(0) != platform_head_seen { - debug!(run_id = %run_id, "platform records landed during the pass; running it again"); drop(tx); - return Ok(PassReport { - run_id, - skipped: false, - contended: true, - petri_events: 0, - platform_records: 0, - stream_seq: 0, - positions: Positions::default(), - health: RecordHealth::default(), - }); + return Ok(false); } let projection_json = serde_json::to_string(&view.projection).map_err(ProjectError::Encode)?; let fold_json = serde_json::to_string(&view.state).map_err(ProjectError::Encode)?; - let positions_json = serde_json::to_string(&positions).map_err(ProjectError::Encode)?; + let positions_json = serde_json::to_string(positions).map_err(ProjectError::Encode)?; sqlx::query( "INSERT INTO petri_projection (run_id, projection_json, fold_json, positions_json, \ stream_seq, updated_at_ms) VALUES (?, ?, ?, ?, ?, ?) ON CONFLICT(run_id) DO UPDATE \ @@ -527,7 +659,7 @@ impl Projector { .execute(&mut *tx) .await .map_err(ProjectError::Database)?; - for row in &rows { + for row in rows { sqlx::query( "INSERT INTO petri_stream (run_id, stream_seq, item_kind, item_id, event_json) \ VALUES (?, ?, ?, ?, ?)", @@ -547,27 +679,7 @@ impl Projector { .map_err(ProjectError::Store)?; } tx.commit().await.map_err(ProjectError::Database)?; - debug!( - run_id = %run_id, - petri_events = events.len(), - platform_records = platform_records.len(), - stream_seq, - "Petri projection pass committed" - ); - if !rows.is_empty() { - // No receiver is not an error: nobody follows the stream. - let _ = self.committed.send(run_id); - } - Ok(PassReport { - run_id, - skipped: false, - contended: false, - petri_events: events.len(), - platform_records: platform_records.len(), - stream_seq, - positions, - health: view.state.health.clone(), - }) + Ok(true) } /// The stored view of the run, or an empty one. @@ -729,6 +841,14 @@ impl petri_execution::RunLogs for SignallingLogs { self.inner.read(log).await } + async fn read_from( + &self, + log: &petri_execution::LogId, + seq: u64, + ) -> Result, StoreError> { + self.inner.read_from(log, seq).await + } + async fn put_blob(&self, bytes: &[u8]) -> Result { self.inner.put_blob(bytes).await } diff --git a/lib/components/fabro-petri/src/projector/cache.rs b/lib/components/fabro-petri/src/projector/cache.rs new file mode 100644 index 000000000..39f242f17 --- /dev/null +++ b/lib/components/fabro-petri/src/projector/cache.rs @@ -0,0 +1,132 @@ +//! The state a live run's passes continue from, kept in memory between +//! passes: Petri's replay of the run (the coordinator state, each +//! execution's engine state, the projection) and Fabro's view as the last +//! committed pass left it. With it a pass reads and folds only the records +//! past the ones the view holds, so its cost is the new records', not the +//! run's. +//! +//! The cache is never a source of facts. It is dropped when the run +//! records its finish, when it has not been used for [`IDLE`], when the +//! stored view moves under it (another projector committed a pass), when +//! the run is deleted, and with the process; the first pass after that +//! rebuilds it by a full replay, which is what every pass did before the +//! cache existed. A replay that fails (a torn tail) stands still and is +//! retried by the next pass. Every pass, cached or not, commits the same +//! rows: the rebuild test in `tests/projection.rs` compares the two. + +use std::collections::HashMap; +use std::sync::{Arc, Mutex, MutexGuard, PoisonError}; +use std::time::{Duration, Instant}; + +use fabro_types::RunId; +use petri_execution::events::{RunEvent, RunReplay}; +use tokio::sync::Mutex as AsyncMutex; + +use super::{Positions, StoredView}; +use crate::projection::RunView; + +/// How long a run's cache is kept after its last pass. A run blocked on a +/// question for longer pays one full replay when its next record lands. +pub(super) const IDLE: Duration = Duration::from_secs(10 * 60); + +/// One live run's cache. +pub(super) struct RunCache { + pub(super) replay: RunReplay, + /// Events derived by an earlier pass that committed nothing: a pass + /// that found a platform record landing under it, or that failed + /// before its view transaction. They lead the next pass's events. + pub(super) pending: Vec, + /// The view as the last committed pass left it, with its positions. + pub(super) view: StoredView, +} + +impl RunCache { + /// A cache over the stored view, with a replay that has consumed + /// nothing: the first advance replays the run whole. + pub(super) fn over(view: StoredView) -> Self { + Self { + replay: RunReplay::new(), + pending: Vec::new(), + view, + } + } + + /// Whether the cache still describes the stored view: its positions and + /// delivery sequence are the ones the view tables hold. + pub(super) fn matches(&self, positions: &Positions, stream_seq: u64) -> bool { + self.view.stream_seq == stream_seq + && self.view.positions.platform_seq == positions.platform_seq + && self.view.positions.held() == positions.held() + } + + /// The pass committed: the view moved on, and nothing is pending. + pub(super) fn committed(&mut self, view: RunView, positions: Positions, stream_seq: u64) { + self.view = StoredView { + view, + positions, + stream_seq, + }; + self.pending.clear(); + } +} + +/// The caches of every run the projector passed over, each behind the +/// run's pass lock, so a pass and a sweep never race over one cache. +#[derive(Default)] +pub(crate) struct Caches { + runs: Mutex>, +} + +struct Entry { + pass: Arc>>, + touched: Instant, +} + +impl Caches { + /// The run's pass lock, holding its cache if one is kept; the run counts + /// as used now. + pub(super) fn pass_of(&self, run_id: RunId) -> Arc>> { + let mut runs = lock(&self.runs); + let entry = runs.entry(run_id).or_insert_with(|| Entry { + pass: Arc::default(), + touched: Instant::now(), + }); + entry.touched = Instant::now(); + Arc::clone(&entry.pass) + } + + /// Drop the cache of every run not used for `idle`, and forget the runs + /// with no cache and no pass under way. A run whose pass is running is + /// in use and left alone. How many caches were dropped. + pub(crate) fn sweep(&self, idle: Duration) -> usize { + let mut runs = lock(&self.runs); + let mut dropped = 0; + runs.retain(|_, entry| { + if entry.touched.elapsed() < idle { + return true; + } + let Ok(mut cache) = entry.pass.try_lock() else { + return true; + }; + if cache.take().is_some() { + dropped += 1; + } + drop(cache); + // An `Arc` held elsewhere is a pass about to take the lock: the + // entry stays so the run keeps one lock. + Arc::strong_count(&entry.pass) > 1 + }); + dropped + } + + /// Whether a cache is kept for the run: a test's view of the cache. + pub(crate) fn holds(&self, run_id: RunId) -> bool { + let runs = lock(&self.runs); + runs.get(&run_id) + .is_some_and(|entry| entry.pass.try_lock().is_ok_and(|cache| cache.is_some())) + } +} + +fn lock(mutex: &Mutex) -> MutexGuard<'_, T> { + mutex.lock().unwrap_or_else(PoisonError::into_inner) +} diff --git a/lib/components/fabro-petri/src/run_store.rs b/lib/components/fabro-petri/src/run_store.rs index 89d55e21d..d1eae1bdb 100644 --- a/lib/components/fabro-petri/src/run_store.rs +++ b/lib/components/fabro-petri/src/run_store.rs @@ -513,11 +513,17 @@ impl RunLogs for SqliteRunLogs { } async fn read(&self, log: &LogId) -> Result, StoreError> { + self.read_from(log, 0).await + } + + async fn read_from(&self, log: &LogId, seq: u64) -> Result, StoreError> { let rows: Vec = sqlx::query_scalar( - "SELECT record_json FROM petri_records WHERE run_id = ? AND log = ? ORDER BY seq", + "SELECT record_json FROM petri_records WHERE run_id = ? AND log = ? AND seq >= ? ORDER \ + BY seq", ) .bind(self.key.as_str()) .bind(log_id_text(log)) + .bind(i64::try_from(seq).unwrap_or(i64::MAX)) .fetch_all(&self.shared.pool) .await .map_err(|cause| self.backend("read a log", cause))?; diff --git a/lib/components/fabro-petri/src/test_support.rs b/lib/components/fabro-petri/src/test_support.rs index 9779fd2ea..7174cfbc5 100644 --- a/lib/components/fabro-petri/src/test_support.rs +++ b/lib/components/fabro-petri/src/test_support.rs @@ -6,6 +6,7 @@ use std::collections::HashMap; use std::sync::{Mutex, MutexGuard, PoisonError}; +use std::time::Duration; use async_trait::async_trait; use bytes::Bytes; @@ -16,6 +17,20 @@ pub use petri_testkit::run_store; use crate::blobs::Blobs; use crate::platform_records::{PlatformRecordError, PlatformRecords}; +use crate::projector::Projector; + +/// Whether the projector keeps a cache for the run: the replay and the +/// view its passes continue from. +#[must_use] +pub fn cache_held(projector: &Projector, run_id: RunId) -> bool { + projector.caches.holds(run_id) +} + +/// Drop the projector's caches not used for `idle`, as its passes do +/// after the documented idle period; how many were dropped. +pub fn drop_idle_caches(projector: &Projector, idle: Duration) -> usize { + projector.caches.sweep(idle) +} /// A blob table in memory. #[derive(Debug, Default)] diff --git a/lib/components/fabro-petri/tests/projection.rs b/lib/components/fabro-petri/tests/projection.rs index 2e6ab7a71..ca9906676 100644 --- a/lib/components/fabro-petri/tests/projection.rs +++ b/lib/components/fabro-petri/tests/projection.rs @@ -3,7 +3,9 @@ //! missed its wake-ups catches up on the next signal; a crash between the //! record commit and the view transaction is recovered by applying only the //! missing suffix; two projectors over one store agree over nested child -//! executions; and a torn tail holds the view where it stands. +//! executions; a torn tail holds the view where it stands; and a pass over +//! a live run costs its new records, with the cache that makes it so +//! dropped at a restart, after the idle period and at the run's finish. //! //! Every run here takes its scope's environment through the sandbox-driver //! host plugin, so the tests skip, and say why, when the executable is not @@ -25,13 +27,13 @@ use std::time::{Duration, Instant}; use fabro_db::DbPool; use fabro_interview::ControlInterviewer; -use fabro_petri::SqliteRunStore; use fabro_petri::blobs::{Blobs, RunBlobs}; use fabro_petri::check::Launch; use fabro_petri::engine::{self, RunStatus as EngineRunStatus}; use fabro_petri::interview::{Approval, FabroInterviewer}; use fabro_petri::projector::{self, Projector}; use fabro_petri::runtime::RuntimeSpec; +use fabro_petri::{SqliteRunStore, test_support as petri_support}; use fabro_store::platform_records::{ PlatformRecord, PlatformRecordStore, RunCreatedRecord, RunLifecycleKind, RunLifecycleRecord, }; @@ -968,6 +970,181 @@ async fn a_torn_tail_holds_the_view_and_reports_the_run_incomplete() { ); } +/// The run's records the projection reads (the coordinator log and the +/// execution logs; the sandbox ledger has no events) in the order they +/// were recorded: by `recorded_at`, the coordinator log first on a tie, +/// each log's own order kept. The ledger's records are copied to `staged` +/// first, since they are no part of any batch. +async fn in_recorded_order<'a>( + rows: &'a [(String, i64, i64, String)], + staged: &DbPool, + run_id: RunId, +) -> Vec<&'a (String, i64, i64, String)> { + let (ledger, projected): (Vec<_>, Vec<_>) = rows.iter().partition(|row| row.0 == "resources"); + for row in ledger { + insert_petri_row(staged, run_id, row).await; + } + let mut ordered = projected; + ordered.sort_by_key(|row| (row.2, row.0 != "coordinator", row.0.clone(), row.1)); + ordered +} + +/// One committed pass over the run, run again while a platform record +/// contends it. +async fn committed_pass(projector: &Projector, run_id: RunId) -> projector::PassReport { + loop { + let report = projector + .project_run(run_id) + .await + .expect("the pass commits"); + if !report.contended { + return report; + } + } +} + +/// The records land in batches and a pass follows each: every pass feeds +/// only its batch through Petri's derivation, never the run so far, and +/// the view the batches build is the rebuild. The finished run's cache is +/// dropped, and a pass over it is skipped. +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn a_pass_over_a_live_run_costs_its_new_records_not_the_run() { + if host_plugin().is_none() { + return; + } + let scenario = parallel_scenario().await; + run_unobserved(&scenario).await; + let rows = petri_rows(&scenario.pool, scenario.run_id).await; + let staged = copy_run_without_records(&scenario.pool, scenario.run_id).await; + let ordered = in_recorded_order(&rows, &staged, scenario.run_id).await; + const BATCH: usize = 7; + assert!( + ordered.len() > 4 * BATCH, + "enough records for several batches: {}", + ordered.len() + ); + let projector = Projector::new(staged.clone(), staged.clone()); + let mut replayed = Vec::new(); + for batch in ordered.chunks(BATCH) { + for row in batch { + insert_petri_row(&staged, scenario.run_id, row).await; + } + let report = committed_pass(&projector, scenario.run_id).await; + assert!(!report.skipped, "a batch is folded: {report:?}"); + assert!( + report.replayed_records <= batch.len(), + "pass {}: {} records replayed for a batch of {}", + replayed.len(), + report.replayed_records, + batch.len() + ); + replayed.push(report.replayed_records); + } + assert_eq!( + replayed.iter().sum::(), + ordered.len(), + "every record was fed once: {replayed:?}" + ); + assert_view_equals_rebuild(&staged, scenario.run_id).await; + assert!( + !petri_support::cache_held(&projector, scenario.run_id), + "a finished run's cache is dropped" + ); + let again = committed_pass(&projector, scenario.run_id).await; + assert!(again.skipped, "nothing is left to fold: {again:?}"); + assert!(again.health.complete, "{:?}", again.health.incomplete); +} + +/// The cache is dropped with the process and after the idle period, and +/// rebuilt by one full replay: the first pass over new records after +/// either feeds the run so far through Petri's derivation, the next only +/// its new records. Nothing is checkpointed for it, and the view it +/// continues is the rebuild. +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn a_restart_and_the_idle_period_drop_the_cache_and_one_full_replay_rebuilds_it() { + if host_plugin().is_none() { + return; + } + let scenario = parallel_scenario().await; + run_unobserved(&scenario).await; + let rows = petri_rows(&scenario.pool, scenario.run_id).await; + let staged = copy_run_without_records(&scenario.pool, scenario.run_id).await; + let ordered = in_recorded_order(&rows, &staged, scenario.run_id).await; + const BATCH: usize = 5; + let half = ordered.len() / 2; + assert!(half > 3 * BATCH, "enough records: {}", ordered.len()); + let mut fed = 0; + let mut feed = |count: usize| { + let rows: Vec<_> = ordered[fed..(fed + count).min(ordered.len())].to_vec(); + fed += rows.len(); + rows + }; + + for row in feed(half) { + insert_petri_row(&staged, scenario.run_id, row).await; + } + let before = Projector::new(staged.clone(), staged.clone()); + let first = committed_pass(&before, scenario.run_id).await; + assert_eq!( + first.replayed_records, half, + "the first pass replays the run so far" + ); + assert!(!first.health.complete, "the run has not finished"); + assert!( + petri_support::cache_held(&before, scenario.run_id), + "a live run's cache is kept" + ); + drop(before); + + // A restarted server builds a new projector: no cache, and the next + // pass replays the run whole once. + let after = Projector::new(staged.clone(), staged.clone()); + assert!( + !petri_support::cache_held(&after, scenario.run_id), + "a restart holds no cache" + ); + for row in feed(BATCH) { + insert_petri_row(&staged, scenario.run_id, row).await; + } + let rebuilt = committed_pass(&after, scenario.run_id).await; + assert_eq!( + rebuilt.replayed_records, + half + BATCH, + "the first pass after a restart replays the run so far" + ); + assert!(petri_support::cache_held(&after, scenario.run_id)); + for row in feed(BATCH) { + insert_petri_row(&staged, scenario.run_id, row).await; + } + let live = committed_pass(&after, scenario.run_id).await; + assert_eq!( + live.replayed_records, BATCH, + "the next pass replays its batch" + ); + + // The idle period passes: the cache is dropped, and rebuilt the same way. + assert_eq!( + petri_support::drop_idle_caches(&after, Duration::ZERO), + 1, + "the run's cache was idle" + ); + assert!(!petri_support::cache_held(&after, scenario.run_id)); + for row in feed(BATCH) { + insert_petri_row(&staged, scenario.run_id, row).await; + } + let idle = committed_pass(&after, scenario.run_id).await; + assert_eq!(idle.replayed_records, half + 3 * BATCH); + let rest = feed(ordered.len()); + let rest_len = rest.len(); + for row in rest { + insert_petri_row(&staged, scenario.run_id, row).await; + } + let last = committed_pass(&after, scenario.run_id).await; + assert_eq!(last.replayed_records, rest_len); + assert!(last.health.complete, "{:?}", last.health.incomplete); + assert_view_equals_rebuild(&staged, scenario.run_id).await; +} + /// A gate scenario runs through the engine assembly with the interview /// adapter, as a Fabro run does, over a store that signals the projector. struct GateRun {