diff --git a/lib/apps/fabro-server/tests/it/scenario/petri.rs b/lib/apps/fabro-server/tests/it/scenario/petri.rs index 9e53c3381..9ad651660 100644 --- a/lib/apps/fabro-server/tests/it/scenario/petri.rs +++ b/lib/apps/fabro-server/tests/it/scenario/petri.rs @@ -28,7 +28,7 @@ use axum::body::Body; use axum::http::{Request, StatusCode}; use fabro_petri::engine::{self, RunStatus}; use fabro_petri::petri::{Access, OwnerId, RunKey, RunStore as _}; -use fabro_petri::{SqliteRunStore, projector}; +use fabro_petri::{SqliteRunStore, test_support}; use fabro_server::server::AppState; use fabro_server::test_support::{ TestAppStateBuilder, llm_overlay_with_provider_base_url, test_app_db_pool, @@ -237,7 +237,7 @@ pub(super) async fn settled_state( /// How many items the run's projected stream holds. async fn petri_stream_len(state: &AppState, run_id: &str) -> usize { let id: RunId = run_id.parse().expect("the run id parses"); - projector::stored_stream(&state.test_petri_view_pool(), id) + test_support::stored_stream(&state.test_petri_view_pool(), id) .await .expect("the stream reads") .len() diff --git a/lib/components/fabro-petri/src/projector/cache.rs b/lib/components/fabro-petri/src/projector/cache.rs index c2f423443..abe4cf5fd 100644 --- a/lib/components/fabro-petri/src/projector/cache.rs +++ b/lib/components/fabro-petri/src/projector/cache.rs @@ -121,6 +121,7 @@ impl Caches { } /// Whether a cache is kept for the run: a test's view of the cache. + #[cfg(any(test, feature = "test-support"))] pub(crate) fn holds(&self, run_id: RunId) -> bool { let runs = sync::lock(&self.runs); runs.get(&run_id) diff --git a/lib/components/fabro-petri/src/projector.rs b/lib/components/fabro-petri/src/projector/mod.rs similarity index 52% rename from lib/components/fabro-petri/src/projector.rs rename to lib/components/fabro-petri/src/projector/mod.rs index 533a590c5..cada3395c 100644 --- a/lib/components/fabro-petri/src/projector.rs +++ b/lib/components/fabro-petri/src/projector/mod.rs @@ -54,21 +54,24 @@ //! completeness once the run has recorded its finish. mod cache; +pub(crate) mod order; +mod signalling; +pub(crate) mod stream; -use std::collections::{BTreeMap, BTreeSet, HashMap}; +use std::collections::{BTreeMap, HashMap}; +#[cfg(any(test, feature = "test-support"))] use std::sync::atomic::{AtomicBool, Ordering}; use std::sync::{Arc, Mutex}; use std::time::Duration; use fabro_db::DbPool; -use fabro_store::platform_records::{PlatformRecordStore, StoredPlatformRecord, now_ms}; +use fabro_store::platform_records::{PlatformRecordStore, now_ms}; use fabro_store::{RunProjection, RunSummaryStore}; -use fabro_types::{RunId, RunStreamItem, RunStreamItemKind}; +use fabro_types::{RunId, RunStreamItem}; use fabro_util::error::collect_chain; use fabro_util::sync; -use petri_execution::events::{self, EventId, EventSource, RunEvent}; +use petri_execution::events::{EventId, EventSource, RunEvent}; 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::broadcast; @@ -76,8 +79,9 @@ use tokio::time; use tracing::{debug, info, warn}; use self::cache::{Caches, IDLE, RunCache}; +use self::stream::StreamRow; use crate::SqliteRunStore; -use crate::projection::{self, FiringKey, FoldState, Item, RecordHealth, RunView}; +use crate::projection::{self, FoldState, RecordHealth, RunView}; /// The positions a view committed: the last event consumed per Petri log, /// and the last platform record consumed. @@ -128,6 +132,48 @@ pub struct PassReport { pub health: RecordHealth, } +impl PassReport { + /// A pass that found the view at the head of every log and wrote + /// nothing. + fn skipped(run_id: RunId, stored: &StoredView) -> Self { + Self { + run_id, + skipped: true, + contended: false, + petri_events: 0, + platform_records: 0, + replayed_records: 0, + stream_seq: stored.stream_seq, + positions: stored.positions.clone(), + health: stored.view.state.health.clone(), + } + } + + /// A pass that left the view alone because a platform record landed + /// under it; the projector runs it again. + fn contended(run_id: RunId, replayed_records: usize) -> Self { + Self { + run_id, + skipped: false, + contended: true, + petri_events: 0, + platform_records: 0, + replayed_records, + stream_seq: 0, + positions: Positions::default(), + health: RecordHealth::default(), + } + } +} + +/// What a pass read past its cache: the events new to the view, the +/// replay failure that held it, and how many records the replay cost. +struct NewEvents { + events: Vec, + replay_failure: Option, + replayed_records: usize, +} + /// What the startup pass did. #[derive(Clone, Debug, Default, PartialEq, Eq)] pub struct StartupReport { @@ -182,8 +228,9 @@ pub struct Projector { /// 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. + /// A test's fault: stop the next pass after its reads, before its + /// view transaction, as a crash there would. + #[cfg(any(test, feature = "test-support"))] 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 @@ -210,6 +257,7 @@ impl Projector { pool: views, slots: Mutex::default(), caches: Caches::default(), + #[cfg(any(test, feature = "test-support"))] fault: AtomicBool::new(false), committed: broadcast::channel(COMMIT_SIGNAL_CAPACITY).0, }) @@ -232,7 +280,7 @@ impl Projector { after: u64, limit: usize, ) -> Result, ProjectError> { - stream_after(&self.pool, run_id, after, limit).await + stream::stream_after(&self.pool, run_id, after, limit).await } /// Delete everything the store and the view tables hold for the run: @@ -344,10 +392,24 @@ impl Projector { /// Stop the next pass after its reads and before its view transaction, /// as a crash there would, once. + #[cfg(any(test, feature = "test-support"))] pub fn fail_before_view(&self) { self.fault.store(true, Ordering::SeqCst); } + /// Whether a test asked this pass to stop before its view transaction; + /// never outside tests. + fn take_fault(&self) -> bool { + #[cfg(any(test, feature = "test-support"))] + { + self.fault.swap(false, Ordering::SeqCst) + } + #[cfg(not(any(test, feature = "test-support")))] + { + false + } + } + /// One pass over every Petri run the database holds: the runs with a /// Petri record, and the runs with platform records. Runs whose view /// already covers every committed record are skipped cheaply. @@ -430,67 +492,20 @@ impl Projector { /// 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 - .head(&run_id) - .await - .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 - .positions - .petri - .iter() - .any(|held| log_text(&held.source) == *log && held.seq == *head) - }); - if at_head && stored.view.projection.is_some() { - return Ok(PassReport { - run_id, - skipped: true, - contended: false, - petri_events: 0, - platform_records: 0, - replayed_records: 0, - stream_seq: stored.stream_seq, - positions: stored.positions.clone(), - health: stored.view.state.health.clone(), - }); + if self.at_head(run_id, &run.view).await? { + return Ok(PassReport::skipped(run_id, &run.view)); } let platform_records = self .platform - .read_after(&run_id, stored.positions.platform_seq) + .read_after(&run_id, run.view.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 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"); - (Vec::new(), Some(chain)) - } - }, - Err(StoreError::NotFound { .. }) => (Vec::new(), None), - Err(error) => return Err(ProjectError::Open(error)), - }; + let NewEvents { + events, + replay_failure, + replayed_records, + } = self.read_new_events(run_id, &key, run).await?; let mut view = run.view.view.clone(); let mut positions = run.view.positions.clone(); @@ -505,7 +520,7 @@ impl Projector { let platform_head_seen = platform_records .last() .map_or(positions.platform_seq, |record| record.seq); - let (items, held) = order_items( + let (items, held) = order::order_items( &events, &platform_records, &view.state.finished_firings, @@ -518,37 +533,11 @@ impl Projector { "platform records held back until their firing's finish is in the stream" ); } - - let mut rows: Vec = Vec::with_capacity(items.len()); - for item in &items { - stream_seq += 1; - view.fold(item, stream_seq); - let row = match item { - Item::Petri(event) => { - positions.advance(event.id); - StreamRow { - stream_seq, - item_kind: "petri", - item_id: event_id_text(&event.id), - event_json: serde_json::to_string(event).map_err(ProjectError::Encode)?, - } - } - Item::Platform(record) => { - positions.platform_seq = record.seq; - StreamRow { - stream_seq, - item_kind: "platform", - item_id: record.seq.to_string(), - event_json: serde_json::to_string(record).map_err(ProjectError::Encode)?, - } - } - }; - rows.push(row); - } + let rows = stream::stream_rows(&items, &mut view, &mut positions, &mut stream_seq)?; drop(items); view.state.health = self.health(&key, &view.state, replay_failure).await?; - if self.fault.swap(false, Ordering::SeqCst) { + if self.take_fault() { run.pending = events; return Err(ProjectError::Injected); } @@ -567,17 +556,7 @@ impl Projector { 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(), - }); + return Ok(PassReport::contended(run_id, replayed_records)); } Err(error) => { run.pending = events; @@ -612,6 +591,77 @@ impl Projector { }) } + /// Whether the stored view already covers every committed record: the + /// platform head and every Petri log's head are the positions it holds, + /// and it has a projection to serve. + async fn at_head(&self, run_id: RunId, stored: &StoredView) -> Result { + let platform_head = self + .platform + .head(&run_id) + .await + .map_err(ProjectError::Store)? + .unwrap_or(0); + let petri_heads = self.petri_heads(&run_id).await?; + Ok(platform_head == stored.positions.platform_seq + && petri_heads.iter().all(|(log, head)| { + stored + .positions + .petri + .iter() + .any(|held| stream::log_text(&held.source) == *log && held.seq == *head) + }) + && stored.view.projection.is_some()) + } + + /// The events past the view's positions: the cache's replay advanced + /// over the records committed since, led by the events an earlier pass + /// derived and did not commit. A rebuilt replay derives the run whole, + /// so only the events past the view's positions are new to it. A + /// replay that fails (a torn tail) stands still, yields nothing, and + /// names its reason; what it derived before stays pending for the next + /// pass. + async fn read_new_events( + &self, + run_id: RunId, + key: &RunKey, + run: &mut RunCache, + ) -> Result { + let nothing = NewEvents { + events: Vec::new(), + replay_failure: None, + replayed_records: 0, + }; + let logs = match self.store.open(key, Access::Read).await { + Ok(logs) => logs, + Err(StoreError::NotFound { .. }) => return Ok(nothing), + Err(error) => return Err(ProjectError::Open(error)), + }; + match run.replay.advance(&*logs).await { + Ok(new) => { + let replayed_records = new.iter().filter(|event| event.id.index == 0).count(); + 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) + })); + Ok(NewEvents { + events, + replay_failure: None, + replayed_records, + }) + } + Err(error) => { + let chain = collect_chain(&error).join(": "); + warn!(run_id = %run_id, error = %chain, "Petri run does not replay; the view holds"); + Ok(NewEvents { + replay_failure: Some(chain), + ..nothing + }) + } + } + } + /// 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). @@ -655,8 +705,8 @@ impl Projector { .bind(projection_json) .bind(fold_json) .bind(positions_json) - .bind(column(stream_seq)) - .bind(column(now_ms())) + .bind(stream::column(stream_seq)) + .bind(stream::column(now_ms())) .execute(&mut *tx) .await .map_err(ProjectError::Database)?; @@ -666,7 +716,7 @@ impl Projector { VALUES (?, ?, ?, ?, ?)", ) .bind(run_id.to_string()) - .bind(column(row.stream_seq)) + .bind(stream::column(row.stream_seq)) .bind(row.item_kind) .bind(&row.item_id) .bind(&row.event_json) @@ -772,335 +822,13 @@ impl Projector { } } -impl Projector { - /// A run store whose appends signal this projector: for a run that - /// executes in the same process as the projector, over the SQLite store - /// directly, where no append endpoint is there to signal. The signal is - /// sent after the store's append returned, so the records it covers are - /// durable before the view sees them. - pub fn observe_store( - self: &Arc, - inner: Arc, - ) -> Arc { - Arc::new(SignallingStore { - inner, - projector: Arc::clone(self), - }) - } -} - -/// A run store that signals a projector after each append. -struct SignallingStore { - inner: Arc, - projector: Arc, -} - -#[async_trait::async_trait] -impl petri_execution::RunStore for SignallingStore { - async fn open( - &self, - key: &RunKey, - access: Access, - ) -> Result, StoreError> { - let logs = self.inner.open(key, access).await?; - Ok(Arc::new(SignallingLogs { - inner: logs, - run_id: projection::run_id_of(key.as_str()), - projector: Arc::clone(&self.projector), - })) - } -} - -struct SignallingLogs { - inner: Arc, - run_id: Option, - projector: Arc, -} - -#[async_trait::async_trait] -impl petri_execution::RunLogs for SignallingLogs { - fn locator(&self) -> String { - self.inner.locator() - } - - async fn append( - &self, - log: &petri_execution::LogId, - records: &[petri_execution::Record], - ) -> Result<(), StoreError> { - self.inner.append(log, records).await?; - if let Some(run_id) = self.run_id { - self.projector.signal(run_id); - } - Ok(()) - } - - async fn read( - &self, - log: &petri_execution::LogId, - ) -> Result, StoreError> { - 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 - } - - async fn get_blob(&self, digest: petri_store::Digest) -> Result>, StoreError> { - self.inner.get_blob(digest).await - } -} - -/// The order one pass streams its new items in, and how many platform -/// records it holds back for a later pass. -/// -/// Every item is first ordered by `recorded_at` (stable: the coordinator -/// log before an execution log before a platform record on a tie, and each -/// log's own order kept). A platform record that carries a Petri position -/// (a checkpoint, keyed on `(execution, firing)`) is then placed by that -/// position, not by its clock, because the server stamps the record and the -/// worker stamps Petri's records and the two clocks can tie or invert: -/// -/// - before the firing's first `routing.resolved` event in the pass, which is -/// right after the firing's finish (its `step.finished` and the -/// `visit.completed` attached to it) and before the next firing's -/// `visit.started`, which is attached to that routing record; -/// - else after the last event of the firing in the pass; -/// - else, when the firing finished in an earlier pass, before the first event -/// of a later firing (a larger firing id) in the same execution, or where its -/// `recorded_at` put it; -/// - else the record is held back, with every platform record after it, and the -/// pass consumes platform records only up to it. The hook that writes a -/// checkpoint record runs after the driver appended the attempt's finish, but -/// the driver's store writer flushes that record on its own schedule, so the -/// platform record can be committed before its firing's `step.finished`; -/// holding it keeps the stream's order the same live and on a rebuild. -/// Nothing is held once the run has recorded its finish. -/// -/// The rule reads only the pass's own items and the firings already -/// finished, so a record is never streamed before its firing's finish and -/// never after the firing's routes. -fn order_items<'a>( - events: &'a [RunEvent], - platform_records: &'a [StoredPlatformRecord], - finished_before: &BTreeSet, - run_finished: bool, -) -> (Vec>, usize) { - let finished_in_pass = |at: FiringKey| { - events.iter().any(|event| { - FiringKey::of_event(event) == Some(at) - && matches!(event.engine(), Some(Event::StepFinished { .. })) - }) - }; - let finished = |at: FiringKey| finished_in_pass(at) || finished_before.contains(&at); - // Platform records are consumed in seq order: the first one whose firing - // has not finished holds itself and everything after it. - let consumed = if run_finished { - platform_records.len() - } else { - platform_records - .iter() - .position(|record| { - record - .position - .is_some_and(|position| !finished(FiringKey::from(position))) - }) - .unwrap_or(platform_records.len()) - }; - let held = platform_records.len() - consumed; - let platform_records = &platform_records[..consumed]; - - let mut items: Vec<(u64, u8, Item<'a>)> = - Vec::with_capacity(events.len() + platform_records.len()); - for event in events { - let rank = match event.id.source { - EventSource::Coordinator => 0, - EventSource::Execution { .. } => 1, - }; - items.push((event.recorded_at, rank, Item::Petri(event))); - } - for record in platform_records { - items.push((record.recorded_at, 2, Item::Platform(record))); - } - items.sort_by_key(|(recorded_at, rank, _)| (*recorded_at, *rank)); - - let item_firing = |item: &Item<'a>| match item { - Item::Petri(event) => FiringKey::of_event(event), - Item::Platform(_) => None, - }; - let is_routing = |item: &Item<'a>| { - matches!( - item, - Item::Petri(event) if matches!(event.engine(), Some(Event::RoutingResolved { .. })) - ) - }; - // The key of each item: its index in clock order, and whether it sits - // before (0), at (1) or after (2) that index. - let mut keys: Vec<(usize, u8)> = (0..items.len()).map(|index| (index, 1)).collect(); - for (index, (_, _, item)) in items.iter().enumerate() { - let Item::Platform(record) = item else { - continue; - }; - let Some(position) = record.position else { - continue; - }; - let at = FiringKey::from(position); - let first_routing = items - .iter() - .position(|(_, _, other)| item_firing(other) == Some(at) && is_routing(other)); - let last_of_firing = items - .iter() - .rposition(|(_, _, other)| item_firing(other) == Some(at)); - let first_later = items.iter().position(|(_, _, other)| { - item_firing(other) - .is_some_and(|key| key.execution == at.execution && key.firing > at.firing) - }); - keys[index] = if let Some(before) = first_routing { - (before, 0) - } else if let Some(after) = last_of_firing { - (after, 2) - } else if let Some(before) = first_later { - (before, 0) - } else { - (index, 1) - }; - } - let mut order: Vec = (0..items.len()).collect(); - order.sort_by_key(|index| keys[*index]); - let mut ordered: Vec>> = - items.into_iter().map(|(_, _, item)| Some(item)).collect(); - let items = order - .into_iter() - .map(|index| ordered[index].take().expect("each item is placed once")) - .collect(); - (items, held) -} - -struct StreamRow { - stream_seq: u64, - item_kind: &'static str, - item_id: String, - event_json: String, -} - /// How many commit signals a slow reader may fall behind before it is told /// it lagged and re-reads from its cursor. const COMMIT_SIGNAL_CAPACITY: usize = 1024; -/// The run's stream past the cursor, read from the view tables: up to -/// `limit` rows with `stream_seq > after`, in order, in Fabro's envelope. -pub async fn stream_after( - views: &DbPool, - run_id: RunId, - after: u64, - limit: usize, -) -> Result, ProjectError> { - let rows: Vec<(i64, String, String, String)> = sqlx::query_as( - "SELECT stream_seq, item_kind, item_id, event_json FROM petri_stream WHERE run_id = ? AND \ - stream_seq > ? ORDER BY stream_seq LIMIT ?", - ) - .bind(run_id.to_string()) - .bind(column(after)) - .bind(i64::try_from(limit).unwrap_or(i64::MAX)) - .fetch_all(views) - .await - .map_err(ProjectError::Database)?; - rows.into_iter() - .map(|(stream_seq, item_kind, item_id, event_json)| { - let item: serde_json::Value = - serde_json::from_str(&event_json).map_err(ProjectError::Encode)?; - let kind = match item_kind.as_str() { - "platform" => RunStreamItemKind::Platform, - _ => RunStreamItemKind::Petri, - }; - let recorded_at = item - .get("recorded_at") - .and_then(serde_json::Value::as_u64) - .unwrap_or(0); - Ok(RunStreamItem { - run_id, - stream_seq: u64::try_from(stream_seq).unwrap_or(0), - kind, - id: item_id, - recorded_at, - item, - }) - }) - .collect() -} - -/// A Petri event id as the stream names it: `//`. -#[must_use] -pub fn event_id_text(id: &EventId) -> String { - format!("{}/{}/{}", log_text(&id.source), id.seq, id.index) -} - -fn log_text(source: &EventSource) -> String { - match source { - EventSource::Coordinator => "coordinator".to_string(), - EventSource::Execution { execution } => format!("execution {execution}"), - } -} - -fn column(value: u64) -> i64 { - i64::try_from(value).unwrap_or(i64::MAX) -} - -/// The run's projection rebuilt from its records alone, with nothing -/// stored: what a fresh projector would commit over the same records. A test -/// compares it with the live view. `records` and `views` are the two pools -/// [`Projector::new`] takes. -pub async fn rebuild( - records: &DbPool, - views: &DbPool, - run_id: RunId, -) -> Result<(Option, Positions, u64), ProjectError> { - let store = SqliteRunStore::new(records.clone()); - let platform = PlatformRecordStore::new(views.clone()); - let key = RunKey::new(run_id.to_string()); - let platform_records = platform.read(&run_id).await.map_err(ProjectError::Store)?; - let events = match store.open(&key, Access::Read).await { - Ok(logs) => events::replay_run(&*logs) - .await - .inspect_err(|error| { - warn!(error = %collect_chain(error).join(": "), "rebuild: the run does not replay"); - }) - .unwrap_or_default(), - Err(StoreError::NotFound { .. }) => Vec::new(), - Err(error) => return Err(ProjectError::Open(error)), - }; - let run_finished = events.iter().any(|event| { - matches!( - event.coordinator(), - Some(CoordinatorEvent::RunFinished { .. }) - ) - }); - let (items, _held) = order_items(&events, &platform_records, &BTreeSet::new(), run_finished); - let mut view = RunView::new(); - let mut positions = Positions::default(); - let mut stream_seq = 0; - for item in &items { - stream_seq += 1; - view.fold(item, stream_seq); - match item { - Item::Petri(event) => positions.advance(event.id), - Item::Platform(record) => positions.platform_seq = record.seq, - } - } - Ok((view.projection, positions, stream_seq)) -} - -/// The stored view's positions and stream sequence, for a test; `views` is -/// the pool the view tables live in. -pub async fn stored_positions( +/// The stored view's positions and stream sequence; `views` is the pool the +/// view tables live in. +pub(crate) async fn stored_positions( views: &DbPool, run_id: RunId, ) -> Result, ProjectError> { @@ -1118,261 +846,3 @@ pub async fn stored_positions( }) .transpose() } - -/// The stored view's projection, for a test or a reader outside the store. -pub async fn stored_projection( - views: &DbPool, - run_id: RunId, -) -> Result, ProjectError> { - let json: Option = - sqlx::query_scalar("SELECT projection_json FROM petri_projection WHERE run_id = ?") - .bind(run_id.to_string()) - .fetch_optional(views) - .await - .map_err(ProjectError::Database)?; - json.map(|json| serde_json::from_str(&json).map_err(ProjectError::Encode)) - .transpose() -} - -/// The stream rows of a run: `(stream_seq, item_kind, item_id)`, in order. -pub async fn stored_stream( - views: &DbPool, - run_id: RunId, -) -> Result, ProjectError> { - let rows: Vec<(i64, String, String)> = sqlx::query_as( - "SELECT stream_seq, item_kind, item_id FROM petri_stream WHERE run_id = ? ORDER BY stream_seq", - ) - .bind(run_id.to_string()) - .fetch_all(views) - .await - .map_err(ProjectError::Database)?; - Ok(rows - .into_iter() - .map(|(seq, kind, id)| (u64::try_from(seq).unwrap_or(0), kind, id)) - .collect()) -} - -/// Every stored platform record of a run, for a reader outside the store. -pub async fn stored_platform_records( - views: &DbPool, - run_id: RunId, -) -> Result, ProjectError> { - PlatformRecordStore::new(views.clone()) - .read(&run_id) - .await - .map_err(ProjectError::Store) -} - -/// A recorded event's projection is what `RunEvent` serializes to. -#[must_use] -pub fn event_json(event: &RunEvent) -> serde_json::Value { - serde_json::to_value(event).unwrap_or_default() -} - -#[cfg(test)] -mod tests { - use fabro_store::PlatformRecord; - use fabro_store::platform_records::{CheckpointRecord, StagePosition}; - use petri_execution::events::{Context, NodeRef, Record, RecordOrigin, Subject}; - use petri_execution::{ExecutionId, StoredEngineRecord}; - use petri_runtime::driver::BranchRole; - use petri_runtime::engine::{DecisionId, EventOrigin, RouteApplied}; - use petri_runtime::ir::{Attempt, FiringId, NodeId, Outcome, Status}; - - use super::*; - - /// A firing's engine event at `seq`, recorded at `at`. - fn engine_event(seq: u64, firing: u64, at: u64, body: Event) -> RunEvent { - RunEvent { - id: EventId { - source: EventSource::Execution { - execution: ExecutionId::new(0), - }, - seq, - index: 0, - }, - origin: RecordOrigin::External, - context: Context { - invocation: None, - execution: Some(ExecutionId::new(0)), - parent: None, - }, - subject: Some(Subject { - node: NodeRef { - id: NodeId::new(1), - name: format!("n{firing}").into(), - kind: "attractor/command".into(), - meta: serde_json::Value::Null, - }, - firing: Some(FiringId::new(firing)), - visit: Some(1), - attempt: Some(Attempt::FIRST), - generation: None, - branch: BranchRole::None, - }), - observed_at: None, - recorded_at: at, - record: Some(Record::Engine(StoredEngineRecord { - seq, - origin: EventOrigin::External, - recorded_at: at, - body, - })), - derived: None, - } - } - - fn finished(seq: u64, firing: u64, at: u64) -> RunEvent { - engine_event(seq, firing, at, Event::StepFinished { - firing: FiringId::new(firing), - attempt: Attempt::FIRST, - outcome: Outcome::new(Status::Success, serde_json::Value::Null), - }) - } - - fn routing(seq: u64, firing: u64, at: u64) -> RunEvent { - engine_event(seq, firing, at, Event::RoutingResolved { - decision_id: DecisionId::route(FiringId::new(firing), Attempt::FIRST), - groups: Vec::new(), - }) - } - - fn applied(seq: u64, firing: u64, at: u64) -> RunEvent { - engine_event(seq, firing, at, Event::RouteApplied { - applied: RouteApplied::None { - firing: FiringId::new(firing), - group: 0, - }, - }) - } - - fn started(seq: u64, firing: u64, at: u64) -> RunEvent { - engine_event(seq, firing, at, Event::StepStarted { - firing: FiringId::new(firing), - attempt: Attempt::FIRST, - }) - } - - fn checkpoint(seq: u64, firing: u64, at: u64) -> StoredPlatformRecord { - StoredPlatformRecord { - seq, - recorded_at: at, - record: PlatformRecord::Checkpoint(CheckpointRecord { - execution: 0, - firing, - attempt: Some(1), - workspace: None, - git_commit_sha: Some("abc".to_string()), - diff_summary: None, - patch_blob: None, - operation: None, - }), - position: Some(StagePosition { - execution: 0, - firing, - }), - } - } - - fn names(items: &[Item<'_>]) -> Vec { - items - .iter() - .map(|item| match item { - Item::Petri(event) => event_id_text(&event.id), - Item::Platform(record) => format!("platform {}", record.seq), - }) - .collect() - } - - /// A firing's events, then a later firing's events, then a checkpoint - /// for the first firing stamped later than all of them: the stream puts - /// the checkpoint right after the first firing's finish, before its - /// routes and before the later firing. - #[test] - fn a_positioned_record_follows_its_firings_finish_whatever_its_clock_says() { - let events = vec![ - finished(10, 1, 100), - routing(11, 1, 101), - applied(12, 1, 102), - started(13, 2, 103), - finished(14, 2, 104), - ]; - let records = vec![checkpoint(1, 1, 250)]; - let (items, held) = order_items(&events, &records, &BTreeSet::new(), false); - assert_eq!(held, 0); - assert_eq!(names(&items), vec![ - "execution 0/10/0", - "platform 1", - "execution 0/11/0", - "execution 0/12/0", - "execution 0/13/0", - "execution 0/14/0", - ]); - } - - /// With the firing finished in an earlier pass, the record goes before - /// the first event of a later firing; a record with no position keeps - /// its clock order. - #[test] - fn a_positioned_record_precedes_later_firings_and_an_unpositioned_one_keeps_its_clock() { - let events = vec![started(13, 2, 103), finished(14, 2, 104)]; - let records = vec![checkpoint(1, 1, 250)]; - let finished_before: BTreeSet = [FiringKey::new(0, 1)].into_iter().collect(); - let (items, held) = order_items(&events, &records, &finished_before, false); - assert_eq!(held, 0); - assert_eq!(names(&items), vec![ - "platform 1", - "execution 0/13/0", - "execution 0/14/0", - ]); - - let unpositioned = StoredPlatformRecord { - position: None, - ..checkpoint(2, 1, 250) - }; - let unpositioned = [unpositioned]; - let (items, held) = order_items(&events, &unpositioned, &BTreeSet::new(), false); - assert_eq!(held, 0); - assert_eq!(names(&items), vec![ - "execution 0/13/0", - "execution 0/14/0", - "platform 2", - ]); - } - - /// A record whose firing has no finish yet, in the stream or in the pass, - /// is held back with everything after it until the finish arrives, or - /// until the run has finished. - #[test] - fn a_positioned_record_is_held_until_its_firings_finish_is_in_the_stream() { - let events = vec![started(13, 2, 103), finished(14, 2, 104)]; - let records = vec![ - checkpoint(1, 2, 50), - checkpoint(2, 3, 60), - checkpoint(3, 2, 70), - ]; - let (items, held) = order_items(&events, &records, &BTreeSet::new(), false); - assert_eq!( - held, 2, - "the record for firing 3 holds itself and the one after it" - ); - assert_eq!(names(&items), vec![ - "execution 0/13/0", - "execution 0/14/0", - "platform 1", - ]); - - // Once the run finished, a record for a firing that never finished - // keeps its clock order; the firing's own records still follow its - // finish, in their seq order. - let (items, held) = order_items(&events, &records, &BTreeSet::new(), true); - assert_eq!(held, 0, "nothing is held once the run finished"); - assert_eq!(names(&items), vec![ - "platform 2", - "execution 0/13/0", - "execution 0/14/0", - "platform 1", - "platform 3", - ]); - } -} diff --git a/lib/components/fabro-petri/src/projector/order.rs b/lib/components/fabro-petri/src/projector/order.rs new file mode 100644 index 000000000..2b27aaff4 --- /dev/null +++ b/lib/components/fabro-petri/src/projector/order.rs @@ -0,0 +1,343 @@ +//! The order one pass streams its new items in. + +use std::collections::BTreeSet; + +use fabro_store::platform_records::StoredPlatformRecord; +use petri_execution::events::{EventSource, RunEvent}; +use petri_runtime::engine::Event; + +use crate::projection::{FiringKey, Item}; + +/// The order one pass streams its new items in, and how many platform +/// records it holds back for a later pass. +/// +/// Every item is first ordered by `recorded_at` (stable: the coordinator +/// log before an execution log before a platform record on a tie, and each +/// log's own order kept). A platform record that carries a Petri position +/// (a checkpoint, keyed on `(execution, firing)`) is then placed by that +/// position, not by its clock, because the server stamps the record and the +/// worker stamps Petri's records and the two clocks can tie or invert: +/// +/// - before the firing's first `routing.resolved` event in the pass, which is +/// right after the firing's finish (its `step.finished` and the +/// `visit.completed` attached to it) and before the next firing's +/// `visit.started`, which is attached to that routing record; +/// - else after the last event of the firing in the pass; +/// - else, when the firing finished in an earlier pass, before the first event +/// of a later firing (a larger firing id) in the same execution, or where its +/// `recorded_at` put it; +/// - else the record is held back, with every platform record after it, and the +/// pass consumes platform records only up to it. The hook that writes a +/// checkpoint record runs after the driver appended the attempt's finish, but +/// the driver's store writer flushes that record on its own schedule, so the +/// platform record can be committed before its firing's `step.finished`; +/// holding it keeps the stream's order the same live and on a rebuild. +/// Nothing is held once the run has recorded its finish. +/// +/// The rule reads only the pass's own items and the firings already +/// finished, so a record is never streamed before its firing's finish and +/// never after the firing's routes. +pub(crate) fn order_items<'a>( + events: &'a [RunEvent], + platform_records: &'a [StoredPlatformRecord], + finished_before: &BTreeSet, + run_finished: bool, +) -> (Vec>, usize) { + let finished_in_pass = |at: FiringKey| { + events.iter().any(|event| { + FiringKey::of_event(event) == Some(at) + && matches!(event.engine(), Some(Event::StepFinished { .. })) + }) + }; + let finished = |at: FiringKey| finished_in_pass(at) || finished_before.contains(&at); + // Platform records are consumed in seq order: the first one whose firing + // has not finished holds itself and everything after it. + let consumed = if run_finished { + platform_records.len() + } else { + platform_records + .iter() + .position(|record| { + record + .position + .is_some_and(|position| !finished(FiringKey::from(position))) + }) + .unwrap_or(platform_records.len()) + }; + let held = platform_records.len() - consumed; + let platform_records = &platform_records[..consumed]; + + let mut items: Vec<(u64, u8, Item<'a>)> = + Vec::with_capacity(events.len() + platform_records.len()); + for event in events { + let rank = match event.id.source { + EventSource::Coordinator => 0, + EventSource::Execution { .. } => 1, + }; + items.push((event.recorded_at, rank, Item::Petri(event))); + } + for record in platform_records { + items.push((record.recorded_at, 2, Item::Platform(record))); + } + items.sort_by_key(|(recorded_at, rank, _)| (*recorded_at, *rank)); + + let item_firing = |item: &Item<'a>| match item { + Item::Petri(event) => FiringKey::of_event(event), + Item::Platform(_) => None, + }; + let is_routing = |item: &Item<'a>| { + matches!( + item, + Item::Petri(event) if matches!(event.engine(), Some(Event::RoutingResolved { .. })) + ) + }; + // The key of each item: its index in clock order, and whether it sits + // before (0), at (1) or after (2) that index. + let mut keys: Vec<(usize, u8)> = (0..items.len()).map(|index| (index, 1)).collect(); + for (index, (_, _, item)) in items.iter().enumerate() { + let Item::Platform(record) = item else { + continue; + }; + let Some(position) = record.position else { + continue; + }; + let at = FiringKey::from(position); + let first_routing = items + .iter() + .position(|(_, _, other)| item_firing(other) == Some(at) && is_routing(other)); + let last_of_firing = items + .iter() + .rposition(|(_, _, other)| item_firing(other) == Some(at)); + let first_later = items.iter().position(|(_, _, other)| { + item_firing(other) + .is_some_and(|key| key.execution == at.execution && key.firing > at.firing) + }); + keys[index] = if let Some(before) = first_routing { + (before, 0) + } else if let Some(after) = last_of_firing { + (after, 2) + } else if let Some(before) = first_later { + (before, 0) + } else { + (index, 1) + }; + } + let mut order: Vec = (0..items.len()).collect(); + order.sort_by_key(|index| keys[*index]); + let mut ordered: Vec>> = + items.into_iter().map(|(_, _, item)| Some(item)).collect(); + let items = order + .into_iter() + .map(|index| ordered[index].take().expect("each item is placed once")) + .collect(); + (items, held) +} + +#[cfg(test)] +mod tests { + use fabro_store::PlatformRecord; + use fabro_store::platform_records::{CheckpointRecord, StagePosition}; + use petri_execution::events::{Context, EventId, NodeRef, Record, RecordOrigin, Subject}; + use petri_execution::{ExecutionId, StoredEngineRecord}; + use petri_runtime::driver::BranchRole; + use petri_runtime::engine::{DecisionId, EventOrigin, RouteApplied}; + use petri_runtime::ir::{Attempt, FiringId, NodeId, Outcome, Status}; + + use super::*; + use crate::projector::stream; + + /// A firing's engine event at `seq`, recorded at `at`. + fn engine_event(seq: u64, firing: u64, at: u64, body: Event) -> RunEvent { + RunEvent { + id: EventId { + source: EventSource::Execution { + execution: ExecutionId::new(0), + }, + seq, + index: 0, + }, + origin: RecordOrigin::External, + context: Context { + invocation: None, + execution: Some(ExecutionId::new(0)), + parent: None, + }, + subject: Some(Subject { + node: NodeRef { + id: NodeId::new(1), + name: format!("n{firing}").into(), + kind: "attractor/command".into(), + meta: serde_json::Value::Null, + }, + firing: Some(FiringId::new(firing)), + visit: Some(1), + attempt: Some(Attempt::FIRST), + generation: None, + branch: BranchRole::None, + }), + observed_at: None, + recorded_at: at, + record: Some(Record::Engine(StoredEngineRecord { + seq, + origin: EventOrigin::External, + recorded_at: at, + body, + })), + derived: None, + } + } + + fn finished(seq: u64, firing: u64, at: u64) -> RunEvent { + engine_event(seq, firing, at, Event::StepFinished { + firing: FiringId::new(firing), + attempt: Attempt::FIRST, + outcome: Outcome::new(Status::Success, serde_json::Value::Null), + }) + } + + fn routing(seq: u64, firing: u64, at: u64) -> RunEvent { + engine_event(seq, firing, at, Event::RoutingResolved { + decision_id: DecisionId::route(FiringId::new(firing), Attempt::FIRST), + groups: Vec::new(), + }) + } + + fn applied(seq: u64, firing: u64, at: u64) -> RunEvent { + engine_event(seq, firing, at, Event::RouteApplied { + applied: RouteApplied::None { + firing: FiringId::new(firing), + group: 0, + }, + }) + } + + fn started(seq: u64, firing: u64, at: u64) -> RunEvent { + engine_event(seq, firing, at, Event::StepStarted { + firing: FiringId::new(firing), + attempt: Attempt::FIRST, + }) + } + + fn checkpoint(seq: u64, firing: u64, at: u64) -> StoredPlatformRecord { + StoredPlatformRecord { + seq, + recorded_at: at, + record: PlatformRecord::Checkpoint(CheckpointRecord { + execution: 0, + firing, + attempt: Some(1), + workspace: None, + git_commit_sha: Some("abc".to_string()), + diff_summary: None, + patch_blob: None, + operation: None, + }), + position: Some(StagePosition { + execution: 0, + firing, + }), + } + } + + fn names(items: &[Item<'_>]) -> Vec { + items + .iter() + .map(|item| match item { + Item::Petri(event) => stream::event_id_text(&event.id), + Item::Platform(record) => format!("platform {}", record.seq), + }) + .collect() + } + + /// A firing's events, then a later firing's events, then a checkpoint + /// for the first firing stamped later than all of them: the stream puts + /// the checkpoint right after the first firing's finish, before its + /// routes and before the later firing. + #[test] + fn a_positioned_record_follows_its_firings_finish_whatever_its_clock_says() { + let events = vec![ + finished(10, 1, 100), + routing(11, 1, 101), + applied(12, 1, 102), + started(13, 2, 103), + finished(14, 2, 104), + ]; + let records = vec![checkpoint(1, 1, 250)]; + let (items, held) = order_items(&events, &records, &BTreeSet::new(), false); + assert_eq!(held, 0); + assert_eq!(names(&items), vec![ + "execution 0/10/0", + "platform 1", + "execution 0/11/0", + "execution 0/12/0", + "execution 0/13/0", + "execution 0/14/0", + ]); + } + + /// With the firing finished in an earlier pass, the record goes before + /// the first event of a later firing; a record with no position keeps + /// its clock order. + #[test] + fn a_positioned_record_precedes_later_firings_and_an_unpositioned_one_keeps_its_clock() { + let events = vec![started(13, 2, 103), finished(14, 2, 104)]; + let records = vec![checkpoint(1, 1, 250)]; + let finished_before: BTreeSet = [FiringKey::new(0, 1)].into_iter().collect(); + let (items, held) = order_items(&events, &records, &finished_before, false); + assert_eq!(held, 0); + assert_eq!(names(&items), vec![ + "platform 1", + "execution 0/13/0", + "execution 0/14/0", + ]); + + let unpositioned = StoredPlatformRecord { + position: None, + ..checkpoint(2, 1, 250) + }; + let unpositioned = [unpositioned]; + let (items, held) = order_items(&events, &unpositioned, &BTreeSet::new(), false); + assert_eq!(held, 0); + assert_eq!(names(&items), vec![ + "execution 0/13/0", + "execution 0/14/0", + "platform 2", + ]); + } + + /// A record whose firing has no finish yet, in the stream or in the pass, + /// is held back with everything after it until the finish arrives, or + /// until the run has finished. + #[test] + fn a_positioned_record_is_held_until_its_firings_finish_is_in_the_stream() { + let events = vec![started(13, 2, 103), finished(14, 2, 104)]; + let records = vec![ + checkpoint(1, 2, 50), + checkpoint(2, 3, 60), + checkpoint(3, 2, 70), + ]; + let (items, held) = order_items(&events, &records, &BTreeSet::new(), false); + assert_eq!( + held, 2, + "the record for firing 3 holds itself and the one after it" + ); + assert_eq!(names(&items), vec![ + "execution 0/13/0", + "execution 0/14/0", + "platform 1", + ]); + + // Once the run finished, a record for a firing that never finished + // keeps its clock order; the firing's own records still follow its + // finish, in their seq order. + let (items, held) = order_items(&events, &records, &BTreeSet::new(), true); + assert_eq!(held, 0, "nothing is held once the run finished"); + assert_eq!(names(&items), vec![ + "platform 2", + "execution 0/13/0", + "execution 0/14/0", + "platform 1", + "platform 3", + ]); + } +} diff --git a/lib/components/fabro-petri/src/projector/signalling.rs b/lib/components/fabro-petri/src/projector/signalling.rs new file mode 100644 index 000000000..c81a97587 --- /dev/null +++ b/lib/components/fabro-petri/src/projector/signalling.rs @@ -0,0 +1,99 @@ +//! The in-process wake-up: a run store whose appends signal the projector, +//! for a run that executes in the same process as the projector, over the +//! SQLite store directly, where no append endpoint is there to signal. + +use std::sync::Arc; + +use fabro_types::RunId; +use petri_execution::{Access, RunKey}; +use petri_store::StoreError; + +use super::Projector; +use crate::projection; + +impl Projector { + /// A run store whose appends signal this projector: for a run that + /// executes in the same process as the projector, over the SQLite store + /// directly, where no append endpoint is there to signal. The signal is + /// sent after the store's append returned, so the records it covers are + /// durable before the view sees them. + pub fn observe_store( + self: &Arc, + inner: Arc, + ) -> Arc { + Arc::new(SignallingStore { + inner, + projector: Arc::clone(self), + }) + } +} + +/// A run store that signals a projector after each append. +struct SignallingStore { + inner: Arc, + projector: Arc, +} + +#[async_trait::async_trait] +impl petri_execution::RunStore for SignallingStore { + async fn open( + &self, + key: &RunKey, + access: Access, + ) -> Result, StoreError> { + let logs = self.inner.open(key, access).await?; + Ok(Arc::new(SignallingLogs { + inner: logs, + run_id: projection::run_id_of(key.as_str()), + projector: Arc::clone(&self.projector), + })) + } +} + +struct SignallingLogs { + inner: Arc, + run_id: Option, + projector: Arc, +} + +#[async_trait::async_trait] +impl petri_execution::RunLogs for SignallingLogs { + fn locator(&self) -> String { + self.inner.locator() + } + + async fn append( + &self, + log: &petri_execution::LogId, + records: &[petri_execution::Record], + ) -> Result<(), StoreError> { + self.inner.append(log, records).await?; + if let Some(run_id) = self.run_id { + self.projector.signal(run_id); + } + Ok(()) + } + + async fn read( + &self, + log: &petri_execution::LogId, + ) -> Result, StoreError> { + 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 + } + + async fn get_blob(&self, digest: petri_store::Digest) -> Result>, StoreError> { + self.inner.get_blob(digest).await + } +} diff --git a/lib/components/fabro-petri/src/projector/stream.rs b/lib/components/fabro-petri/src/projector/stream.rs new file mode 100644 index 000000000..f35e63d10 --- /dev/null +++ b/lib/components/fabro-petri/src/projector/stream.rs @@ -0,0 +1,113 @@ +//! The run's stream as the view tables hold it: one row per folded item, +//! written by a pass and read back in `stream_seq` order. + +use fabro_db::DbPool; +use fabro_types::{RunId, RunStreamItem, RunStreamItemKind}; +use petri_execution::events::{EventId, EventSource}; + +use super::{Positions, ProjectError}; +use crate::projection::{Item, RunView}; + +pub(crate) struct StreamRow { + pub(super) stream_seq: u64, + pub(super) item_kind: &'static str, + pub(super) item_id: String, + pub(super) event_json: String, +} + +/// Fold the items into the view in order, each at the next delivery +/// sequence, advancing the positions with each, and produce its stream +/// row. +pub(crate) fn stream_rows( + items: &[Item<'_>], + view: &mut RunView, + positions: &mut Positions, + stream_seq: &mut u64, +) -> Result, ProjectError> { + let mut rows = Vec::with_capacity(items.len()); + for item in items { + *stream_seq += 1; + view.fold(item, *stream_seq); + let row = match item { + Item::Petri(event) => { + positions.advance(event.id); + StreamRow { + stream_seq: *stream_seq, + item_kind: "petri", + item_id: event_id_text(&event.id), + event_json: serde_json::to_string(event).map_err(ProjectError::Encode)?, + } + } + Item::Platform(record) => { + positions.platform_seq = record.seq; + StreamRow { + stream_seq: *stream_seq, + item_kind: "platform", + item_id: record.seq.to_string(), + event_json: serde_json::to_string(record).map_err(ProjectError::Encode)?, + } + } + }; + rows.push(row); + } + Ok(rows) +} + +/// The run's stream past the cursor, read from the view tables: up to +/// `limit` rows with `stream_seq > after`, in order, in Fabro's envelope. +pub(super) async fn stream_after( + views: &DbPool, + run_id: RunId, + after: u64, + limit: usize, +) -> Result, ProjectError> { + let rows: Vec<(i64, String, String, String)> = sqlx::query_as( + "SELECT stream_seq, item_kind, item_id, event_json FROM petri_stream WHERE run_id = ? AND \ + stream_seq > ? ORDER BY stream_seq LIMIT ?", + ) + .bind(run_id.to_string()) + .bind(column(after)) + .bind(i64::try_from(limit).unwrap_or(i64::MAX)) + .fetch_all(views) + .await + .map_err(ProjectError::Database)?; + rows.into_iter() + .map(|(stream_seq, item_kind, item_id, event_json)| { + let item: serde_json::Value = + serde_json::from_str(&event_json).map_err(ProjectError::Encode)?; + let kind = match item_kind.as_str() { + "platform" => RunStreamItemKind::Platform, + _ => RunStreamItemKind::Petri, + }; + let recorded_at = item + .get("recorded_at") + .and_then(serde_json::Value::as_u64) + .unwrap_or(0); + Ok(RunStreamItem { + run_id, + stream_seq: u64::try_from(stream_seq).unwrap_or(0), + kind, + id: item_id, + recorded_at, + item, + }) + }) + .collect() +} + +/// A Petri event id as the stream names it: `//`. +#[must_use] +pub(crate) fn event_id_text(id: &EventId) -> String { + format!("{}/{}/{}", log_text(&id.source), id.seq, id.index) +} + +pub(super) fn log_text(source: &EventSource) -> String { + match source { + EventSource::Coordinator => "coordinator".to_string(), + EventSource::Execution { execution } => format!("execution {execution}"), + } +} + +pub(super) fn column(value: u64) -> i64 { + i64::try_from(value).unwrap_or(i64::MAX) +} diff --git a/lib/components/fabro-petri/src/test_support.rs b/lib/components/fabro-petri/src/test_support.rs index 38b37ce07..91793f60e 100644 --- a/lib/components/fabro-petri/src/test_support.rs +++ b/lib/components/fabro-petri/src/test_support.rs @@ -1,24 +1,35 @@ //! Petri's test kit, for Fabro crates that check a store implementation -//! against Petri's contract from their own tests, an in-memory platform +//! against Petri's contract from their own tests; an in-memory platform //! record store and an in-memory blob table for tests of the hooks and -//! recovery. Compiled only with the `test-support` feature, which a +//! recovery; and the readers over the view tables a test compares a live +//! view with. Compiled only with the `test-support` feature, which a //! dev-dependency turns on. -use std::collections::HashMap; +use std::collections::{BTreeSet, HashMap}; use std::sync::Mutex; use std::time::Duration; use async_trait::async_trait; use bytes::Bytes; -use fabro_store::platform_records::now_ms; -use fabro_store::{PlatformRecord, PlatformRecordKind, StagePosition, StoredPlatformRecord}; +use fabro_db::DbPool; +use fabro_store::platform_records::{PlatformRecordStore, now_ms}; +use fabro_store::{ + PlatformRecord, PlatformRecordKind, RunProjection, StagePosition, StoredPlatformRecord, +}; use fabro_types::{BlobHash, RunId}; +use fabro_util::error::collect_chain; use fabro_util::sync; +use petri_execution::events::{self, RunEvent}; +use petri_execution::{Access, CoordinatorEvent, RunKey, RunStore as _}; +use petri_store::StoreError; pub use petri_testkit::run_store; +use tracing::warn; +use crate::SqliteRunStore; use crate::blobs::Blobs; use crate::platform_records::{PlatformRecordError, PlatformRecords}; -use crate::projector::Projector; +use crate::projection::RunView; +use crate::projector::{self, Positions, ProjectError, Projector, order, stream}; /// Whether the projector keeps a cache for the run: the replay and the /// view its passes continue from. @@ -126,3 +137,100 @@ impl PlatformRecords for MemoryPlatformRecords { .collect()) } } + +/// The run's projection rebuilt from its records alone, with nothing +/// stored: what a fresh projector would commit over the same records. A test +/// compares it with the live view. `records` and `views` are the two pools +/// [`Projector::new`] takes. +pub async fn rebuild( + records: &DbPool, + views: &DbPool, + run_id: RunId, +) -> Result<(Option, Positions, u64), ProjectError> { + let store = SqliteRunStore::new(records.clone()); + let platform = PlatformRecordStore::new(views.clone()); + let key = RunKey::new(run_id.to_string()); + let platform_records = platform.read(&run_id).await.map_err(ProjectError::Store)?; + let events = match store.open(&key, Access::Read).await { + Ok(logs) => events::replay_run(&*logs) + .await + .inspect_err(|error| { + warn!(error = %collect_chain(error).join(": "), "rebuild: the run does not replay"); + }) + .unwrap_or_default(), + Err(StoreError::NotFound { .. }) => Vec::new(), + Err(error) => return Err(ProjectError::Open(error)), + }; + let run_finished = events.iter().any(|event| { + matches!( + event.coordinator(), + Some(CoordinatorEvent::RunFinished { .. }) + ) + }); + let (items, _held) = + order::order_items(&events, &platform_records, &BTreeSet::new(), run_finished); + let mut view = RunView::new(); + let mut positions = Positions::default(); + let mut stream_seq = 0; + stream::stream_rows(&items, &mut view, &mut positions, &mut stream_seq)?; + Ok((view.projection, positions, stream_seq)) +} + +/// The stored view's positions and stream sequence; `views` is the pool +/// the view tables live in. +pub async fn stored_positions( + views: &DbPool, + run_id: RunId, +) -> Result, ProjectError> { + projector::stored_positions(views, run_id).await +} + +/// The stored view's projection, for a test or a reader outside the store. +pub async fn stored_projection( + views: &DbPool, + run_id: RunId, +) -> Result, ProjectError> { + let json: Option = + sqlx::query_scalar("SELECT projection_json FROM petri_projection WHERE run_id = ?") + .bind(run_id.to_string()) + .fetch_optional(views) + .await + .map_err(ProjectError::Database)?; + json.map(|json| serde_json::from_str(&json).map_err(ProjectError::Encode)) + .transpose() +} + +/// The stream rows of a run: `(stream_seq, item_kind, item_id)`, in order. +pub async fn stored_stream( + views: &DbPool, + run_id: RunId, +) -> Result, ProjectError> { + let rows: Vec<(i64, String, String)> = sqlx::query_as( + "SELECT stream_seq, item_kind, item_id FROM petri_stream WHERE run_id = ? ORDER BY stream_seq", + ) + .bind(run_id.to_string()) + .fetch_all(views) + .await + .map_err(ProjectError::Database)?; + Ok(rows + .into_iter() + .map(|(seq, kind, id)| (u64::try_from(seq).unwrap_or(0), kind, id)) + .collect()) +} + +/// Every stored platform record of a run, for a reader outside the store. +pub async fn stored_platform_records( + views: &DbPool, + run_id: RunId, +) -> Result, ProjectError> { + PlatformRecordStore::new(views.clone()) + .read(&run_id) + .await + .map_err(ProjectError::Store) +} + +/// A recorded event's projection is what `RunEvent` serializes to. +#[must_use] +pub fn event_json(event: &RunEvent) -> serde_json::Value { + serde_json::to_value(event).unwrap_or_default() +} diff --git a/lib/components/fabro-petri/tests/projection.rs b/lib/components/fabro-petri/tests/projection.rs index 31f6e6ca1..aa7ff9914 100644 --- a/lib/components/fabro-petri/tests/projection.rs +++ b/lib/components/fabro-petri/tests/projection.rs @@ -418,15 +418,15 @@ fn json(value: &T) -> serde_json::Value { /// The stored view equals the view rebuilt from the records alone: the /// projection, the positions and the delivery sequence. async fn assert_view_equals_rebuild(pool: &DbPool, run_id: RunId) { - let stored = projector::stored_projection(pool, run_id) + let stored = petri_support::stored_projection(pool, run_id) .await .expect("the stored projection reads") .expect("the run has a stored projection"); - let (stored_positions, stored_stream_seq) = projector::stored_positions(pool, run_id) + let (stored_positions, stored_stream_seq) = petri_support::stored_positions(pool, run_id) .await .expect("the positions read") .expect("the run has positions"); - let (rebuilt, positions, stream_seq) = projector::rebuild(pool, pool, run_id) + let (rebuilt, positions, stream_seq) = petri_support::rebuild(pool, pool, run_id) .await .expect("the run rebuilds"); let rebuilt = rebuilt.expect("the rebuild has a projection"); @@ -442,7 +442,7 @@ async fn assert_view_equals_rebuild(pool: &DbPool, run_id: RunId) { positions.petri.sort(); assert_eq!(stored_positions, positions); assert_eq!(stored_stream_seq, stream_seq); - let stream = projector::stored_stream(pool, run_id) + let stream = petri_support::stored_stream(pool, run_id) .await .expect("the stream reads"); let seqs: Vec = stream.iter().map(|(seq, _, _)| *seq).collect(); @@ -454,7 +454,7 @@ async fn assert_view_equals_rebuild(pool: &DbPool, run_id: RunId) { } async fn stage_states(pool: &DbPool, run_id: RunId) -> Vec<(String, StageState)> { - let stored = projector::stored_projection(pool, run_id) + let stored = petri_support::stored_projection(pool, run_id) .await .expect("the stored projection reads") .expect("the run has a stored projection"); @@ -472,7 +472,7 @@ async fn the_hello_bundle_projects_live_as_it_rebuilds() { let scenario = hello_scenario().await; run_live(&scenario).await; assert_view_equals_rebuild(&scenario.pool, scenario.run_id).await; - let stored = projector::stored_projection(&scenario.pool, scenario.run_id) + let stored = petri_support::stored_projection(&scenario.pool, scenario.run_id) .await .expect("reads") .expect("stored"); @@ -501,7 +501,7 @@ async fn a_large_output_projects_as_its_blob_reference() { let scenario = large_output_scenario().await; run_live(&scenario).await; assert_view_equals_rebuild(&scenario.pool, scenario.run_id).await; - let stored = projector::stored_projection(&scenario.pool, scenario.run_id) + let stored = petri_support::stored_projection(&scenario.pool, scenario.run_id) .await .expect("reads") .expect("stored"); @@ -525,7 +525,7 @@ async fn a_command_workflow_projects_live_as_it_rebuilds() { let scenario = command_scenario().await; run_live(&scenario).await; assert_view_equals_rebuild(&scenario.pool, scenario.run_id).await; - let stored = projector::stored_projection(&scenario.pool, scenario.run_id) + let stored = petri_support::stored_projection(&scenario.pool, scenario.run_id) .await .expect("reads") .expect("stored"); @@ -554,7 +554,7 @@ async fn a_parallel_workflow_projects_its_branches_as_child_executions() { let scenario = parallel_scenario().await; run_live(&scenario).await; assert_view_equals_rebuild(&scenario.pool, scenario.run_id).await; - let stored = projector::stored_projection(&scenario.pool, scenario.run_id) + let stored = petri_support::stored_projection(&scenario.pool, scenario.run_id) .await .expect("reads") .expect("stored"); @@ -593,7 +593,7 @@ async fn dropped_wake_ups_are_caught_up_by_the_next_signal() { let scenario = command_scenario().await; run_unobserved(&scenario).await; assert!( - projector::stored_projection(&scenario.pool, scenario.run_id) + petri_support::stored_projection(&scenario.pool, scenario.run_id) .await .expect("reads") .is_none(), @@ -780,12 +780,13 @@ async fn a_crash_between_the_record_commit_and_the_view_applies_only_the_suffix( .expect("the first pass commits"); assert!(!pass.skipped); assert!(!pass.health.complete, "the run has not finished"); - let (positions_before, stream_before) = projector::stored_positions(&replayed, scenario.run_id) - .await - .expect("reads") - .expect("positions"); + let (positions_before, stream_before) = + petri_support::stored_positions(&replayed, scenario.run_id) + .await + .expect("reads") + .expect("positions"); assert_eq!(pass.stream_seq, stream_before); - let stream_rows_before = projector::stored_stream(&replayed, scenario.run_id) + let stream_rows_before = petri_support::stored_stream(&replayed, scenario.run_id) .await .expect("reads") .len(); @@ -801,7 +802,7 @@ async fn a_crash_between_the_record_commit_and_the_view_applies_only_the_suffix( "{crashed:?}" ); assert_eq!( - projector::stored_positions(&replayed, scenario.run_id) + petri_support::stored_positions(&replayed, scenario.run_id) .await .expect("reads") .expect("positions"), @@ -813,11 +814,12 @@ async fn a_crash_between_the_record_commit_and_the_view_applies_only_the_suffix( let after = Projector::new(replayed.clone(), replayed.clone()); let report = after.startup_pass().await.expect("the restart catches up"); assert_eq!((report.runs, report.projected), (1, 1)); - let (positions_after, stream_after) = projector::stored_positions(&replayed, scenario.run_id) - .await - .expect("reads") - .expect("positions"); - let stream_rows_after = projector::stored_stream(&replayed, scenario.run_id) + let (positions_after, stream_after) = + petri_support::stored_positions(&replayed, scenario.run_id) + .await + .expect("reads") + .expect("positions"); + let stream_rows_after = petri_support::stored_stream(&replayed, scenario.run_id) .await .expect("reads"); // Only the suffix was applied: the stream grew by the suffix's events, @@ -845,11 +847,11 @@ async fn a_crash_between_the_record_commit_and_the_view_applies_only_the_suffix( // And the copy agrees with the run projected in one go over the source. let source = Projector::new(scenario.pool.clone(), scenario.pool.clone()); source.startup_pass().await.expect("the source projects"); - let whole = projector::stored_projection(&scenario.pool, scenario.run_id) + let whole = petri_support::stored_projection(&scenario.pool, scenario.run_id) .await .expect("reads") .expect("stored"); - let pieced = projector::stored_projection(&replayed, scenario.run_id) + let pieced = petri_support::stored_projection(&replayed, scenario.run_id) .await .expect("reads") .expect("stored"); @@ -904,11 +906,11 @@ async fn a_restarted_projector_agrees_over_nested_child_executions() { let whole = Projector::new(scenario.pool.clone(), scenario.pool.clone()); whole.startup_pass().await.expect("the source projects"); - let one_go = projector::stored_projection(&scenario.pool, scenario.run_id) + let one_go = petri_support::stored_projection(&scenario.pool, scenario.run_id) .await .expect("reads") .expect("stored"); - let restarted = projector::stored_projection(&staged, scenario.run_id) + let restarted = petri_support::stored_projection(&staged, scenario.run_id) .await .expect("reads") .expect("stored"); @@ -937,11 +939,11 @@ async fn a_torn_tail_holds_the_view_and_reports_the_run_incomplete() { .await .expect("the clean pass commits"); assert!(clean.health.complete, "{:?}", clean.health.incomplete); - let (positions, stream_seq) = projector::stored_positions(&scenario.pool, scenario.run_id) + let (positions, stream_seq) = petri_support::stored_positions(&scenario.pool, scenario.run_id) .await .expect("reads") .expect("positions"); - let before = projector::stored_projection(&scenario.pool, scenario.run_id) + let before = petri_support::stored_projection(&scenario.pool, scenario.run_id) .await .expect("reads") .expect("stored"); @@ -973,7 +975,7 @@ async fn a_torn_tail_holds_the_view_and_reports_the_run_incomplete() { held.health ); let (positions_after, stream_after) = - projector::stored_positions(&scenario.pool, scenario.run_id) + petri_support::stored_positions(&scenario.pool, scenario.run_id) .await .expect("reads") .expect("positions"); @@ -982,7 +984,7 @@ async fn a_torn_tail_holds_the_view_and_reports_the_run_incomplete() { "the view did not advance past the tear" ); assert_eq!(stream_after, stream_seq); - let after = projector::stored_projection(&scenario.pool, scenario.run_id) + let after = petri_support::stored_projection(&scenario.pool, scenario.run_id) .await .expect("reads") .expect("stored"); @@ -1240,9 +1242,10 @@ impl GateRun { async fn pending(&self) -> fabro_types::RunProjection { let deadline = Instant::now() + Duration::from_secs(30); loop { - let stored = projector::stored_projection(&self.scenario.pool, self.scenario.run_id) - .await - .expect("the stored projection reads"); + let stored = + petri_support::stored_projection(&self.scenario.pool, self.scenario.run_id) + .await + .expect("the stored projection reads"); if let Some(stored) = stored.filter(|stored| !stored.pending_interviews.is_empty()) { return stored; } @@ -1255,7 +1258,7 @@ impl GateRun { } async fn stored(&self) -> fabro_types::RunProjection { - projector::stored_projection(&self.scenario.pool, self.scenario.run_id) + petri_support::stored_projection(&self.scenario.pool, self.scenario.run_id) .await .expect("the stored projection reads") .expect("the run has a stored projection") @@ -1282,7 +1285,7 @@ async fn an_expired_question_is_pending_while_the_gate_waits_and_closes_on_the_e let run_id = gate.scenario.run_id; let deadline = Instant::now() + Duration::from_secs(30); loop { - let stored = projector::stored_projection(&pool, run_id) + let stored = petri_support::stored_projection(&pool, run_id) .await .expect("the stored projection reads"); if let Some(stored) = stored.filter(|stored| !stored.pending_interviews.is_empty()) {