Skip to main content

weavatrix_memory/store/
file.rs

1use super::{
2    EventStore, ExpectedVersion, InMemoryStore,
3    frame::{self, ScanOutcome},
4};
5use crate::{AppendReceipt, Codec, MemoryError, NewEvent, Result, StoredEvent, StreamId};
6use std::{
7    fs::{File, OpenOptions, TryLockError},
8    io::{Seek, SeekFrom, Write},
9    path::{Path, PathBuf},
10};
11
12#[derive(Debug, Clone, Copy, PartialEq, Eq)]
13pub enum RecoveryPolicy {
14    Strict,
15    TruncatePartialTail,
16}
17
18#[derive(Debug, Clone, Copy, PartialEq, Eq)]
19pub enum Durability {
20    Flush,
21    SyncData,
22    SyncAll,
23}
24
25#[derive(Debug, Clone, Copy, PartialEq, Eq)]
26pub struct FileStoreOptions {
27    pub recovery: RecoveryPolicy,
28    pub durability: Durability,
29    pub max_frame_bytes: usize,
30}
31
32impl Default for FileStoreOptions {
33    fn default() -> Self {
34        Self {
35            recovery: RecoveryPolicy::Strict,
36            durability: Durability::SyncAll,
37            max_frame_bytes: 256 * 1024 * 1024,
38        }
39    }
40}
41
42pub struct FileEventStore<E, C> {
43    path: PathBuf,
44    file: File,
45    codec: C,
46    options: FileStoreOptions,
47    inner: InMemoryStore<E>,
48    durable_len: u64,
49}
50
51impl<E, C> FileEventStore<E, C>
52where
53    E: Clone,
54    C: Codec<StoredEvent<E>>,
55{
56    /// Opens or creates a framed append-only event journal.
57    ///
58    /// # Errors
59    ///
60    /// Rejects invalid headers, checksum corruption, invalid restored event
61    /// sequences, and partial tails in strict mode.
62    pub fn open(path: impl AsRef<Path>, codec: C, options: FileStoreOptions) -> Result<Self> {
63        if options.max_frame_bytes < 4 {
64            return Err(MemoryError::InvalidValue {
65                field: "max_frame_bytes",
66                reason: "must fit at least the event-count field",
67            });
68        }
69        let path = path.as_ref().to_path_buf();
70        let mut file = OpenOptions::new()
71            .read(true)
72            .write(true)
73            .create(true)
74            .truncate(false)
75            .open(&path)
76            .map_err(|error| io("open event log", error))?;
77        file.try_lock().map_err(|error| match error {
78            TryLockError::WouldBlock => MemoryError::ExternalModification,
79            TryLockError::Error(error) => io("lock event log", error),
80        })?;
81        if file
82            .metadata()
83            .map_err(|error| io("read event log metadata", error))?
84            .len()
85            == 0
86        {
87            file.write_all(frame::FILE_HEADER)
88                .map_err(|error| io("write event log header", error))?;
89            sync(&file, options.durability)?;
90        }
91        let outcome = frame::scan(&mut file, &codec, options.max_frame_bytes)?;
92        let (events, durable_len, partial) = match outcome {
93            ScanOutcome::Complete {
94                events,
95                durable_len,
96            } => (events, durable_len, false),
97            ScanOutcome::PartialTail {
98                events,
99                durable_len,
100            } => (events, durable_len, true),
101        };
102        if partial && options.recovery == RecoveryPolicy::Strict {
103            return Err(MemoryError::CorruptLog {
104                offset: durable_len,
105                reason: "partial trailing batch".to_owned(),
106            });
107        }
108        if partial {
109            file.set_len(durable_len)
110                .map_err(|error| io("truncate partial event batch", error))?;
111            sync(&file, options.durability)?;
112        }
113        file.seek(SeekFrom::End(0))
114            .map_err(|error| io("seek event log end", error))?;
115        Ok(Self {
116            path,
117            file,
118            codec,
119            options,
120            inner: InMemoryStore::restore(events)?,
121            durable_len,
122        })
123    }
124
125    #[must_use]
126    pub fn path(&self) -> &Path {
127        &self.path
128    }
129
130    fn persist(&mut self, events: &[StoredEvent<E>]) -> Result<()> {
131        if events.is_empty() {
132            return Ok(());
133        }
134        let actual_len = self
135            .file
136            .metadata()
137            .map_err(|error| io("read event log metadata", error))?
138            .len();
139        if actual_len != self.durable_len {
140            return Err(MemoryError::ExternalModification);
141        }
142        let frame = frame::encode_batch(events, &self.codec, self.options.max_frame_bytes)?;
143        self.file
144            .seek(SeekFrom::Start(self.durable_len))
145            .map_err(|error| io("seek append position", error))?;
146        if let Err(error) = self.file.write_all(&frame) {
147            self.rollback_tail()?;
148            return Err(io("append event batch", error));
149        }
150        if let Err(error) = sync(&self.file, self.options.durability) {
151            self.rollback_tail()?;
152            return Err(error);
153        }
154        self.durable_len = self
155            .durable_len
156            .checked_add(u64::try_from(frame.len()).map_err(|_| MemoryError::CapacityOverflow)?)
157            .ok_or(MemoryError::CapacityOverflow)?;
158        Ok(())
159    }
160
161    fn rollback_tail(&mut self) -> Result<()> {
162        self.file
163            .set_len(self.durable_len)
164            .map_err(|error| io("rollback partial event batch", error))?;
165        self.file
166            .seek(SeekFrom::Start(self.durable_len))
167            .map_err(|error| io("seek after rollback", error))?;
168        sync(&self.file, self.options.durability)
169    }
170}
171
172impl<E, C> EventStore<E> for FileEventStore<E, C>
173where
174    E: Clone,
175    C: Codec<StoredEvent<E>>,
176{
177    fn append(
178        &mut self,
179        stream: &StreamId,
180        expected: ExpectedVersion,
181        events: &[NewEvent<E>],
182    ) -> Result<Vec<StoredEvent<E>>> {
183        let committed = self.inner.prepare_append(stream, expected, events)?;
184        self.persist(&committed)?;
185        self.inner.commit_prepared(&committed);
186        Ok(committed)
187    }
188
189    fn append_owned(
190        &mut self,
191        stream: &StreamId,
192        expected: ExpectedVersion,
193        events: Vec<NewEvent<E>>,
194    ) -> Result<Vec<StoredEvent<E>>> {
195        let committed = self.inner.prepare_append_owned(stream, expected, events)?;
196        self.persist(&committed)?;
197        self.inner.commit_prepared(&committed);
198        Ok(committed)
199    }
200
201    fn append_owned_receipt(
202        &mut self,
203        stream: &StreamId,
204        expected: ExpectedVersion,
205        events: Vec<NewEvent<E>>,
206    ) -> Result<AppendReceipt> {
207        let committed = self.inner.prepare_append_owned(stream, expected, events)?;
208        self.persist(&committed)?;
209        let receipt = AppendReceipt::from_events(&committed);
210        self.inner.commit_prepared_owned(committed);
211        Ok(receipt)
212    }
213
214    fn load_stream(&self, stream: &StreamId, after: Option<u64>) -> Vec<StoredEvent<E>> {
215        self.inner.load_stream(stream, after)
216    }
217
218    fn load_all(&self, after: Option<u64>, limit: usize) -> Vec<StoredEvent<E>> {
219        self.inner.load_all(after, limit)
220    }
221
222    fn stream_version(&self, stream: &StreamId) -> Option<u64> {
223        self.inner.stream_version(stream)
224    }
225
226    fn len(&self) -> usize {
227        self.inner.len()
228    }
229}
230
231fn sync(file: &File, durability: Durability) -> Result<()> {
232    match durability {
233        Durability::Flush => Ok(()),
234        Durability::SyncData => file
235            .sync_data()
236            .map_err(|error| io("sync event log", error)),
237        Durability::SyncAll => file
238            .sync_all()
239            .map_err(|error| io("sync event log data and metadata", error)),
240    }
241}
242
243#[allow(clippy::needless_pass_by_value)]
244fn io(operation: &'static str, error: std::io::Error) -> MemoryError {
245    MemoryError::Io {
246        operation,
247        message: error.to_string(),
248    }
249}