Skip to main content

behest_runtime/
event_store.rs

1//! Reliable runtime event log.
2//!
3//! [`RuntimeEventStore`] is the **authoritative replay source** for runtime
4//! events. Unlike [`RuntimeStreamAdapter`](super::stream_adapter::RuntimeStreamAdapter),
5//! which only performs best-effort live fanout, the store guarantees that any
6//! event accepted by [`RuntimeEventStore::append`] can be replayed later via
7//! [`RuntimeEventStore::list_after`].
8//!
9//! Delivery semantics are at-least-once: a consumer reconnecting with
10//! `run_id + after_seq` may receive duplicates of events it already observed
11//! live; deduplicate via [`RuntimeEventEnvelope::event_id`](super::stream::RuntimeEventEnvelope::event_id)
12//! or `seq`.
13
14use std::collections::HashMap;
15use std::sync::Arc;
16
17use async_trait::async_trait;
18use chrono::Utc;
19use thiserror::Error;
20use tokio::sync::Mutex;
21
22use super::event::AgentEvent;
23use super::run::RunId;
24use super::stream::{RuntimeEventEnvelope, RuntimeEventId};
25
26#[cfg(feature = "redis")]
27#[path = "event_store/redis.rs"]
28pub mod redis;
29
30#[cfg(feature = "sqlx-postgres")]
31#[path = "event_store/postgres.rs"]
32pub mod postgres;
33
34/// Single locked state for [`MemoryRuntimeEventStore`].
35///
36/// Merging `events`, `seq`, and `sessions` under one [`Mutex`] guarantees that
37/// sequence assignment, session propagation, and event insertion are atomic
38/// per append — no interleaved racing between the three maps.
39#[derive(Debug, Default)]
40struct StoreState {
41    events: HashMap<RunId, Vec<RuntimeEventEnvelope>>,
42    seq: HashMap<RunId, u64>,
43    sessions: HashMap<RunId, Option<uuid::Uuid>>,
44}
45
46/// Errors raised by a [`RuntimeEventStore`].
47#[derive(Debug, Error)]
48pub enum RuntimeEventStoreError {
49    /// An append could not be persisted.
50    #[error("runtime event store append failed: {message}")]
51    Append {
52        /// Human-readable diagnostic.
53        message: String,
54    },
55    /// The requested run has no recorded events.
56    #[error("runtime event store has no events for run {run_id}")]
57    NotFound {
58        /// Run that was queried.
59        run_id: RunId,
60    },
61}
62
63/// Authoritative replay source for runtime events.
64///
65/// Implementations are responsible for minting [`RuntimeEventId`] and the
66/// per-run `seq` counter on [`RuntimeEventStore::append`].
67#[async_trait]
68pub trait RuntimeEventStore: Send + Sync {
69    /// Appends an event and returns the resulting envelope with identity and
70    /// sequence assigned.
71    ///
72    /// On failure the event MUST NOT be considered persisted; callers (such as
73    /// [`RuntimeEventBridge`](super::subscription::RuntimeEventBridge)) rely on
74    /// this contract to avoid publishing live events whose replay source is
75    /// incomplete.
76    async fn append(
77        &self,
78        event: AgentEvent,
79    ) -> Result<RuntimeEventEnvelope, RuntimeEventStoreError>;
80
81    /// Replays events for `run_id` with `seq > after_seq`.
82    ///
83    /// `after_seq = None` replays from the beginning. `limit` caps the page
84    /// size to avoid unbounded memory use.
85    async fn list_after(
86        &self,
87        run_id: RunId,
88        after_seq: Option<u64>,
89        limit: usize,
90    ) -> Result<Vec<RuntimeEventEnvelope>, RuntimeEventStoreError>;
91}
92
93/// In-memory [`RuntimeEventStore`] for tests and single-instance development.
94///
95/// `seq` is monotonic per `run_id`. When a [`AgentEvent::RunStarted`] is
96/// appended, its `session_id` is cached and attached to subsequent events of
97/// the same run. All state is guarded by a single [`Mutex`] so sequence
98/// assignment, session propagation, and event insertion are atomic per
99/// [`RuntimeEventStore::append`].
100#[derive(Debug, Default)]
101pub struct MemoryRuntimeEventStore {
102    state: Mutex<StoreState>,
103}
104
105impl MemoryRuntimeEventStore {
106    /// Creates an empty store.
107    #[must_use]
108    pub fn new() -> Self {
109        Self::default()
110    }
111}
112
113#[async_trait]
114impl RuntimeEventStore for MemoryRuntimeEventStore {
115    async fn append(
116        &self,
117        event: AgentEvent,
118    ) -> Result<RuntimeEventEnvelope, RuntimeEventStoreError> {
119        let run_id = event.run_id();
120        let mut state = self.state.lock().await;
121
122        let session_id = if let AgentEvent::RunStarted(started) = &event {
123            state.sessions.insert(run_id, Some(started.session_id));
124            Some(started.session_id)
125        } else {
126            state.sessions.get(&run_id).copied().flatten()
127        };
128
129        let next_seq = {
130            let entry = state.seq.entry(run_id).or_default();
131            *entry += 1;
132            *entry
133        };
134
135        let envelope = RuntimeEventEnvelope {
136            event_id: RuntimeEventId::new(),
137            seq: next_seq,
138            run_id,
139            session_id,
140            event,
141            emitted_at: Utc::now(),
142        };
143
144        state
145            .events
146            .entry(run_id)
147            .or_default()
148            .push(envelope.clone());
149
150        Ok(envelope)
151    }
152
153    async fn list_after(
154        &self,
155        run_id: RunId,
156        after_seq: Option<u64>,
157        limit: usize,
158    ) -> Result<Vec<RuntimeEventEnvelope>, RuntimeEventStoreError> {
159        let state = self.state.lock().await;
160        let Some(run_events) = state.events.get(&run_id) else {
161            return Ok(Vec::new());
162        };
163
164        let filtered: Vec<RuntimeEventEnvelope> = run_events
165            .iter()
166            .filter(|env| match after_seq {
167                None => true,
168                Some(seq) => env.seq > seq,
169            })
170            .take(limit)
171            .cloned()
172            .collect();
173
174        Ok(filtered)
175    }
176}
177
178/// [`RuntimeEventStore`] that always fails. Used by tests that assert a failed
179/// append does not propagate to the live adapter.
180#[derive(Debug, Default, Clone, Copy)]
181pub struct FailingRuntimeEventStore;
182
183impl FailingRuntimeEventStore {
184    /// Creates a new failing store.
185    #[must_use]
186    pub fn new() -> Self {
187        Self
188    }
189}
190
191#[async_trait]
192impl RuntimeEventStore for FailingRuntimeEventStore {
193    async fn append(
194        &self,
195        _event: AgentEvent,
196    ) -> Result<RuntimeEventEnvelope, RuntimeEventStoreError> {
197        Err(RuntimeEventStoreError::Append {
198            message: "failing runtime event store always rejects appends".to_owned(),
199        })
200    }
201
202    async fn list_after(
203        &self,
204        run_id: RunId,
205        _after_seq: Option<u64>,
206        _limit: usize,
207    ) -> Result<Vec<RuntimeEventEnvelope>, RuntimeEventStoreError> {
208        Err(RuntimeEventStoreError::NotFound { run_id })
209    }
210}
211
212/// Convenience alias for shared, trait-object event stores.
213pub type DynRuntimeEventStore = Arc<dyn RuntimeEventStore>;
214
215#[cfg(test)]
216mod tests {
217    #![allow(clippy::unwrap_used, clippy::expect_used)]
218
219    use chrono::Utc;
220    use uuid::Uuid;
221
222    use super::*;
223    use crate::event::{RunCancelled, RunCompleted, RunFailed, RunStarted};
224    use behest_provider::{ModelName, ProviderId};
225
226    fn started(run_id: RunId, session_id: Uuid) -> AgentEvent {
227        AgentEvent::RunStarted(RunStarted {
228            run_id,
229            session_id,
230            provider: ProviderId::new("acme"),
231            model: ModelName::new("gpt-test"),
232            timestamp: Utc::now(),
233        })
234    }
235
236    fn terminal(run_id: RunId) -> AgentEvent {
237        AgentEvent::RunCompleted(RunCompleted {
238            run_id,
239            finish_reason: behest_provider::FinishReason::Stop,
240            iterations: 1,
241            timestamp: Utc::now(),
242        })
243    }
244
245    fn failed(run_id: RunId) -> AgentEvent {
246        AgentEvent::RunFailed(RunFailed {
247            run_id,
248            error: "boom".to_owned(),
249            timestamp: Utc::now(),
250        })
251    }
252
253    fn cancelled(run_id: RunId) -> AgentEvent {
254        AgentEvent::RunCancelled(RunCancelled {
255            run_id,
256            timestamp: Utc::now(),
257        })
258    }
259
260    #[tokio::test]
261    async fn append_assigns_monotonic_seq_per_run() {
262        let store = MemoryRuntimeEventStore::new();
263        let run = RunId::new();
264        let sid = Uuid::now_v7();
265
266        let e1 = store.append(started(run, sid)).await.unwrap();
267        let e2 = store.append(terminal(run)).await.unwrap();
268        let e3 = store.append(failed(run)).await.unwrap();
269
270        assert_eq!(e1.seq, 1);
271        assert_eq!(e2.seq, 2);
272        assert_eq!(e3.seq, 3);
273    }
274
275    #[tokio::test]
276    async fn append_propagates_session_id_from_run_started() {
277        let store = MemoryRuntimeEventStore::new();
278        let run = RunId::new();
279        let sid = Uuid::now_v7();
280
281        let started_env = store.append(started(run, sid)).await.unwrap();
282        assert_eq!(started_env.session_id, Some(sid));
283
284        let terminal_env = store.append(terminal(run)).await.unwrap();
285        assert_eq!(terminal_env.session_id, Some(sid));
286    }
287
288    #[tokio::test]
289    async fn list_after_filters_by_seq() {
290        let store = MemoryRuntimeEventStore::new();
291        let run = RunId::new();
292        let sid = Uuid::now_v7();
293
294        store.append(started(run, sid)).await.unwrap();
295        let e2 = store.append(terminal(run)).await.unwrap();
296        let e3 = store.append(failed(run)).await.unwrap();
297
298        let page = store.list_after(run, Some(e2.seq), 10).await.unwrap();
299        assert_eq!(page.len(), 1);
300        assert_eq!(page[0].seq, e3.seq);
301    }
302
303    #[tokio::test]
304    async fn list_after_respects_limit() {
305        let store = MemoryRuntimeEventStore::new();
306        let run = RunId::new();
307        let sid = Uuid::now_v7();
308
309        store.append(started(run, sid)).await.unwrap();
310        store.append(terminal(run)).await.unwrap();
311        store.append(failed(run)).await.unwrap();
312
313        let page = store.list_after(run, None, 2).await.unwrap();
314        assert_eq!(page.len(), 2);
315    }
316
317    #[tokio::test]
318    async fn list_after_unknown_run_returns_empty() {
319        let store = MemoryRuntimeEventStore::new();
320        let run = RunId::new();
321        let page = store.list_after(run, None, 10).await.unwrap();
322        assert!(page.is_empty());
323    }
324
325    #[tokio::test]
326    async fn envelope_is_terminal_recognizes_terminal_variants() {
327        let store = MemoryRuntimeEventStore::new();
328        let run = RunId::new();
329        let sid = Uuid::now_v7();
330
331        let non_terminal = store.append(started(run, sid)).await.unwrap();
332        assert!(!non_terminal.is_terminal());
333
334        let completed = store.append(terminal(run)).await.unwrap();
335        let failed_env = store.append(failed(run)).await.unwrap();
336        let cancelled_env = store.append(cancelled(run)).await.unwrap();
337
338        assert!(completed.is_terminal());
339        assert!(failed_env.is_terminal());
340        assert!(cancelled_env.is_terminal());
341    }
342}