1use corium_core::{
4 Datom, EntityId,
5 encoding::{decode_value, encode_value},
6};
7use std::{
8 fs::{self, File, OpenOptions},
9 io::{self, Read, Write},
10 path::{Path, PathBuf},
11 sync::{Arc, RwLock},
12};
13use thiserror::Error;
14
15#[derive(Clone, Debug, Eq, PartialEq)]
17pub struct TxRecord {
18 pub t: u64,
20 pub tx_instant: i64,
22 pub datoms: Vec<Datom>,
24}
25
26#[derive(Debug, Error)]
28pub enum LogError {
29 #[error("log I/O failed: {0}")]
31 Io(#[from] io::Error),
32 #[error("corrupt transaction log")]
34 Corrupt,
35}
36
37pub trait TransactionLog: Send + Sync {
39 fn append(&self, record: &TxRecord) -> Result<(), LogError>;
44 fn tx_range(&self, start: u64, end: Option<u64>) -> Result<Vec<TxRecord>, LogError>;
49 fn replay(&self) -> Result<Vec<TxRecord>, LogError> {
54 self.tx_range(0, None)
55 }
56}
57
58#[derive(Clone, Default)]
60pub struct MemoryLog(Arc<RwLock<Vec<TxRecord>>>);
61impl TransactionLog for MemoryLog {
62 fn append(&self, record: &TxRecord) -> Result<(), LogError> {
63 let mut records = self.0.write().expect("poisoned log lock");
64 if records.last().map_or(1, |r| r.t + 1) != record.t {
65 return Err(LogError::Corrupt);
66 }
67 records.push(record.clone());
68 Ok(())
69 }
70 fn tx_range(&self, start: u64, end: Option<u64>) -> Result<Vec<TxRecord>, LogError> {
71 Ok(self
72 .0
73 .read()
74 .expect("poisoned log lock")
75 .iter()
76 .filter(|r| r.t >= start && end.is_none_or(|e| r.t < e))
77 .cloned()
78 .collect())
79 }
80}
81
82pub struct FileLog {
88 path: PathBuf,
89 next_t: RwLock<u64>,
90}
91impl FileLog {
92 pub fn open(path: impl AsRef<Path>) -> Result<Self, LogError> {
98 let path = path.as_ref().to_path_buf();
99 if let Some(parent) = path.parent() {
100 fs::create_dir_all(parent)?;
101 }
102 OpenOptions::new().create(true).append(true).open(&path)?;
103 let (records, durable_len) = read_records(&path)?;
104 if fs::metadata(&path)?.len() > durable_len {
105 let file = OpenOptions::new().write(true).open(&path)?;
106 file.set_len(durable_len)?;
107 file.sync_all()?;
108 }
109 Ok(Self {
110 path,
111 next_t: RwLock::new(records.last().map_or(1, |r| r.t + 1)),
112 })
113 }
114}
115impl TransactionLog for FileLog {
116 fn append(&self, record: &TxRecord) -> Result<(), LogError> {
117 let mut next_t = self.next_t.write().expect("poisoned log lock");
118 if *next_t != record.t {
119 return Err(LogError::Corrupt);
120 }
121 let payload = encode_record(record);
122 let mut file = OpenOptions::new().append(true).open(&self.path)?;
123 file.write_all(
124 &u64::try_from(payload.len())
125 .map_err(|_| LogError::Corrupt)?
126 .to_be_bytes(),
127 )?;
128 file.write_all(&payload)?;
129 file.sync_all()?;
130 *next_t += 1;
131 Ok(())
132 }
133 fn tx_range(&self, start: u64, end: Option<u64>) -> Result<Vec<TxRecord>, LogError> {
134 let _guard = self.next_t.read().expect("poisoned log lock");
135 Ok(read_records(&self.path)?
136 .0
137 .into_iter()
138 .filter(|r| r.t >= start && end.is_none_or(|e| r.t < e))
139 .collect())
140 }
141}
142
143pub struct VersionedLog {
155 dir: PathBuf,
156 name: String,
157 write_path: PathBuf,
158 next_t: RwLock<u64>,
159}
160
161impl VersionedLog {
162 pub fn open(dir: impl AsRef<Path>, name: &str, write_version: u64) -> Result<Self, LogError> {
170 let dir = dir.as_ref().to_path_buf();
171 fs::create_dir_all(&dir)?;
172 let write_path = version_path(&dir, name, write_version);
173 OpenOptions::new()
174 .create(true)
175 .append(true)
176 .open(&write_path)?;
177 let (_, durable_len) = read_records(&write_path)?;
178 if fs::metadata(&write_path)?.len() > durable_len {
179 let file = OpenOptions::new().write(true).open(&write_path)?;
180 file.set_len(durable_len)?;
181 file.sync_all()?;
182 }
183 let records = read_merged(&dir, name)?;
184 Ok(Self {
185 dir,
186 name: name.to_owned(),
187 write_path,
188 next_t: RwLock::new(records.last().map_or(1, |r| r.t + 1)),
189 })
190 }
191
192 pub fn open_read_only(dir: impl AsRef<Path>, name: &str) -> Result<Self, LogError> {
198 let dir = dir.as_ref().to_path_buf();
199 Ok(Self {
200 write_path: PathBuf::new(),
201 name: name.to_owned(),
202 next_t: RwLock::new(u64::MAX),
203 dir,
204 })
205 }
206
207 #[must_use]
209 pub fn exists(dir: impl AsRef<Path>, name: &str) -> bool {
210 !version_files(dir.as_ref(), name).is_empty()
211 }
212
213 pub fn delete_all(dir: impl AsRef<Path>, name: &str) -> Result<(), LogError> {
218 for (_, path) in version_files(dir.as_ref(), name) {
219 match fs::remove_file(&path) {
220 Ok(()) => {}
221 Err(error) if error.kind() == io::ErrorKind::NotFound => {}
222 Err(error) => return Err(error.into()),
223 }
224 }
225 Ok(())
226 }
227}
228
229impl TransactionLog for VersionedLog {
230 fn append(&self, record: &TxRecord) -> Result<(), LogError> {
231 let mut next_t = self.next_t.write().expect("poisoned log lock");
232 if *next_t != record.t {
233 return Err(LogError::Corrupt);
234 }
235 let payload = encode_record(record);
236 let mut file = OpenOptions::new().append(true).open(&self.write_path)?;
237 file.write_all(
238 &u64::try_from(payload.len())
239 .map_err(|_| LogError::Corrupt)?
240 .to_be_bytes(),
241 )?;
242 file.write_all(&payload)?;
243 file.sync_all()?;
244 *next_t += 1;
245 Ok(())
246 }
247
248 fn tx_range(&self, start: u64, end: Option<u64>) -> Result<Vec<TxRecord>, LogError> {
249 let _guard = self.next_t.read().expect("poisoned log lock");
250 Ok(read_merged(&self.dir, &self.name)?
251 .into_iter()
252 .filter(|r| r.t >= start && end.is_none_or(|e| r.t < e))
253 .collect())
254 }
255}
256
257fn version_path(dir: &Path, name: &str, version: u64) -> PathBuf {
258 if version == 0 {
259 dir.join(format!("{name}.log"))
260 } else {
261 dir.join(format!("{name}.v{version}.log"))
262 }
263}
264
265fn version_files(dir: &Path, name: &str) -> Vec<(u64, PathBuf)> {
267 let mut files = Vec::new();
268 let legacy = version_path(dir, name, 0);
269 if legacy.is_file() {
270 files.push((0, legacy));
271 }
272 let prefix = format!("{name}.v");
273 if let Ok(entries) = fs::read_dir(dir) {
274 for entry in entries.flatten() {
275 let file_name = entry.file_name();
276 let Some(text) = file_name.to_str() else {
277 continue;
278 };
279 if let Some(version) = text
280 .strip_prefix(&prefix)
281 .and_then(|rest| rest.strip_suffix(".log"))
282 .and_then(|v| v.parse::<u64>().ok())
283 {
284 if version > 0 {
285 files.push((version, entry.path()));
286 }
287 }
288 }
289 }
290 files.sort_by_key(|(version, _)| *version);
291 files
292}
293
294fn read_merged(dir: &Path, name: &str) -> Result<Vec<TxRecord>, LogError> {
297 let files = version_files(dir, name);
298 let mut per_file: Vec<Vec<TxRecord>> = Vec::with_capacity(files.len());
299 for (_, path) in &files {
300 per_file.push(read_records(path)?.0);
301 }
302 let mut cutoff = u64::MAX;
307 for records in per_file.iter_mut().rev() {
308 let first = records.first().map(|r| r.t);
309 records.retain(|r| r.t < cutoff);
310 if let Some(first) = first {
311 cutoff = cutoff.min(first);
312 }
313 }
314 let merged: Vec<TxRecord> = per_file.into_iter().flatten().collect();
315 for pair in merged.windows(2) {
316 if pair[1].t != pair[0].t + 1 {
317 return Err(LogError::Corrupt);
318 }
319 }
320 Ok(merged)
321}
322
323fn encode_record(record: &TxRecord) -> Vec<u8> {
324 let mut out = Vec::new();
325 out.extend_from_slice(&record.t.to_be_bytes());
326 out.extend_from_slice(&record.tx_instant.to_be_bytes());
327 out.extend_from_slice(&(record.datoms.len() as u64).to_be_bytes());
328 for d in &record.datoms {
329 out.extend_from_slice(&d.e.raw().to_be_bytes());
330 out.extend_from_slice(&d.a.raw().to_be_bytes());
331 out.extend_from_slice(&d.tx.raw().to_be_bytes());
332 out.push(u8::from(d.added));
333 let v = encode_value(&d.v);
334 out.extend_from_slice(&(v.len() as u64).to_be_bytes());
335 out.extend_from_slice(&v);
336 }
337 out
338}
339fn decode_record(mut bytes: &[u8]) -> Result<TxRecord, LogError> {
340 fn take<'a>(bytes: &mut &'a [u8], n: usize) -> Result<&'a [u8], LogError> {
341 let value = bytes.get(..n).ok_or(LogError::Corrupt)?;
342 *bytes = &bytes[n..];
343 Ok(value)
344 }
345 fn u64_be(bytes: &mut &[u8]) -> Result<u64, LogError> {
346 Ok(u64::from_be_bytes(
347 take(bytes, 8)?.try_into().map_err(|_| LogError::Corrupt)?,
348 ))
349 }
350 let t = u64_be(&mut bytes)?;
351 let tx_instant = i64::from_be_bytes(
352 take(&mut bytes, 8)?
353 .try_into()
354 .map_err(|_| LogError::Corrupt)?,
355 );
356 let count = u64_be(&mut bytes)?;
357 let mut datoms = Vec::new();
358 for _ in 0..count {
359 let e = EntityId::from_raw(u64_be(&mut bytes)?);
360 let a = EntityId::from_raw(u64_be(&mut bytes)?);
361 let tx = EntityId::from_raw(u64_be(&mut bytes)?);
362 let added = take(&mut bytes, 1)?[0] != 0;
363 let len = usize::try_from(u64_be(&mut bytes)?).map_err(|_| LogError::Corrupt)?;
364 let raw = take(&mut bytes, len)?;
365 let (v, used) = decode_value(raw).map_err(|_| LogError::Corrupt)?;
366 if used != len {
367 return Err(LogError::Corrupt);
368 }
369 datoms.push(Datom { e, a, v, tx, added });
370 }
371 if !bytes.is_empty() {
372 return Err(LogError::Corrupt);
373 }
374 Ok(TxRecord {
375 t,
376 tx_instant,
377 datoms,
378 })
379}
380fn read_records(path: &Path) -> Result<(Vec<TxRecord>, u64), LogError> {
386 let mut file = File::open(path)?;
387 let mut records = Vec::new();
388 let mut durable_len = 0_u64;
389 loop {
390 let mut len = [0; 8];
391 match file.read_exact(&mut len) {
392 Ok(()) => {}
393 Err(e) if e.kind() == io::ErrorKind::UnexpectedEof => break,
394 Err(e) => return Err(e.into()),
395 }
396 let len = usize::try_from(u64::from_be_bytes(len)).map_err(|_| LogError::Corrupt)?;
397 let mut payload = vec![0; len];
398 match file.read_exact(&mut payload) {
399 Ok(()) => {}
400 Err(e) if e.kind() == io::ErrorKind::UnexpectedEof => break,
401 Err(e) => return Err(e.into()),
402 }
403 records.push(decode_record(&payload)?);
404 durable_len += 8 + len as u64;
405 }
406 Ok((records, durable_len))
407}