Skip to main content

weavatrix_memory/store/
file.rs

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