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();