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