From 4c7d13aff0f8188bc1d967880ff5541c5de39fbc Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Fri, 24 Jul 2026 09:18:43 -0400 Subject: [PATCH] fix(store): enforce event sequence key limit --- lib/components/fabro-store/src/error.rs | 2 + lib/components/fabro-store/src/keys.rs | 9 ++-- .../fabro-store/src/slate/run_store.rs | 48 ++++++++++++++++++- 3 files changed, 54 insertions(+), 5 deletions(-) diff --git a/lib/components/fabro-store/src/error.rs b/lib/components/fabro-store/src/error.rs index 33126d474..43c26b44a 100644 --- a/lib/components/fabro-store/src/error.rs +++ b/lib/components/fabro-store/src/error.rs @@ -24,6 +24,8 @@ pub enum Error { SessionAlreadyExists(String), #[error("run store is read-only")] ReadOnly, + #[error("event sequence limit of {max_seq} reached")] + EventSequenceExhausted { max_seq: u32 }, #[error("invalid key segment: {segment:?}")] InvalidKeySegment { segment: String }, #[error("failed to parse key: {0}")] diff --git a/lib/components/fabro-store/src/keys.rs b/lib/components/fabro-store/src/keys.rs index f63bb27b7..343cdc1d0 100644 --- a/lib/components/fabro-store/src/keys.rs +++ b/lib/components/fabro-store/src/keys.rs @@ -3,6 +3,8 @@ use std::ops::Range; use fabro_types::{RunBlobId, RunId, SessionId}; +pub(crate) const MAX_EVENT_SEQ: u32 = 999_999; + #[derive(Debug, PartialEq, Eq)] pub(crate) struct SlateKey(String); @@ -61,8 +63,9 @@ pub(crate) fn run_events_prefix(run_id: &RunId) -> SlateKey { } // Sequence keys zero-pad `seq` to six digits so lexicographic key order -// matches numeric seq order for up to 999,999 events per run. Seek-based -// event listing (`run_events_range`) depends on this invariant. +// matches numeric seq order through `MAX_EVENT_SEQ`. Seek-based event listing +// (`run_events_range`) depends on this invariant, so event allocation rejects +// larger sequences. pub(crate) fn run_event_key(run_id: &RunId, seq: u32, epoch_ms: i64) -> SlateKey { SlateKey::new("runs") .with(run_id) @@ -184,7 +187,7 @@ mod tests { assert!(!contains(&run_event_key(&run_id, 1, 123))); assert!(contains(&run_event_key(&run_id, 2, 123))); - assert!(contains(&run_event_key(&run_id, 999_999, 123))); + assert!(contains(&run_event_key(&run_id, MAX_EVENT_SEQ, 123))); // Sibling namespaces of the same run sort outside the range. assert!(!contains(&SlateKey::new("runs").with(run_id).with("state"))); assert!(!contains( diff --git a/lib/components/fabro-store/src/slate/run_store.rs b/lib/components/fabro-store/src/slate/run_store.rs index 664ce3162..7efa68aba 100644 --- a/lib/components/fabro-store/src/slate/run_store.rs +++ b/lib/components/fabro-store/src/slate/run_store.rs @@ -292,7 +292,7 @@ impl RunDatabase { async fn append_event_envelope_locked(&self, payload: &EventPayload) -> Result { let event_seq = self.inner.event_seq.as_ref().ok_or(Error::ReadOnly)?; - let seq = event_seq.fetch_add(1, Ordering::SeqCst); + let seq = allocate_event_seq(event_seq)?; let event = EventEnvelope { seq, event: RunEvent::try_from(payload)?, @@ -507,6 +507,16 @@ impl RunDatabase { } } +fn allocate_event_seq(event_seq: &AtomicU32) -> Result { + event_seq + .fetch_update(Ordering::SeqCst, Ordering::SeqCst, |seq| { + (seq <= keys::MAX_EVENT_SEQ).then_some(seq + 1) + }) + .map_err(|_| Error::EventSequenceExhausted { + max_seq: keys::MAX_EVENT_SEQ, + }) +} + fn apply_cached_projection_event( state: &mut Option, event: &EventEnvelope, @@ -736,13 +746,14 @@ fn key_to_str(key: &Bytes) -> Result<&str> { #[cfg(test)] mod tests { use std::sync::Arc; + use std::sync::atomic::Ordering; use std::time::Duration; use fabro_types::{Graph, RunId, SessionId, StageId, WorkflowSettings, test_support}; use object_store::memory::InMemory; use serde_json::json; - use crate::{Database, EventPayload, keys}; + use crate::{Database, Error, EventPayload, keys}; #[tokio::test] async fn list_blobs_reads_global_cas_namespace() { @@ -894,6 +905,39 @@ mod tests { assert_eq!(seqs, vec![3]); } + #[tokio::test] + async fn append_event_rejects_sequences_beyond_key_order_limit() { + let run = fresh_run().await; + let run_id = run.run_id(); + run.inner + .event_seq + .as_ref() + .unwrap() + .store(keys::MAX_EVENT_SEQ, Ordering::SeqCst); + + let seq = run + .append_event(&stage_prompt_payload(&run_id, 1, Some("alpha"))) + .await + .unwrap(); + assert_eq!(seq, keys::MAX_EVENT_SEQ); + + let err = run + .append_event(&stage_prompt_payload(&run_id, 2, Some("beta"))) + .await + .unwrap_err(); + assert!(matches!( + err, + Error::EventSequenceExhausted { max_seq } + if max_seq == keys::MAX_EVENT_SEQ + )); + assert!( + run.get_event(keys::MAX_EVENT_SEQ + 1) + .await + .unwrap() + .is_none() + ); + } + #[tokio::test] async fn list_events_for_stage_returns_only_matching_events_in_seq_order() { let run = fresh_run().await;