fabro(01KYQMV1VW6139EGNHEM1RGF2G): implement (succeeded)

Fabro-Run: 01KYQMV1VW6139EGNHEM1RGF2G
Fabro-Completed: 5
Fabro-Checkpoint: c84d4147ba

⚒️ Generated with [Fabro](https://fabro.sh)
This commit is contained in:
Fabro 2026-07-29 20:35:54 +00:00
parent ed59858bd6
commit 4746d143fd
6 changed files with 295 additions and 163 deletions

View file

@ -14,6 +14,8 @@ pub enum Error {
Io(#[from] std::io::Error),
#[error("Invalid event payload: {0}")]
InvalidEvent(String),
#[error("event rejected by run projection: {reason}")]
EventRejected { reason: String },
#[error("Run not found: {0}")]
RunNotFound(String),
#[error("Run already exists: {0}")]

View file

@ -154,6 +154,11 @@ impl RunSummaryStore {
Ok(())
}
#[cfg(test)]
pub(crate) async fn close_pool(&self) {
self.pool.close().await;
}
pub(crate) async fn reconcile(&self, entries: &[CachedRunProjection]) -> Result<()> {
let mut transaction = self.pool.begin().await?;
let stored_seqs: HashMap<String, i64> =

View file

@ -684,6 +684,57 @@ mod tests {
.unwrap();
}
async fn append_runnable(run: &RunDatabase, label: &str, created_at: DateTime<Utc>) {
append_created(run, label, created_at).await;
run.append_event(&event_payload(
label,
"2026-03-27T12:00:01Z",
"run.submitted",
&serde_json::json!({}),
))
.await
.unwrap();
run.append_event(&event_payload(
label,
"2026-03-27T12:00:02Z",
"run.start_requested",
&serde_json::json!({ "resume": false }),
))
.await
.unwrap();
run.append_event(&event_payload(
label,
"2026-03-27T12:00:03Z",
"run.runnable",
&serde_json::json!({ "source": "start_requested" }),
))
.await
.unwrap();
}
fn workflow_failure_payload(label: &str) -> EventPayload {
event_payload(
label,
"2026-03-27T12:00:04Z",
"run.failed",
&serde_json::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 append_completed(run: &RunDatabase, label: &str, created_at: DateTime<Utc>) {
append_running(run, label, created_at).await;
run.append_event(&event_payload(
@ -849,6 +900,160 @@ mod tests {
assert_eq!(run.list_events().await.unwrap().len(), 2);
}
#[tokio::test]
async fn rejected_transition_writes_nothing_and_preserves_projection_cache() {
let (_object_store, store) = make_store();
let run_id = test_run_id("run-1");
let run = store.create_run(&run_id).await.unwrap();
append_runnable(&run, "run-1", dt("2026-03-27T12:00:00Z")).await;
let events_before = run.list_events().await.unwrap();
let err = run
.append_event(&workflow_failure_payload("run-1"))
.await
.unwrap_err();
let Error::EventRejected { reason } = err else {
panic!("expected event rejection");
};
assert_eq!(
reason,
"invalid status transition: runnable -> failed(workflow_error)"
);
assert_eq!(run.list_events().await.unwrap(), events_before);
assert_eq!(run.state().await.unwrap().status, RunStatus::Runnable);
let cached = store.get_cached_run(&run_id).await.unwrap().unwrap();
assert_eq!(cached.last_seq, 4);
assert_eq!(cached.projection.status, RunStatus::Runnable);
}
#[tokio::test]
async fn rejected_transition_leaves_reconciled_summary_present() {
let (_object_store, store) = make_store();
let (_directory, summaries) = make_summary_store().await;
store.attach_run_summary_store(Arc::clone(&summaries));
let run_id = test_run_id("run-1");
let run = store.create_run(&run_id).await.unwrap();
append_runnable(&run, "run-1", dt("2026-03-27T12:00:00Z")).await;
let err = run
.append_event(&workflow_failure_payload("run-1"))
.await
.unwrap_err();
assert!(matches!(err, Error::EventRejected { .. }));
let entries = store
.list_cached_runs(&ListRunsQuery::default(), Utc::now())
.await
.unwrap();
summaries.reconcile(&entries).await.unwrap();
let summary = summaries.get(&run_id, Utc::now()).await.unwrap().unwrap();
assert_eq!(summary.lifecycle.status, RunStatus::Runnable);
}
#[tokio::test]
async fn committed_append_succeeds_when_summary_update_fails_and_is_repairable() {
let (_object_store, store) = make_store();
let (directory, summaries) = make_summary_store().await;
store.attach_run_summary_store(Arc::clone(&summaries));
let run_id = test_run_id("run-1");
let run = store.create_run(&run_id).await.unwrap();
append_created(&run, "run-1", dt("2026-03-27T12:00:00Z")).await;
summaries.close_pool().await;
let result = run
.append_event_envelope(&event_payload(
"run-1",
"2026-03-27T12:00:01Z",
"run.title.updated",
&serde_json::json!({ "title": "Committed title" }),
))
.await;
assert!(result.is_ok(), "committed append returned {result:?}");
assert_eq!(run.list_events().await.unwrap().len(), 2);
let cached = store.get_cached_run(&run_id).await.unwrap().unwrap();
assert_eq!(cached.last_seq, 2);
assert_eq!(cached.summary.title, "Committed title");
let stored = run.get_event(2).await.unwrap().unwrap();
assert_eq!(stored.event, result.unwrap().event);
let repaired_database = fabro_db::Database::connect(directory.path().join("fabro.sqlite3"))
.await
.unwrap();
repaired_database.migrate().await.unwrap();
let repaired_summaries = RunSummaryStore::new(repaired_database.clone_pool());
let stale = repaired_summaries
.get(&run_id, Utc::now())
.await
.unwrap()
.unwrap();
assert_ne!(stale.title, "Committed title");
let entries = store
.list_cached_runs(&ListRunsQuery::default(), Utc::now())
.await
.unwrap();
repaired_summaries.reconcile(&entries).await.unwrap();
let repaired = repaired_summaries
.get(&run_id, Utc::now())
.await
.unwrap()
.unwrap();
assert_eq!(repaired.title, "Committed title");
}
#[tokio::test]
async fn first_event_is_validated_before_write() {
let (_object_store, store) = make_store();
let run_id = test_run_id("run-1");
let run = store.create_run(&run_id).await.unwrap();
let invalid_first = event_payload(
"run-1",
"2026-03-27T12:00:00Z",
"run.title.updated",
&serde_json::json!({ "title": "Too early" }),
);
let err = run.append_event(&invalid_first).await.unwrap_err();
assert!(matches!(err, Error::EventRejected { .. }));
assert!(run.list_events().await.unwrap().is_empty());
append_created(&run, "run-1", dt("2026-03-27T12:00:01Z")).await;
assert_eq!(run.list_events().await.unwrap().len(), 1);
assert!(run.state().await.is_ok());
}
#[tokio::test]
async fn malformed_optional_envelope_field_is_rejected_before_write() {
let (_object_store, store) = make_store();
let run_id = test_run_id("run-1");
let run = store.create_run(&run_id).await.unwrap();
let malformed = EventPayload::new(
serde_json::json!({
"id": "evt-created",
"ts": "2026-03-27T12:00:00Z",
"run_id": run_id.to_string(),
"event": "run.created",
"node_id": 42,
"properties": {
"settings": WorkflowSettings::default(),
"graph": Graph::new("test"),
"run_dir": "/tmp/test",
"provenance": test_support::test_run_provenance(),
},
}),
&run_id,
)
.unwrap();
let err = run.append_event(&malformed).await.unwrap_err();
assert!(matches!(err, Error::InvalidEvent(_)));
assert!(run.list_events().await.unwrap().is_empty());
}
#[tokio::test]
async fn control_request_events_set_pending_control_without_overwriting_status() {
let (_object_store, store) = make_store();

View file

@ -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<RunId>,
parent_id: Option<RunId>,
) {
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<CachedRunProjection> {
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);
}

View file

@ -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<Arc<RunProjection>> {
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<Option<Arc<RunProjection>>> {
let next_seq = {
let cache = self.inner.projection_cache.lock().await;
cache.last_seq.saturating_add(1)
@ -202,55 +211,42 @@ 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 install_derived_state_after_append(
&self,
event: &EventEnvelope,
cached: CachedRunProjection,
) {
{
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(Arc::clone(&cached.projection));
projection_cache.last_seq = event.seq;
}
self.inner
.shared_projection_cache
.replace(cached.clone())
.await;
let mut recent_events = self.inner.recent_events.lock().await;
recent_events.push_back(event.clone());
while recent_events.len() > self.inner.recent_event_limit {
recent_events.pop_front();
}
drop(recent_events);
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
)));
if let Some(store) = self.inner.run_summary_store.get() {
if let Err(err) = store.upsert_projection(&cached).await {
warn!(
run_id = %self.inner.run_id,
source_last_seq = event.seq,
error = ?err,
"failed to update SQLite run summary after committed append"
);
}
}
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<Vec<EventEnvelope>> {
@ -270,12 +266,26 @@ impl RunDatabase {
}
impl RunDatabase {
/// Appends an event after validating it against the current run projection.
///
/// A rejected event writes nothing. Every returned error means the event
/// was not committed and is safe to retry. Once the SlateDB write succeeds,
/// the append returns success even if a derived cache or SQLite summary
/// update fails; those failures are logged and repaired by later updates or
/// startup reconciliation.
pub async fn append_event(&self, payload: &EventPayload) -> Result<u32> {
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. An invalid transition is also rejected before write, and every
/// returned error means the event was not committed and is safe to retry.
/// After the SlateDB write succeeds, derived cache and SQLite summary
/// updates are best-effort and cannot turn the committed append into an
/// error.
pub async fn append_event_if(
&self,
payload: &EventPayload,
@ -293,6 +303,12 @@ impl RunDatabase {
Ok(Some(self.append_event_envelope_locked(payload).await?.seq))
}
/// Appends and returns the stored event envelope after pre-write reduction.
///
/// A rejected event writes nothing. Every returned error means the event
/// was not committed and is safe to retry. Once the SlateDB write succeeds,
/// derived cache and SQLite summary updates are best-effort: failures are
/// logged, and this method still returns the committed envelope.
pub async fn append_event_envelope(&self, payload: &EventPayload) -> Result<EventEnvelope> {
if self.read_only {
return Err(Error::ReadOnly);
@ -309,73 +325,32 @@ impl RunDatabase {
seq,
event: RunEvent::try_from(payload)?,
};
let current_projection = self.projected_state_option_locked().await?;
let next_projection = match current_projection {
Some(projection) => {
let mut projection = (*projection).clone();
projection.apply_event(&event).map_err(event_rejected)?;
projection
}
None => {
RunProjection::apply_events(std::slice::from_ref(&event)).map_err(event_rejected)?
}
};
let cached = CachedRunProjection::from_projection(self.inner.run_id, next_projection, seq);
let event_bytes = serde_json::to_vec(payload)?;
self.inner
.db
.put(
keys::run_event_key(&self.inner.run_id, seq, Utc::now().timestamp_millis()),
serde_json::to_vec(payload)?,
event_bytes,
)
.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?;
// Box the derived-update future so this frequently awaited append API
// does not pass a large state machine into every caller.
Box::pin(self.install_derived_state_after_append(&event, cached)).await;
Ok(event)
}
async fn update_summary_projection_after_append(&self, event: &EventEnvelope) -> Result<()> {
let cached = match 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);
}
}
}
};
if let Some(store) = self.inner.run_summary_store.get() {
if let Err(err) = store.upsert_projection(&cached).await {
error!(
run_id = %self.inner.run_id,
source_last_seq = cached.last_seq,
error = %err,
"Failed to update SQLite run summary after append"
);
return Err(err);
}
}
Ok(())
}
pub async fn list_events(&self) -> Result<Vec<EventEnvelope>> {
self.list_events_from_with_limit(1, usize::MAX).await
}
@ -589,6 +564,14 @@ impl RunDatabase {
}
}
fn event_rejected(error: Error) -> Error {
let reason = match error {
Error::InvalidTransition(transition) => transition.to_string(),
error => error.to_string(),
};
Error::EventRejected { reason }
}
fn allocate_event_seq(event_seq: &AtomicU32) -> Result<u32> {
event_seq
.fetch_update(Ordering::SeqCst, Ordering::SeqCst, |seq| {
@ -1304,6 +1287,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 +1297,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

View file

@ -57,7 +57,7 @@ impl TryFrom<&EventPayload> for RunEvent {
type Error = Error;
fn try_from(value: &EventPayload) -> Result<Self> {
Self::from_ref(value.as_value())
Self::from_value(value.as_value().clone())
.map_err(|err| Error::InvalidEvent(format!("invalid stored event: {err}")))
}
}