Centralize run projection replay on ProjectedRun

Move ProjectedRun next to EventProjectionCache in run_state, since the
summary store both produces and consumes it, and give it a replay
constructor that owns the events-to-head derivation. load_projection now
returns RunNotFound directly instead of erasing it to None and having
callers rebuild it; load_run_projection is the single Option translation
point. install_in_memory_state reuses the existing From impl, and the
commit path passes its Arc through instead of unwrapping and
reallocating it.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
This commit is contained in:
Scott Werner 2026-09-02 16:33:14 -04:00
parent 2360e8046b
commit 0f1e5e1c9c
5 changed files with 104 additions and 146 deletions

View file

@ -8,7 +8,7 @@ use std::collections::HashSet;
use std::error::Error as StdError;
use std::fmt;
use fabro_types::{EventEnvelope, RunEvent, RunId, RunProjection};
use fabro_types::{EventEnvelope, RunEvent, RunId};
use sha2::{Digest as _, Sha256};
use sqlx::SqlitePool;
#[cfg(test)]
@ -16,8 +16,8 @@ use tokio::sync::Barrier;
use tracing::debug;
use crate::keys::SlateKey;
use crate::slate::ProjectedRun;
use crate::{Database, EventPayload, RunProjectionReducer, RunSummaryStore, keys};
use crate::run_state::ProjectedRun;
use crate::{Database, EventPayload, RunSummaryStore, keys};
/// Count-only observations about the legacy catalog and session indexes.
#[derive(Clone, Copy, Debug, Default, Eq, PartialEq)]
@ -527,14 +527,8 @@ impl LegacyRunHistorySource {
.iter()
.map(|event| event.envelope.clone())
.collect::<Vec<_>>();
let projection = RunProjection::apply_events(&envelopes)
let current = ProjectedRun::replay(run_id, &envelopes)
.map_err(LegacyRunHistorySourceFailure::Replay)?;
let last_seq = events
.last()
.expect("a validated history contains at least one event")
.envelope
.seq;
let current = ProjectedRun::new(run_id, projection, last_seq);
Ok(Some(ValidatedLegacyRunHistory {
run_id,
events,
@ -1160,28 +1154,20 @@ fn replay_destination(
run_id: &RunId,
events: &[(EventEnvelope, String)],
) -> crate::Result<ProjectedRun> {
let Some((first, _event_json)) = events.first() else {
return Err(crate::Error::InvalidEvent(
"run projection requires an event".to_owned(),
));
};
if first.seq != 1 {
return Err(crate::Error::RunEventMismatch {
run_id: run_id.to_string(),
seq: first.seq,
field: "seq",
});
if let Some((first, _event_json)) = events.first() {
if first.seq != 1 {
return Err(crate::Error::RunEventMismatch {
run_id: run_id.to_string(),
seq: first.seq,
field: "seq",
});
}
}
let envelopes = events
.iter()
.map(|(envelope, _event_json)| envelope.clone())
.collect::<Vec<_>>();
let projection = RunProjection::apply_events(&envelopes)?;
let last_seq = envelopes
.last()
.expect("a destination history validated as nonempty")
.seq;
Ok(ProjectedRun::new(*run_id, projection, last_seq))
ProjectedRun::replay(*run_id, &envelopes)
}
fn usize_to_import_count(value: usize) -> Result<u64, LegacyRunHistoryImportFailure> {
@ -1256,9 +1242,7 @@ mod tests {
use std::time::Duration;
use chrono::{TimeZone as _, Utc};
use fabro_types::{
Graph, RunEvent, RunId, RunProjection, SessionId, WorkflowSettings, test_support,
};
use fabro_types::{Graph, RunEvent, RunId, SessionId, WorkflowSettings, test_support};
use fabro_util::error;
use object_store::memory::InMemory;
use tokio::sync::Barrier;
@ -1271,9 +1255,9 @@ mod tests {
parse_source_event,
};
use crate::keys::SlateKey;
use crate::slate::ProjectedRun;
use crate::run_state::ProjectedRun;
use crate::{
Database, EventEnvelope, EventPayload, RunProjectionReducer, RunSummaryStore, keys,
Database, EventEnvelope, EventPayload, RunSummaryStore, keys,
test_support as store_test_support,
};
@ -1476,8 +1460,7 @@ mod tests {
.iter()
.map(|(_payload, envelope)| envelope.clone())
.collect::<Vec<_>>();
let projection = RunProjection::apply_events(&envelopes)?;
let current = ProjectedRun::new(*run_id, projection, events.last().unwrap().0);
let current = ProjectedRun::replay(*run_id, &envelopes)?;
let mut transaction = pool.begin().await?;
RunSummaryStore::insert_imported_run_on_connection(&mut transaction, &current).await?;
for ((_, event_json), (payload, envelope)) in events.iter().zip(&decoded) {
@ -1510,8 +1493,7 @@ mod tests {
.iter()
.map(|(_payload, envelope)| envelope.clone())
.collect::<Vec<_>>();
let projection = RunProjection::apply_events(&envelopes)?;
let current = ProjectedRun::new(*run_id, projection, next.0);
let current = ProjectedRun::replay(*run_id, &envelopes)?;
let mut transaction = pool.begin().await?;
RunSummaryStore::append_event_on_connection(
&mut transaction,

View file

@ -33,6 +33,45 @@ pub(crate) struct EventProjectionCache {
pub state: Option<Arc<RunProjection>>,
}
/// A run's projection at its committed head. `RunSummaryStore` replays it
/// from SQLite history and derives summary rows from it; `RunDatabase` builds
/// it from a newly committed event and seeds its in-memory cache with it.
#[derive(Debug, Clone)]
pub(crate) struct ProjectedRun {
pub(crate) run_id: RunId,
pub(crate) projection: Arc<RunProjection>,
pub(crate) last_seq: u32,
}
impl ProjectedRun {
pub(crate) fn new(run_id: RunId, projection: Arc<RunProjection>, last_seq: u32) -> Self {
Self {
run_id,
projection,
last_seq,
}
}
/// Replays a run's full history; the last event's `seq` becomes the head.
pub(crate) fn replay(run_id: RunId, events: &[EventEnvelope]) -> Result<Self> {
let projection = RunProjection::apply_events(events)?;
let last_seq = events
.last()
.expect("a successfully replayed history contains at least one event")
.seq;
Ok(Self::new(run_id, Arc::new(projection), last_seq))
}
}
impl From<ProjectedRun> for EventProjectionCache {
fn from(projected: ProjectedRun) -> Self {
Self {
last_seq: projected.last_seq,
state: Some(projected.projection),
}
}
}
pub trait RunProjectionReducer {
fn apply_events(events: &[EventEnvelope]) -> Result<Self>
where

View file

@ -3,8 +3,8 @@ use std::sync::LazyLock;
use chrono::{DateTime, Utc};
use fabro_types::{
BilledTokenCounts, EventEnvelope, Run, RunEvent, RunId, RunProjection, RunSize, RunStatusKind,
RunTiming, SessionId, StageId, timing,
BilledTokenCounts, EventEnvelope, Run, RunEvent, RunId, RunSize, RunStatusKind, RunTiming,
SessionId, StageId, timing,
};
use sqlx::pool::PoolConnection;
use sqlx::query::Query;
@ -12,8 +12,7 @@ use sqlx::sqlite::{SqliteArguments, SqliteConnection, SqliteRow};
use sqlx::{Connection as _, QueryBuilder, Row as _, Sqlite, SqlitePool, Transaction};
use strum::VariantArray as _;
use crate::run_state::{RunProjectionReducer, build_summary, projected_billing};
use crate::slate::ProjectedRun;
use crate::run_state::{ProjectedRun, build_summary, projected_billing};
use crate::{Error, EventPayload, Result, keys};
const INSERT_RUN_SQL: &str = r"
@ -351,20 +350,11 @@ ON CONFLICT(singleton) DO NOTHING
select_run_head(&mut connection, run_id).await
}
/// Replays one run's canonical history from one validated SQLite snapshot,
/// returning `None` when the run does not exist.
pub(crate) async fn load_projection(&self, run_id: &RunId) -> Result<Option<ProjectedRun>> {
let events = match self.list_events_for_run(run_id).await {
Ok(events) => events,
Err(Error::RunNotFound(_)) => return Ok(None),
Err(error) => return Err(error),
};
let last_seq = events
.last()
.map(|event| event.seq)
.ok_or_else(|| Error::InvalidEvent(format!("run {run_id} has no run.created event")))?;
let projection = RunProjection::apply_events(&events)?;
Ok(Some(ProjectedRun::new(*run_id, projection, last_seq)))
/// Replays one run's canonical history from one validated SQLite snapshot.
/// Fails with `RunNotFound` when the run does not exist.
pub(crate) async fn load_projection(&self, run_id: &RunId) -> Result<ProjectedRun> {
let events = self.list_events_for_run(run_id).await?;
ProjectedRun::replay(*run_id, &events)
}
pub(crate) async fn list_events_for_run(&self, run_id: &RunId) -> Result<Vec<EventEnvelope>> {
@ -1543,6 +1533,7 @@ fn overlay_live_wall_time(run: &mut Run, now: DateTime<Utc>) {
#[cfg(test)]
mod tests {
use std::collections::HashMap;
use std::sync::Arc;
use std::time::Duration;
use chrono::{DateTime, Utc};
@ -1560,7 +1551,7 @@ mod tests {
INSERT_EVENT_SQL, RunSummaryListQuery, RunSummarySort, RunSummarySortDirection,
RunSummaryStore, RunSummaryVisibility, decode_event_row,
};
use crate::slate::ProjectedRun;
use crate::run_state::ProjectedRun;
use crate::{Error, EventPayload, RunProjectionReducer, test_support as store_test_support};
fn dt(value: &str) -> DateTime<Utc> {
@ -1597,7 +1588,7 @@ mod tests {
}
fn entry(projection: RunProjection, last_seq: u32) -> ProjectedRun {
ProjectedRun::new(projection.spec.run_id, projection, last_seq)
ProjectedRun::new(projection.spec.run_id, Arc::new(projection), last_seq)
}
async fn store() -> (tempfile::TempDir, RunSummaryStore) {
@ -1976,7 +1967,7 @@ mod tests {
let expected = RunProjection::apply_events(&[first_envelope, second_envelope]).unwrap();
let loaded = store.load_projection(&id).await.unwrap().unwrap();
let loaded = store.load_projection(&id).await.unwrap();
assert_eq!(loaded.run_id, id);
assert_eq!(loaded.last_seq, 2);
assert_eq!(
@ -1985,43 +1976,37 @@ mod tests {
);
let missing = run_id(created_at.timestamp_millis().cast_unsigned() + 1, 32);
assert!(store.load_projection(&missing).await.unwrap().is_none());
assert!(matches!(
store.load_projection(&missing).await,
Err(Error::RunNotFound(text)) if text == missing.to_string()
));
}
#[tokio::test]
async fn load_projection_reports_removed_events_without_poisoning_following_reads() {
async fn load_projection_reports_removed_events() {
let (_directory, store) = store().await;
let created_at = dt("2026-08-27T12:00:00Z");
let broken_id = run_id(created_at.timestamp_millis().cast_unsigned(), 33);
let healthy_id = run_id(created_at.timestamp_millis().cast_unsigned() + 1, 34);
for id in [broken_id, healthy_id] {
let current = entry(projection(id, "created", created_at), 1);
let mut transaction = store.pool.begin().await.unwrap();
RunSummaryStore::insert_first_event_on_connection(
&mut transaction,
&current,
&created_payload(&id),
)
.await
.unwrap();
transaction.commit().await.unwrap();
}
store.test_delete_run_events(&broken_id).await.unwrap();
let id = run_id(created_at.timestamp_millis().cast_unsigned(), 33);
let current = entry(projection(id, "created", created_at), 1);
let mut transaction = store.pool.begin().await.unwrap();
RunSummaryStore::insert_first_event_on_connection(
&mut transaction,
&current,
&created_payload(&id),
)
.await
.unwrap();
transaction.commit().await.unwrap();
store.test_delete_run_events(&id).await.unwrap();
assert!(matches!(
store.load_projection(&broken_id).await,
store.load_projection(&id).await,
Err(Error::RunHeadMismatch {
expected_last_seq: 1,
actual_last_seq: None,
..
})
));
let loaded = store.load_projection(&healthy_id).await.unwrap().unwrap();
assert_eq!(loaded.run_id, healthy_id);
assert_eq!(loaded.last_seq, 1);
assert_eq!(loaded.projection.title, "created");
}
#[tokio::test]

View file

@ -8,7 +8,6 @@ use std::time::Duration;
use chrono::{DateTime, Utc};
use fabro_types::{RunId, SessionId};
use object_store::ObjectStore;
pub(crate) use run_store::ProjectedRun;
pub use run_store::RunDatabase;
use run_store::RunDatabaseInner;
use slatedb::config::{CompressionCodec, Settings};
@ -129,7 +128,7 @@ impl Database {
) -> Result<RunDatabase> {
let (mut active_runs, run_store) = self.reserve_new_run(run_id).await?;
let (envelope, projected) = run_store.commit_first_event(payload).await?;
run_store.install_in_memory_state(&projected);
run_store.install_in_memory_state(projected);
Self::cache_active_run(&mut active_runs, &run_store);
run_store.publish(&envelope);
Ok(run_store)
@ -186,18 +185,12 @@ impl Database {
let run_ids = self.run_summary_store.list_run_ids().await?;
let mut unreadable = Vec::new();
for run_id in run_ids {
match self.run_summary_store.load_projection(&run_id).await {
Ok(Some(_)) => {}
Ok(None) => unreadable.push(UnreadableRun {
run_id,
created_at: run_id.created_at(),
error: "run has no events".to_string(),
}),
Err(err) => unreadable.push(UnreadableRun {
if let Err(err) = self.run_summary_store.load_projection(&run_id).await {
unreadable.push(UnreadableRun {
run_id,
created_at: run_id.created_at(),
error: err.to_string(),
}),
});
}
}
unreadable.sort_by(|left, right| {
@ -243,11 +236,11 @@ impl Database {
if let Some(active) = self.get_active_run(run_id).await {
return active.projection_snapshot().await.map(Some);
}
Ok(self
.run_summary_store
.load_projection(run_id)
.await?
.map(|projected| projected.projection))
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),
}
}
/// Resolves the run that owns `session_id` from the canonical typed
@ -332,6 +325,7 @@ mod tests {
use object_store::path::Path;
use super::*;
use crate::run_state::ProjectedRun;
use crate::{EventPayload, keys, test_support as store_test_support};
fn dt(value: &str) -> DateTime<Utc> {
@ -1025,11 +1019,7 @@ mod tests {
let projection = store.load_run_projection(&run_id).await.unwrap().unwrap();
let last_seq = run.last_event_seq().await.unwrap().unwrap();
let entries = [ProjectedRun::new(
run_id,
Arc::unwrap_or_clone(projection),
last_seq,
)];
let entries = [ProjectedRun::new(run_id, projection, last_seq)];
summaries.reconcile(&entries).await.unwrap();
let summary = summaries.get(&run_id, Utc::now()).await.unwrap().unwrap();
assert_eq!(summary.lifecycle.status, RunStatus::Runnable);

View file

@ -6,7 +6,7 @@ use futures::Stream;
use tokio::sync::{Mutex as AsyncMutex, broadcast, mpsc};
use tokio_stream::wrappers::UnboundedReceiverStream;
use crate::run_state::{EventProjectionCache, RunProjectionReducer};
use crate::run_state::{EventProjectionCache, ProjectedRun, RunProjectionReducer};
use crate::{
BlobStore, Error, EventEnvelope, EventPayload, Result, RunProjection, RunSummaryStore, StageId,
run_summary_store,
@ -16,35 +16,6 @@ use crate::{
/// from SQLite.
const EVENT_BROADCAST_CAPACITY: usize = 1024;
/// A run's projection as of its last committed event. Produced by replaying
/// SQLite history or by applying a newly committed event, and consumed by
/// `RunSummaryStore` writes that must stay in step with the event log.
#[derive(Debug, Clone)]
pub(crate) struct ProjectedRun {
pub(crate) run_id: RunId,
pub(crate) projection: Arc<RunProjection>,
pub(crate) last_seq: u32,
}
impl ProjectedRun {
pub(crate) fn new(run_id: RunId, projection: RunProjection, last_seq: u32) -> Self {
Self {
run_id,
projection: Arc::new(projection),
last_seq,
}
}
}
impl From<ProjectedRun> for EventProjectionCache {
fn from(projected: ProjectedRun) -> Self {
Self {
last_seq: projected.last_seq,
state: Some(projected.projection),
}
}
}
#[derive(Clone)]
pub struct RunDatabase {
inner: Arc<RunDatabaseInner>,
@ -85,10 +56,7 @@ impl RunDatabase {
blob_store: Arc<BlobStore>,
run_summary_store: Arc<RunSummaryStore>,
) -> Result<Self> {
let projected = run_summary_store
.load_projection(&run_id)
.await?
.ok_or_else(|| Error::RunNotFound(run_id.to_string()))?;
let projected = run_summary_store.load_projection(&run_id).await?;
Ok(Self::from_event_projection_cache(
run_id,
read_only,
@ -177,10 +145,8 @@ impl RunDatabase {
})
}
pub(crate) fn install_in_memory_state(&self, projected: &ProjectedRun) {
let mut projection_cache = self.inner.lock_projection_cache();
projection_cache.state = Some(Arc::clone(&projected.projection));
projection_cache.last_seq = projected.last_seq;
pub(crate) fn install_in_memory_state(&self, projected: ProjectedRun) {
*self.inner.lock_projection_cache() = projected.into();
}
pub(crate) fn publish(&self, event: &EventEnvelope) {
@ -254,7 +220,7 @@ 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.
self.install_in_memory_state(&projected);
self.install_in_memory_state(projected);
self.publish(&envelope);
Ok(envelope)
}
@ -273,11 +239,7 @@ impl RunDatabase {
apply_cached_projection_event(&mut next_state, &prospective).map_err(event_rejected)?;
let next_projection =
next_state.expect("applying a valid event should always produce a projection");
let projected = ProjectedRun::new(
self.inner.run_id,
Arc::unwrap_or_clone(next_projection),
seq,
);
let projected = ProjectedRun::new(self.inner.run_id, next_projection, seq);
let mut transaction = self.inner.run_summary_store.begin().await?;
let envelope = if expected_last_seq == 0 {