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 line in entries.lines {
19                state.log_entries.push(LogEntry {
20                    id: Uuid::now_v7(),
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            run_id,
88            step_id,
89            step_name: "build".to_string(),
90            stream,
91            lines: vec!["line 1".to_string(), "line 2".to_string()],
92        }
93    }
94
95    #[tokio::test]
96    async fn append_and_get_golden_path() {
97        let store = InMemoryStore::new();
98        let run_id = Uuid::now_v7();
99        let step_id = Uuid::now_v7();
100
101        store
102            .append_logs(new_entries(run_id, step_id, LogStream::Stdout))
103            .await
104            .unwrap();
105
106        let logs = store
107            .get_logs(run_id, LogFilter::default(), None, 100)
108            .await
109            .unwrap();
110
111        assert_eq!(logs.len(), 2);
112        assert_eq!(logs[0].line, "line 1");
113        assert_eq!(logs[1].line, "line 2");
114        assert_eq!(logs[0].run_id, run_id);
115        assert_eq!(logs[0].step_id, step_id);
116        assert_eq!(logs[0].step_name, "build");
117        assert_eq!(logs[0].stream, LogStream::Stdout);
118    }
119
120    #[tokio::test]
121    async fn get_empty_returns_empty_vec() {
122        let store = InMemoryStore::new();
123        let run_id = Uuid::now_v7();
124
125        let logs = store
126            .get_logs(run_id, LogFilter::default(), None, 100)
127            .await
128            .unwrap();
129
130        assert!(logs.is_empty());
131    }
132
133    #[tokio::test]
134    async fn cursor_based_pagination() {
135        let store = InMemoryStore::new();
136        let run_id = Uuid::now_v7();
137        let step_id = Uuid::now_v7();
138
139        store
140            .append_logs(NewLogEntries {
141                run_id,
142                step_id,
143                step_name: "build".to_string(),
144                stream: LogStream::Stdout,
145                lines: (0..5).map(|i| format!("line {i}")).collect(),
146            })
147            .await
148            .unwrap();
149
150        let page1 = store
151            .get_logs(run_id, LogFilter::default(), None, 2)
152            .await
153            .unwrap();
154        assert_eq!(page1.len(), 2);
155        assert_eq!(page1[0].line, "line 0");
156        assert_eq!(page1[1].line, "line 1");
157
158        let cursor = page1.last().unwrap().id;
159        let page2 = store
160            .get_logs(run_id, LogFilter::default(), Some(cursor), 2)
161            .await
162            .unwrap();
163        assert_eq!(page2.len(), 2);
164        assert_eq!(page2[0].line, "line 2");
165        assert_eq!(page2[1].line, "line 3");
166
167        let cursor = page2.last().unwrap().id;
168        let page3 = store
169            .get_logs(run_id, LogFilter::default(), Some(cursor), 2)
170            .await
171            .unwrap();
172        assert_eq!(page3.len(), 1);
173        assert_eq!(page3[0].line, "line 4");
174    }
175
176    #[tokio::test]
177    async fn filter_by_step_id() {
178        let store = InMemoryStore::new();
179        let run_id = Uuid::now_v7();
180        let step_a = Uuid::now_v7();
181        let step_b = Uuid::now_v7();
182
183        store
184            .append_logs(new_entries(run_id, step_a, LogStream::Stdout))
185            .await
186            .unwrap();
187        store
188            .append_logs(new_entries(run_id, step_b, LogStream::Stdout))
189            .await
190            .unwrap();
191
192        let filter = LogFilter {
193            step_id: Some(step_a),
194            ..LogFilter::default()
195        };
196        let logs = store.get_logs(run_id, filter, None, 100).await.unwrap();
197
198        assert_eq!(logs.len(), 2);
199        assert!(logs.iter().all(|e| e.step_id == step_a));
200    }
201
202    #[tokio::test]
203    async fn filter_by_stream() {
204        let store = InMemoryStore::new();
205        let run_id = Uuid::now_v7();
206        let step_id = Uuid::now_v7();
207
208        store
209            .append_logs(new_entries(run_id, step_id, LogStream::Stdout))
210            .await
211            .unwrap();
212        store
213            .append_logs(new_entries(run_id, step_id, LogStream::Stderr))
214            .await
215            .unwrap();
216
217        let filter = LogFilter {
218            stream: Some(LogStream::Stderr),
219            ..LogFilter::default()
220        };
221        let logs = store.get_logs(run_id, filter, None, 100).await.unwrap();
222
223        assert_eq!(logs.len(), 2);
224        assert!(logs.iter().all(|e| e.stream == LogStream::Stderr));
225    }
226
227    #[tokio::test]
228    async fn filter_by_step_id_and_stream() {
229        let store = InMemoryStore::new();
230        let run_id = Uuid::now_v7();
231        let step_a = Uuid::now_v7();
232        let step_b = Uuid::now_v7();
233
234        store
235            .append_logs(new_entries(run_id, step_a, LogStream::Stdout))
236            .await
237            .unwrap();
238        store
239            .append_logs(new_entries(run_id, step_a, LogStream::Stderr))
240            .await
241            .unwrap();
242        store
243            .append_logs(new_entries(run_id, step_b, LogStream::Stdout))
244            .await
245            .unwrap();
246
247        let filter = LogFilter {
248            step_id: Some(step_a),
249            stream: Some(LogStream::Stderr),
250        };
251        let logs = store.get_logs(run_id, filter, None, 100).await.unwrap();
252
253        assert_eq!(logs.len(), 2);
254        assert!(logs.iter().all(|e| e.step_id == step_a));
255        assert!(logs.iter().all(|e| e.stream == LogStream::Stderr));
256    }
257
258    #[tokio::test]
259    async fn different_runs_are_isolated() {
260        let store = InMemoryStore::new();
261        let run_a = Uuid::now_v7();
262        let run_b = Uuid::now_v7();
263        let step_id = Uuid::now_v7();
264
265        store
266            .append_logs(new_entries(run_a, step_id, LogStream::Stdout))
267            .await
268            .unwrap();
269        store
270            .append_logs(new_entries(run_b, step_id, LogStream::Stdout))
271            .await
272            .unwrap();
273
274        let logs_a = store
275            .get_logs(run_a, LogFilter::default(), None, 100)
276            .await
277            .unwrap();
278        assert_eq!(logs_a.len(), 2);
279        assert!(logs_a.iter().all(|e| e.run_id == run_a));
280
281        let logs_b = store
282            .get_logs(run_b, LogFilter::default(), None, 100)
283            .await
284            .unwrap();
285        assert_eq!(logs_b.len(), 2);
286        assert!(logs_b.iter().all(|e| e.run_id == run_b));
287    }
288
289    #[tokio::test]
290    async fn entries_have_unique_ids() {
291        let store = InMemoryStore::new();
292        let run_id = Uuid::now_v7();
293        let step_id = Uuid::now_v7();
294
295        store
296            .append_logs(new_entries(run_id, step_id, LogStream::Stdout))
297            .await
298            .unwrap();
299
300        let logs = store
301            .get_logs(run_id, LogFilter::default(), None, 100)
302            .await
303            .unwrap();
304
305        assert_ne!(logs[0].id, logs[1].id);
306    }
307}