Add platform records and the Petri projection tables

A Petri run's own Fabro facts (its lifecycle before and after the
engine, a checkpoint commit, a pull request, a notification, a pairing)
are platform records in a table beside Petri's records, one typed enum
of kinds tagged on the wire, each keyed to a Petri stage where it
belongs to one and carrying the operation identity of the effect it
records. The run summary store derives the lifecycle kinds from the
legacy run events a Petri run still appends, in the event's
transaction, and calls a hook after the commit so the run's projector
can wake up.

Two more tables serve the projection that follows: the per-run
projection document with its committed positions, and the ordered
stream of everything the view consumed. The run summary store reads the
Petri projection back for the API and writes the narrowed runs row from
it without touching the legacy concurrency guard.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
This commit is contained in:
Bryan Helmkamp 2026-09-17 21:46:21 -04:00
parent 16354186fa
commit abbc7ca11d
No known key found for this signature in database
9 changed files with 1278 additions and 14 deletions

View file

@ -7,6 +7,7 @@ mod keyed_mutex;
mod keys;
mod legacy_blob_import;
mod legacy_run_history_import;
pub mod platform_records;
#[cfg(test)]
mod record;
mod run_session_record_store;
@ -43,9 +44,13 @@ pub use legacy_run_history_import::{
LegacyRunHistorySourceIdentity, LegacyRunHistorySourceIdentityError,
LegacyRunHistoryVerificationError, LegacyRunHistoryVerificationReport,
};
pub use platform_records::{
PlatformRecord, PlatformRecordHook, PlatformRecordKind, PlatformRecordStore, StagePosition,
StoredPlatformRecord,
};
pub use run_session_record_store::{RunSessionRecordStore, StoredSessionRecord};
pub use run_sessions::{ProjectedRunSession, project_run_session, project_run_sessions};
pub use run_state::RunProjectionReducer;
pub use run_state::{RunProjectionReducer, build_summary, projected_usage};
pub use run_summary_store::{
RunSummaryIdentity, RunSummaryListQuery, RunSummaryPage, RunSummarySort,
RunSummarySortDirection, RunSummaryStore, RunSummaryVisibility,

File diff suppressed because it is too large Load diff

View file

@ -1145,7 +1145,10 @@ fn stage_at_completed_visit<'a>(
Some(state.stage_entry(node_id, visit, first_event_seq(seq)))
}
pub(crate) fn build_summary(state: &RunProjection, run_id: &RunId) -> Run {
/// The run summary (`Run`) a projection stands for: what the run list, the
/// board and the scheduler read.
#[must_use]
pub fn build_summary(state: &RunProjection, run_id: &RunId) -> Run {
let goal = state.spec.graph.goal().to_string();
let diff_summary = state
.conclusion
@ -1241,7 +1244,8 @@ pub(crate) fn build_summary(state: &RunProjection, run_id: &RunId) -> Run {
/// The run's usage: the conclusion's total once the run ended, else the sum
/// of every non-boundary stage's usage so far.
pub(crate) fn projected_usage(state: &RunProjection) -> Usage {
#[must_use]
pub fn projected_usage(state: &RunProjection) -> Usage {
if let Some(usage) = state
.conclusion
.as_ref()

View file

@ -1,5 +1,5 @@
use std::fmt::Write as _;
use std::sync::LazyLock;
use std::sync::{Arc, LazyLock, RwLock};
use chrono::{DateTime, Utc};
use fabro_types::{
@ -12,8 +12,9 @@ use sqlx::sqlite::{SqliteArguments, SqliteConnection, SqliteRow};
use sqlx::{Connection as _, QueryBuilder, Row as _, Sqlite, SqlitePool, Transaction};
use strum::VariantArray as _;
use crate::platform_records::{self, PlatformRecordHook, PlatformRecordStore};
use crate::run_state::{ProjectedRun, build_summary, projected_usage};
use crate::{Error, EventPayload, Result, keys};
use crate::{Error, EventPayload, Result, RunProjection, keys};
const INSERT_RUN_SQL: &str = r"
INSERT INTO runs (
@ -65,6 +66,38 @@ ON CONFLICT(id) DO UPDATE SET
WHERE excluded.source_last_seq > runs.source_last_seq
";
/// The `runs` row of a Petri run, written by its projector: every column the
/// list views and the scheduler read, and never `source_last_seq`, which the
/// legacy event path owns while it still writes the row.
const UPSERT_PETRI_RUN_SQL: &str = r"
INSERT INTO runs (
id, source_last_seq, created_at_ms, started_at_ms, last_event_at_ms, completed_at_ms,
status, archived_at_ms, parent_id, title, workflow_slug, workflow_name,
repository_name, automation_id, diff_files_changed, diff_additions, diff_deletions,
input_tokens, output_tokens, reasoning_tokens, cache_read_tokens, cache_write_tokens,
total_usd_micros, summary_json
) VALUES (
?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?
)
ON CONFLICT(id) DO UPDATE SET
created_at_ms = excluded.created_at_ms,
started_at_ms = excluded.started_at_ms,
last_event_at_ms = excluded.last_event_at_ms,
completed_at_ms = excluded.completed_at_ms,
status = excluded.status,
archived_at_ms = excluded.archived_at_ms,
parent_id = excluded.parent_id,
title = excluded.title,
workflow_slug = excluded.workflow_slug,
workflow_name = excluded.workflow_name,
repository_name = excluded.repository_name,
automation_id = excluded.automation_id,
diff_additions = excluded.diff_additions,
diff_deletions = excluded.diff_deletions,
total_usd_micros = excluded.total_usd_micros,
summary_json = excluded.summary_json
";
const UPDATE_RUN_SQL: &str = r"
UPDATE runs SET
source_last_seq = ?,
@ -184,7 +217,10 @@ pub struct RunSummaryPage {
#[derive(Clone)]
pub struct RunSummaryStore {
pool: SqlitePool,
pool: SqlitePool,
/// Called after a platform record for a Petri run is committed beside
/// its legacy event: the projector's wake-up.
platform_hook: Arc<RwLock<Option<PlatformRecordHook>>>,
}
impl std::fmt::Debug for RunSummaryStore {
@ -196,7 +232,73 @@ impl std::fmt::Debug for RunSummaryStore {
impl RunSummaryStore {
#[must_use]
pub fn new(pool: SqlitePool) -> Self {
Self { pool }
Self {
pool,
platform_hook: Arc::new(RwLock::new(None)),
}
}
/// The platform records over the same pool.
#[must_use]
pub fn platform_records(&self) -> PlatformRecordStore {
PlatformRecordStore::new(self.pool.clone())
}
/// Install the wake-up called after a platform record of a Petri run is
/// committed beside its legacy event.
pub fn set_platform_record_hook(&self, hook: PlatformRecordHook) {
*self
.platform_hook
.write()
.unwrap_or_else(std::sync::PoisonError::into_inner) = Some(hook);
}
pub(crate) fn notify_platform_record(&self, run_id: RunId) {
let hook = self
.platform_hook
.read()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.clone();
if let Some(hook) = hook {
hook(run_id);
}
}
/// The stored projection of a Petri run, as the run's projector last
/// committed it, or `None` when no view pass has run for it yet.
pub async fn load_petri_projection(
&self,
run_id: &RunId,
) -> Result<Option<Arc<RunProjection>>> {
let json: Option<String> =
sqlx::query_scalar("SELECT projection_json FROM petri_projection WHERE run_id = ?")
.bind(run_id.to_string())
.fetch_optional(&self.pool)
.await?;
json.map(|json| Ok(Arc::new(serde_json::from_str(&json)?)))
.transpose()
}
/// Write the `runs` row of a Petri run from its projection, on a
/// connection the caller holds a transaction on: the columns the list
/// views and the scheduler read, and the summary JSON. The legacy
/// concurrency guard `source_last_seq` is left as the legacy path set it
/// (or `1` when this write creates the row), so both writers keep
/// working until the legacy events go.
pub async fn write_petri_run_row_on_connection(
connection: &mut SqliteConnection,
run_id: &RunId,
projection: &RunProjection,
) -> Result<()> {
let entry = ProjectedRun::new(*run_id, Arc::new(projection.clone()), 1);
let record = PreparedRunSummary::from_entry(&entry);
bind_run_columns(
sqlx::query(UPSERT_PETRI_RUN_SQL).bind(run_id.to_string()),
&record,
)?
.execute(connection)
.await?;
Ok(())
}
#[cfg(test)]
@ -657,6 +759,7 @@ impl RunSummaryStore {
insert_run_on_connection(connection, &record).await?;
insert_event_on_connection(connection, &record, payload, &envelope).await?;
insert_platform_record_on_connection(connection, entry, &envelope).await?;
Ok(envelope)
}
@ -674,6 +777,7 @@ impl RunSummaryStore {
update_run_on_connection(connection, &record, expected_last_seq).await?;
insert_event_on_connection(connection, &record, payload, &envelope).await?;
insert_platform_record_on_connection(connection, entry, &envelope).await?;
Ok(envelope)
}
@ -1113,6 +1217,43 @@ async fn insert_event_json_on_connection(
Ok(())
}
/// For a Petri run, the platform record the legacy event stands for, stored
/// in the event's transaction so the projection over Petri's records reads
/// the lifecycle from platform records alone. Whether one was written is
/// what [`platform_record_written`] answers after the commit.
async fn insert_platform_record_on_connection(
connection: &mut SqliteConnection,
entry: &ProjectedRun,
envelope: &EventEnvelope,
) -> Result<()> {
let Some(record) = platform_record_written(entry, envelope) else {
return Ok(());
};
let recorded_at = u64::try_from(envelope.event.ts.timestamp_millis()).unwrap_or(0);
PlatformRecordStore::append_on_connection(
connection,
&entry.run_id,
recorded_at,
&record,
None,
)
.await?;
Ok(())
}
/// The platform record a committed legacy event of a Petri run produced,
/// if any: the same derivation the insert makes, for the caller that
/// notifies after the commit.
pub(crate) fn platform_record_written(
entry: &ProjectedRun,
envelope: &EventEnvelope,
) -> Option<platform_records::PlatformRecord> {
if !entry.projection.spec.engine.is_petri() {
return None;
}
platform_records::platform_record_for(&envelope.event)
}
fn sql_limit(limit: usize) -> i64 {
i64::try_from(limit.saturating_add(1)).unwrap_or(i64::MAX)
}

View file

@ -232,15 +232,32 @@ impl Database {
Ok(())
}
/// The run's projection: for a legacy run the reducer's fold of its
/// events; for a Petri run the projection its projector last committed
/// over Petri's records and the platform records, falling back to the
/// legacy fold (the lifecycle alone) until the first view pass commits.
pub async fn load_run_projection(&self, run_id: &RunId) -> Result<Option<Arc<RunProjection>>> {
if let Some(active) = self.get_active_run(run_id).await {
return active.projection_snapshot().await.map(Some);
}
match self.run_summary_store.load_projection(run_id).await {
Ok(projected) => Ok(Some(projected.projection)),
Err(Error::RunNotFound(_)) => Ok(None),
Err(error) => Err(error),
let legacy = if let Some(active) = self.get_active_run(run_id).await {
active.projection_snapshot().await?
} else {
match self.run_summary_store.load_projection(run_id).await {
Ok(projected) => projected.projection,
Err(Error::RunNotFound(_)) => return Ok(None),
Err(error) => return Err(error),
}
};
if legacy.spec.engine.is_petri() {
if let Some(petri) = self.run_summary_store.load_petri_projection(run_id).await? {
return Ok(Some(petri));
}
}
Ok(Some(legacy))
}
/// Install the wake-up called after a platform record of a Petri run is
/// committed beside its legacy event.
pub fn set_platform_record_hook(&self, hook: crate::PlatformRecordHook) {
self.run_summary_store.set_platform_record_hook(hook);
}
/// Resolves the run that owns `session_id` from the canonical typed

View file

@ -220,8 +220,14 @@ impl RunDatabase {
let (envelope, projected) = self.commit_event_locked(payload, event).await?;
// Keep post-commit propagation await-free: cancellation after SQLite
// commits must not leave in-memory state stale or omit the broadcast.
let platform_record = run_summary_store::platform_record_written(&projected, &envelope);
self.install_in_memory_state(projected);
self.publish(&envelope);
if platform_record.is_some() {
self.inner
.run_summary_store
.notify_platform_record(self.inner.run_id);
}
Ok(envelope)
}

View file

@ -34,6 +34,7 @@ pub fn test_run_summary_store() -> Arc<RunSummaryStore> {
fabro_db::RUN_EVENTS_MIGRATION_SQL,
fabro_db::RUN_HISTORY_ACTIVATION_MIGRATION_SQL,
fabro_db::RUN_EVENT_SESSION_OWNER_MIGRATION_SQL,
fabro_db::PETRI_PROJECTION_MIGRATION_SQL,
])))
}
@ -115,6 +116,7 @@ pub fn test_run_summary_store_at(store_dir: &Path) -> Arc<RunSummaryStore> {
fabro_db::RUN_EVENTS_MIGRATION_SQL,
fabro_db::RUN_HISTORY_ACTIVATION_MIGRATION_SQL,
fabro_db::RUN_EVENT_SESSION_OWNER_MIGRATION_SQL,
fabro_db::PETRI_PROJECTION_MIGRATION_SQL,
],
)))
}

View file

@ -0,0 +1,60 @@
-- Fabro's own facts about a Petri run, beside Petri's records.
--
-- `platform_records` holds every fact Fabro records about a run that Petri
-- does not: the lifecycle before and after the engine, a checkpoint commit,
-- a pull request, a notification, a pairing. `seq` is per run and assigned
-- by the store; `kind` is the record's kind and `record_json` the typed
-- record with its kind tag; `execution` and `firing` name the Petri stage a
-- record belongs to, when it belongs to one.
CREATE TABLE platform_records (
run_id TEXT NOT NULL,
seq INTEGER NOT NULL,
recorded_at INTEGER NOT NULL,
kind TEXT NOT NULL,
record_json TEXT NOT NULL,
execution INTEGER NULL,
firing INTEGER NULL,
PRIMARY KEY (run_id, seq),
CHECK (seq >= 1),
CHECK (json_valid(record_json))
);
CREATE INDEX platform_records_by_kind
ON platform_records(run_id, kind, seq);
-- The projection of a Petri run: the view document Fabro's read side serves,
-- derived from the run's Petri records and platform records, rewritten in
-- one transaction per view pass together with the positions it covers.
-- `projection_json` is the `RunProjection`; `fold_json` is the projector's
-- own bookkeeping; `positions_json` is the last event consumed per Petri
-- log and the last platform record consumed; `stream_seq` is the last
-- delivery sequence assigned to `petri_stream`.
CREATE TABLE petri_projection (
run_id TEXT PRIMARY KEY NOT NULL,
projection_json TEXT NOT NULL,
fold_json TEXT NOT NULL,
positions_json TEXT NOT NULL,
stream_seq INTEGER NOT NULL,
updated_at_ms INTEGER NOT NULL,
CHECK (stream_seq >= 0),
CHECK (json_valid(projection_json)),
CHECK (json_valid(fold_json)),
CHECK (json_valid(positions_json))
);
-- One ordered stream per run of everything the projection consumed: each
-- Petri event and each platform record, in the order the view committed
-- them. `stream_seq` is the cursor a client resumes from; `item_kind` and
-- `item_id` are the item's own identity (a Petri event id as
-- `<log>/<seq>/<index>`, or a platform record's `seq`), for deduplication.
CREATE TABLE petri_stream (
run_id TEXT NOT NULL,
stream_seq INTEGER NOT NULL,
item_kind TEXT NOT NULL,
item_id TEXT NOT NULL,
event_json TEXT NOT NULL,
PRIMARY KEY (run_id, stream_seq),
CHECK (stream_seq >= 1),
CHECK (item_kind IN ('petri', 'platform')),
CHECK (json_valid(event_json))
);

View file

@ -47,6 +47,12 @@ pub const RUN_SESSION_RECORDS_MIGRATION_SQL: &str =
pub const PETRI_RECORDS_MIGRATION_SQL: &str =
include_str!("../migrations/2026091701_petri_records.sql");
/// The Petri projection migration (`platform_records`, `petri_projection`,
/// `petri_stream`), exposed so fixtures in other crates can install the
/// production schema without a filesystem path into this crate.
pub const PETRI_PROJECTION_MIGRATION_SQL: &str =
include_str!("../migrations/2026091801_petri_projection.sql");
/// The temporary run-history activation migration, exposed so fixtures in
/// other crates can install the production compatibility schema.
pub const RUN_HISTORY_ACTIVATION_MIGRATION_SQL: &str =