diff --git a/lib/components/fabro-store/src/error.rs b/lib/components/fabro-store/src/error.rs index 43c26b44a..e5af1dbc9 100644 --- a/lib/components/fabro-store/src/error.rs +++ b/lib/components/fabro-store/src/error.rs @@ -14,6 +14,11 @@ pub enum Error { Io(#[from] std::io::Error), #[error("Invalid event payload: {0}")] InvalidEvent(String), + #[error("Event rejected by run state: {source}")] + EventRejected { + #[source] + source: Box, + }, #[error("Run not found: {0}")] RunNotFound(String), #[error("Run already exists: {0}")] diff --git a/lib/components/fabro-store/src/run_summary_store.rs b/lib/components/fabro-store/src/run_summary_store.rs index d80936dc7..57d1a126d 100644 --- a/lib/components/fabro-store/src/run_summary_store.rs +++ b/lib/components/fabro-store/src/run_summary_store.rs @@ -147,6 +147,11 @@ impl RunSummaryStore { Self { pool } } + #[cfg(test)] + pub(crate) fn test_pool(&self) -> &SqlitePool { + &self.pool + } + pub(crate) async fn upsert_projection(&self, entry: &CachedRunProjection) -> Result<()> { let record = ProjectedRunSummary::from_entry(entry); let mut connection = self.pool.acquire().await?; diff --git a/lib/components/fabro-store/src/slate/projection_cache.rs b/lib/components/fabro-store/src/slate/projection_cache.rs index 7cec9a204..79ff0c03b 100644 --- a/lib/components/fabro-store/src/slate/projection_cache.rs +++ b/lib/components/fabro-store/src/slate/projection_cache.rs @@ -5,8 +5,8 @@ use chrono::{DateTime, Utc}; use fabro_types::{Run, RunId, RunProjection}; use tokio::sync::Mutex; -use crate::run_state::{RunProjectionReducer, build_summary}; -use crate::{Error, EventEnvelope, ListRunsQuery, Result}; +use crate::ListRunsQuery; +use crate::run_state::build_summary; #[derive(Debug, Clone)] pub struct CachedRunProjection { @@ -85,26 +85,6 @@ impl RunProjectionCacheState { } } - fn update_parent_index( - &mut self, - run_id: RunId, - previous_parent_id: Option, - parent_id: Option, - ) { - if previous_parent_id == parent_id { - return; - } - if let Some(previous_parent_id) = previous_parent_id { - self.remove_parent_link(&previous_parent_id, &run_id); - } - if let Some(parent_id) = parent_id { - self.children_by_parent - .entry(parent_id) - .or_default() - .insert(run_id); - } - } - fn count_children(&self, run_id: &RunId) -> u64 { self.children_by_parent .get(run_id) @@ -222,51 +202,6 @@ impl RunProjectionCache { Some(entry.summary) } - pub(crate) async fn apply_event( - &self, - run_id: &RunId, - event: &EventEnvelope, - ) -> Result { - let mut state = self.state.lock().await; - let Some(entry) = state.entries.get(run_id) else { - if event.seq == 1 { - let projection = RunProjection::apply_events(std::slice::from_ref(event))?; - let entry = CachedRunProjection::from_projection(*run_id, projection, event.seq); - state.insert(entry.clone()); - return Ok(entry); - } - return Err(Error::InvalidEvent(format!( - "projection cache cannot initialize run {run_id} from event seq {}", - event.seq - ))); - }; - - let last_seq = entry.last_seq; - if event.seq <= last_seq { - return Ok(entry.clone()); - } - if event.seq != last_seq.saturating_add(1) { - return Err(Error::Other(format!( - "projection cache sequence gap for run {run_id}: last_seq={}, event_seq={}", - last_seq, event.seq - ))); - } - - let (previous_parent_id, parent_id, entry) = { - let entry = state - .entries - .get_mut(run_id) - .expect("entry was read from the same locked map"); - let previous_parent_id = entry.summary.parent_id; - Arc::make_mut(&mut entry.projection).apply_event(event)?; - entry.summary = build_summary(&entry.projection, run_id); - entry.last_seq = event.seq; - (previous_parent_id, entry.summary.parent_id, entry.clone()) - }; - state.update_parent_index(*run_id, previous_parent_id, parent_id); - Ok(entry) - } - pub(crate) async fn remove(&self, run_id: &RunId) { self.state.lock().await.remove(run_id); } diff --git a/lib/components/fabro-store/src/slate/run_store.rs b/lib/components/fabro-store/src/slate/run_store.rs index 9058627a6..c71fe7078 100644 --- a/lib/components/fabro-store/src/slate/run_store.rs +++ b/lib/components/fabro-store/src/slate/run_store.rs @@ -9,7 +9,7 @@ use futures::Stream; use slatedb::{Db, DbIterator, DbRead}; use tokio::sync::{Mutex, broadcast, mpsc}; use tokio_stream::wrappers::UnboundedReceiverStream; -use tracing::{error, warn}; +use tracing::warn; use super::blob_store::BlobStore; use super::projection_cache::{CachedRunProjection, RunProjectionCache}; @@ -192,6 +192,15 @@ impl RunDatabase { } async fn projected_state_locked(&self) -> Result> { + self.projected_state_option_locked().await?.ok_or_else(|| { + Error::InvalidEvent(format!( + "run {} has no run.created event", + self.inner.run_id + )) + }) + } + + async fn projected_state_option_locked(&self) -> Result>> { let next_seq = { let cache = self.inner.projection_cache.lock().await; cache.last_seq.saturating_add(1) @@ -202,25 +211,14 @@ impl RunDatabase { apply_cached_projection_event(&mut cache.state, event)?; cache.last_seq = event.seq; } - cache.state.clone().ok_or_else(|| { - Error::InvalidEvent(format!( - "run {} has no run.created event", - self.inner.run_id - )) - }) + Ok(cache.state.clone()) } - async fn cache_event(&self, event: &EventEnvelope) -> Result<()> { + async fn cache_event(&self, event: &EventEnvelope, projection: Arc) { { let mut projection_cache = self.inner.projection_cache.lock().await; - if projection_cache.state.is_none() && event.seq > 1 { - drop(projection_cache); - self.rebuild_local_projection_cache_through(event.seq) - .await?; - } else { - apply_cached_projection_event(&mut projection_cache.state, event)?; - projection_cache.last_seq = event.seq; - } + projection_cache.state = Some(projection); + projection_cache.last_seq = event.seq; } let mut recent_events = self.inner.recent_events.lock().await; recent_events.push_back(event.clone()); @@ -228,29 +226,6 @@ impl RunDatabase { recent_events.pop_front(); } let _ = self.inner.event_tx.send(event.clone()); - Ok(()) - } - - async fn rebuild_local_projection_cache_through(&self, seq: u32) -> Result<()> { - let events = list_events_from(&self.inner.db, &self.inner.run_id, 1).await?; - let Some(last_seq) = events.last().map(|event| event.seq) else { - return Err(Error::InvalidEvent(format!( - "run {} has no events while rebuilding projection cache", - self.inner.run_id - ))); - }; - if last_seq < seq { - return Err(Error::InvalidEvent(format!( - "run {} projection cache rebuild stopped at seq {last_seq}, before appended seq {seq}", - self.inner.run_id - ))); - } - - let state = RunProjection::apply_events(&events)?; - let mut projection_cache = self.inner.projection_cache.lock().await; - projection_cache.state = Some(Arc::new(state)); - projection_cache.last_seq = last_seq; - Ok(()) } async fn cached_events_from(&self, start_seq: u32, limit: usize) -> Option> { @@ -270,12 +245,40 @@ impl RunDatabase { } impl RunDatabase { + /// Appends `payload` to the run's authoritative event log. + /// + /// A payload rejected by the current run state is not written. Any other + /// error also means the event was not committed and is safe to retry. Once + /// the authoritative write succeeds, derived cache and SQLite updates are + /// best-effort and this method returns success even when one of them fails. + /// + /// # Errors + /// + /// Returns [`Error::EventRejected`] when the run reducer rejects the + /// candidate event. Other errors report validation, sequence allocation, + /// projection loading, serialization, read-only access, or an + /// authoritative storage failure. No returned error represents a committed + /// append. pub async fn append_event(&self, payload: &EventPayload) -> Result { Ok(self.append_event_envelope(payload).await?.seq) } /// Atomically appends `payload` when `predicate` matches the latest run /// projection. + /// + /// `Ok(None)` means the predicate rejected the append and nothing was + /// written. A state-machine rejection or any other error also leaves the + /// event uncommitted and safe to retry. Once the authoritative write + /// succeeds, derived cache and SQLite updates are best-effort and this + /// method returns the committed sequence. + /// + /// # Errors + /// + /// Returns [`Error::EventRejected`] when the predicate accepts but the run + /// reducer rejects the candidate event. Other errors report validation, + /// sequence allocation, projection loading, serialization, read-only + /// access, or an authoritative storage failure. No returned error + /// represents a committed append. pub async fn append_event_if( &self, payload: &EventPayload, @@ -290,90 +293,91 @@ impl RunDatabase { if !predicate(&projection) { return Ok(None); } - Ok(Some(self.append_event_envelope_locked(payload).await?.seq)) + Ok(Some( + self.append_event_envelope_locked(payload, Some(projection)) + .await? + .seq, + )) } + /// Appends `payload` and returns its committed envelope. + /// + /// A payload rejected by the current run state is not written. Any other + /// error also means the event was not committed and is safe to retry. Once + /// the authoritative write succeeds, derived cache and SQLite updates are + /// best-effort and this method returns success even when one of them fails. + /// + /// # Errors + /// + /// Returns [`Error::EventRejected`] when the run reducer rejects the + /// candidate event. Other errors report validation, sequence allocation, + /// projection loading, serialization, read-only access, or an + /// authoritative storage failure. No returned error represents a committed + /// append. 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 + self.append_event_envelope_locked(payload, None).await } - async fn append_event_envelope_locked(&self, payload: &EventPayload) -> Result { + async fn append_event_envelope_locked( + &self, + payload: &EventPayload, + current_projection: Option>, + ) -> Result { let event_seq = self.inner.event_seq.as_ref().ok_or(Error::ReadOnly)?; - let seq = allocate_event_seq(event_seq)?; + let seq = event_seq.load(Ordering::SeqCst); + if seq > keys::MAX_EVENT_SEQ { + return Err(Error::EventSequenceExhausted { + max_seq: keys::MAX_EVENT_SEQ, + }); + } let event = EventEnvelope { seq, event: RunEvent::try_from(payload)?, }; + let projection = match current_projection { + Some(projection) => validate_projected_event(&projection, &event)?, + None => match self.projected_state_option_locked().await? { + Some(projection) => validate_projected_event(&projection, &event)?, + None => validate_first_event(&event)?, + }, + }; + let encoded = serde_json::to_vec(payload)?; + allocate_event_seq(event_seq)?; self.inner .db .put( keys::run_event_key(&self.inner.run_id, seq, Utc::now().timestamp_millis()), - serde_json::to_vec(payload)?, + encoded, ) .await?; - self.cache_event(&event).await?; - // Box::pin keeps append_event_envelope's future small enough for the - // clippy::large_futures budget of its many callers. - Box::pin(self.update_summary_projection_after_append(&event)).await?; - Ok(event) - } - async fn update_summary_projection_after_append(&self, event: &EventEnvelope) -> Result<()> { - let cached = match self - .inner + let cached = CachedRunProjection::from_projection( + self.inner.run_id, + Arc::unwrap_or_clone(projection), + event.seq, + ); + self.inner .shared_projection_cache - .apply_event(&self.inner.run_id, event) - .await - { - Ok(entry) => entry, - Err(err) => { - match Self::build_cached_projection(&self.inner.db, &self.inner.run_id).await { - Ok(Some(entry)) => { - self.inner - .shared_projection_cache - .replace(entry.clone()) - .await; - entry - } - rebuild => { - self.inner - .shared_projection_cache - .remove(&self.inner.run_id) - .await; - if let Err(rebuild_err) = rebuild { - warn!( - run_id = %self.inner.run_id, - error = %rebuild_err, - "Failed to rebuild run projection cache after append" - ); - } - warn!( - run_id = %self.inner.run_id, - error = %err, - "Failed to update run projection cache after append" - ); - return Err(err); - } - } - } - }; + .replace(cached.clone()) + .await; + self.cache_event(&event, Arc::clone(&cached.projection)) + .await; if let Some(store) = self.inner.run_summary_store.get() { if let Err(err) = store.upsert_projection(&cached).await { - error!( + warn!( run_id = %self.inner.run_id, source_last_seq = cached.last_seq, error = %err, - "Failed to update SQLite run summary after append" + "Failed to update SQLite run summary after committed append" ); - return Err(err); } } - Ok(()) + Ok(event) } pub async fn list_events(&self) -> Result> { @@ -599,6 +603,27 @@ fn allocate_event_seq(event_seq: &AtomicU32) -> Result { }) } +fn validate_first_event(event: &EventEnvelope) -> Result> { + RunProjection::apply_events(std::slice::from_ref(event)) + .map(Arc::new) + .map_err(event_rejected) +} + +fn validate_projected_event( + current: &RunProjection, + event: &EventEnvelope, +) -> Result> { + let mut projection = current.clone(); + projection.apply_event(event).map_err(event_rejected)?; + Ok(Arc::new(projection)) +} + +fn event_rejected(source: Error) -> Error { + Error::EventRejected { + source: Box::new(source), + } +} + fn apply_cached_projection_event( state: &mut Option>, event: &EventEnvelope, @@ -909,11 +934,15 @@ mod tests { use std::sync::atomic::Ordering; use std::time::Duration; - use fabro_types::{Graph, RunId, SessionId, StageId, WorkflowSettings, test_support}; + use chrono::Utc; + use fabro_types::{ + Graph, RunId, RunStatus, SessionId, StageId, WorkflowSettings, test_support, + }; + use fabro_util::error; use object_store::memory::InMemory; use serde_json::json; - use crate::{Database, Error, EventPayload, keys}; + use crate::{Database, Error, EventPayload, keys, test_util}; #[tokio::test] async fn list_blobs_reads_global_cas_namespace() { @@ -973,6 +1002,79 @@ mod tests { .unwrap() } + fn lifecycle_payload( + run_id: &RunId, + id: &str, + ts: &str, + event: &str, + properties: &serde_json::Value, + ) -> EventPayload { + EventPayload::new( + json!({ + "id": id, + "ts": ts, + "run_id": run_id.to_string(), + "event": event, + "properties": properties, + }), + run_id, + ) + .unwrap() + } + + fn run_failed_payload(run_id: &RunId) -> EventPayload { + lifecycle_payload( + run_id, + "evt-failed", + "2026-04-09T12:00:04Z", + "run.failed", + &json!({ + "failure": { + "reason": "workflow_error", + "detail": { + "message": "workflow failed", + "category": "deterministic", + }, + }, + "timing": { + "wall_time_ms": 1, + "inference_time_ms": 0, + "tool_time_ms": 0, + "active_time_ms": 0, + }, + }), + ) + } + + async fn drive_to_runnable(run: &super::RunDatabase) { + let run_id = run.run_id(); + for payload in [ + lifecycle_payload( + &run_id, + "evt-submitted", + "2026-04-09T12:00:00Z", + "run.submitted", + &json!({ "definition_blob": null }), + ), + lifecycle_payload( + &run_id, + "evt-start-requested", + "2026-04-09T12:00:01Z", + "run.start_requested", + &json!({ "resume": false }), + ), + lifecycle_payload( + &run_id, + "evt-runnable", + "2026-04-09T12:00:02Z", + "run.runnable", + &json!({ "source": "start_requested" }), + ), + ] { + run.append_event(&payload).await.unwrap(); + } + } + fn stage_prompt_payload_for_stage( run_id: &RunId, idx: u32, @@ -1015,6 +1117,214 @@ mod tests { run } + #[tokio::test] + async fn invalid_transition_is_rejected_without_changing_durable_or_cached_state() { + let run = fresh_run().await; + let run_id = run.run_id(); + drive_to_runnable(&run).await; + let events_before = run.list_events().await.unwrap(); + + let err = run + .append_event(&run_failed_payload(&run_id)) + .await + .unwrap_err(); + + assert!(matches!( + &err, + Error::EventRejected { source } + if matches!(source.as_ref(), Error::InvalidTransition(_)) + )); + let error_chain = error::collect_chain(&err); + assert!(error_chain.len() >= 2); + assert!( + error_chain + .iter() + .skip(1) + .any(|cause| cause.contains("invalid status transition")) + ); + assert_eq!(run.list_events().await.unwrap(), events_before); + assert!( + run.get_event(events_before.last().unwrap().seq + 1) + .await + .unwrap() + .is_none() + ); + assert_eq!(run.state().await.unwrap().status, RunStatus::Runnable); + let cached = run + .inner + .shared_projection_cache + .get(&run_id) + .await + .unwrap(); + assert_eq!(cached.last_seq, events_before.last().unwrap().seq); + assert_eq!(cached.projection.status, RunStatus::Runnable); + + let next_seq = run + .append_event(&stage_prompt_payload(&run_id, 1, Some("alpha"))) + .await + .unwrap(); + assert_eq!(next_seq, events_before.last().unwrap().seq + 1); + } + + #[tokio::test] + async fn rejected_transition_leaves_summary_available_after_reconcile() { + let object_store = Arc::new(InMemory::new()); + let database = Database::new(object_store, "", Duration::from_millis(1), None); + let (_directory, summary_store) = test_util::sqlite_summary_store().await; + let summary_store = Arc::new(summary_store); + database.attach_run_summary_store(Arc::clone(&summary_store)); + let run_id: RunId = "01JT56VE4Z5NZ814GZN2JZD65A".parse().unwrap(); + let run = database.create_run(&run_id).await.unwrap(); + run.append_event(&run_created_payload(&run_id)) + .await + .unwrap(); + drive_to_runnable(&run).await; + + let err = run + .append_event(&run_failed_payload(&run_id)) + .await + .unwrap_err(); + assert!(matches!(err, Error::EventRejected { .. })); + let cached = run + .inner + .shared_projection_cache + .get(&run_id) + .await + .unwrap(); + assert_eq!(cached.last_seq, 4); + summary_store.reconcile(&[cached]).await.unwrap(); + + assert!( + summary_store + .get(&run_id, Utc::now()) + .await + .unwrap() + .is_some() + ); + } + + #[tokio::test] + async fn committed_append_succeeds_when_summary_update_fails_and_is_repairable() { + let object_store = Arc::new(InMemory::new()); + let database = Database::new(object_store, "", Duration::from_millis(1), None); + let (_directory, summary_store) = test_util::sqlite_summary_store().await; + let summary_store = Arc::new(summary_store); + database.attach_run_summary_store(Arc::clone(&summary_store)); + let run_id: RunId = "01JT56VE4Z5NZ814GZN2JZD65A".parse().unwrap(); + let run = database.create_run(&run_id).await.unwrap(); + run.append_event(&run_created_payload(&run_id)) + .await + .unwrap(); + sqlx::query( + "CREATE TRIGGER reject_run_summary_update BEFORE UPDATE ON runs BEGIN SELECT \ + RAISE(ABORT, 'forced run summary update failure'); END", + ) + .execute(summary_store.test_pool()) + .await + .unwrap(); + let first_update = lifecycle_payload( + &run_id, + "evt-title-1", + "2026-04-09T12:00:01Z", + "run.title.updated", + &json!({ "title": "first update" }), + ); + + let seq = run.append_event(&first_update).await.unwrap(); + + assert_eq!(run.get_event(seq).await.unwrap().unwrap().seq, seq); + assert_ne!( + summary_store + .get(&run_id, Utc::now()) + .await + .unwrap() + .unwrap() + .title, + "first update" + ); + sqlx::query("DROP TRIGGER reject_run_summary_update") + .execute(summary_store.test_pool()) + .await + .unwrap(); + let repaired_seq = run + .append_event(&lifecycle_payload( + &run_id, + "evt-title-2", + "2026-04-09T12:00:02Z", + "run.title.updated", + &json!({ "title": "repaired" }), + )) + .await + .unwrap(); + let repaired = summary_store + .get(&run_id, Utc::now()) + .await + .unwrap() + .unwrap(); + assert_eq!(repaired.title, "repaired"); + assert_eq!(repaired_seq, seq + 1); + } + + #[tokio::test] + async fn first_event_must_initialize_a_projection_before_any_write() { + let object_store = Arc::new(InMemory::new()); + let database = Database::new(object_store, "", Duration::from_millis(1), None); + let run_id: RunId = "01JT56VE4Z5NZ814GZN2JZD65A".parse().unwrap(); + let run = database.create_run(&run_id).await.unwrap(); + + let err = run + .append_event(&stage_prompt_payload(&run_id, 1, Some("alpha"))) + .await + .unwrap_err(); + + assert!(matches!( + &err, + Error::EventRejected { source } + if matches!(source.as_ref(), Error::InvalidEvent(_)) + )); + assert!(run.list_events().await.unwrap().is_empty()); + assert!( + run.inner + .shared_projection_cache + .projection_snapshot(&run_id) + .await + .is_none() + ); + assert_eq!( + run.inner.event_seq.as_ref().unwrap().load(Ordering::SeqCst), + 1 + ); + + assert_eq!( + run.append_event(&run_created_payload(&run_id)) + .await + .unwrap(), + 1 + ); + assert_eq!(run.list_events().await.unwrap().len(), 1); + } + + #[tokio::test] + async fn append_event_if_false_writes_nothing() { + let run = fresh_run().await; + let run_id = run.run_id(); + let events_before = run.list_events().await.unwrap(); + + let result = run + .append_event_if(&stage_prompt_payload(&run_id, 1, Some("alpha")), |_| false) + .await + .unwrap(); + + assert_eq!(result, None); + assert_eq!(run.list_events().await.unwrap(), events_before); + assert!( + run.get_event(events_before.last().unwrap().seq + 1) + .await + .unwrap() + .is_none() + ); + } + #[tokio::test] async fn list_events_from_with_limit_does_not_read_past_limit_plus_one() { let run = fresh_run().await; @@ -1304,6 +1614,7 @@ mod tests { .unwrap(); assert_eq!(seq, keys::MAX_EVENT_SEQ); + let events_before_error = run.list_events().await.unwrap(); let err = run .append_event(&stage_prompt_payload(&run_id, 2, Some("beta"))) .await @@ -1313,6 +1624,7 @@ mod tests { Error::EventSequenceExhausted { max_seq } if max_seq == keys::MAX_EVENT_SEQ )); + assert_eq!(run.list_events().await.unwrap(), events_before_error); assert!( run.get_event(keys::MAX_EVENT_SEQ + 1) .await