diff --git a/lib/crates/fabro-server/src/server/handler/runs.rs b/lib/crates/fabro-server/src/server/handler/runs.rs index 57562fc4c..614f2d719 100644 --- a/lib/crates/fabro-server/src/server/handler/runs.rs +++ b/lib/crates/fabro-server/src/server/handler/runs.rs @@ -743,23 +743,6 @@ fn spawn_generated_title_task(task: GeneratedTitleTask) { return; } - let current = match task - .state - .stores - .runs - .get_cached_summary(&task.run_id, Utc::now()) - .await - { - Ok(Some(summary)) => summary, - Ok(None) => return, - Err(err) => { - tracing::debug!(run_id = %task.run_id, error = %err, "Failed to re-read run summary for title update"); - return; - } - }; - if current.title != task.deterministic_title { - return; - } let run_store = match task.state.stores.runs.open_run(&task.run_id).await { Ok(store) => store, Err(err) => { @@ -767,7 +750,8 @@ fn spawn_generated_title_task(task: GeneratedTitleTask) { return; } }; - if let Err(err) = workflow_event::append_event( + let expected_title = task.deterministic_title; + if let Err(err) = workflow_event::append_event_if( &run_store, &task.run_id, &workflow_event::Event::RunTitleUpdated { @@ -776,6 +760,7 @@ fn spawn_generated_title_task(task: GeneratedTitleTask) { system_kind: SystemActorKind::Engine, }), }, + move |projection| projection.title().as_ref() == expected_title, ) .await { diff --git a/lib/crates/fabro-server/src/server/tests.rs b/lib/crates/fabro-server/src/server/tests.rs index 7aafcec03..c90e1d10f 100644 --- a/lib/crates/fabro-server/src/server/tests.rs +++ b/lib/crates/fabro-server/src/server/tests.rs @@ -3495,7 +3495,7 @@ async fn generated_title_does_not_overwrite_user_title_edit() { response_json!(response, StatusCode::OK).await; wait_for_mock_hits(&title_mock, 1).await; - tokio::time::sleep(std::time::Duration::from_millis(50)).await; + tokio::time::sleep(std::time::Duration::from_millis(250)).await; assert_eq!( state diff --git a/lib/crates/fabro-store/src/slate/mod.rs b/lib/crates/fabro-store/src/slate/mod.rs index 60c40d1c4..7546a741f 100644 --- a/lib/crates/fabro-store/src/slate/mod.rs +++ b/lib/crates/fabro-store/src/slate/mod.rs @@ -803,6 +803,40 @@ mod tests { assert!(matches!(err, Error::ReadOnly)); } + #[tokio::test] + async fn append_event_if_evaluates_latest_projection_before_appending() { + let (_object_store, store) = make_store(); + let run = store.create_run(&test_run_id("run-1")).await.unwrap(); + append_created(&run, "run-1", dt("2026-03-27T12:00:00Z")).await; + let initial_title = run.state().await.unwrap().title().into_owned(); + + run.append_event(&event_payload( + "run-1", + "2026-03-27T12:00:01Z", + "run.title.updated", + &serde_json::json!({ "title": "User title" }), + )) + .await + .unwrap(); + + let generated_update = event_payload( + "run-1", + "2026-03-27T12:00:02Z", + "run.title.updated", + &serde_json::json!({ "title": "Generated title" }), + ); + let appended = run + .append_event_if(&generated_update, |projection| { + projection.title() == initial_title + }) + .await + .unwrap(); + + assert_eq!(appended, None); + assert_eq!(run.state().await.unwrap().title(), "User title"); + assert_eq!(run.list_events().await.unwrap().len(), 2); + } + #[tokio::test] async fn control_request_events_set_pending_control_without_overwriting_status() { let (_object_store, store) = make_store(); diff --git a/lib/crates/fabro-store/src/slate/run_store.rs b/lib/crates/fabro-store/src/slate/run_store.rs index 2b67a066a..a1d7f8441 100644 --- a/lib/crates/fabro-store/src/slate/run_store.rs +++ b/lib/crates/fabro-store/src/slate/run_store.rs @@ -169,6 +169,10 @@ impl RunDatabase { async fn projected_state(&self) -> Result { let _state_guard = self.inner.state_lock.lock().await; + self.projected_state_locked().await + } + + async fn projected_state_locked(&self) -> Result { let next_seq = { let cache = self.inner.projection_cache.lock().await; cache.last_seq.saturating_add(1) @@ -254,12 +258,35 @@ impl RunDatabase { Ok(self.append_event_envelope(payload).await?.seq) } + /// Atomically appends `payload` when `predicate` matches the latest run + /// projection. + pub async fn append_event_if( + &self, + payload: &EventPayload, + predicate: impl FnOnce(&RunProjection) -> bool, + ) -> Result> { + if self.read_only { + return Err(Error::ReadOnly); + } + payload.validate(&self.inner.run_id)?; + let _state_guard = self.inner.state_lock.lock().await; + let projection = self.projected_state_locked().await?; + if !predicate(&projection) { + return Ok(None); + } + Ok(Some(self.append_event_envelope_locked(payload).await?.seq)) + } + pub async fn append_event_envelope(&self, payload: &EventPayload) -> Result { if self.read_only { return Err(Error::ReadOnly); } payload.validate(&self.inner.run_id)?; let _state_guard = self.inner.state_lock.lock().await; + self.append_event_envelope_locked(payload).await + } + + async fn append_event_envelope_locked(&self, payload: &EventPayload) -> Result { let seq = self.inner.event_seq.fetch_add(1, Ordering::SeqCst); let event = EventEnvelope { seq, diff --git a/lib/crates/fabro-workflow/src/event.rs b/lib/crates/fabro-workflow/src/event.rs index e64b2ae30..a5c1f583e 100644 --- a/lib/crates/fabro-workflow/src/event.rs +++ b/lib/crates/fabro-workflow/src/event.rs @@ -18,6 +18,7 @@ pub use self::redaction::{ build_redacted_event_payload, event_payload_from_redacted_json, redacted_event_json, }; pub use self::sink::{ - RunEventLogger, RunEventSink, StoreProgressLogger, append_event, append_event_to_sink, + RunEventLogger, RunEventSink, StoreProgressLogger, append_event, append_event_if, + append_event_to_sink, }; pub use crate::stage_scope::StageScope; diff --git a/lib/crates/fabro-workflow/src/event/sink.rs b/lib/crates/fabro-workflow/src/event/sink.rs index 1017b2b29..7d62b7151 100644 --- a/lib/crates/fabro-workflow/src/event/sink.rs +++ b/lib/crates/fabro-workflow/src/event/sink.rs @@ -2,7 +2,7 @@ use std::future::Future; use std::pin::Pin; use std::sync::Arc; -use ::fabro_types::{RunEvent, RunId}; +use ::fabro_types::{RunEvent, RunId, RunProjection}; use anyhow::Result; use fabro_store::RunDatabase; use tokio::io::{AsyncWrite, AsyncWriteExt}; @@ -23,6 +23,21 @@ pub async fn append_event(run_store: &RunDatabase, run_id: &RunId, event: &Event .map_err(anyhow::Error::from) } +pub async fn append_event_if( + run_store: &RunDatabase, + run_id: &RunId, + event: &Event, + predicate: impl FnOnce(&RunProjection) -> bool, +) -> Result { + let stored = to_run_event(run_id, event); + let payload = build_redacted_event_payload(&stored, run_id)?; + run_store + .append_event_if(&payload, predicate) + .await + .map(|seq| seq.is_some()) + .map_err(anyhow::Error::from) +} + pub async fn append_event_to_sink( sink: &RunEventSink, run_id: &RunId,