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    fs::{self, File, OpenOptions},
9    io::{self, Read, Write},
10    path::{Path, PathBuf},
11    sync::{Arc, RwLock},
12};
13use thiserror::Error;
14
15/// One committed transaction record.
16#[derive(Clone, Debug, Eq, PartialEq)]
17pub struct TxRecord {
18    /// Monotonic transaction number.
19    pub t: u64,
20    /// Monotonic UTC millisecond timestamp.
21    pub tx_instant: i64,
22    /// Facts asserted/retracted by the transaction.
23    pub datoms: Vec<Datom>,
24}
25
26/// Log errors.
27#[derive(Debug, Error)]
28pub enum LogError {
29    /// Filesystem error.
30    #[error("log I/O failed: {0}")]
31    Io(#[from] io::Error),
32    /// Malformed or incomplete log data.
33    #[error("corrupt transaction log")]
34    Corrupt,
35}
36
37/// Common transaction log interface.
38pub trait TransactionLog: Send + Sync {
39    /// Durably appends exactly the next transaction.
40    ///
41    /// # Errors
42    /// Returns an error for I/O failure, corruption, or a non-contiguous `t`.
43    fn append(&self, record: &TxRecord) -> Result<(), LogError>;
44    /// Returns records in the half-open transaction range `[start, end)`.
45    ///
46    /// # Errors
47    /// Returns an error when stored records cannot be read or decoded.
48    fn tx_range(&self, start: u64, end: Option<u64>) -> Result<Vec<TxRecord>, LogError>;
49    /// Replays every committed record.
50    ///
51    /// # Errors
52    /// Returns an error when stored records cannot be read or decoded.
53    fn replay(&self) -> Result<Vec<TxRecord>, LogError> {
54        self.tx_range(0, None)
55    }
56}
57
58/// In-memory log implementation.
59#[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
82/// Filesystem append log. Each append is flushed and `fsync`ed before returning.
83///
84/// A crash mid-append leaves a torn, never-acked record at the tail; `open`
85/// truncates it away so replay stops at the durability point of the last
86/// acked transaction and later appends extend a clean tail.
87pub struct FileLog {
88    path: PathBuf,
89    next_t: RwLock<u64>,
90}
91impl FileLog {
92    /// Opens or creates a log file, dropping any torn tail left by a crash.
93    ///
94    /// # Errors
95    /// Returns an error if the file cannot be created or a fully written
96    /// record is corrupt.
97    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
143fn encode_record(record: &TxRecord) -> Vec<u8> {
144    let mut out = Vec::new();
145    out.extend_from_slice(&record.t.to_be_bytes());
146    out.extend_from_slice(&record.tx_instant.to_be_bytes());
147    out.extend_from_slice(&(record.datoms.len() as u64).to_be_bytes());
148    for d in &record.datoms {
149        out.extend_from_slice(&d.e.raw().to_be_bytes());
150        out.extend_from_slice(&d.a.raw().to_be_bytes());
151        out.extend_from_slice(&d.tx.raw().to_be_bytes());
152        out.push(u8::from(d.added));
153        let v = encode_value(&d.v);
154        out.extend_from_slice(&(v.len() as u64).to_be_bytes());
155        out.extend_from_slice(&v);
156    }
157    out
158}
159fn decode_record(mut bytes: &[u8]) -> Result<TxRecord, LogError> {
160    fn take<'a>(bytes: &mut &'a [u8], n: usize) -> Result<&'a [u8], LogError> {
161        let value = bytes.get(..n).ok_or(LogError::Corrupt)?;
162        *bytes = &bytes[n..];
163        Ok(value)
164    }
165    fn u64_be(bytes: &mut &[u8]) -> Result<u64, LogError> {
166        Ok(u64::from_be_bytes(
167            take(bytes, 8)?.try_into().map_err(|_| LogError::Corrupt)?,
168        ))
169    }
170    let t = u64_be(&mut bytes)?;
171    let tx_instant = i64::from_be_bytes(
172        take(&mut bytes, 8)?
173            .try_into()
174            .map_err(|_| LogError::Corrupt)?,
175    );
176    let count = u64_be(&mut bytes)?;
177    let mut datoms = Vec::new();
178    for _ in 0..count {
179        let e = EntityId::from_raw(u64_be(&mut bytes)?);
180        let a = EntityId::from_raw(u64_be(&mut bytes)?);
181        let tx = EntityId::from_raw(u64_be(&mut bytes)?);
182        let added = take(&mut bytes, 1)?[0] != 0;
183        let len = usize::try_from(u64_be(&mut bytes)?).map_err(|_| LogError::Corrupt)?;
184        let raw = take(&mut bytes, len)?;
185        let (v, used) = decode_value(raw).map_err(|_| LogError::Corrupt)?;
186        if used != len {
187            return Err(LogError::Corrupt);
188        }
189        datoms.push(Datom { e, a, v, tx, added });
190    }
191    if !bytes.is_empty() {
192        return Err(LogError::Corrupt);
193    }
194    Ok(TxRecord {
195        t,
196        tx_instant,
197        datoms,
198    })
199}
200/// Reads fully written records plus the byte length of that durable prefix.
201///
202/// A record cut short by a crash mid-append (truncated length prefix or
203/// payload) ends the scan; a fully present record that fails to decode is
204/// genuine corruption and errors.
205fn read_records(path: &Path) -> Result<(Vec<TxRecord>, u64), LogError> {
206    let mut file = File::open(path)?;
207    let mut records = Vec::new();
208    let mut durable_len = 0_u64;
209    loop {
210        let mut len = [0; 8];
211        match file.read_exact(&mut len) {
212            Ok(()) => {}
213            Err(e) if e.kind() == io::ErrorKind::UnexpectedEof => break,
214            Err(e) => return Err(e.into()),
215        }
216        let len = usize::try_from(u64::from_be_bytes(len)).map_err(|_| LogError::Corrupt)?;
217        let mut payload = vec![0; len];
218        match file.read_exact(&mut payload) {
219            Ok(()) => {}
220            Err(e) if e.kind() == io::ErrorKind::UnexpectedEof => break,
221            Err(e) => return Err(e.into()),
222        }
223        records.push(decode_record(&payload)?);
224        durable_len += 8 + len as u64;
225    }
226    Ok((records, durable_len))
227}