Project a Petri run's records into Fabro's run view

The projection folds Petri's public events (replay_since over the run's
stored records) and Fabro's platform records into the RunProjection the
API serves, row by row as VIEWS.md maps them. The stage key is the
execution and firing; the StageId label is node@visit, made unique with
the execution when two child invocations would share one. A stage's
first_event_seq is the milliseconds from the run's creation to its
visit.started, so the view built live equals the view rebuilt from the
records whatever order two logs' records were committed in.

The projector is the view pass and its wake-up. Records first: an append
returns before any view work; a pass reads what is committed, folds the
items past the committed positions, and writes the projection document,
the ordered stream (one stream_seq per Petri event or platform record,
with the item's own identity beside it) and the narrowed runs row in one
later transaction. Signals coalesce per run, a lost signal costs only
latency, the startup pass folds every run the view trails, and a pass
that races a platform record leaves the view alone and runs again. A
torn tail holds the view where it stands and reports the run incomplete
with the replay's error; inspect_run decides completeness once the run
recorded its finish.

The tests build the view live for the hello bundle, a command workflow
and a two-branch parallel workflow and compare it with the rebuild; drop
every wake-up and catch up by a signal and by the startup pass; crash
between the record commit and the view transaction and apply only the
suffix; restart the projector over child executions; and hold at a torn
tail.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
This commit is contained in:
Bryan Helmkamp 2026-09-17 22:32:23 -04:00
parent 162791979c
commit e2bf05c0f0
No known key found for this signature in database
8 changed files with 3053 additions and 6 deletions

2
Cargo.lock generated
View file

@ -2891,6 +2891,7 @@ dependencies = [
"anyhow",
"async-trait",
"bytes",
"chrono",
"fabro-api",
"fabro-auth",
"fabro-client",
@ -2899,6 +2900,7 @@ dependencies = [
"fabro-llm",
"fabro-store",
"fabro-types",
"fabro-util",
"lithos-llm",
"petri-attractor-steps",
"petri-execution",

View file

@ -25,6 +25,7 @@ fabro-db = { path = "../../foundation/fabro-db" }
fabro-http.workspace = true
fabro-store = { path = "../fabro-store" }
fabro-types = { path = "../../foundation/fabro-types" }
fabro-util = { path = "../../foundation/fabro-util" }
petri_runtime.workspace = true
petri_execution.workspace = true
petri_store.workspace = true
@ -36,6 +37,7 @@ petri_testkit = { workspace = true, optional = true }
anyhow.workspace = true
bytes.workspace = true
async-trait.workspace = true
chrono = { workspace = true, features = ["serde"] }
serde.workspace = true
serde_json.workspace = true
sqlx.workspace = true
@ -49,5 +51,6 @@ tracing.workspace = true
fabro-auth = { path = "../../foundation/fabro-auth", features = ["test-support"] }
fabro-llm = { path = "../fabro-llm", features = ["test-support"] }
fabro-store = { path = "../fabro-store", features = ["test-support"] }
fabro-types = { path = "../../foundation/fabro-types", features = ["test-support"] }
petri_testkit.workspace = true
tokio = { workspace = true, features = ["macros", "rt-multi-thread"] }

View file

@ -46,8 +46,39 @@ Every adapter the integration plan describes lands here.
- `petri`: the Petri store vocabulary re-exported for the server, which
answers the worker endpoints from a `SqliteRunStore` without naming a Petri
package in its own manifest.
- `projection`: the fold of a Petri run's public events (`replay_since` over
its records) and Fabro's platform records (`fabro-store`'s
`platform_records`) into the `RunProjection` the API serves, row by row as
`VIEWS.md` maps them. The stage key is `(execution, firing)`; the
`StageId` label is `node@visit`, made unique with the execution when two
child invocations would share one.
- `projector`: the view pass and its wake-up. Records first: Petri's append
and a platform record's insert return before any view work; a pass reads
what is committed, folds the items past the committed positions, and
writes the projection document (`petri_projection`), the ordered stream
(`petri_stream`, one `stream_seq` per Petri event or platform record) and
the narrowed `runs` row in one later transaction. The server signals the
projector after each committed worker append, after each committed
platform record (the run summary store's hook), at worker exit and, over
every Petri run, at startup. A run that executes in the server process
goes through `Projector::observe_store`, which signals after each append.
A torn tail (a record Petri cannot read) holds the view where it stands
and reports the run incomplete with the reason.
- The platform adapters the plan adds after it: hooks, interviews over
Fabro's API, secrets, output storage, run tools, the event projection.
Fabro's API, secrets, output storage, run tools.
### What the projection leaves default
`VIEWS.md` rows with no source yet, or whose source this crate does not read
yet, keep their default value in the projection: `StageProjection.diff` and
`Conclusion.diff.patch` (the checkpoint's `patch_blob` is not resolved),
`Checkpoint`'s engine-derived maps (`completed_nodes`, `node_retries`,
`context_values`, `node_outcomes`, `next_node_id`), `agent_tools`,
`permission_level`, `script_invocation` and `script_timing`, a stage's
`notes`, `StageCompletion` details for a `parsed.note`, the sandbox instance
(the matrix's two gaps), `Run.ask_fabro`, an interview option's
`description` and `preview`, the pull request `creation` state, and the
run's notices, notifications and pairings (recorded, not shown).
A run goes to Petri when its workflow version's `workflow.toml` names
`engine = "petri"` in `[workflow]`, or when the server's
@ -79,6 +110,16 @@ Integration tests live under `tests/`:
operator release, lease exclusivity, a crash between appends, and blob
interoperation with Fabro's `BlobStore`.
- `projection.rs` builds the view live (every append signals the
projector) for the `hello` bundle on the stub registry, a command-only
workflow and a two-branch parallel workflow, and checks it equals the view
rebuilt from the records alone (`projector::rebuild`); catches a view up
after every wake-up was dropped, by a signal and by the startup pass;
recovers a crash between the record commit and the view transaction by
applying only the missing suffix, with the positions and `stream_seq`
continuing; runs two projectors over one store with child executions; and
holds the view at a torn tail. All skip without the host plugin.
The conformance suite over `HttpRunStore` needs a server to talk to, so it
lives with the server's integration tests
(`lib/apps/fabro-server/tests/it/api/petri_store.rs`), which reach the suite
@ -91,10 +132,12 @@ ulimit -n 4096 && cargo nextest run -p fabro-petri
```
The server's end-to-end coverage is `lib/apps/fabro-server/tests/it/scenario/petri.rs`:
the `hello` bundle on the OpenAI twin and a command-only bundle run to
completion through the create handler and the scheduler, in the server
process under its test override, under the version flag and under the
server setting, and Petri's diagnostics refuse a run at create. The
the `hello` bundle on the OpenAI twin, a command-only bundle and a
two-branch parallel bundle run to completion through the create handler and
the scheduler, in the server process under its test override, under the
version flag and under the server setting, with `GET /runs/{id}/state`
serving the projection over Petri's records, and Petri's diagnostics refuse
a run at create. The
server's `petri_runs` unit tests cover the lease ending at worker exit and
the restart reconcile that relaunches a worker in resume mode.

View file

@ -6,6 +6,11 @@ comes from once Petri's records are the store. It is written before any view
changes. F2.2 (the projection), F2.3 (platform records) and F2.4 (API, CLI,
web) build from it.
F2.2 and F2.3 implement this matrix: `src/projection.rs` is the fold,
`src/projector.rs` the view pass and its wake-up, and `fabro-store`'s
`platform_records` module the platform record kinds and their table. The
crate README names the rows the fold still leaves default.
Sources are named three ways:
- A Petri event, by its `<subject>.<verb>` name from

View file

@ -23,8 +23,11 @@
//! - [`HttpRunStore`]: the same store as a run's worker process reaches it,
//! over the server's API with the worker's token and its launch id as the
//! lease owner;
//! - [`projection`] and [`projector`]: the view of a Petri run, folded from its
//! records and Fabro's platform records, and the pass that writes it after
//! each committed record;
//! - the platform adapters still to come: hooks, interviews over Fabro's API,
//! secrets, output storage, the run tools, the event projection.
//! secrets, output storage, the run tools.
//!
//! The Petri packages are pinned by revision in the workspace `Cargo.toml`
//! under `petri_*` keys.
@ -35,6 +38,8 @@ pub mod engine;
pub mod http_store;
pub mod interviewer;
pub mod petri;
pub mod projection;
pub mod projector;
pub mod run_store;
pub mod runtime;
#[cfg(feature = "test-support")]

File diff suppressed because it is too large Load diff

View file

@ -0,0 +1,761 @@
//! The projector: the view pass that folds a Petri run's committed records
//! into its stored projection, and the wake-up that drives it.
//!
//! # Commit rule
//!
//! Records first. Petri's append (the worker's append endpoint, then
//! `SqliteRunStore::append`) and a platform record's insert are the
//! durability boundaries, and both return before any view work. A view pass
//! then reads what is committed, folds the items past the positions the
//! view last committed, and writes the derived rows in one later
//! transaction together with the new positions: the last event consumed per
//! 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.
//!
//! 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.
//!
//! # Where it runs
//!
//! In the server. [`Projector::signal`] schedules a pass for a run: the
//! server calls it after each committed worker append and, through the run
//! summary store's hook, after each committed platform record; signals
//! that arrive while a pass runs coalesce into one more pass. A signal is a
//! wake-up only, never a source of facts: a signal that is lost costs
//! nothing but latency, because the next signal or the startup pass
//! ([`Projector::startup_pass`]) folds everything the view still trails.
//!
//! # A torn tail
//!
//! A record the store holds that Petri cannot read (a gap in a log, a line
//! that does not decode) fails the replay. The pass then advances no Petri
//! position, folds only the platform records, and reports the run's record
//! as incomplete with the replay's error; `inspect_run` decides
//! completeness once the run has recorded its finish.
use std::collections::{BTreeMap, HashMap};
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::{Arc, Mutex, MutexGuard, PoisonError};
use std::time::Duration;
use fabro_db::DbPool;
use fabro_store::platform_records::{PlatformRecordStore, StoredPlatformRecord, now_ms};
use fabro_store::{RunProjection, RunSummaryStore};
use fabro_types::RunId;
use fabro_util::error::collect_chain;
use petri_execution::events::{self, EventId, EventSource, RunEvent};
use petri_execution::{Access, RunKey, RunStore as _, inspect};
use petri_store::StoreError;
use serde::{Deserialize, Serialize};
use tokio::time;
use tracing::{debug, info, warn};
use crate::SqliteRunStore;
use crate::projection::{self, FoldState, Item, RecordHealth, RunView};
/// The positions a view committed: the last event consumed per Petri log,
/// and the last platform record consumed.
#[derive(Clone, Debug, Default, PartialEq, Eq, Serialize, Deserialize)]
pub struct Positions {
#[serde(default)]
pub petri: Vec<EventId>,
#[serde(default)]
pub platform_seq: u64,
}
impl Positions {
fn held(&self) -> BTreeMap<EventSource, EventId> {
self.petri.iter().map(|id| (id.source, *id)).collect()
}
fn advance(&mut self, id: EventId) {
match self.petri.iter_mut().find(|held| held.source == id.source) {
Some(held) => {
if id > *held {
*held = id;
}
}
None => self.petri.push(id),
}
}
}
/// What one pass did.
#[derive(Clone, Debug, PartialEq, Eq)]
pub struct PassReport {
pub run_id: RunId,
/// The pass found nothing past the committed positions and wrote nothing.
pub skipped: bool,
/// The view was left alone because a platform record landed during the
/// pass; the projector runs the pass again.
pub contended: bool,
pub petri_events: usize,
pub platform_records: usize,
/// The last delivery sequence the view holds.
pub stream_seq: u64,
pub positions: Positions,
pub health: RecordHealth,
}
/// What the startup pass did.
#[derive(Clone, Debug, Default, PartialEq, Eq)]
pub struct StartupReport {
pub runs: usize,
pub projected: usize,
}
/// Why a pass could not run or commit.
#[derive(Debug, thiserror::Error)]
pub enum ProjectError {
#[error("the run's Petri record could not be opened")]
Open(#[source] StoreError),
#[error("the projection tables could not be read or written")]
Database(#[source] sqlx::Error),
#[error("the platform records could not be read or written")]
Store(#[source] fabro_store::Error),
#[error("the view could not be encoded")]
Encode(#[source] serde_json::Error),
#[error("the pass was stopped before its view transaction (injected)")]
Injected,
}
/// The stored view of a run, as the projection tables hold it.
struct StoredView {
view: RunView,
positions: Positions,
stream_seq: u64,
}
/// A run's pass state under the projector's lock.
#[derive(Default)]
struct Slot {
running: bool,
pending: bool,
}
/// The projector over one database.
pub struct Projector {
pool: DbPool,
store: SqliteRunStore,
platform: PlatformRecordStore,
slots: Mutex<HashMap<RunId, Slot>>,
/// Test-only: stop the next pass after its reads, before its view
/// transaction, as a crash there would.
fault: AtomicBool,
}
impl std::fmt::Debug for Projector {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("Projector").finish_non_exhaustive()
}
}
impl Projector {
/// A projector over a pool whose migrations have run.
#[must_use]
pub fn new(pool: DbPool) -> Arc<Self> {
Arc::new(Self {
store: SqliteRunStore::new(pool.clone()),
platform: PlatformRecordStore::new(pool.clone()),
pool,
slots: Mutex::default(),
fault: AtomicBool::new(false),
})
}
/// Schedule a pass for the run. A pass already running for it runs once
/// more when it ends; any number of signals in between coalesce.
pub fn signal(self: &Arc<Self>, run_id: RunId) {
{
let mut slots = lock(&self.slots);
let slot = slots.entry(run_id).or_default();
if slot.running {
slot.pending = true;
return;
}
slot.running = true;
}
let projector = Arc::clone(self);
tokio::spawn(async move {
loop {
let again = match projector.project_run(run_id).await {
Ok(report) => report.contended,
Err(error) => {
warn!(
run_id = %run_id,
error = %collect_chain(&error).join(": "),
"Petri projection pass failed; the next signal retries it"
);
false
}
};
let mut slots = lock(&projector.slots);
let slot = slots.entry(run_id).or_default();
if again || slot.pending {
slot.pending = false;
continue;
}
slot.running = false;
return;
}
});
}
/// Wait until no pass is running or pending for the run: a test's way
/// to observe the view after its signals.
pub async fn settle(&self, run_id: RunId) {
loop {
let idle = {
let slots = lock(&self.slots);
slots
.get(&run_id)
.is_none_or(|slot| !slot.running && !slot.pending)
};
if idle {
return;
}
time::sleep(Duration::from_millis(5)).await;
}
}
/// Stop the next pass after its reads and before its view transaction,
/// as a crash there would, once.
pub fn fail_before_view(&self) {
self.fault.store(true, Ordering::SeqCst);
}
/// 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.
pub async fn startup_pass(&self) -> Result<StartupReport, ProjectError> {
let ids: Vec<String> = sqlx::query_scalar(
"SELECT run_id FROM petri_runs UNION SELECT run_id FROM platform_records ORDER BY 1",
)
.fetch_all(&self.pool)
.await
.map_err(ProjectError::Database)?;
let mut report = StartupReport::default();
for id in ids {
let Some(run_id) = projection::run_id_of(&id) else {
debug!(run_key = %id, "Petri run key is not a Fabro run id; not projected");
continue;
};
report.runs += 1;
let pass = self.project_run(run_id).await?;
if !pass.skipped {
report.projected += 1;
}
}
if report.projected > 0 {
info!(
runs = report.runs,
projected = report.projected,
"Petri projections caught up at startup"
);
}
Ok(report)
}
/// One view pass for the run.
pub async fn project_run(&self, run_id: RunId) -> Result<PassReport, ProjectError> {
let stored = self.load_view(&run_id).await?;
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 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,
stream_seq: stored.stream_seq,
positions: stored.positions,
health: stored.view.state.health,
});
}
let StoredView {
mut view,
mut positions,
mut stream_seq,
} = stored;
let platform_records = self
.platform
.read_after(&run_id, positions.platform_seq)
.await
.map_err(ProjectError::Store)?;
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),
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 mut items: Vec<(u64, u8, Item<'_>)> =
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 mut rows: Vec<StreamRow> = 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);
}
view.state.health = self.health(&key, &view.state, replay_failure).await?;
if self.fault.swap(false, Ordering::SeqCst) {
return Err(ProjectError::Injected);
}
let mut tx = self
.pool
.begin_with("BEGIN IMMEDIATE")
.await
.map_err(ProjectError::Database)?;
let head_now: i64 = sqlx::query_scalar(
"SELECT COALESCE(MAX(seq), 0) FROM platform_records WHERE run_id = ?",
)
.bind(run_id.to_string())
.fetch_one(&mut *tx)
.await
.map_err(ProjectError::Database)?;
if u64::try_from(head_now).unwrap_or(0) != positions.platform_seq {
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(),
});
}
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)?;
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 \
SET projection_json = excluded.projection_json, fold_json = excluded.fold_json, \
positions_json = excluded.positions_json, stream_seq = excluded.stream_seq, \
updated_at_ms = excluded.updated_at_ms",
)
.bind(run_id.to_string())
.bind(projection_json)
.bind(fold_json)
.bind(positions_json)
.bind(column(stream_seq))
.bind(column(now_ms()))
.execute(&mut *tx)
.await
.map_err(ProjectError::Database)?;
for row in &rows {
sqlx::query(
"INSERT INTO petri_stream (run_id, stream_seq, item_kind, item_id, event_json) \
VALUES (?, ?, ?, ?, ?)",
)
.bind(run_id.to_string())
.bind(column(row.stream_seq))
.bind(row.item_kind)
.bind(&row.item_id)
.bind(&row.event_json)
.execute(&mut *tx)
.await
.map_err(ProjectError::Database)?;
}
if let Some(projection) = view.projection.as_ref() {
RunSummaryStore::write_petri_run_row_on_connection(&mut tx, &run_id, projection)
.await
.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"
);
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(),
})
}
/// The stored view of the run, or an empty one.
async fn load_view(&self, run_id: &RunId) -> Result<StoredView, ProjectError> {
let row: Option<(String, String, String, i64)> = sqlx::query_as(
"SELECT projection_json, fold_json, positions_json, stream_seq FROM petri_projection \
WHERE run_id = ?",
)
.bind(run_id.to_string())
.fetch_optional(&self.pool)
.await
.map_err(ProjectError::Database)?;
let Some((projection_json, fold_json, positions_json, stream_seq)) = row else {
return Ok(StoredView {
view: RunView::new(),
positions: Positions::default(),
stream_seq: 0,
});
};
let projection: Option<RunProjection> =
serde_json::from_str(&projection_json).map_err(ProjectError::Encode)?;
let state: FoldState = serde_json::from_str(&fold_json).map_err(ProjectError::Encode)?;
let positions: Positions =
serde_json::from_str(&positions_json).map_err(ProjectError::Encode)?;
Ok(StoredView {
view: RunView { projection, state },
positions,
stream_seq: u64::try_from(stream_seq).unwrap_or(0),
})
}
/// The last seq of every Petri log of the run, by the log column's text.
async fn petri_heads(&self, run_id: &RunId) -> Result<Vec<(String, u64)>, ProjectError> {
// The coordinator log and the execution logs are what the projection
// reads; the resources log is the sandbox ledger and has no events.
let rows: Vec<(String, i64)> = sqlx::query_as(
"SELECT log, MAX(seq) FROM petri_records WHERE run_id = ? AND (log = 'coordinator' \
OR log LIKE 'execution %') GROUP BY log",
)
.bind(run_id.to_string())
.fetch_all(&self.pool)
.await
.map_err(ProjectError::Database)?;
Ok(rows
.into_iter()
.map(|(log, seq)| (log, u64::try_from(seq).unwrap_or(0)))
.collect())
}
/// Whether the run's record is whole: a replay failure says no with its
/// reason; a run that has not recorded its finish is not yet; a finished
/// run is what `inspect_run` says, checked until it says complete.
async fn health(
&self,
key: &RunKey,
state: &FoldState,
replay_failure: Option<String>,
) -> Result<RecordHealth, ProjectError> {
if let Some(failure) = replay_failure {
return Ok(RecordHealth {
complete: false,
incomplete: vec![failure],
});
}
if state.finished.is_none() {
return Ok(RecordHealth {
complete: false,
incomplete: vec!["the run has not recorded its finish".to_string()],
});
}
if state.health.complete {
return Ok(state.health.clone());
}
let logs = match self.store.open(key, Access::Read).await {
Ok(logs) => logs,
Err(StoreError::NotFound { .. }) => return Ok(state.health.clone()),
Err(error) => return Err(ProjectError::Open(error)),
};
match inspect::inspect_run(&*logs).await {
Ok(inspection) => Ok(RecordHealth {
complete: inspection.complete,
incomplete: inspection.incomplete,
}),
Err(error) => Ok(RecordHealth {
complete: false,
incomplete: vec![collect_chain(&error).join(": ")],
}),
}
}
}
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<Self>,
inner: Arc<dyn petri_execution::RunStore>,
) -> Arc<dyn petri_execution::RunStore> {
Arc::new(SignallingStore {
inner,
projector: Arc::clone(self),
})
}
}
/// A run store that signals a projector after each append.
struct SignallingStore {
inner: Arc<dyn petri_execution::RunStore>,
projector: Arc<Projector>,
}
#[async_trait::async_trait]
impl petri_execution::RunStore for SignallingStore {
async fn open(
&self,
key: &RunKey,
access: Access,
) -> Result<Arc<dyn petri_execution::RunLogs>, 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<dyn petri_execution::RunLogs>,
run_id: Option<RunId>,
projector: Arc<Projector>,
}
#[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<Vec<petri_execution::Record>, StoreError> {
self.inner.read(log).await
}
async fn put_blob(&self, bytes: &[u8]) -> Result<petri_store::Digest, StoreError> {
self.inner.put_blob(bytes).await
}
async fn get_blob(&self, digest: petri_store::Digest) -> Result<Option<Vec<u8>>, StoreError> {
self.inner.get_blob(digest).await
}
}
struct StreamRow {
stream_seq: u64,
item_kind: &'static str,
item_id: String,
event_json: String,
}
/// A Petri event id as the stream names it: `<log>/<seq>/<index>`.
#[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)
}
fn lock<T>(mutex: &Mutex<T>) -> MutexGuard<'_, T> {
mutex.lock().unwrap_or_else(PoisonError::into_inner)
}
/// 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.
pub async fn rebuild(
pool: &DbPool,
run_id: RunId,
) -> Result<(Option<RunProjection>, Positions, u64), ProjectError> {
let store = SqliteRunStore::new(pool.clone());
let platform = PlatformRecordStore::new(pool.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 mut items: Vec<(u64, u8, Item<'_>)> = Vec::new();
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 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.
pub async fn stored_positions(
pool: &DbPool,
run_id: RunId,
) -> Result<Option<(Positions, u64)>, ProjectError> {
let row: Option<(String, i64)> =
sqlx::query_as("SELECT positions_json, stream_seq FROM petri_projection WHERE run_id = ?")
.bind(run_id.to_string())
.fetch_optional(pool)
.await
.map_err(ProjectError::Database)?;
row.map(|(positions, stream_seq)| {
Ok((
serde_json::from_str(&positions).map_err(ProjectError::Encode)?,
u64::try_from(stream_seq).unwrap_or(0),
))
})
.transpose()
}
/// The stored view's projection, for a test or a reader outside the store.
pub async fn stored_projection(
pool: &DbPool,
run_id: RunId,
) -> Result<Option<RunProjection>, ProjectError> {
let json: Option<String> =
sqlx::query_scalar("SELECT projection_json FROM petri_projection WHERE run_id = ?")
.bind(run_id.to_string())
.fetch_optional(pool)
.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(
pool: &DbPool,
run_id: RunId,
) -> Result<Vec<(u64, String, String)>, 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(pool)
.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(
pool: &DbPool,
run_id: RunId,
) -> Result<Vec<StoredPlatformRecord>, ProjectError> {
PlatformRecordStore::new(pool.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()
}

View file

@ -0,0 +1,853 @@
//! The projection of a Petri run: the view built live, as the run appends
//! its records, equals the view rebuilt from the records alone; a view that
//! 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.
//!
//! 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
//! found, unless `FABRO_REQUIRE_SANDBOX_PLUGINS` is set.
#![expect(
clippy::disallowed_methods,
reason = "the tests locate the plugin executable through the process environment"
)]
#![expect(clippy::print_stderr, reason = "a skipped test says why on its stderr")]
use std::collections::BTreeSet;
use std::env;
use std::path::{Path, PathBuf};
use std::sync::Arc;
use std::time::Duration;
use fabro_db::DbPool;
use fabro_petri::SqliteRunStore;
use fabro_petri::projector::{self, Projector};
use fabro_store::platform_records::{
PlatformRecord, PlatformRecordStore, RunCreatedRecord, RunLifecycleKind, RunLifecycleRecord,
};
use fabro_store::test_support;
use fabro_types::{
BlobHash, PetriAdmission, PetriGraphRef, RunEngine, RunId, RunStatus, StageHandler, StageId,
StageState, test_support as types_support,
};
use petri_execution::host::{self, HostRun};
use petri_frontend_fabro::Fabro;
use petri_runtime::executor::Retention;
use petri_runtime::frontend::CompileInputs;
use petri_runtime::ir::RunStatus as PetriRunStatus;
use petri_runtime::{RunOptions, Runtime};
use petri_store::{RunKey, RunStore};
use tokio::fs;
const HOST_PLUGIN: &str = "sandbox-driver-host";
const HOST_PLUGIN_OVERRIDE: &str = "PETRI_SANDBOX_HOST_PLUGIN";
const REQUIRE_ENV: &str = "FABRO_REQUIRE_SANDBOX_PLUGINS";
const COMMAND_WORKFLOW: &str = r#"digraph Command {
graph [goal="Run one command"]
start [shape=Mdiamond]
exit [shape=Msquare]
say [shape=parallelogram, script="echo hello from petri"]
start -> say -> exit
}"#;
/// Two branches, each a command, joined by a fan-in.
const PARALLEL_WORKFLOW: &str = r#"digraph Parallel {
graph [goal="Run two branches"]
start [shape=Mdiamond]
exit [shape=Msquare]
fork [shape=component]
a [shape=parallelogram, script="echo a"]
b [shape=parallelogram, script="echo b"]
merge [shape=tripleoctagon]
report [shape=parallelogram, script="echo done"]
start -> fork
fork -> a
fork -> b
a -> merge
b -> merge
merge -> report -> exit
}"#;
const SETTINGS: &str = "_version = 1\n\n[workflow]\ngraph = \"workflow.fabro\"\n";
fn host_plugin() -> Option<PathBuf> {
let found = env::var_os(HOST_PLUGIN_OVERRIDE)
.map(PathBuf::from)
.or_else(|| {
env::split_paths(&env::var_os("PATH")?)
.map(|dir| dir.join(HOST_PLUGIN))
.find(|candidate| candidate.is_file())
});
if found.is_none() {
assert!(
env::var_os(REQUIRE_ENV).is_none(),
"{REQUIRE_ENV} is set, but {HOST_PLUGIN} is not on PATH and {HOST_PLUGIN_OVERRIDE} is unset"
);
eprintln!("skipping: {HOST_PLUGIN} is not on PATH and {HOST_PLUGIN_OVERRIDE} is unset");
}
found
}
/// A fresh in-memory database with every table the projection touches.
fn pool() -> DbPool {
test_support::in_memory_pool_with(&[
fabro_db::BLOBS_MIGRATION_SQL,
fabro_db::RUNS_MIGRATION_SQL,
fabro_db::PETRI_RECORDS_MIGRATION_SQL,
fabro_db::PETRI_PROJECTION_MIGRATION_SQL,
])
}
fn hello_bundle() -> PathBuf {
Path::new(env!("CARGO_MANIFEST_DIR")).join("../../../.fabro/workflows/hello")
}
async fn install_bundle(root: &Path, name: &str, files: &[(&str, &str)]) -> PathBuf {
let bundle = root.join(".fabro").join("workflows").join(name);
fs::create_dir_all(&bundle)
.await
.expect("the bundle directory is creatable");
for (file, text) in files {
fs::write(bundle.join(file), text)
.await
.expect("the bundle file is writable");
}
bundle.join("workflow.fabro")
}
fn run_options(run_dir: &Path, run_id: RunId) -> RunOptions {
let mut options = RunOptions::new(run_dir);
options.grace = Duration::from_secs(2);
options.retention = Retention::Never;
options.echo = false;
options.run_key = Some(RunKey::new(run_id.to_string()));
options
}
/// The run's `run.created` platform record, as the create handler writes it,
/// and the `running` lifecycle record the execute path writes.
async fn create_run(pool: &DbPool, run_id: RunId, goal: &str) {
let store = PlatformRecordStore::new(pool.clone());
let mut spec = types_support::test_run_spec();
spec.run_id = run_id;
spec.engine = RunEngine::Petri(PetriAdmission {
graph: PetriGraphRef {
blob: BlobHash::new(b"graph"),
digest: "digest".to_string(),
},
children: Vec::new(),
});
store
.append(
&run_id,
&PlatformRecord::RunCreated(RunCreatedRecord {
spec,
title: Some(goal.to_string()),
parent_id: None,
retried_from: None,
web_url: None,
}),
None,
)
.await
.expect("the created record stores");
for (transition, status) in [
(RunLifecycleKind::Runnable, RunStatus::Runnable),
(RunLifecycleKind::Starting, RunStatus::Starting),
(RunLifecycleKind::Running, RunStatus::Running),
] {
store
.append(
&run_id,
&PlatformRecord::RunLifecycle(
RunLifecycleRecord::new(transition).with_status(status),
),
None,
)
.await
.expect("the lifecycle record stores");
}
}
/// Run `workflow` to completion on the real registry over `store`.
async fn run_workflow(
store: Arc<dyn RunStore>,
run_dir: &Path,
run_id: RunId,
workflow: &Path,
stubs: bool,
) {
let runtime = Runtime::standard().frontend(Fabro::new());
let runtime = if stubs {
petri_attractor_steps::register_stubs(runtime)
} else {
petri_attractor_steps::register(runtime)
};
let rt = runtime.store(store).options(run_options(run_dir, run_id));
let lowered = rt
.check(workflow, None, None, &CompileInputs::new())
.expect("the workflow file loads");
let graph = lowered
.graph
.unwrap_or_else(|| panic!("the workflow lowers: {:?}", lowered.diagnostics));
let host_run = HostRun::new(graph).with_children(lowered.children);
let report = host::run_configured(&rt, host_run, |_, _| {})
.await
.expect("the run completes");
assert_eq!(
report.status,
PetriRunStatus::Success,
"errors: {:?}",
report.state.errors()
);
}
/// A scenario: its bundle installed, its run created in the database.
struct Scenario {
pool: DbPool,
run_id: RunId,
workflow: PathBuf,
run_dir: PathBuf,
stubs: bool,
_root: tempfile::TempDir,
}
async fn scenario(name: &str, files: &[(&str, &str)], stubs: bool) -> Scenario {
let root = tempfile::tempdir().expect("a temp dir");
let workflow = install_bundle(root.path(), name, files).await;
let pool = pool();
let run_id = RunId::new();
create_run(&pool, run_id, name).await;
Scenario {
pool,
run_id,
workflow,
run_dir: root.path().join("run"),
stubs,
_root: root,
}
}
async fn hello_scenario() -> Scenario {
let bundle = hello_bundle();
let workflow = fs::read_to_string(bundle.join("workflow.fabro"))
.await
.expect("the hello workflow is checked in");
let settings = fs::read_to_string(bundle.join("workflow.toml"))
.await
.expect("the hello settings are checked in");
scenario(
"hello",
&[("workflow.fabro", &workflow), ("workflow.toml", &settings)],
true,
)
.await
}
async fn command_scenario() -> Scenario {
scenario(
"command",
&[
("workflow.fabro", COMMAND_WORKFLOW),
("workflow.toml", SETTINGS),
],
false,
)
.await
}
async fn parallel_scenario() -> Scenario {
scenario(
"parallel",
&[
("workflow.fabro", PARALLEL_WORKFLOW),
("workflow.toml", SETTINGS),
],
false,
)
.await
}
/// Run the scenario live: every append signals the projector, and the view
/// settles before the run is compared with its rebuild.
async fn run_live(scenario: &Scenario) -> Arc<Projector> {
let projector = Projector::new(scenario.pool.clone());
projector.signal(scenario.run_id);
let store = projector.observe_store(Arc::new(SqliteRunStore::new(scenario.pool.clone())));
run_workflow(
store,
&scenario.run_dir,
scenario.run_id,
&scenario.workflow,
scenario.stubs,
)
.await;
projector.settle(scenario.run_id).await;
projector
}
/// Run the scenario with no projector attached: the records land and
/// nothing wakes the view.
async fn run_unobserved(scenario: &Scenario) {
run_workflow(
Arc::new(SqliteRunStore::new(scenario.pool.clone())),
&scenario.run_dir,
scenario.run_id,
&scenario.workflow,
scenario.stubs,
)
.await;
}
/// Every path where two JSON values differ, with both sides.
fn diff_json(path: &str, left: &serde_json::Value, right: &serde_json::Value) -> Vec<String> {
use serde_json::Value;
match (left, right) {
(Value::Object(left), Value::Object(right)) => {
let keys: BTreeSet<&String> = left.keys().chain(right.keys()).collect();
keys.into_iter()
.flat_map(|key| {
diff_json(
&format!("{path}/{key}"),
left.get(key).unwrap_or(&Value::Null),
right.get(key).unwrap_or(&Value::Null),
)
})
.collect()
}
(Value::Array(left), Value::Array(right)) if left.len() == right.len() => left
.iter()
.zip(right)
.enumerate()
.flat_map(|(index, (left, right))| diff_json(&format!("{path}[{index}]"), left, right))
.collect(),
_ if left == right => Vec::new(),
_ => vec![format!("{path}: live {left} != rebuilt {right}")],
}
}
fn json<T: serde::Serialize>(value: &T) -> serde_json::Value {
serde_json::to_value(value).expect("the value serializes")
}
/// 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)
.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)
.await
.expect("the positions read")
.expect("the run has positions");
let (rebuilt, positions, stream_seq) = projector::rebuild(pool, run_id)
.await
.expect("the run rebuilds");
let rebuilt = rebuilt.expect("the rebuild has a projection");
let differences = diff_json("", &json(&stored), &json(&rebuilt));
assert!(
differences.is_empty(),
"live view differs from the rebuild at:\n{}",
differences.join("\n")
);
let mut stored_positions = stored_positions;
let mut positions = positions;
stored_positions.petri.sort();
positions.petri.sort();
assert_eq!(stored_positions, positions);
assert_eq!(stored_stream_seq, stream_seq);
let stream = projector::stored_stream(pool, run_id)
.await
.expect("the stream reads");
let seqs: Vec<u64> = stream.iter().map(|(seq, _, _)| *seq).collect();
assert_eq!(
seqs,
(1..=stream_seq).collect::<Vec<_>>(),
"contiguous stream"
);
}
async fn stage_states(pool: &DbPool, run_id: RunId) -> Vec<(String, StageState)> {
let stored = projector::stored_projection(pool, run_id)
.await
.expect("the stored projection reads")
.expect("the run has a stored projection");
stored
.iter_stages()
.map(|(id, stage)| (id.to_string(), stage.state))
.collect()
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn the_hello_bundle_projects_live_as_it_rebuilds() {
if host_plugin().is_none() {
return;
}
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)
.await
.expect("reads")
.expect("stored");
assert!(
matches!(stored.status, RunStatus::Succeeded { .. }),
"{:?}",
stored.status
);
assert!(stored.conclusion.is_some(), "the run concluded");
let states = stage_states(&scenario.pool, scenario.run_id).await;
assert!(
states
.iter()
.any(|(label, state)| label.starts_with("start@") && *state == StageState::Succeeded),
"{states:?}"
);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn a_command_workflow_projects_live_as_it_rebuilds() {
if host_plugin().is_none() {
return;
}
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)
.await
.expect("reads")
.expect("stored");
let say = stored
.stage(&StageId::new("say", 1))
.expect("the command stage is shown");
assert_eq!(say.state, StageState::Succeeded);
assert_eq!(say.handler, Some(StageHandler::Command));
assert!(
say.output
.as_deref()
.is_some_and(|output| output.contains("hello from petri")),
"{:?}",
say.output
);
assert!(say.timing.is_some());
let states = stage_states(&scenario.pool, scenario.run_id).await;
assert_eq!(states.len(), 3, "start, say, exit: {states:?}");
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn a_parallel_workflow_projects_its_branches_as_child_executions() {
if host_plugin().is_none() {
return;
}
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)
.await
.expect("reads")
.expect("stored");
let fork = StageId::new("fork", 1);
for branch in ["a", "b"] {
let stage = stored
.stage(&StageId::new(branch, 1))
.unwrap_or_else(|| panic!("branch {branch} is a stage"));
assert_eq!(stage.state, StageState::Succeeded);
let branch_id = stage
.parallel_branch_id
.as_ref()
.unwrap_or_else(|| panic!("branch {branch} is grouped under the fork"));
assert_eq!(branch_id.group(), &fork);
}
let fork_stage = stored.stage(&fork).expect("the fork is a stage");
let results = fork_stage
.parallel_results
.as_ref()
.expect("the fork carries its branch results");
assert_eq!(results.len(), 2, "{results:?}");
let labels: Vec<String> = stored.iter_stages().map(|(id, _)| id.to_string()).collect();
assert!(
!labels.iter().any(|label| label.contains("fan_in")),
"synthetic nodes stay off the list: {labels:?}"
);
}
/// The projector is not signalled for any append; one signal at the end
/// folds everything.
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn dropped_wake_ups_are_caught_up_by_the_next_signal() {
if host_plugin().is_none() {
return;
}
let scenario = command_scenario().await;
run_unobserved(&scenario).await;
assert!(
projector::stored_projection(&scenario.pool, scenario.run_id)
.await
.expect("reads")
.is_none(),
"nothing woke the view"
);
let projector = Projector::new(scenario.pool.clone());
projector.signal(scenario.run_id);
projector.settle(scenario.run_id).await;
assert_view_equals_rebuild(&scenario.pool, scenario.run_id).await;
let report = projector
.project_run(scenario.run_id)
.await
.expect("a pass over a caught-up view");
assert!(report.skipped, "nothing is left to fold: {report:?}");
assert!(report.health.complete, "{:?}", report.health.incomplete);
}
/// The same, through the startup pass.
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn the_startup_pass_catches_up_a_view_nobody_signalled() {
if host_plugin().is_none() {
return;
}
let scenario = command_scenario().await;
run_unobserved(&scenario).await;
let projector = Projector::new(scenario.pool.clone());
let report = projector
.startup_pass()
.await
.expect("the startup pass runs");
assert_eq!((report.runs, report.projected), (1, 1));
assert_view_equals_rebuild(&scenario.pool, scenario.run_id).await;
let again = projector
.startup_pass()
.await
.expect("a second startup pass");
assert_eq!((again.runs, again.projected), (1, 0), "nothing left to do");
}
/// Every Petri record of the run, as `(log, seq, recorded_at, record_json)`.
async fn petri_rows(pool: &DbPool, run_id: RunId) -> Vec<(String, i64, i64, String)> {
sqlx::query_as(
"SELECT log, seq, recorded_at, record_json FROM petri_records WHERE run_id = ? ORDER BY \
log, seq",
)
.bind(run_id.to_string())
.fetch_all(pool)
.await
.expect("the records read")
}
async fn insert_petri_row(pool: &DbPool, run_id: RunId, row: &(String, i64, i64, String)) {
sqlx::query(
"INSERT INTO petri_records (run_id, log, seq, recorded_at, record_json) VALUES (?, ?, ?, \
?, ?)",
)
.bind(run_id.to_string())
.bind(&row.0)
.bind(row.1)
.bind(row.2)
.bind(&row.3)
.execute(pool)
.await
.expect("the record inserts");
}
/// A copy of the run in a fresh database: its blobs, its Petri run row and
/// its platform records, but none of its Petri records yet.
async fn copy_run_without_records(source: &DbPool, run_id: RunId) -> DbPool {
let target = pool();
let blobs: Vec<(String, Vec<u8>)> = sqlx::query_as("SELECT hash, data FROM blobs")
.fetch_all(source)
.await
.expect("the blobs read");
for (hash, data) in blobs {
sqlx::query("INSERT INTO blobs (hash, data) VALUES (?, ?)")
.bind(hash)
.bind(data)
.execute(&target)
.await
.expect("the blob inserts");
}
sqlx::query("INSERT INTO petri_runs (run_id, created_at_ms, owner_id, acquired_at_ms) VALUES (?, 0, NULL, NULL)")
.bind(run_id.to_string())
.execute(&target)
.await
.expect("the run row inserts");
let platform: Vec<(i64, i64, String, String)> = sqlx::query_as(
"SELECT seq, recorded_at, kind, record_json FROM platform_records WHERE run_id = ? ORDER \
BY seq",
)
.bind(run_id.to_string())
.fetch_all(source)
.await
.expect("the platform records read");
for (seq, recorded_at, kind, record_json) in platform {
sqlx::query(
"INSERT INTO platform_records (run_id, seq, recorded_at, kind, record_json) VALUES \
(?, ?, ?, ?, ?)",
)
.bind(run_id.to_string())
.bind(seq)
.bind(recorded_at)
.bind(kind)
.bind(record_json)
.execute(&target)
.await
.expect("the platform record inserts");
}
target
}
/// The records commit in two halves and the view runs between them, then
/// the process dies before the view catches the second half: the rebuilt
/// view applies only the suffix, with the positions and the delivery
/// sequence continuing from where the committed view stood.
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn a_crash_between_the_record_commit_and_the_view_applies_only_the_suffix() {
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 replayed = copy_run_without_records(&scenario.pool, scenario.run_id).await;
// The first half of every log: a prefix per log, the coordinator log
// short of its finish.
let mut first: Vec<&(String, i64, i64, String)> = Vec::new();
let mut second: Vec<&(String, i64, i64, String)> = Vec::new();
for row in &rows {
let head = rows
.iter()
.filter(|other| other.0 == row.0)
.map(|other| other.1)
.max()
.expect("the log has a head");
if row.1 <= head / 2 {
first.push(row);
} else {
second.push(row);
}
}
for row in &first {
insert_petri_row(&replayed, scenario.run_id, row).await;
}
let before = Projector::new(replayed.clone());
let pass = before
.project_run(scenario.run_id)
.await
.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");
assert_eq!(pass.stream_seq, stream_before);
let stream_rows_before = projector::stored_stream(&replayed, scenario.run_id)
.await
.expect("reads")
.len();
// The rest of the records commit; the view transaction never runs.
for row in &second {
insert_petri_row(&replayed, scenario.run_id, row).await;
}
before.fail_before_view();
let crashed = before.project_run(scenario.run_id).await;
assert!(
matches!(crashed, Err(projector::ProjectError::Injected)),
"{crashed:?}"
);
assert_eq!(
projector::stored_positions(&replayed, scenario.run_id)
.await
.expect("reads")
.expect("positions"),
(positions_before.clone(), stream_before),
"the crash left the committed view alone"
);
// A new projector, as a restarted server builds one.
let after = Projector::new(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)
.await
.expect("reads");
// Only the suffix was applied: the stream grew by the suffix's events,
// numbered on from the committed sequence, and every earlier row stayed.
assert_eq!(
stream_rows_after.len(),
stream_rows_before + usize::try_from(stream_after - stream_before).expect("a small count")
);
assert!(stream_after > stream_before);
assert_eq!(
stream_rows_after[stream_rows_before].0,
stream_before + 1,
"the suffix starts right after the committed sequence"
);
for held in &positions_before.petri {
let now = positions_after
.petri
.iter()
.find(|after| after.source == held.source)
.expect("a held log is still held");
assert!(now >= held, "{now:?} >= {held:?}");
}
assert_eq!(positions_after.platform_seq, positions_before.platform_seq);
assert_view_equals_rebuild(&replayed, scenario.run_id).await;
// And the copy agrees with the run projected in one go over the source.
let source = Projector::new(scenario.pool.clone());
source.startup_pass().await.expect("the source projects");
let whole = projector::stored_projection(&scenario.pool, scenario.run_id)
.await
.expect("reads")
.expect("stored");
let pieced = projector::stored_projection(&replayed, scenario.run_id)
.await
.expect("reads")
.expect("stored");
assert_eq!(json(&whole), json(&pieced));
}
/// Two projectors over one store, one after the other, over a run with
/// child executions: the second continues where the first stopped and both
/// agree with a projector that saw the run whole.
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn a_restarted_projector_agrees_over_nested_child_executions() {
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;
assert!(
rows.iter()
.filter(|row| row.0.starts_with("execution "))
.map(|row| &row.0)
.collect::<BTreeSet<_>>()
.len()
>= 3,
"the parallel run has child executions: {:?}",
rows.iter().map(|row| &row.0).collect::<BTreeSet<_>>()
);
let staged = copy_run_without_records(&scenario.pool, scenario.run_id).await;
let first = Projector::new(staged.clone());
// The parent execution and the coordinator log up to the first child's
// declaration go in first; a restart then sees the children.
let (early, late): (Vec<_>, Vec<_>) = rows
.iter()
.partition(|row| row.0 == "execution 0" || (row.0 == "coordinator" && row.1 < 6));
for row in &early {
insert_petri_row(&staged, scenario.run_id, row).await;
}
first
.startup_pass()
.await
.expect("the first projector passes");
for row in &late {
insert_petri_row(&staged, scenario.run_id, row).await;
}
drop(first);
let second = Projector::new(staged.clone());
second
.startup_pass()
.await
.expect("the second projector passes");
assert_view_equals_rebuild(&staged, scenario.run_id).await;
let whole = Projector::new(scenario.pool.clone());
whole.startup_pass().await.expect("the source projects");
let one_go = projector::stored_projection(&scenario.pool, scenario.run_id)
.await
.expect("reads")
.expect("stored");
let restarted = projector::stored_projection(&staged, scenario.run_id)
.await
.expect("reads")
.expect("stored");
assert_eq!(json(&one_go), json(&restarted));
let states = stage_states(&staged, scenario.run_id).await;
assert!(
states.iter().any(|(label, _)| label == "a@1")
&& states.iter().any(|(label, _)| label == "b@1"),
"{states:?}"
);
}
/// A record at seq n+2 of an execution log, past a gap: Petri cannot read
/// the log, the view does not advance past what it held, and the run is
/// reported incomplete with the reason.
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn a_torn_tail_holds_the_view_and_reports_the_run_incomplete() {
if host_plugin().is_none() {
return;
}
let scenario = command_scenario().await;
run_unobserved(&scenario).await;
let projector = Projector::new(scenario.pool.clone());
let clean = projector
.project_run(scenario.run_id)
.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)
.await
.expect("reads")
.expect("positions");
let before = projector::stored_projection(&scenario.pool, scenario.run_id)
.await
.expect("reads")
.expect("stored");
// A record two past the head of the execution log.
let rows = petri_rows(&scenario.pool, scenario.run_id).await;
let last = rows
.iter()
.filter(|row| row.0 == "execution 0")
.max_by_key(|row| row.1)
.expect("the execution log has records");
let mut torn: serde_json::Value = serde_json::from_str(&last.3).expect("the record is JSON");
torn["seq"] = serde_json::json!(last.1 + 2);
insert_petri_row(
&scenario.pool,
scenario.run_id,
&(last.0.clone(), last.1 + 2, last.2, torn.to_string()),
)
.await;
let held = projector
.project_run(scenario.run_id)
.await
.expect("the pass over the torn log still commits its health");
assert!(!held.health.complete, "the torn log is incomplete");
assert!(
!held.health.incomplete.is_empty(),
"the reason is reported: {:?}",
held.health
);
let (positions_after, stream_after) =
projector::stored_positions(&scenario.pool, scenario.run_id)
.await
.expect("reads")
.expect("positions");
assert_eq!(
positions_after, positions,
"the view did not advance past the tear"
);
assert_eq!(stream_after, stream_seq);
let after = projector::stored_projection(&scenario.pool, scenario.run_id)
.await
.expect("reads")
.expect("stored");
assert_eq!(
json(&before),
json(&after),
"the projection stands where it was"
);
}