Skip to main content

harn_vm/event_log/
file.rs

1use std::path::{Path, PathBuf};
2use std::sync::{Arc, Mutex};
3use std::time::Duration;
4
5use futures::stream::BoxStream;
6use serde::{Deserialize, Serialize};
7
8use super::util::{
9    dir_size_bytes, now_ms, prepare_event_after, sanitize_filename, stream_from_broadcast,
10    sync_tree, write_json_atomically, BroadcastMap,
11};
12use super::{
13    require_expected_topic_head, AppendHeadExpectation, AppendOutcome, CompactReport, ConsumerId,
14    EventId, EventLog, EventLogBackendKind, EventLogDescription, LogError, LogEvent, LogEventBytes,
15    Topic,
16};
17
18const TOPIC_LOCK_TIMEOUT: Duration = Duration::from_secs(30);
19
20#[derive(Serialize, Deserialize)]
21struct FileRecord {
22    id: EventId,
23    event: LogEvent,
24}
25
26pub struct FileEventLog {
27    root: PathBuf,
28    write_lock: Mutex<()>,
29    pub(super) broadcasts: BroadcastMap,
30    pub(super) queue_depth: usize,
31}
32
33impl FileEventLog {
34    pub fn open(root: PathBuf, queue_depth: usize) -> Result<Self, LogError> {
35        std::fs::create_dir_all(root.join("topics"))
36            .map_err(|error| LogError::Io(format!("event log mkdir error: {error}")))?;
37        std::fs::create_dir_all(root.join("consumers"))
38            .map_err(|error| LogError::Io(format!("event log mkdir error: {error}")))?;
39        std::fs::create_dir_all(root.join("locks"))
40            .map_err(|error| LogError::Io(format!("event log mkdir error: {error}")))?;
41        Ok(Self {
42            root,
43            write_lock: Mutex::new(()),
44            broadcasts: BroadcastMap::default(),
45            queue_depth: queue_depth.max(1),
46        })
47    }
48
49    fn topic_path(&self, topic: &Topic) -> PathBuf {
50        self.root
51            .join("topics")
52            .join(format!("{}.jsonl", topic.as_str()))
53    }
54
55    fn consumer_path(&self, topic: &Topic, consumer: &ConsumerId) -> PathBuf {
56        self.root.join("consumers").join(format!(
57            "{}__{}.json",
58            topic.as_str(),
59            sanitize_filename(consumer.as_str())
60        ))
61    }
62
63    fn lock_topic(&self, topic: &Topic) -> Result<std::fs::File, LogError> {
64        let path = self
65            .root
66            .join("locks")
67            .join(format!("{}.lock", topic.as_str()));
68        let file = std::fs::OpenOptions::new()
69            .create(true)
70            .truncate(false)
71            .read(true)
72            .write(true)
73            .open(&path)
74            .map_err(|error| LogError::Io(format!("event log topic lock open error: {error}")))?;
75        harn_flock::lock_with_deadline(
76            &file,
77            &path,
78            harn_flock::LockMode::Exclusive,
79            TOPIC_LOCK_TIMEOUT,
80        )
81        .map_err(|error| LogError::Io(format!("event log topic lock error: {error}")))?;
82        Ok(file)
83    }
84
85    fn latest_id_from_disk(&self, topic: &Topic) -> Result<EventId, LogError> {
86        let mut latest = 0;
87        let path = self.topic_path(topic);
88        if path.is_file() {
89            for record in read_file_records(&path)? {
90                latest = record.id;
91            }
92        }
93        Ok(latest)
94    }
95
96    fn read_range_sync(
97        &self,
98        topic: &Topic,
99        from: Option<EventId>,
100        limit: usize,
101    ) -> Result<Vec<(EventId, LogEvent)>, LogError> {
102        let path = self.topic_path(topic);
103        if !path.is_file() {
104            return Ok(Vec::new());
105        }
106        let from = from.unwrap_or(0);
107        let mut events = Vec::new();
108        for record in read_file_records(&path)? {
109            if record.id > from {
110                events.push((record.id, record.event));
111            }
112            if events.len() >= limit {
113                break;
114            }
115        }
116        Ok(events)
117    }
118
119    pub(super) fn topics(&self) -> Result<Vec<Topic>, LogError> {
120        let topics_dir = self.root.join("topics");
121        if !topics_dir.is_dir() {
122            return Ok(Vec::new());
123        }
124        let mut topics = Vec::new();
125        for entry in std::fs::read_dir(&topics_dir)
126            .map_err(|error| LogError::Io(format!("event log topics read error: {error}")))?
127        {
128            let entry = entry
129                .map_err(|error| LogError::Io(format!("event log topic entry error: {error}")))?;
130            let path = entry.path();
131            if path.extension().and_then(|ext| ext.to_str()) != Some("jsonl") {
132                continue;
133            }
134            let Some(stem) = path.file_stem().and_then(|stem| stem.to_str()) else {
135                continue;
136            };
137            topics.push(Topic::new(stem.to_string())?);
138        }
139        topics.sort_by(|left, right| left.as_str().cmp(right.as_str()));
140        Ok(topics)
141    }
142
143    pub(super) fn append_idempotent_by_header(
144        &self,
145        topic: &Topic,
146        header: &str,
147        value: &str,
148        event: LogEvent,
149    ) -> Result<AppendOutcome, LogError> {
150        self.append_idempotent_by_header_with_expectation(
151            topic,
152            header,
153            value,
154            AppendHeadExpectation::Any,
155            event,
156        )
157    }
158
159    pub(super) fn append_idempotent_chained_by_header(
160        &self,
161        topic: &Topic,
162        header: &str,
163        value: &str,
164        expected_head: Option<&str>,
165        event: LogEvent,
166    ) -> Result<AppendOutcome, LogError> {
167        self.append_idempotent_by_header_with_expectation(
168            topic,
169            header,
170            value,
171            AppendHeadExpectation::Exact(expected_head),
172            event,
173        )
174    }
175
176    fn append_idempotent_by_header_with_expectation(
177        &self,
178        topic: &Topic,
179        header: &str,
180        value: &str,
181        expectation: AppendHeadExpectation<'_>,
182        event: LogEvent,
183    ) -> Result<AppendOutcome, LogError> {
184        let _guard = self
185            .write_lock
186            .lock()
187            .expect("file event log write lock poisoned");
188        let _topic_lock = self.lock_topic(topic)?;
189        let existing_events = self.read_range_sync(topic, None, usize::MAX)?;
190        if let Some((event_id, existing)) = existing_events.iter().find(|(_, event)| {
191            event
192                .headers
193                .get(header)
194                .is_some_and(|found| found == value)
195        }) {
196            return Ok(AppendOutcome {
197                event_id: *event_id,
198                event: existing.clone(),
199                inserted: false,
200            });
201        }
202
203        let next_id = existing_events
204            .last()
205            .map(|(event_id, _)| *event_id)
206            .unwrap_or(0)
207            + 1;
208        let previous = existing_events
209            .last()
210            .map(|(previous_id, previous_event)| (*previous_id, previous_event));
211        require_expected_topic_head(topic, previous, expectation)?;
212        let event = prepare_event_after(topic, next_id, previous, event)?;
213        self.append_record_locked(topic, next_id, event)
214    }
215
216    /// Read counterpart of [`Self::append_idempotent_by_header`]. The file
217    /// backend has no header index, so this scans the topic — acceptable for a
218    /// dev/test backend (SQLite is the durable default).
219    pub(super) fn read_idempotent_by_header(
220        &self,
221        topic: &Topic,
222        header: &str,
223        value: &str,
224    ) -> Result<Option<(EventId, LogEvent)>, LogError> {
225        let events = self.read_range_sync(topic, None, usize::MAX)?;
226        Ok(events.into_iter().find(|(_, event)| {
227            event
228                .headers
229                .get(header)
230                .is_some_and(|found| found == value)
231        }))
232    }
233
234    fn append_record_locked(
235        &self,
236        topic: &Topic,
237        event_id: EventId,
238        event: LogEvent,
239    ) -> Result<AppendOutcome, LogError> {
240        let record = FileRecord {
241            id: event_id,
242            event: event.clone(),
243        };
244        let path = self.topic_path(topic);
245        if let Some(parent) = path.parent() {
246            std::fs::create_dir_all(parent)
247                .map_err(|error| LogError::Io(format!("event log mkdir error: {error}")))?;
248        }
249        let line = serde_json::to_string(&record)
250            .map_err(|error| LogError::Serde(format!("event log encode error: {error}")))?;
251        use std::io::Write as _;
252        let mut file = std::fs::OpenOptions::new()
253            .create(true)
254            .append(true)
255            .open(&path)
256            .map_err(|error| LogError::Io(format!("event log open error: {error}")))?;
257        writeln!(file, "{line}")
258            .map_err(|error| LogError::Io(format!("event log write error: {error}")))?;
259        self.broadcasts
260            .publish(topic, self.queue_depth, (event_id, event.clone()));
261        Ok(AppendOutcome {
262            event_id,
263            event,
264            inserted: true,
265        })
266    }
267}
268
269fn read_file_records(path: &Path) -> Result<Vec<FileRecord>, LogError> {
270    let file = std::fs::File::open(path)
271        .map_err(|error| LogError::Io(format!("event log open error: {error}")))?;
272    let mut reader = std::io::BufReader::new(file);
273    let mut records = Vec::new();
274    let mut line = Vec::new();
275    loop {
276        line.clear();
277        let bytes_read = std::io::BufRead::read_until(&mut reader, b'\n', &mut line)
278            .map_err(|error| LogError::Io(format!("event log read error: {error}")))?;
279        if bytes_read == 0 {
280            break;
281        }
282        if line.iter().all(u8::is_ascii_whitespace) {
283            continue;
284        }
285        let complete_line = line.ends_with(b"\n");
286        match serde_json::from_slice::<FileRecord>(&line) {
287            Ok(record) => records.push(record),
288            Err(_) if !complete_line => break,
289            Err(error) => {
290                return Err(LogError::Serde(format!("event log parse error: {error}")));
291            }
292        }
293    }
294    Ok(records)
295}
296
297impl EventLog for FileEventLog {
298    fn describe(&self) -> EventLogDescription {
299        EventLogDescription {
300            backend: EventLogBackendKind::File,
301            location: Some(self.root.clone()),
302            size_bytes: Some(dir_size_bytes(&self.root)),
303            queue_depth: self.queue_depth,
304        }
305    }
306
307    async fn append(&self, topic: &Topic, event: LogEvent) -> Result<EventId, LogError> {
308        let _guard = self
309            .write_lock
310            .lock()
311            .expect("file event log write lock poisoned");
312        let _topic_lock = self.lock_topic(topic)?;
313        let existing_events = self.read_range_sync(topic, None, usize::MAX)?;
314        let next_id = existing_events
315            .last()
316            .map(|(event_id, _)| *event_id)
317            .unwrap_or(0)
318            + 1;
319        let previous = existing_events
320            .last()
321            .map(|(previous_id, previous_event)| (*previous_id, previous_event));
322        let event = prepare_event_after(topic, next_id, previous, event)?;
323        self.append_record_locked(topic, next_id, event)
324            .map(|outcome| outcome.event_id)
325    }
326
327    async fn flush(&self) -> Result<(), LogError> {
328        sync_tree(&self.root)
329    }
330
331    async fn read_range(
332        &self,
333        topic: &Topic,
334        from: Option<EventId>,
335        limit: usize,
336    ) -> Result<Vec<(EventId, LogEvent)>, LogError> {
337        self.read_range_sync(topic, from, limit)
338    }
339
340    async fn read_range_bytes(
341        &self,
342        topic: &Topic,
343        from: Option<EventId>,
344        limit: usize,
345    ) -> Result<Vec<(EventId, LogEventBytes)>, LogError> {
346        self.read_range_sync(topic, from, limit)?
347            .into_iter()
348            .map(|(event_id, event)| Ok((event_id, event.try_into()?)))
349            .collect()
350    }
351
352    async fn subscribe(
353        self: Arc<Self>,
354        topic: &Topic,
355        from: Option<EventId>,
356    ) -> Result<BoxStream<'static, Result<(EventId, LogEvent), LogError>>, LogError> {
357        let rx = self.broadcasts.subscribe(topic, self.queue_depth);
358        let history = self.read_range_sync(topic, from, usize::MAX)?;
359        Ok(stream_from_broadcast(history, from, rx, self.queue_depth))
360    }
361
362    async fn ack(
363        &self,
364        topic: &Topic,
365        consumer: &ConsumerId,
366        up_to: EventId,
367    ) -> Result<(), LogError> {
368        let path = self.consumer_path(topic, consumer);
369        let payload = serde_json::json!({
370            "topic": topic.as_str(),
371            "consumer_id": consumer.as_str(),
372            "cursor": up_to,
373            "updated_at_ms": now_ms(),
374        });
375        write_json_atomically(&path, &payload)
376    }
377
378    async fn consumer_cursor(
379        &self,
380        topic: &Topic,
381        consumer: &ConsumerId,
382    ) -> Result<Option<EventId>, LogError> {
383        let path = self.consumer_path(topic, consumer);
384        if !path.is_file() {
385            return Ok(None);
386        }
387        let raw = std::fs::read_to_string(&path)
388            .map_err(|error| LogError::Io(format!("event log consumer read error: {error}")))?;
389        let payload: serde_json::Value = serde_json::from_str(&raw)
390            .map_err(|error| LogError::Serde(format!("event log consumer parse error: {error}")))?;
391        let cursor = payload
392            .get("cursor")
393            .and_then(serde_json::Value::as_u64)
394            .ok_or_else(|| {
395                LogError::Serde("event log consumer record missing numeric cursor".to_string())
396            })?;
397        Ok(Some(cursor))
398    }
399
400    async fn latest(&self, topic: &Topic) -> Result<Option<EventId>, LogError> {
401        let _topic_lock = self.lock_topic(topic)?;
402        let latest = self.latest_id_from_disk(topic)?;
403        if latest == 0 {
404            Ok(None)
405        } else {
406            Ok(Some(latest))
407        }
408    }
409
410    async fn compact(&self, topic: &Topic, before: EventId) -> Result<CompactReport, LogError> {
411        let _guard = self
412            .write_lock
413            .lock()
414            .expect("file event log write lock poisoned");
415        let _topic_lock = self.lock_topic(topic)?;
416        let path = self.topic_path(topic);
417        if !path.is_file() {
418            return Ok(CompactReport::default());
419        }
420        let retained = self.read_range_sync(topic, Some(before), usize::MAX)?;
421        let removed = self.read_range_sync(topic, None, usize::MAX)?.len() - retained.len();
422        if retained.is_empty() {
423            let _ = std::fs::remove_file(&path);
424        } else {
425            crate::atomic_io::atomic_write_with(&path, |writer| {
426                use std::io::Write as _;
427                for (event_id, event) in &retained {
428                    let line = serde_json::to_string(&FileRecord {
429                        id: *event_id,
430                        event: event.clone(),
431                    })
432                    .map_err(|error| std::io::Error::other(error.to_string()))?;
433                    writeln!(writer, "{line}")?;
434                }
435                Ok(())
436            })
437            .map_err(|error| LogError::Io(format!("event log compact finalize error: {error}")))?;
438        }
439        let latest = retained.last().map(|(event_id, _)| *event_id);
440        Ok(CompactReport {
441            removed,
442            remaining: retained.len(),
443            latest,
444            checkpointed: false,
445        })
446    }
447}