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