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
143/// A transaction log split into per-lease-version files for HA append
144/// isolation (see `docs/design/log-and-transactor.md`).
145///
146/// The active writer under lease version `V` appends only to
147/// `{name}.v{V}.log` (the pre-HA `{name}.log` reads as version 0). Readers
148/// merge the files in version order and drop any record in an older file
149/// whose `t` is at or past the first record of a later file: such records
150/// were appended by a deposed writer after a takeover and were never
151/// acknowledged, because acknowledgement re-verifies lease ownership after
152/// the durable append. A deposed writer therefore cannot corrupt or fork
153/// the log — its stale appends land in a file nobody considers current.
154pub struct VersionedLog {
155    dir: PathBuf,
156    name: String,
157    write_path: PathBuf,
158    next_t: RwLock<u64>,
159}
160
161impl VersionedLog {
162    /// Opens the log for writing under `write_version`, creating the
163    /// version file if needed and dropping any torn tail it carries.
164    /// Files of other versions are never modified.
165    ///
166    /// # Errors
167    /// Returns an error if files cannot be read/created or a fully written
168    /// record is corrupt.
169    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    /// Opens the log read-only (offline inspection, backup); appends fail.
193    ///
194    /// # Errors
195    /// Returns an error when the directory cannot be read or a fully
196    /// written record is corrupt.
197    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    /// Reports whether any log file exists for this database.
208    #[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    /// Deletes every version file for this database.
214    ///
215    /// # Errors
216    /// Returns an error when a file cannot be removed.
217    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
265/// Existing version files for `name`, sorted by version.
266fn 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
294/// Merges every version file, applying the takeover cutoff rule, and
295/// verifies the surviving sequence is contiguous.
296fn 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    // A record in an older file is dead once any later file starts at or
303    // below its t: every record acked under version v precedes the first
304    // record of every later version (the successor replayed it before
305    // choosing its own first t), so only never-acked stale appends die.
306    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}
380/// Reads fully written records plus the byte length of that durable prefix.
381///
382/// A record cut short by a crash mid-append (truncated length prefix or
383/// payload) ends the scan; a fully present record that fails to decode is
384/// genuine corruption and errors.
385fn 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}