From 162791979c5eb09d38761f3f96753341757bab8d Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Thu, 17 Sep 2026 22:32:23 -0400 Subject: [PATCH] Wake the projector after a Petri run's first event too The run's creation commits its first event on the create path, not the append path, so the platform record hook never fired for run.created. The store now notifies after that commit as well, and a test proves a Petri run's lifecycle events leave platform records beside them with one wake-up per record while a legacy run leaves none. Co-Authored-By: Claude Fable 5.1 --- .../fabro-store/src/platform_records.rs | 8 +- .../fabro-store/src/run_summary_store.rs | 6 +- lib/components/fabro-store/src/slate/mod.rs | 8 +- .../fabro-store/src/slate/run_store.rs | 80 +++++++++++++++++++ 4 files changed, 94 insertions(+), 8 deletions(-) diff --git a/lib/components/fabro-store/src/platform_records.rs b/lib/components/fabro-store/src/platform_records.rs index d2ae7625f..83a436180 100644 --- a/lib/components/fabro-store/src/platform_records.rs +++ b/lib/components/fabro-store/src/platform_records.rs @@ -766,7 +766,7 @@ fn pull_request_created_record(props: &PullRequestCreatedProps) -> PullRequestCr fn run_paired_record(props: &RunPairStartedProps) -> RunPairedRecord { RunPairedRecord { - pair_id: props.pair_id.clone(), + pair_id: props.pair_id, target: props.target.clone(), } } @@ -784,7 +784,7 @@ fn interview_answered_record( #[cfg(test)] mod tests { - use fabro_types::{FailureReason, RunStatus, fixtures}; + use fabro_types::{FailureReason, RunStatus, fixtures, test_support as types_support}; use serde_json::json; use super::*; @@ -803,7 +803,7 @@ mod tests { fn sample(kind: PlatformRecordKind) -> PlatformRecord { match kind { PlatformRecordKind::RunCreated => PlatformRecord::RunCreated(RunCreatedRecord { - spec: fabro_types::test_support::test_run_spec(), + spec: types_support::test_run_spec(), title: Some("A run".to_string()), parent_id: None, retried_from: None, @@ -957,7 +957,7 @@ mod tests { assert_eq!(json(&stored), json(&[first, second.clone()])); assert_eq!( json(&store.read_after(&run, 1).await.expect("the tail reads")), - json(&[second.clone()]) + json(std::slice::from_ref(&second)) ); assert_eq!( json( diff --git a/lib/components/fabro-store/src/run_summary_store.rs b/lib/components/fabro-store/src/run_summary_store.rs index 8b51a41fd..802415cef 100644 --- a/lib/components/fabro-store/src/run_summary_store.rs +++ b/lib/components/fabro-store/src/run_summary_store.rs @@ -1,5 +1,5 @@ use std::fmt::Write as _; -use std::sync::{Arc, LazyLock, RwLock}; +use std::sync::{Arc, LazyLock, PoisonError, RwLock}; use chrono::{DateTime, Utc}; use fabro_types::{ @@ -250,14 +250,14 @@ impl RunSummaryStore { *self .platform_hook .write() - .unwrap_or_else(std::sync::PoisonError::into_inner) = Some(hook); + .unwrap_or_else(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) + .unwrap_or_else(PoisonError::into_inner) .clone(); if let Some(hook) = hook { hook(run_id); diff --git a/lib/components/fabro-store/src/slate/mod.rs b/lib/components/fabro-store/src/slate/mod.rs index 93330261d..b55e047e5 100644 --- a/lib/components/fabro-store/src/slate/mod.rs +++ b/lib/components/fabro-store/src/slate/mod.rs @@ -13,7 +13,9 @@ use run_store::RunDatabaseInner; use slatedb::config::{CompressionCodec, Settings}; use tokio::sync::{Mutex, MutexGuard, OnceCell}; -use crate::{BlobStore, Error, EventPayload, Result, RunProjection, RunSummaryStore, keys}; +use crate::{ + BlobStore, Error, EventPayload, Result, RunProjection, RunSummaryStore, keys, run_summary_store, +}; #[derive(Debug, Clone, PartialEq, Eq)] pub struct UnreadableRun { @@ -128,9 +130,13 @@ impl Database { ) -> Result { let (mut active_runs, run_store) = self.reserve_new_run(run_id).await?; let (envelope, projected) = run_store.commit_first_event(payload).await?; + let platform_record = run_summary_store::platform_record_written(&projected, &envelope); run_store.install_in_memory_state(projected); Self::cache_active_run(&mut active_runs, &run_store); run_store.publish(&envelope); + if platform_record.is_some() { + self.run_summary_store.notify_platform_record(*run_id); + } Ok(run_store) } diff --git a/lib/components/fabro-store/src/slate/run_store.rs b/lib/components/fabro-store/src/slate/run_store.rs index 19e15a938..ca70766aa 100644 --- a/lib/components/fabro-store/src/slate/run_store.rs +++ b/lib/components/fabro-store/src/slate/run_store.rs @@ -510,6 +510,86 @@ mod tests { ); } + /// A Petri run's legacy lifecycle events leave platform records beside + /// them, in the same commit, and the hook fires after each; a legacy + /// run's events leave none. + #[tokio::test] + async fn a_petri_runs_lifecycle_events_become_platform_records_and_wake_the_hook() { + use std::sync::atomic::{AtomicUsize, Ordering}; + + use crate::platform_records::{PlatformRecord, PlatformRecordKind, RunLifecycleKind}; + + let store = store(); + let woken = Arc::new(AtomicUsize::new(0)); + let counter = Arc::clone(&woken); + store.set_platform_record_hook(Arc::new(move |_| { + counter.fetch_add(1, Ordering::SeqCst); + })); + let run_id: RunId = "01JT56VE4Z5NZ814GZN2JZD65E".parse().unwrap(); + let mut created = run_created_payload(&run_id); + let mut petri = serde_json::to_value(&created).unwrap(); + petri["properties"]["engine"] = json!({ + "kind": "petri", + "graph": { "blob": fabro_types::BlobHash::new(b"graph").to_string(), "digest": "d" }, + }); + created = EventPayload::new(petri, &run_id).unwrap(); + let run = store + .create_run_with_first_event(&run_id, &created) + .await + .unwrap(); + run.append_event( + &EventPayload::new( + json!({ + "id": "evt-starting", + "ts": "2026-04-09T12:00:00Z", + "run_id": run_id.to_string(), + "event": "run.start_requested", + "properties": { "resume": false }, + }), + &run_id, + ) + .unwrap(), + ) + .await + .unwrap(); + run.append_event(&stage_payload(&run_id, 3)).await.unwrap(); + + let records = store + .run_summary_store() + .platform_records() + .read(&run_id) + .await + .unwrap(); + let kinds: Vec = records.iter().map(|r| r.record.kind()).collect(); + assert_eq!(kinds, vec![ + PlatformRecordKind::RunCreated, + PlatformRecordKind::RunLifecycle + ]); + let PlatformRecord::RunLifecycle(lifecycle) = &records[1].record else { + panic!("the second record is the lifecycle"); + }; + assert_eq!(lifecycle.transition, RunLifecycleKind::StartRequested); + assert_eq!(lifecycle.source.as_deref(), Some("start")); + assert_eq!(woken.load(Ordering::SeqCst), 2, "one wake-up per record"); + + let legacy_id: RunId = "01JT56VE4Z5NZ814GZN2JZD65F".parse().unwrap(); + store + .create_run_with_first_event(&legacy_id, &run_created_payload(&legacy_id)) + .await + .unwrap(); + assert!( + store + .run_summary_store() + .platform_records() + .read(&legacy_id) + .await + .unwrap() + .is_empty(), + "a legacy run leaves no platform records" + ); + assert_eq!(woken.load(Ordering::SeqCst), 2); + } + #[tokio::test] async fn watcher_catches_up_from_sql_without_duplicates() { let store = store();