mirror of
https://github.com/fabro-sh/fabro.git
synced 2026-10-09 03:20:56 +00:00
fabro(01KYNBXZ4PAGMNZGHVHPNAQ341): implement (succeeded)
Fabro-Run: 01KYNBXZ4PAGMNZGHVHPNAQ341
Fabro-Completed: 5
Fabro-Checkpoint: 85b5e079f4
⚒️ Generated with [Fabro](https://fabro.sh)
This commit is contained in:
parent
12cdaee5da
commit
e35c5f726f
4 changed files with 418 additions and 161 deletions
|
|
@ -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<Self>,
|
||||
},
|
||||
#[error("Run not found: {0}")]
|
||||
RunNotFound(String),
|
||||
#[error("Run already exists: {0}")]
|
||||
|
|
|
|||
|
|
@ -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?;
|
||||
|
|
|
|||
|
|
@ -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);
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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,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<RunProjection>) {
|
||||
{
|
||||
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<Vec<EventEnvelope>> {
|
||||
|
|
@ -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<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. 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<EventEnvelope> {
|
||||
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<EventEnvelope> {
|
||||
async fn append_event_envelope_locked(
|
||||
&self,
|
||||
payload: &EventPayload,
|
||||
current_projection: Option<Arc<RunProjection>>,
|
||||
) -> Result<EventEnvelope> {
|
||||
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<Vec<EventEnvelope>> {
|
||||
|
|
@ -599,6 +603,27 @@ fn allocate_event_seq(event_seq: &AtomicU32) -> Result<u32> {
|
|||
})
|
||||
}
|
||||
|
||||
fn validate_first_event(event: &EventEnvelope) -> Result<Arc<RunProjection>> {
|
||||
RunProjection::apply_events(std::slice::from_ref(event))
|
||||
.map(Arc::new)
|
||||
.map_err(event_rejected)
|
||||
}
|
||||
|
||||
fn validate_projected_event(
|
||||
current: &RunProjection,
|
||||
event: &EventEnvelope,
|
||||
) -> Result<Arc<RunProjection>> {
|
||||
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<Arc<RunProjection>>,
|
||||
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
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue