Skip to main content

ironflow_store/memory/
log_store.rs

1//! [`LogStore`] trait implementation for [`InMemoryStore`].
2
3use chrono::Utc;
4use uuid::Uuid;
5
6use crate::entities::{LogEntry, LogFilter, NewLogEntries};
7use crate::log_store::LogStore;
8use crate::store::StoreFuture;
9
10use super::InMemoryStore;
11
12impl LogStore for InMemoryStore {
13    fn append_logs(&self, entries: NewLogEntries) -> StoreFuture<'_, ()> {
14        Box::pin(async move {
15            let now = Utc::now();
16            let mut state = self.state.write().await;
17
18            for (id, line) in entries.ids.into_iter().zip(entries.lines) {
19                state.log_entries.push(LogEntry {
20                    id,
21                    run_id: entries.run_id,
22                    step_id: entries.step_id,
23                    step_name: entries.step_name.clone(),
24                    stream: entries.stream,
25                    line,
26                    created_at: now,
27                });
28            }
29
30            Ok(())
31        })
32    }
33
34    fn get_logs(
35        &self,
36        run_id: Uuid,
37        filter: LogFilter,
38        cursor: Option<Uuid>,
39        limit: u32,
40    ) -> StoreFuture<'_, Vec<LogEntry>> {
41        Box::pin(async move {
42            let state = self.state.read().await;
43
44            let entries: Vec<LogEntry> = state
45                .log_entries
46                .iter()
47                .filter(|e| {
48                    if e.run_id != run_id {
49                        return false;
50                    }
51                    if let Some(step_id) = filter.step_id
52                        && e.step_id != step_id
53                    {
54                        return false;
55                    }
56                    if let Some(stream) = filter.stream
57                        && e.stream != stream
58                    {
59                        return false;
60                    }
61                    if let Some(cursor) = cursor
62                        && e.id <= cursor
63                    {
64                        return false;
65                    }
66                    true
67                })
68                .take(limit as usize)
69                .cloned()
70                .collect();
71
72            Ok(entries)
73        })
74    }
75}
76
77#[cfg(test)]
78mod tests {
79    use uuid::Uuid;
80
81    use crate::entities::{LogFilter, LogStream, NewLogEntries};
82    use crate::log_store::LogStore;
83    use crate::memory::InMemoryStore;
84
85    fn new_entries(run_id: Uuid, step_id: Uuid, stream: LogStream) -> NewLogEntries {
86        NewLogEntries {
87            ids: vec![Uuid::now_v7(), Uuid::now_v7()],
88            run_id,
89            step_id,
90            step_name: "build".to_string(),
91            stream,
92            lines: vec!["line 1".to_string(), "line 2".to_string()],
93        }
94    }
95
96    #[tokio::test]
97    async fn append_and_get_golden_path() {
98        let store = InMemoryStore::new();
99        let run_id = Uuid::now_v7();
100        let step_id = Uuid::now_v7();
101
102        store
103            .append_logs(new_entries(run_id, step_id, LogStream::Stdout))
104            .await
105            .unwrap();
106
107        let logs = store
108            .get_logs(run_id, LogFilter::default(), None, 100)
109            .await
110            .unwrap();
111
112        assert_eq!(logs.len(), 2);
113        assert_eq!(logs[0].line, "line 1");
114        assert_eq!(logs[1].line, "line 2");
115        assert_eq!(logs[0].run_id, run_id);
116        assert_eq!(logs[0].step_id, step_id);
117        assert_eq!(logs[0].step_name, "build");
118        assert_eq!(logs[0].stream, LogStream::Stdout);
119    }
120
121    #[tokio::test]
122    async fn get_empty_returns_empty_vec() {
123        let store = InMemoryStore::new();
124        let run_id = Uuid::now_v7();
125
126        let logs = store
127            .get_logs(run_id, LogFilter::default(), None, 100)
128            .await
129            .unwrap();
130
131        assert!(logs.is_empty());
132    }
133
134    #[tokio::test]
135    async fn cursor_based_pagination() {
136        let store = InMemoryStore::new();
137        let run_id = Uuid::now_v7();
138        let step_id = Uuid::now_v7();
139
140        store
141            .append_logs(NewLogEntries {
142                ids: (0..5).map(|_| Uuid::now_v7()).collect(),
143                run_id,
144                step_id,
145                step_name: "build".to_string(),
146                stream: LogStream::Stdout,
147                lines: (0..5).map(|i| format!("line {i}")).collect(),
148            })
149            .await
150            .unwrap();
151
152        let page1 = store
153            .get_logs(run_id, LogFilter::default(), None, 2)
154            .await
155            .unwrap();
156        assert_eq!(page1.len(), 2);
157        assert_eq!(page1[0].line, "line 0");
158        assert_eq!(page1[1].line, "line 1");
159
160        let cursor = page1.last().unwrap().id;
161        let page2 = store
162            .get_logs(run_id, LogFilter::default(), Some(cursor), 2)
163            .await
164            .unwrap();
165        assert_eq!(page2.len(), 2);
166        assert_eq!(page2[0].line, "line 2");
167        assert_eq!(page2[1].line, "line 3");
168
169        let cursor = page2.last().unwrap().id;
170        let page3 = store
171            .get_logs(run_id, LogFilter::default(), Some(cursor), 2)
172            .await
173            .unwrap();
174        assert_eq!(page3.len(), 1);
175        assert_eq!(page3[0].line, "line 4");
176    }
177
178    #[tokio::test]
179    async fn filter_by_step_id() {
180        let store = InMemoryStore::new();
181        let run_id = Uuid::now_v7();
182        let step_a = Uuid::now_v7();
183        let step_b = Uuid::now_v7();
184
185        store
186            .append_logs(new_entries(run_id, step_a, LogStream::Stdout))
187            .await
188            .unwrap();
189        store
190            .append_logs(new_entries(run_id, step_b, LogStream::Stdout))
191            .await
192            .unwrap();
193
194        let filter = LogFilter {
195            step_id: Some(step_a),
196            ..LogFilter::default()
197        };
198        let logs = store.get_logs(run_id, filter, None, 100).await.unwrap();
199
200        assert_eq!(logs.len(), 2);
201        assert!(logs.iter().all(|e| e.step_id == step_a));
202    }
203
204    #[tokio::test]
205    async fn filter_by_stream() {
206        let store = InMemoryStore::new();
207        let run_id = Uuid::now_v7();
208        let step_id = Uuid::now_v7();
209
210        store
211            .append_logs(new_entries(run_id, step_id, LogStream::Stdout))
212            .await
213            .unwrap();
214        store
215            .append_logs(new_entries(run_id, step_id, LogStream::Stderr))
216            .await
217            .unwrap();
218
219        let filter = LogFilter {
220            stream: Some(LogStream::Stderr),
221            ..LogFilter::default()
222        };
223        let logs = store.get_logs(run_id, filter, None, 100).await.unwrap();
224
225        assert_eq!(logs.len(), 2);
226        assert!(logs.iter().all(|e| e.stream == LogStream::Stderr));
227    }
228
229    #[tokio::test]
230    async fn filter_by_step_id_and_stream() {
231        let store = InMemoryStore::new();
232        let run_id = Uuid::now_v7();
233        let step_a = Uuid::now_v7();
234        let step_b = Uuid::now_v7();
235
236        store
237            .append_logs(new_entries(run_id, step_a, LogStream::Stdout))
238            .await
239            .unwrap();
240        store
241            .append_logs(new_entries(run_id, step_a, LogStream::Stderr))
242            .await
243            .unwrap();
244        store
245            .append_logs(new_entries(run_id, step_b, LogStream::Stdout))
246            .await
247            .unwrap();
248
249        let filter = LogFilter {
250            step_id: Some(step_a),
251            stream: Some(LogStream::Stderr),
252        };
253        let logs = store.get_logs(run_id, filter, None, 100).await.unwrap();
254
255        assert_eq!(logs.len(), 2);
256        assert!(logs.iter().all(|e| e.step_id == step_a));
257        assert!(logs.iter().all(|e| e.stream == LogStream::Stderr));
258    }
259
260    #[tokio::test]
261    async fn different_runs_are_isolated() {
262        let store = InMemoryStore::new();
263        let run_a = Uuid::now_v7();
264        let run_b = Uuid::now_v7();
265        let step_id = Uuid::now_v7();
266
267        store
268            .append_logs(new_entries(run_a, step_id, LogStream::Stdout))
269            .await
270            .unwrap();
271        store
272            .append_logs(new_entries(run_b, step_id, LogStream::Stdout))
273            .await
274            .unwrap();
275
276        let logs_a = store
277            .get_logs(run_a, LogFilter::default(), None, 100)
278            .await
279            .unwrap();
280        assert_eq!(logs_a.len(), 2);
281        assert!(logs_a.iter().all(|e| e.run_id == run_a));
282
283        let logs_b = store
284            .get_logs(run_b, LogFilter::default(), None, 100)
285            .await
286            .unwrap();
287        assert_eq!(logs_b.len(), 2);
288        assert!(logs_b.iter().all(|e| e.run_id == run_b));
289    }
290
291    #[tokio::test]
292    async fn entries_have_unique_ids() {
293        let store = InMemoryStore::new();
294        let run_id = Uuid::now_v7();
295        let step_id = Uuid::now_v7();
296
297        store
298            .append_logs(new_entries(run_id, step_id, LogStream::Stdout))
299            .await
300            .unwrap();
301
302        let logs = store
303            .get_logs(run_id, LogFilter::default(), None, 100)
304            .await
305            .unwrap();
306
307        assert_ne!(logs[0].id, logs[1].id);
308    }
309}