1use async_trait::async_trait;
2use chrono::{DateTime, Utc};
3use std::path::{Path, PathBuf};
4use std::sync::Arc;
5use tokio::fs::{File, OpenOptions};
6use tokio::io::{AsyncBufReadExt, AsyncWriteExt, BufReader};
7use tokio::sync::Mutex;
8use uuid::Uuid;
9
10use crate::error::{FlowError, Result};
11use crate::model::{project_run, FlowEvent, FlowEventEnvelope};
12
13use super::FlowEventStore;
14
15#[derive(Debug, Clone)]
25pub struct LocalFileEventStore {
26 root: PathBuf,
27 lock: Arc<Mutex<()>>,
28}
29
30#[derive(Debug, Clone, Copy, Eq, PartialEq)]
31enum TailRepair {
32 None,
33 AppendDelimiter,
34 Truncate(u64),
35}
36
37#[derive(Debug)]
38struct LoadedEventLog {
39 events: Vec<FlowEventEnvelope>,
40 tail_repair: TailRepair,
41}
42
43impl LocalFileEventStore {
44 pub fn new(root: impl Into<PathBuf>) -> Self {
45 Self {
46 root: root.into(),
47 lock: Arc::new(Mutex::new(())),
48 }
49 }
50
51 pub fn root(&self) -> &Path {
52 &self.root
53 }
54
55 fn run_path(&self, run_id: &str) -> Result<PathBuf> {
56 if !is_safe_run_id(run_id) {
57 return Err(FlowError::Store(format!(
58 "run id {run_id:?} is not safe for local file storage"
59 )));
60 }
61 Ok(self.root.join(format!("{run_id}.jsonl")))
62 }
63
64 async fn load_inner(&self, run_id: &str, missing_is_empty: bool) -> Result<LoadedEventLog> {
65 let path = self.run_path(run_id)?;
66 let file = match File::open(&path).await {
67 Ok(file) => file,
68 Err(err) if err.kind() == std::io::ErrorKind::NotFound && missing_is_empty => {
69 return Ok(LoadedEventLog {
70 events: Vec::new(),
71 tail_repair: TailRepair::None,
72 });
73 }
74 Err(err) if err.kind() == std::io::ErrorKind::NotFound => {
75 return Err(FlowError::RunNotFound(run_id.to_string()));
76 }
77 Err(err) => return Err(FlowError::Io(err)),
78 };
79
80 let mut reader = BufReader::new(file);
81 let mut events = Vec::new();
82 let mut line_no = 0usize;
83 let mut valid_prefix_len = 0u64;
84 let mut buffer = Vec::new();
85 loop {
86 buffer.clear();
87 let bytes_read = reader.read_until(b'\n', &mut buffer).await?;
88 if bytes_read == 0 {
89 break;
90 }
91 line_no += 1;
92 let terminated = buffer.last() == Some(&b'\n');
93 let line = if terminated {
94 &buffer[..buffer.len() - 1]
95 } else {
96 buffer.as_slice()
97 };
98 if line.iter().all(u8::is_ascii_whitespace) {
99 if !terminated {
100 return Ok(LoadedEventLog {
101 events,
102 tail_repair: TailRepair::Truncate(valid_prefix_len),
103 });
104 }
105 valid_prefix_len = checked_file_offset(valid_prefix_len, bytes_read, &path)?;
106 continue;
107 }
108 let envelope: FlowEventEnvelope = match serde_json::from_slice(line) {
109 Ok(envelope) => envelope,
110 Err(_) if !terminated => {
111 return Ok(LoadedEventLog {
112 events,
113 tail_repair: TailRepair::Truncate(valid_prefix_len),
114 });
115 }
116 Err(err) => {
117 return Err(FlowError::Store(format!(
118 "failed to decode event line {line_no} from {}: {err}",
119 path.display()
120 )));
121 }
122 };
123 if envelope.run_id != run_id {
124 return Err(FlowError::Store(format!(
125 "event line {line_no} in {} belongs to run {}, not {run_id}",
126 path.display(),
127 envelope.run_id
128 )));
129 }
130 events.push(envelope);
131 valid_prefix_len = checked_file_offset(valid_prefix_len, bytes_read, &path)?;
132 if !terminated {
133 return Ok(LoadedEventLog {
134 events,
135 tail_repair: TailRepair::AppendDelimiter,
136 });
137 }
138 }
139
140 Ok(LoadedEventLog {
141 events,
142 tail_repair: TailRepair::None,
143 })
144 }
145
146 async fn list_inner(
147 &self,
148 run_id: &str,
149 missing_is_empty: bool,
150 ) -> Result<Vec<FlowEventEnvelope>> {
151 Ok(self.load_inner(run_id, missing_is_empty).await?.events)
152 }
153
154 fn validate_existing_log(&self, run_id: &str, events: &[FlowEventEnvelope]) -> Result<()> {
155 if events.is_empty() {
156 return Ok(());
157 }
158 project_run(run_id, events)?;
159 Ok(())
160 }
161
162 async fn append_inner(&self, run_id: &str, event: FlowEvent) -> Result<FlowEventEnvelope> {
163 tokio::fs::create_dir_all(&self.root).await?;
164
165 let LoadedEventLog {
166 events,
167 tail_repair,
168 } = self.load_inner(run_id, true).await?;
169 self.validate_existing_log(run_id, &events)?;
170 let envelope = FlowEventEnvelope {
171 run_id: run_id.to_string(),
172 sequence: events.last().map_or(1, |event| event.sequence + 1),
173 event_id: Uuid::new_v4(),
174 timestamp: Utc::now(),
175 event,
176 };
177
178 self.repair_tail(run_id, tail_repair).await?;
179 self.write_envelope(&envelope).await?;
180 Ok(envelope)
181 }
182
183 async fn append_if_sequence_inner(
184 &self,
185 run_id: &str,
186 expected_sequence: u64,
187 event: FlowEvent,
188 ) -> Result<FlowEventEnvelope> {
189 tokio::fs::create_dir_all(&self.root).await?;
190
191 let LoadedEventLog {
192 events,
193 tail_repair,
194 } = self.load_inner(run_id, true).await?;
195 self.validate_existing_log(run_id, &events)?;
196 let actual_sequence = events.last().map_or(0, |event| event.sequence);
197 if actual_sequence != expected_sequence {
198 return Err(FlowError::EventConflict {
199 run_id: run_id.to_string(),
200 expected_sequence,
201 actual_sequence,
202 });
203 }
204
205 let envelope = FlowEventEnvelope {
206 run_id: run_id.to_string(),
207 sequence: actual_sequence + 1,
208 event_id: Uuid::new_v4(),
209 timestamp: Utc::now(),
210 event,
211 };
212
213 self.repair_tail(run_id, tail_repair).await?;
214 self.write_envelope(&envelope).await?;
215 Ok(envelope)
216 }
217
218 async fn repair_tail(&self, run_id: &str, repair: TailRepair) -> Result<()> {
219 let path = self.run_path(run_id)?;
220 match repair {
221 TailRepair::None => Ok(()),
222 TailRepair::AppendDelimiter => {
223 let mut file = OpenOptions::new().append(true).open(path).await?;
224 file.write_all(b"\n").await?;
225 file.flush().await?;
226 file.sync_data().await?;
227 Ok(())
228 }
229 TailRepair::Truncate(valid_prefix_len) => {
230 let file = OpenOptions::new().write(true).open(path).await?;
231 file.set_len(valid_prefix_len).await?;
232 file.sync_data().await?;
233 Ok(())
234 }
235 }
236 }
237
238 async fn write_envelope(&self, envelope: &FlowEventEnvelope) -> Result<()> {
239 let path = self.run_path(&envelope.run_id)?;
240 let mut file = OpenOptions::new()
241 .create(true)
242 .append(true)
243 .open(path)
244 .await?;
245 let mut line = serde_json::to_vec(envelope)?;
246 line.push(b'\n');
247 file.write_all(&line).await?;
248 file.flush().await?;
249 file.sync_data().await?;
250 Ok(())
251 }
252
253 async fn list_run_ids_inner(&self) -> Result<Vec<String>> {
254 let mut ids = Vec::new();
255
256 let mut dir = match tokio::fs::read_dir(&self.root).await {
257 Ok(dir) => dir,
258 Err(err) if err.kind() == std::io::ErrorKind::NotFound => return Ok(ids),
259 Err(err) => return Err(FlowError::Io(err)),
260 };
261
262 while let Some(entry) = dir.next_entry().await? {
263 let path = entry.path();
264 if path.extension().and_then(|ext| ext.to_str()) != Some("jsonl") {
265 continue;
266 }
267 let Some(stem) = path.file_stem().and_then(|stem| stem.to_str()) else {
268 continue;
269 };
270 if is_safe_run_id(stem) {
271 ids.push(stem.to_string());
272 }
273 }
274
275 ids.sort();
276 Ok(ids)
277 }
278
279 pub async fn prune_terminal_runs_older_than(
286 &self,
287 terminal_before: DateTime<Utc>,
288 ) -> Result<Vec<String>> {
289 let _guard = self.lock.lock().await;
290 let mut removed = Vec::new();
291
292 for run_id in self.list_run_ids_inner().await? {
293 let events = self.list_inner(&run_id, false).await?;
294 self.validate_existing_log(&run_id, &events)?;
295 let Some(terminal_at) = terminal_event_timestamp(&events) else {
296 continue;
297 };
298 if terminal_at >= terminal_before {
299 continue;
300 }
301
302 let path = self.run_path(&run_id)?;
303 match tokio::fs::remove_file(&path).await {
304 Ok(()) => removed.push(run_id),
305 Err(err) if err.kind() == std::io::ErrorKind::NotFound => {}
306 Err(err) => return Err(FlowError::Io(err)),
307 }
308 }
309
310 removed.sort();
311 Ok(removed)
312 }
313}
314
315#[async_trait]
316impl FlowEventStore for LocalFileEventStore {
317 async fn append(&self, run_id: &str, event: FlowEvent) -> Result<FlowEventEnvelope> {
318 let _guard = self.lock.lock().await;
319 self.append_inner(run_id, event).await
320 }
321
322 async fn append_if_sequence(
323 &self,
324 run_id: &str,
325 expected_sequence: u64,
326 event: FlowEvent,
327 ) -> Result<FlowEventEnvelope> {
328 let _guard = self.lock.lock().await;
329 self.append_if_sequence_inner(run_id, expected_sequence, event)
330 .await
331 }
332
333 async fn list(&self, run_id: &str) -> Result<Vec<FlowEventEnvelope>> {
334 let _guard = self.lock.lock().await;
335 self.list_inner(run_id, false).await
336 }
337
338 async fn list_run_ids(&self) -> Result<Vec<String>> {
339 let _guard = self.lock.lock().await;
340 self.list_run_ids_inner().await
341 }
342}
343
344fn checked_file_offset(current: u64, bytes_read: usize, path: &Path) -> Result<u64> {
345 let bytes_read = u64::try_from(bytes_read).map_err(|_| {
346 FlowError::Store(format!(
347 "event line length from {} exceeds the supported file offset",
348 path.display()
349 ))
350 })?;
351 current.checked_add(bytes_read).ok_or_else(|| {
352 FlowError::Store(format!(
353 "event log {} exceeds the supported file offset",
354 path.display()
355 ))
356 })
357}
358
359fn terminal_event_timestamp(events: &[FlowEventEnvelope]) -> Option<DateTime<Utc>> {
360 events
361 .iter()
362 .rev()
363 .find_map(|envelope| match envelope.event {
364 FlowEvent::RunCompleted { .. }
365 | FlowEvent::RunFailed { .. }
366 | FlowEvent::RunCancelled { .. } => Some(envelope.timestamp),
367 _ => None,
368 })
369}
370
371fn is_safe_run_id(run_id: &str) -> bool {
372 !run_id.is_empty()
373 && run_id
374 .chars()
375 .all(|ch| ch.is_ascii_alphanumeric() || ch == '-' || ch == '_')
376}