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 <noreply@anthropic.com>
This commit is contained in:
Bryan Helmkamp 2026-09-17 22:32:23 -04:00
parent abbc7ca11d
commit 162791979c
No known key found for this signature in database
4 changed files with 94 additions and 8 deletions

View file

@ -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(

View file

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

View file

@ -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<RunDatabase> {
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)
}

View file

@ -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<PlatformRecordKind> = 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();