Skip to main content

corium_log/
lib.rs

1//! Durable append-only transaction logs with replay and range scans.
2
3use corium_core::{
4    Datom, EntityId,
5    encoding::{decode_value, encode_value},
6};
7use std::{
8    collections::HashMap,
9    fs::{self, File, OpenOptions},
10    io::{self, Read, Write},
11    path::{Path, PathBuf},
12    sync::{Arc, Mutex, RwLock},
13};
14use thiserror::Error;
15
16/// One committed transaction record.
17#[derive(Clone, Debug, Eq, PartialEq)]
18pub struct TxRecord {
19    /// Monotonic transaction number.
20    pub t: u64,
21    /// Monotonic UTC millisecond timestamp.
22    pub tx_instant: i64,
23    /// Facts asserted/retracted by the transaction.
24    pub datoms: Vec<Datom>,
25}
26
27/// Log errors.
28#[derive(Debug, Error)]
29pub enum LogError {
30    /// Filesystem error.
31    #[error("log I/O failed: {0}")]
32    Io(#[from] io::Error),
33    /// Malformed or incomplete log data.
34    #[error("corrupt transaction log")]
35    Corrupt,
36}
37
38/// Common transaction log interface.
39pub trait TransactionLog: Send + Sync {
40    /// Durably appends exactly the next transaction.
41    ///
42    /// # Errors
43    /// Returns an error for I/O failure, corruption, or a non-contiguous `t`.
44    fn append(&self, record: &TxRecord) -> Result<(), LogError>;
45    /// Returns records in the half-open transaction range `[start, end)`.
46    ///
47    /// # Errors
48    /// Returns an error when stored records cannot be read or decoded.
49    fn tx_range(&self, start: u64, end: Option<u64>) -> Result<Vec<TxRecord>, LogError>;
50    /// Replays every committed record.
51    ///
52    /// # Errors
53    /// Returns an error when stored records cannot be read or decoded.
54    fn replay(&self) -> Result<Vec<TxRecord>, LogError> {
55        self.tx_range(0, None)
56    }
57}
58
59/// In-memory log implementation.
60#[derive(Clone, Default)]
61pub struct MemoryLog(Arc<RwLock<Vec<TxRecord>>>);
62impl TransactionLog for MemoryLog {
63    fn append(&self, record: &TxRecord) -> Result<(), LogError> {
64        let mut records = self.0.write().expect("poisoned log lock");
65        if records.last().map_or(1, |r| r.t + 1) != record.t {
66            return Err(LogError::Corrupt);
67        }
68        records.push(record.clone());
69        Ok(())
70    }
71    fn tx_range(&self, start: u64, end: Option<u64>) -> Result<Vec<TxRecord>, LogError> {
72        Ok(self
73            .0
74            .read()
75            .expect("poisoned log lock")
76            .iter()
77            .filter(|r| r.t >= start && end.is_none_or(|e| r.t < e))
78            .cloned()
79            .collect())
80    }
81}
82
83/// Filesystem append log. Each append is flushed and `fsync`ed before returning.
84///
85/// A crash mid-append leaves a torn, never-acked record at the tail; `open`
86/// truncates it away so replay stops at the durability point of the last
87/// acked transaction and later appends extend a clean tail.
88pub struct FileLog {
89    path: PathBuf,
90    next_t: RwLock<u64>,
91}
92impl FileLog {
93    /// Opens or creates a log file, dropping any torn tail left by a crash.
94    ///
95    /// # Errors
96    /// Returns an error if the file cannot be created or a fully written
97    /// record is corrupt.
98    pub fn open(path: impl AsRef<Path>) -> Result<Self, LogError> {
99        let path = path.as_ref().to_path_buf();
100        if let Some(parent) = path.parent() {
101            fs::create_dir_all(parent)?;
102        }
103        OpenOptions::new().create(true).append(true).open(&path)?;
104        let (records, durable_len) = read_records(&path)?;
105        if fs::metadata(&path)?.len() > durable_len {
106            let file = OpenOptions::new().write(true).open(&path)?;
107            file.set_len(durable_len)?;
108            file.sync_all()?;
109        }
110        Ok(Self {
111            path,
112            next_t: RwLock::new(records.last().map_or(1, |r| r.t + 1)),
113        })
114    }
115}
116impl TransactionLog for FileLog {
117    fn append(&self, record: &TxRecord) -> Result<(), LogError> {
118        let mut next_t = self.next_t.write().expect("poisoned log lock");
119        if *next_t != record.t {
120            return Err(LogError::Corrupt);
121        }
122        let payload = encode_record(record);
123        let mut file = OpenOptions::new().append(true).open(&self.path)?;
124        file.write_all(
125            &u64::try_from(payload.len())
126                .map_err(|_| LogError::Corrupt)?
127                .to_be_bytes(),
128        )?;
129        file.write_all(&payload)?;
130        file.sync_all()?;
131        *next_t += 1;
132        Ok(())
133    }
134    fn tx_range(&self, start: u64, end: Option<u64>) -> Result<Vec<TxRecord>, LogError> {
135        let _guard = self.next_t.read().expect("poisoned log lock");
136        Ok(read_records(&self.path)?
137            .0
138            .into_iter()
139            .filter(|r| r.t >= start && end.is_none_or(|e| r.t < e))
140            .collect())
141    }
142}
143
144/// A transaction log split into per-lease-version files for HA append
145/// isolation (see `docs/design/log-and-transactor.md`).
146///
147/// The active writer under lease version `V` appends only to
148/// `{name}.v{V}.log` (the pre-HA `{name}.log` reads as version 0). Readers
149/// merge the files in version order and drop any record in an older file
150/// whose `t` is at or past the first record of a later file: such records
151/// were appended by a deposed writer after a takeover and were never
152/// acknowledged, because acknowledgement re-verifies lease ownership after
153/// the durable append. A deposed writer therefore cannot corrupt or fork
154/// the log — its stale appends land in a file nobody considers current.
155pub struct VersionedLog {
156    dir: PathBuf,
157    name: String,
158    write_path: PathBuf,
159    next_t: RwLock<u64>,
160}
161
162impl VersionedLog {
163    /// Opens the log for writing under `write_version`, creating the
164    /// version file if needed and dropping any torn tail it carries.
165    /// Files of other versions are never modified.
166    ///
167    /// # Errors
168    /// Returns an error if files cannot be read/created or a fully written
169    /// record is corrupt.
170    pub fn open(dir: impl AsRef<Path>, name: &str, write_version: u64) -> Result<Self, LogError> {
171        let dir = dir.as_ref().to_path_buf();
172        fs::create_dir_all(&dir)?;
173        let write_path = version_path(&dir, name, write_version);
174        OpenOptions::new()
175            .create(true)
176            .append(true)
177            .open(&write_path)?;
178        let (_, durable_len) = read_records(&write_path)?;
179        if fs::metadata(&write_path)?.len() > durable_len {
180            let file = OpenOptions::new().write(true).open(&write_path)?;
181            file.set_len(durable_len)?;
182            file.sync_all()?;
183        }
184        let records = read_merged(&dir, name)?;
185        Ok(Self {
186            dir,
187            name: name.to_owned(),
188            write_path,
189            next_t: RwLock::new(records.last().map_or(1, |r| r.t + 1)),
190        })
191    }
192
193    /// Opens the log read-only (offline inspection, backup); appends fail.
194    ///
195    /// # Errors
196    /// Returns an error when the directory cannot be read or a fully
197    /// written record is corrupt.
198    pub fn open_read_only(dir: impl AsRef<Path>, name: &str) -> Result<Self, LogError> {
199        let dir = dir.as_ref().to_path_buf();
200        Ok(Self {
201            write_path: PathBuf::new(),
202            name: name.to_owned(),
203            next_t: RwLock::new(u64::MAX),
204            dir,
205        })
206    }
207
208    /// Reports whether any log file exists for this database.
209    #[must_use]
210    pub fn exists(dir: impl AsRef<Path>, name: &str) -> bool {
211        !version_files(dir.as_ref(), name).is_empty()
212    }
213
214    /// Deletes every version file for this database.
215    ///
216    /// # Errors
217    /// Returns an error when a file cannot be removed.
218    pub fn delete_all(dir: impl AsRef<Path>, name: &str) -> Result<(), LogError> {
219        for (_, path) in version_files(dir.as_ref(), name) {
220            match fs::remove_file(&path) {
221                Ok(()) => {}
222                Err(error) if error.kind() == io::ErrorKind::NotFound => {}
223                Err(error) => return Err(error.into()),
224            }
225        }
226        Ok(())
227    }
228}
229
230impl TransactionLog for VersionedLog {
231    fn append(&self, record: &TxRecord) -> Result<(), LogError> {
232        let mut next_t = self.next_t.write().expect("poisoned log lock");
233        if *next_t != record.t {
234            return Err(LogError::Corrupt);
235        }
236        let payload = encode_record(record);
237        let mut file = OpenOptions::new().append(true).open(&self.write_path)?;
238        file.write_all(
239            &u64::try_from(payload.len())
240                .map_err(|_| LogError::Corrupt)?
241                .to_be_bytes(),
242        )?;
243        file.write_all(&payload)?;
244        file.sync_all()?;
245        *next_t += 1;
246        Ok(())
247    }
248
249    fn tx_range(&self, start: u64, end: Option<u64>) -> Result<Vec<TxRecord>, LogError> {
250        let _guard = self.next_t.read().expect("poisoned log lock");
251        Ok(read_merged(&self.dir, &self.name)?
252            .into_iter()
253            .filter(|r| r.t >= start && end.is_none_or(|e| r.t < e))
254            .collect())
255    }
256}
257
258/// Applies the takeover cutoff rule to per-version record lists, in the same
259/// way [`read_merged`] does for on-disk files: a record in an older version
260/// dies once any later version begins at or below its `t`, dropping only the
261/// never-acked stale appends of a deposed writer.
262fn merge_versions(mut per_version: Vec<Vec<TxRecord>>) -> Vec<TxRecord> {
263    let mut cutoff = u64::MAX;
264    for records in per_version.iter_mut().rev() {
265        let first = records.first().map(|r| r.t);
266        records.retain(|r| r.t < cutoff);
267        if let Some(first) = first {
268            cutoff = cutoff.min(first);
269        }
270    }
271    per_version.into_iter().flatten().collect()
272}
273
274/// Shared store of one log's records, each tagged with the lease version it
275/// was appended under.
276type VersionedRecords = Arc<Mutex<Vec<(u64, TxRecord)>>>;
277
278/// Process-shared registry of in-memory transaction logs, keyed by database
279/// name. It plays the role the log directory plays for [`VersionedLog`]:
280/// opening the same name (under any lease version) reaches the same records,
281/// so a mem-backed transactor recovers state across `open`/`create` calls
282/// within one process. Cloning a registry shares its storage.
283#[derive(Clone, Default)]
284pub struct MemLogRegistry {
285    logs: Arc<Mutex<HashMap<String, VersionedRecords>>>,
286}
287
288impl MemLogRegistry {
289    /// Creates an empty registry.
290    #[must_use]
291    pub fn new() -> Self {
292        Self::default()
293    }
294
295    fn entry(&self, name: &str) -> VersionedRecords {
296        Arc::clone(
297            self.logs
298                .lock()
299                .unwrap_or_else(std::sync::PoisonError::into_inner)
300                .entry(name.to_owned())
301                .or_default(),
302        )
303    }
304
305    /// Opens the named log for writing under `write_version`, mirroring
306    /// [`VersionedLog::open`] with in-memory storage.
307    #[must_use]
308    pub fn open(&self, name: &str, write_version: u64) -> MemVersionedLog {
309        let records = self.entry(name);
310        let next_t = {
311            let guard = records
312                .lock()
313                .unwrap_or_else(std::sync::PoisonError::into_inner);
314            MemVersionedLog::merged(&guard)
315                .last()
316                .map_or(1, |r| r.t + 1)
317        };
318        MemVersionedLog {
319            records,
320            write_version,
321            next_t: Mutex::new(next_t),
322        }
323    }
324
325    /// Reports whether any records exist for the named log.
326    #[must_use]
327    pub fn exists(&self, name: &str) -> bool {
328        self.logs
329            .lock()
330            .unwrap_or_else(std::sync::PoisonError::into_inner)
331            .get(name)
332            .is_some_and(|entry| {
333                !entry
334                    .lock()
335                    .unwrap_or_else(std::sync::PoisonError::into_inner)
336                    .is_empty()
337            })
338    }
339
340    /// Discards every record for the named log.
341    pub fn delete_all(&self, name: &str) {
342        self.logs
343            .lock()
344            .unwrap_or_else(std::sync::PoisonError::into_inner)
345            .remove(name);
346    }
347}
348
349/// An in-memory transaction log with the same per-lease-version merge
350/// semantics as [`VersionedLog`], obtained from a [`MemLogRegistry`]. Used by
351/// the mem-backed transactor: fully ephemeral, confined to one process.
352pub struct MemVersionedLog {
353    records: VersionedRecords,
354    write_version: u64,
355    /// The next `t` this writer will accept, tracked per opened instance
356    /// exactly as [`VersionedLog`] does — a deposed writer keeps appending
357    /// under its own stale count, and the merge cutoff discards those records.
358    next_t: Mutex<u64>,
359}
360
361impl MemVersionedLog {
362    fn merged(records: &[(u64, TxRecord)]) -> Vec<TxRecord> {
363        let mut versions: Vec<u64> = records.iter().map(|(version, _)| *version).collect();
364        versions.sort_unstable();
365        versions.dedup();
366        let per_version = versions
367            .into_iter()
368            .map(|version| {
369                records
370                    .iter()
371                    .filter(|(record_version, _)| *record_version == version)
372                    .map(|(_, record)| record.clone())
373                    .collect::<Vec<_>>()
374            })
375            .collect();
376        merge_versions(per_version)
377    }
378}
379
380impl TransactionLog for MemVersionedLog {
381    fn append(&self, record: &TxRecord) -> Result<(), LogError> {
382        let mut next_t = self
383            .next_t
384            .lock()
385            .unwrap_or_else(std::sync::PoisonError::into_inner);
386        if *next_t != record.t {
387            return Err(LogError::Corrupt);
388        }
389        self.records
390            .lock()
391            .unwrap_or_else(std::sync::PoisonError::into_inner)
392            .push((self.write_version, record.clone()));
393        *next_t += 1;
394        Ok(())
395    }
396
397    fn tx_range(&self, start: u64, end: Option<u64>) -> Result<Vec<TxRecord>, LogError> {
398        let records = self
399            .records
400            .lock()
401            .unwrap_or_else(std::sync::PoisonError::into_inner);
402        Ok(Self::merged(&records)
403            .into_iter()
404            .filter(|r| r.t >= start && end.is_none_or(|e| r.t < e))
405            .collect())
406    }
407}
408
409fn version_path(dir: &Path, name: &str, version: u64) -> PathBuf {
410    if version == 0 {
411        dir.join(format!("{name}.log"))
412    } else {
413        dir.join(format!("{name}.v{version}.log"))
414    }
415}
416
417/// Existing version files for `name`, sorted by version.
418fn version_files(dir: &Path, name: &str) -> Vec<(u64, PathBuf)> {
419    let mut files = Vec::new();
420    let legacy = version_path(dir, name, 0);
421    if legacy.is_file() {
422        files.push((0, legacy));
423    }
424    let prefix = format!("{name}.v");
425    if let Ok(entries) = fs::read_dir(dir) {
426        for entry in entries.flatten() {
427            let file_name = entry.file_name();
428            let Some(text) = file_name.to_str() else {
429                continue;
430            };
431            if let Some(version) = text
432                .strip_prefix(&prefix)
433                .and_then(|rest| rest.strip_suffix(".log"))
434                .and_then(|v| v.parse::<u64>().ok())
435                && version > 0
436            {
437                files.push((version, entry.path()));
438            }
439        }
440    }
441    files.sort_by_key(|(version, _)| *version);
442    files
443}
444
445/// Merges every version file, applying the takeover cutoff rule, and
446/// verifies the surviving sequence is contiguous.
447fn read_merged(dir: &Path, name: &str) -> Result<Vec<TxRecord>, LogError> {
448    let files = version_files(dir, name);
449    let mut per_file: Vec<Vec<TxRecord>> = Vec::with_capacity(files.len());
450    for (_, path) in &files {
451        per_file.push(read_records(path)?.0);
452    }
453    // A record in an older file is dead once any later file starts at or
454    // below its t: every record acked under version v precedes the first
455    // record of every later version (the successor replayed it before
456    // choosing its own first t), so only never-acked stale appends die.
457    let merged = merge_versions(per_file);
458    for pair in merged.windows(2) {
459        if pair[1].t != pair[0].t + 1 {
460            return Err(LogError::Corrupt);
461        }
462    }
463    Ok(merged)
464}
465
466fn encode_record(record: &TxRecord) -> Vec<u8> {
467    let mut out = Vec::new();
468    out.extend_from_slice(&record.t.to_be_bytes());
469    out.extend_from_slice(&record.tx_instant.to_be_bytes());
470    out.extend_from_slice(&(record.datoms.len() as u64).to_be_bytes());
471    for d in &record.datoms {
472        out.extend_from_slice(&d.e.raw().to_be_bytes());
473        out.extend_from_slice(&d.a.raw().to_be_bytes());
474        out.extend_from_slice(&d.tx.raw().to_be_bytes());
475        out.push(u8::from(d.added));
476        let v = encode_value(&d.v);
477        out.extend_from_slice(&(v.len() as u64).to_be_bytes());
478        out.extend_from_slice(&v);
479    }
480    out
481}
482fn decode_record(mut bytes: &[u8]) -> Result<TxRecord, LogError> {
483    fn take<'a>(bytes: &mut &'a [u8], n: usize) -> Result<&'a [u8], LogError> {
484        let value = bytes.get(..n).ok_or(LogError::Corrupt)?;
485        *bytes = &bytes[n..];
486        Ok(value)
487    }
488    fn u64_be(bytes: &mut &[u8]) -> Result<u64, LogError> {
489        Ok(u64::from_be_bytes(
490            take(bytes, 8)?.try_into().map_err(|_| LogError::Corrupt)?,
491        ))
492    }
493    let t = u64_be(&mut bytes)?;
494    let tx_instant = i64::from_be_bytes(
495        take(&mut bytes, 8)?
496            .try_into()
497            .map_err(|_| LogError::Corrupt)?,
498    );
499    let count = u64_be(&mut bytes)?;
500    let mut datoms = Vec::new();
501    for _ in 0..count {
502        let e = EntityId::from_raw(u64_be(&mut bytes)?);
503        let a = EntityId::from_raw(u64_be(&mut bytes)?);
504        let tx = EntityId::from_raw(u64_be(&mut bytes)?);
505        let added = take(&mut bytes, 1)?[0] != 0;
506        let len = usize::try_from(u64_be(&mut bytes)?).map_err(|_| LogError::Corrupt)?;
507        let raw = take(&mut bytes, len)?;
508        let (v, used) = decode_value(raw).map_err(|_| LogError::Corrupt)?;
509        if used != len {
510            return Err(LogError::Corrupt);
511        }
512        datoms.push(Datom { e, a, v, tx, added });
513    }
514    if !bytes.is_empty() {
515        return Err(LogError::Corrupt);
516    }
517    Ok(TxRecord {
518        t,
519        tx_instant,
520        datoms,
521    })
522}
523/// Reads fully written records plus the byte length of that durable prefix.
524///
525/// A record cut short by a crash mid-append (truncated length prefix or
526/// payload) ends the scan; a fully present record that fails to decode is
527/// genuine corruption and errors.
528fn read_records(path: &Path) -> Result<(Vec<TxRecord>, u64), LogError> {
529    let mut file = File::open(path)?;
530    let mut records = Vec::new();
531    let mut durable_len = 0_u64;
532    loop {
533        let mut len = [0; 8];
534        match file.read_exact(&mut len) {
535            Ok(()) => {}
536            Err(e) if e.kind() == io::ErrorKind::UnexpectedEof => break,
537            Err(e) => return Err(e.into()),
538        }
539        let len = usize::try_from(u64::from_be_bytes(len)).map_err(|_| LogError::Corrupt)?;
540        let mut payload = vec![0; len];
541        match file.read_exact(&mut payload) {
542            Ok(()) => {}
543            Err(e) if e.kind() == io::ErrorKind::UnexpectedEof => break,
544            Err(e) => return Err(e.into()),
545        }
546        records.push(decode_record(&payload)?);
547        durable_len += 8 + len as u64;
548    }
549    Ok((records, durable_len))
550}