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