mirror of
https://github.com/fabro-sh/fabro.git
synced 2026-10-10 03:30:59 +00:00
fix(store): enforce event sequence key limit
This commit is contained in:
parent
0e4244a24a
commit
4c7d13aff0
3 changed files with 54 additions and 5 deletions
|
|
@ -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}")]
|
||||
|
|
|
|||
|
|
@ -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(
|
||||
|
|
|
|||
|
|
@ -292,7 +292,7 @@ impl RunDatabase {
|
|||
|
||||
async fn append_event_envelope_locked(&self, payload: &EventPayload) -> Result<EventEnvelope> {
|
||||
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<u32> {
|
||||
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<RunProjection>,
|
||||
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;
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue