1use 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}