1use 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#[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#[derive(Debug, Error)]
48pub enum RuntimeEventStoreError {
49 #[error("runtime event store append failed: {message}")]
51 Append {
52 message: String,
54 },
55 #[error("runtime event store has no events for run {run_id}")]
57 NotFound {
58 run_id: RunId,
60 },
61}
62
63#[async_trait]
68pub trait RuntimeEventStore: Send + Sync {
69 async fn append(
77 &self,
78 event: AgentEvent,
79 ) -> Result<RuntimeEventEnvelope, RuntimeEventStoreError>;
80
81 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#[derive(Debug, Default)]
101pub struct MemoryRuntimeEventStore {
102 state: Mutex<StoreState>,
103}
104
105impl MemoryRuntimeEventStore {
106 #[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#[derive(Debug, Default, Clone, Copy)]
181pub struct FailingRuntimeEventStore;
182
183impl FailingRuntimeEventStore {
184 #[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
212pub 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}