Skip to main content

a3s_flow/store/
local_file.rs

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/// JSONL-backed event store for local durable runs.
16///
17/// Each workflow run is stored as `<root>/<run_id>.jsonl`; every line is a full
18/// [`FlowEventEnvelope`]. The store serializes appends inside this process, but
19/// it does not provide cross-process locking. Use it for local development,
20/// embedded Rust hosts, and crash/restart durability. An unterminated malformed
21/// tail is treated as a torn append and truncated before the next write;
22/// terminated or interior corruption remains an error. Use a database-backed
23/// store for multi-writer deployments.
24#[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    /// Remove completed, failed, or cancelled local run histories whose terminal
280    /// event timestamp is strictly before `terminal_before`.
281    ///
282    /// Suspended and running runs are never removed by this helper. Corrupt
283    /// histories are returned as errors rather than deleted, so operators can
284    /// inspect them before cleanup.
285    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}