Skip to main content

corium_log/
lib.rs

1//! Durable append-only transaction logs with replay and range scans.
2
3use async_trait::async_trait;
4use corium_core::{
5    Datom, EntityId,
6    encoding::{decode_value, encode_value},
7};
8use std::{
9    collections::HashMap,
10    fs::{self, File, OpenOptions},
11    io::{self, Read, Write},
12    path::{Path, PathBuf},
13    sync::{Arc, Mutex, RwLock},
14};
15use thiserror::Error;
16
17const CHECKSUMMED_FRAME: u64 = 1 << 63;
18const FRAME_CHECKSUM_LEN: usize = size_of::<u32>();
19
20/// One committed transaction record.
21#[derive(Clone, Debug, Eq, PartialEq)]
22pub struct TxRecord {
23    /// Monotonic transaction number.
24    pub t: u64,
25    /// Monotonic UTC millisecond timestamp.
26    pub tx_instant: i64,
27    /// Facts asserted/retracted by the transaction.
28    pub datoms: Vec<Datom>,
29}
30
31/// Log errors.
32#[derive(Debug, Error)]
33pub enum LogError {
34    /// Filesystem error.
35    #[error("log I/O failed: {0}")]
36    Io(#[from] io::Error),
37    /// Malformed or incomplete log data.
38    #[error("corrupt transaction log")]
39    Corrupt,
40    /// Native store backend failure.
41    #[error("native transaction log store failed: {0}")]
42    Native(String),
43    /// The operation requires the asynchronous log interface.
44    #[error("this transaction log requires asynchronous access")]
45    AsyncOnly,
46}
47
48/// Common transaction log interface.
49#[async_trait]
50pub trait TransactionLog: Send + Sync {
51    /// Durably appends exactly the next transaction.
52    ///
53    /// # Errors
54    /// Returns an error for I/O failure, corruption, or a non-contiguous `t`.
55    fn append(&self, record: &TxRecord) -> Result<(), LogError>;
56    /// Durably appends exactly the next transaction without blocking an async
57    /// runtime worker. Synchronous logs use [`Self::append`] by default;
58    /// storage-backed logs override this method and await their backend.
59    ///
60    /// # Errors
61    /// Returns an error for I/O failure, corruption, or a non-contiguous `t`.
62    async fn append_async(&self, record: &TxRecord) -> Result<(), LogError> {
63        self.append(record)
64    }
65    /// Durably appends a contiguous run of transactions under a single
66    /// durability boundary where the backend supports one (one `fsync`, one
67    /// object, one database transaction), so group commit amortizes the
68    /// per-append cost across the batch. `records` must be contiguous in `t`
69    /// starting at the log's next expected `t`; an empty slice is a no-op. The
70    /// default appends them one at a time; batching backends override this.
71    ///
72    /// # Errors
73    /// Returns an error for I/O failure, corruption, or a non-contiguous `t`.
74    async fn append_batch_async(&self, records: &[TxRecord]) -> Result<(), LogError> {
75        for record in records {
76            self.append_async(record).await?;
77        }
78        Ok(())
79    }
80    /// Returns records in the half-open transaction range `[start, end)`.
81    ///
82    /// # Errors
83    /// Returns an error when stored records cannot be read or decoded.
84    fn tx_range(&self, start: u64, end: Option<u64>) -> Result<Vec<TxRecord>, LogError>;
85    /// Asynchronous form of [`Self::tx_range`].
86    ///
87    /// # Errors
88    /// Returns an error when stored records cannot be read or decoded.
89    async fn tx_range_async(
90        &self,
91        start: u64,
92        end: Option<u64>,
93    ) -> Result<Vec<TxRecord>, LogError> {
94        self.tx_range(start, end)
95    }
96    /// Replays every committed record.
97    ///
98    /// # Errors
99    /// Returns an error when stored records cannot be read or decoded.
100    fn replay(&self) -> Result<Vec<TxRecord>, LogError> {
101        self.tx_range(0, None)
102    }
103    /// Asynchronously replays every committed record.
104    ///
105    /// # Errors
106    /// Returns an error when stored records cannot be read or decoded.
107    async fn replay_async(&self) -> Result<Vec<TxRecord>, LogError> {
108        self.tx_range_async(0, None).await
109    }
110}
111
112/// In-memory log implementation.
113#[derive(Clone, Default)]
114pub struct MemoryLog(Arc<RwLock<Vec<TxRecord>>>);
115impl TransactionLog for MemoryLog {
116    fn append(&self, record: &TxRecord) -> Result<(), LogError> {
117        let mut records = self.0.write().expect("poisoned log lock");
118        if records.last().map_or(1, |r| r.t + 1) != record.t {
119            return Err(LogError::Corrupt);
120        }
121        records.push(record.clone());
122        Ok(())
123    }
124    fn tx_range(&self, start: u64, end: Option<u64>) -> Result<Vec<TxRecord>, LogError> {
125        Ok(self
126            .0
127            .read()
128            .expect("poisoned log lock")
129            .iter()
130            .filter(|r| r.t >= start && end.is_none_or(|e| r.t < e))
131            .cloned()
132            .collect())
133    }
134}
135
136/// Filesystem append log. Each append is flushed and `fsync`ed before returning.
137///
138/// A crash mid-append leaves a torn, never-acked record at the tail; `open`
139/// truncates it away so replay stops at the durability point of the last
140/// acked transaction and later appends extend a clean tail.
141pub struct FileLog {
142    path: PathBuf,
143    next_t: RwLock<u64>,
144}
145impl FileLog {
146    /// Opens or creates a log file, dropping any torn tail left by a crash.
147    ///
148    /// # Errors
149    /// Returns an error if the file cannot be created or a fully written
150    /// record is corrupt.
151    pub fn open(path: impl AsRef<Path>) -> Result<Self, LogError> {
152        let path = path.as_ref().to_path_buf();
153        if let Some(parent) = path.parent() {
154            fs::create_dir_all(parent)?;
155        }
156        OpenOptions::new().create(true).append(true).open(&path)?;
157        let (records, durable_len) = read_records(&path)?;
158        if fs::metadata(&path)?.len() > durable_len {
159            let file = OpenOptions::new().write(true).open(&path)?;
160            file.set_len(durable_len)?;
161            file.sync_all()?;
162        }
163        Ok(Self {
164            path,
165            next_t: RwLock::new(records.last().map_or(1, |r| r.t + 1)),
166        })
167    }
168}
169impl TransactionLog for FileLog {
170    fn append(&self, record: &TxRecord) -> Result<(), LogError> {
171        let mut next_t = self.next_t.write().expect("poisoned log lock");
172        if *next_t != record.t {
173            return Err(LogError::Corrupt);
174        }
175        let mut frame = Vec::new();
176        append_framed_record(&mut frame, record)?;
177        let mut file = OpenOptions::new().append(true).open(&self.path)?;
178        file.write_all(&frame)?;
179        file.sync_all()?;
180        *next_t += 1;
181        Ok(())
182    }
183    fn tx_range(&self, start: u64, end: Option<u64>) -> Result<Vec<TxRecord>, LogError> {
184        let _guard = self.next_t.read().expect("poisoned log lock");
185        Ok(read_records(&self.path)?
186            .0
187            .into_iter()
188            .filter(|r| r.t >= start && end.is_none_or(|e| r.t < e))
189            .collect())
190    }
191}
192
193/// A transaction log split into per-lease-version files for HA append
194/// isolation (see `docs/design/log-and-transactor.md`).
195///
196/// The active writer under lease version `V` appends only to
197/// `{name}.v{V}.log` (the pre-HA `{name}.log` reads as version 0). Readers
198/// merge the files in version order and drop any record in an older file
199/// whose `t` is at or past the first record of a later file: such records
200/// were appended by a deposed writer after a takeover and were never
201/// acknowledged, because acknowledgement re-verifies lease ownership after
202/// the durable append. A deposed writer therefore cannot corrupt or fork
203/// the log — its stale appends land in a file nobody considers current.
204pub struct VersionedLog {
205    dir: PathBuf,
206    name: String,
207    write_path: PathBuf,
208    next_t: RwLock<u64>,
209}
210
211impl VersionedLog {
212    /// Opens the log for writing under `write_version`, creating the
213    /// version file if needed and dropping any torn tail it carries.
214    /// Files of other versions are never modified.
215    ///
216    /// # Errors
217    /// Returns an error if files cannot be read/created or a fully written
218    /// record is corrupt.
219    pub fn open(dir: impl AsRef<Path>, name: &str, write_version: u64) -> Result<Self, LogError> {
220        let dir = dir.as_ref().to_path_buf();
221        fs::create_dir_all(&dir)?;
222        let write_path = version_path(&dir, name, write_version);
223        OpenOptions::new()
224            .create(true)
225            .append(true)
226            .open(&write_path)?;
227        let (_, durable_len) = read_records(&write_path)?;
228        if fs::metadata(&write_path)?.len() > durable_len {
229            let file = OpenOptions::new().write(true).open(&write_path)?;
230            file.set_len(durable_len)?;
231            file.sync_all()?;
232        }
233        let records = read_merged(&dir, name)?;
234        Ok(Self {
235            dir,
236            name: name.to_owned(),
237            write_path,
238            next_t: RwLock::new(records.last().map_or(1, |r| r.t + 1)),
239        })
240    }
241
242    /// Opens the log read-only (independent inspection or backup); appends fail.
243    ///
244    /// # Errors
245    /// Returns an error when the directory cannot be read or a fully
246    /// written record is corrupt.
247    pub fn open_read_only(dir: impl AsRef<Path>, name: &str) -> Result<Self, LogError> {
248        let dir = dir.as_ref().to_path_buf();
249        Ok(Self {
250            write_path: PathBuf::new(),
251            name: name.to_owned(),
252            next_t: RwLock::new(u64::MAX),
253            dir,
254        })
255    }
256
257    /// Reports whether any log file exists for this database.
258    #[must_use]
259    pub fn exists(dir: impl AsRef<Path>, name: &str) -> bool {
260        !version_files(dir.as_ref(), name).is_empty()
261    }
262
263    /// Deletes every version file for this database.
264    ///
265    /// # Errors
266    /// Returns an error when a file cannot be removed.
267    pub fn delete_all(dir: impl AsRef<Path>, name: &str) -> Result<(), LogError> {
268        for (_, path) in version_files(dir.as_ref(), name) {
269            match fs::remove_file(&path) {
270                Ok(()) => {}
271                Err(error) if error.kind() == io::ErrorKind::NotFound => {}
272                Err(error) => return Err(error.into()),
273            }
274        }
275        Ok(())
276    }
277}
278
279impl TransactionLog for VersionedLog {
280    fn append(&self, record: &TxRecord) -> Result<(), LogError> {
281        let mut next_t = self.next_t.write().expect("poisoned log lock");
282        if *next_t != record.t {
283            return Err(LogError::Corrupt);
284        }
285        let mut frame = Vec::new();
286        append_framed_record(&mut frame, record)?;
287        let mut file = OpenOptions::new().append(true).open(&self.write_path)?;
288        file.write_all(&frame)?;
289        file.sync_all()?;
290        *next_t += 1;
291        Ok(())
292    }
293
294    fn tx_range(&self, start: u64, end: Option<u64>) -> Result<Vec<TxRecord>, LogError> {
295        let _guard = self.next_t.read().expect("poisoned log lock");
296        Ok(read_merged(&self.dir, &self.name)?
297            .into_iter()
298            .filter(|r| r.t >= start && end.is_none_or(|e| r.t < e))
299            .collect())
300    }
301}
302
303/// Applies the takeover cutoff rule to per-version record lists, in the same
304/// way [`read_merged`] does for on-disk files: a record in an older version
305/// dies once any later version begins at or below its `t`, dropping only the
306/// never-acked stale appends of a deposed writer.
307fn merge_versions(mut per_version: Vec<Vec<TxRecord>>) -> Vec<TxRecord> {
308    let mut cutoff = u64::MAX;
309    for records in per_version.iter_mut().rev() {
310        let first = records.first().map(|r| r.t);
311        records.retain(|r| r.t < cutoff);
312        if let Some(first) = first {
313            cutoff = cutoff.min(first);
314        }
315    }
316    per_version.into_iter().flatten().collect()
317}
318
319/// Asynchronous object store for transaction-log records.
320///
321/// Implementations adapt the same native storage system used for blobs and
322/// roots. The live log is written **one object per transaction** — a
323/// create-only write keyed `(name, version, t)` whose success is the
324/// durability point — so an append is O(1) (a small insert) instead of a
325/// read-modify-write of a growing chunk. On a SQL backend that is a
326/// row-per-commit insert; on an object store, a create-only `PUT`.
327///
328/// Earlier releases wrote a different layout: a sequence of *chunk* objects
329/// `(name, version, chunk)`, each a run of framed records, appended in place
330/// and rolled at a size cap. Those objects are still read, read-only, through
331/// [`Self::list_legacy_chunks`] / [`Self::read_legacy_chunk`], so a log
332/// written by an older binary keeps replaying after an upgrade; new records
333/// are always written in the per-transaction layout.
334#[async_trait]
335pub trait NativeLogStorage: Send + Sync {
336    /// Create-only, atomic write of one contiguous batch of transactions as a
337    /// single object, keyed by the batch's last `t`. Each element is
338    /// `(t, framed_bytes)` in ascending `t`; the object holds their framed
339    /// bytes concatenated (the same encoding a multi-record chunk uses).
340    /// Returns `Ok(true)` when written and `Ok(false)` when an object already
341    /// exists for that last-`t` — a lost create race or a retry of an
342    /// already-durable batch. The create-only condition is the log's fence: a
343    /// given `(version, last t)` is written at most once, and the batch is
344    /// durable in full or not at all.
345    ///
346    /// # Errors
347    /// Returns an error when the native backend cannot publish the object.
348    async fn put_batch(
349        &self,
350        name: &str,
351        version: u64,
352        records: &[(u64, Vec<u8>)],
353    ) -> Result<bool, LogError>;
354    /// Reads the bytes of the batch object keyed by last-`t` `t` for
355    /// `(name, version)`.
356    ///
357    /// # Errors
358    /// Returns an error when the native backend cannot read the object.
359    async fn read_record(
360        &self,
361        name: &str,
362        version: u64,
363        t: u64,
364    ) -> Result<Option<Vec<u8>>, LogError>;
365    /// Lists every `(version, last-t)` batch object present for `name`.
366    ///
367    /// # Errors
368    /// Returns an error when the native backend cannot enumerate log objects
369    /// or returns an invalid identifier.
370    async fn list_records(&self, name: &str) -> Result<Vec<(u64, u64)>, LogError>;
371    /// Reads one legacy chunk object (pre-per-record layout), for read-only
372    /// replay of logs written by older binaries.
373    ///
374    /// # Errors
375    /// Returns an error when the native backend cannot read the chunk object.
376    async fn read_legacy_chunk(
377        &self,
378        name: &str,
379        version: u64,
380        chunk: u64,
381    ) -> Result<Option<Vec<u8>>, LogError>;
382    /// Lists every legacy `(version, chunk)` object present for `name`, for
383    /// read-only replay of older logs. Returns an empty list on a store that
384    /// only ever wrote the per-record layout.
385    ///
386    /// # Errors
387    /// Returns an error when the native backend cannot enumerate log objects
388    /// or returns an invalid identifier.
389    async fn list_legacy_chunks(&self, name: &str) -> Result<Vec<(u64, u64)>, LogError>;
390    /// Deletes every log object for `name`, in both layouts.
391    ///
392    /// # Errors
393    /// Returns an error when the native backend cannot remove an object.
394    async fn delete_all(&self, name: &str) -> Result<(), LogError>;
395}
396
397/// Versioned transaction log backed by a native key/value-style store.
398///
399/// Each transaction is written as its own create-only object keyed
400/// `(name, write_version, t)`, so an append is a single small insert with no
401/// read and no growing buffer. The writer is the sole appender under its lease
402/// version (the fence gives each active owner its own version; a deposed
403/// writer's stale appends land in a version the takeover cutoff discards), so
404/// tracking `next_t` in memory is all the append state required.
405pub struct NativeVersionedLog<S: ?Sized> {
406    storage: Arc<S>,
407    name: String,
408    write_version: u64,
409    read_only: bool,
410    /// Next `t` this writer will accept; also serializes concurrent appends.
411    next_t: tokio::sync::Mutex<u64>,
412}
413
414impl<S: NativeLogStorage + ?Sized + 'static> NativeVersionedLog<S> {
415    /// Opens the log for writing under `write_version`.
416    ///
417    /// # Errors
418    /// Returns an error when stored records cannot be read or decoded.
419    pub async fn open(storage: Arc<S>, name: &str, write_version: u64) -> Result<Self, LogError> {
420        // The merged view across every version and both layouts establishes the
421        // next `t` — the takeover cutoff may place it past this writer's own
422        // last record.
423        let records = read_native_merged(storage.as_ref(), name).await?;
424        let next_t = records.last().map_or(1, |r| r.t + 1);
425        Ok(Self {
426            storage,
427            name: name.to_owned(),
428            write_version,
429            read_only: false,
430            next_t: tokio::sync::Mutex::new(next_t),
431        })
432    }
433
434    /// Opens the log for read-only range replay without first scanning it to
435    /// initialize writer state.
436    #[must_use]
437    pub fn open_read_only(storage: Arc<S>, name: &str) -> Self {
438        Self {
439            storage,
440            name: name.to_owned(),
441            write_version: 0,
442            read_only: true,
443            next_t: tokio::sync::Mutex::new(0),
444        }
445    }
446}
447
448#[async_trait]
449impl<S: NativeLogStorage + ?Sized + 'static> TransactionLog for NativeVersionedLog<S> {
450    fn append(&self, record: &TxRecord) -> Result<(), LogError> {
451        let _ = record;
452        Err(LogError::AsyncOnly)
453    }
454
455    async fn append_async(&self, record: &TxRecord) -> Result<(), LogError> {
456        self.append_batch_async(std::slice::from_ref(record)).await
457    }
458
459    async fn append_batch_async(&self, records: &[TxRecord]) -> Result<(), LogError> {
460        if self.read_only {
461            return Err(LogError::Native("transaction log is read-only".into()));
462        }
463        if records.is_empty() {
464            return Ok(());
465        }
466        let mut next_t = self.next_t.lock().await;
467        // The batch must be exactly the next contiguous run of transactions.
468        for (offset, record) in records.iter().enumerate() {
469            if record.t != *next_t + offset as u64 {
470                return Err(LogError::Corrupt);
471            }
472        }
473        let framed = records
474            .iter()
475            .map(|record| {
476                let mut bytes = Vec::new();
477                append_framed_record(&mut bytes, record)?;
478                Ok((record.t, bytes))
479            })
480            .collect::<Result<Vec<_>, LogError>>()?;
481        // Create-only write of the whole batch as one object. As the sole
482        // appender under this lease version, an object that already exists for
483        // this last-`t` is a duplicate or a racing writer under our version —
484        // never a legitimate append — so reject it rather than overwrite.
485        if !self
486            .storage
487            .put_batch(&self.name, self.write_version, &framed)
488            .await?
489        {
490            return Err(LogError::Corrupt);
491        }
492        *next_t += records.len() as u64;
493        Ok(())
494    }
495
496    fn tx_range(&self, start: u64, end: Option<u64>) -> Result<Vec<TxRecord>, LogError> {
497        let _ = (start, end);
498        Err(LogError::AsyncOnly)
499    }
500
501    async fn tx_range_async(
502        &self,
503        start: u64,
504        end: Option<u64>,
505    ) -> Result<Vec<TxRecord>, LogError> {
506        // Range/replay must merge every version (for the takeover cutoff), so
507        // they read the store; the lock only serializes them with appends.
508        let _guard = self.next_t.lock().await;
509        Ok(read_native_merged(self.storage.as_ref(), &self.name)
510            .await?
511            .into_iter()
512            .filter(|r| r.t >= start && end.is_none_or(|e| r.t < e))
513            .collect())
514    }
515}
516
517async fn read_native_merged<S: NativeLogStorage + ?Sized>(
518    storage: &S,
519    name: &str,
520) -> Result<Vec<TxRecord>, LogError> {
521    use std::collections::BTreeMap;
522
523    // Gather every version's records from both layouts, keyed by version so the
524    // cross-version takeover cutoff below sees them in ascending version order.
525    let mut per_version: BTreeMap<u64, Vec<TxRecord>> = BTreeMap::new();
526
527    // Legacy chunk objects (read-only): a version's chunks concatenate in chunk
528    // order, which is the order they were filled — i.e. transaction order.
529    // Empty on a store that only ever wrote the per-record layout.
530    let mut chunks = storage.list_legacy_chunks(name).await?;
531    chunks.sort_unstable();
532    for (version, chunk) in chunks {
533        let bytes = storage
534            .read_legacy_chunk(name, version, chunk)
535            .await?
536            .unwrap_or_default();
537        per_version
538            .entry(version)
539            .or_default()
540            .extend(decode_framed_records(&bytes)?);
541    }
542
543    // Per-record objects, one framed record each.
544    let mut records = storage.list_records(name).await?;
545    records.sort_unstable();
546    for (version, t) in records {
547        let bytes = storage
548            .read_record(name, version, t)
549            .await?
550            .unwrap_or_default();
551        per_version
552            .entry(version)
553            .or_default()
554            .extend(decode_framed_records(&bytes)?);
555    }
556
557    // Order each version's records by `t` (a version is written in a single
558    // layout in practice; sorting keeps even a version that carries both — a
559    // legacy tail then per-record appends — correct), then apply the takeover
560    // cutoff across versions.
561    let per_version: Vec<Vec<TxRecord>> = per_version
562        .into_values()
563        .map(|mut records| {
564            records.sort_by_key(|record| record.t);
565            records
566        })
567        .collect();
568    let merged = merge_versions(per_version);
569    for pair in merged.windows(2) {
570        if pair[1].t != pair[0].t + 1 {
571            return Err(LogError::Corrupt);
572        }
573    }
574    Ok(merged)
575}
576
577/// Shared store of one log's records, each tagged with the lease version it
578/// was appended under.
579type VersionedRecords = Arc<Mutex<Vec<(u64, TxRecord)>>>;
580
581/// Process-shared registry of in-memory transaction logs, keyed by database
582/// name. It plays the role the log directory plays for [`VersionedLog`]:
583/// opening the same name (under any lease version) reaches the same records,
584/// so a mem-backed transactor recovers state across `open`/`create` calls
585/// within one process. Cloning a registry shares its storage.
586#[derive(Clone, Default)]
587pub struct MemLogRegistry {
588    logs: Arc<Mutex<HashMap<String, VersionedRecords>>>,
589}
590
591impl MemLogRegistry {
592    /// Creates an empty registry.
593    #[must_use]
594    pub fn new() -> Self {
595        Self::default()
596    }
597
598    fn entry(&self, name: &str) -> VersionedRecords {
599        Arc::clone(
600            self.logs
601                .lock()
602                .unwrap_or_else(std::sync::PoisonError::into_inner)
603                .entry(name.to_owned())
604                .or_default(),
605        )
606    }
607
608    /// Opens the named log for writing under `write_version`, mirroring
609    /// [`VersionedLog::open`] with in-memory storage.
610    #[must_use]
611    pub fn open(&self, name: &str, write_version: u64) -> MemVersionedLog {
612        let records = self.entry(name);
613        let next_t = {
614            let guard = records
615                .lock()
616                .unwrap_or_else(std::sync::PoisonError::into_inner);
617            MemVersionedLog::merged(&guard)
618                .last()
619                .map_or(1, |r| r.t + 1)
620        };
621        MemVersionedLog {
622            records,
623            write_version,
624            next_t: Mutex::new(next_t),
625        }
626    }
627
628    /// Reports whether any records exist for the named log.
629    #[must_use]
630    pub fn exists(&self, name: &str) -> bool {
631        self.logs
632            .lock()
633            .unwrap_or_else(std::sync::PoisonError::into_inner)
634            .get(name)
635            .is_some_and(|entry| {
636                !entry
637                    .lock()
638                    .unwrap_or_else(std::sync::PoisonError::into_inner)
639                    .is_empty()
640            })
641    }
642
643    /// Discards every record for the named log.
644    pub fn delete_all(&self, name: &str) {
645        self.logs
646            .lock()
647            .unwrap_or_else(std::sync::PoisonError::into_inner)
648            .remove(name);
649    }
650}
651
652/// An in-memory transaction log with the same per-lease-version merge
653/// semantics as [`VersionedLog`], obtained from a [`MemLogRegistry`]. Used by
654/// the mem-backed transactor: fully ephemeral, confined to one process.
655pub struct MemVersionedLog {
656    records: VersionedRecords,
657    write_version: u64,
658    /// The next `t` this writer will accept, tracked per opened instance
659    /// exactly as [`VersionedLog`] does — a deposed writer keeps appending
660    /// under its own stale count, and the merge cutoff discards those records.
661    next_t: Mutex<u64>,
662}
663
664impl MemVersionedLog {
665    fn merged(records: &[(u64, TxRecord)]) -> Vec<TxRecord> {
666        let mut versions: Vec<u64> = records.iter().map(|(version, _)| *version).collect();
667        versions.sort_unstable();
668        versions.dedup();
669        let per_version = versions
670            .into_iter()
671            .map(|version| {
672                records
673                    .iter()
674                    .filter(|(record_version, _)| *record_version == version)
675                    .map(|(_, record)| record.clone())
676                    .collect::<Vec<_>>()
677            })
678            .collect();
679        merge_versions(per_version)
680    }
681}
682
683impl TransactionLog for MemVersionedLog {
684    fn append(&self, record: &TxRecord) -> Result<(), LogError> {
685        let mut next_t = self
686            .next_t
687            .lock()
688            .unwrap_or_else(std::sync::PoisonError::into_inner);
689        if *next_t != record.t {
690            return Err(LogError::Corrupt);
691        }
692        self.records
693            .lock()
694            .unwrap_or_else(std::sync::PoisonError::into_inner)
695            .push((self.write_version, record.clone()));
696        *next_t += 1;
697        Ok(())
698    }
699
700    fn tx_range(&self, start: u64, end: Option<u64>) -> Result<Vec<TxRecord>, LogError> {
701        let records = self
702            .records
703            .lock()
704            .unwrap_or_else(std::sync::PoisonError::into_inner);
705        Ok(Self::merged(&records)
706            .into_iter()
707            .filter(|r| r.t >= start && end.is_none_or(|e| r.t < e))
708            .collect())
709    }
710}
711
712fn version_path(dir: &Path, name: &str, version: u64) -> PathBuf {
713    if version == 0 {
714        dir.join(format!("{name}.log"))
715    } else {
716        dir.join(format!("{name}.v{version}.log"))
717    }
718}
719
720/// Existing version files for `name`, sorted by version.
721fn version_files(dir: &Path, name: &str) -> Vec<(u64, PathBuf)> {
722    let mut files = Vec::new();
723    let legacy = version_path(dir, name, 0);
724    if legacy.is_file() {
725        files.push((0, legacy));
726    }
727    let prefix = format!("{name}.v");
728    if let Ok(entries) = fs::read_dir(dir) {
729        for entry in entries.flatten() {
730            let file_name = entry.file_name();
731            let Some(text) = file_name.to_str() else {
732                continue;
733            };
734            if let Some(version) = text
735                .strip_prefix(&prefix)
736                .and_then(|rest| rest.strip_suffix(".log"))
737                .and_then(|v| v.parse::<u64>().ok())
738                && version > 0
739            {
740                files.push((version, entry.path()));
741            }
742        }
743    }
744    files.sort_by_key(|(version, _)| *version);
745    files
746}
747
748/// Merges every version file, applying the takeover cutoff rule, and
749/// verifies the surviving sequence is contiguous.
750fn read_merged(dir: &Path, name: &str) -> Result<Vec<TxRecord>, LogError> {
751    let files = version_files(dir, name);
752    let mut per_file: Vec<Vec<TxRecord>> = Vec::with_capacity(files.len());
753    for (_, path) in &files {
754        per_file.push(read_records(path)?.0);
755    }
756    // A record in an older file is dead once any later file starts at or
757    // below its t: every record acked under version v precedes the first
758    // record of every later version (the successor replayed it before
759    // choosing its own first t), so only never-acked stale appends die.
760    let merged = merge_versions(per_file);
761    for pair in merged.windows(2) {
762        if pair[1].t != pair[0].t + 1 {
763            return Err(LogError::Corrupt);
764        }
765    }
766    Ok(merged)
767}
768
769fn encode_record(record: &TxRecord) -> Vec<u8> {
770    let mut out = Vec::new();
771    out.extend_from_slice(&record.t.to_be_bytes());
772    out.extend_from_slice(&record.tx_instant.to_be_bytes());
773    out.extend_from_slice(&(record.datoms.len() as u64).to_be_bytes());
774    for d in &record.datoms {
775        out.extend_from_slice(&d.e.raw().to_be_bytes());
776        out.extend_from_slice(&d.a.raw().to_be_bytes());
777        out.extend_from_slice(&d.tx.raw().to_be_bytes());
778        out.push(u8::from(d.added));
779        let v = encode_value(&d.v);
780        out.extend_from_slice(&(v.len() as u64).to_be_bytes());
781        out.extend_from_slice(&v);
782    }
783    out
784}
785fn decode_record(mut bytes: &[u8]) -> Result<TxRecord, LogError> {
786    fn take<'a>(bytes: &mut &'a [u8], n: usize) -> Result<&'a [u8], LogError> {
787        let value = bytes.get(..n).ok_or(LogError::Corrupt)?;
788        *bytes = &bytes[n..];
789        Ok(value)
790    }
791    fn u64_be(bytes: &mut &[u8]) -> Result<u64, LogError> {
792        Ok(u64::from_be_bytes(
793            take(bytes, 8)?.try_into().map_err(|_| LogError::Corrupt)?,
794        ))
795    }
796    let t = u64_be(&mut bytes)?;
797    let tx_instant = i64::from_be_bytes(
798        take(&mut bytes, 8)?
799            .try_into()
800            .map_err(|_| LogError::Corrupt)?,
801    );
802    let count = u64_be(&mut bytes)?;
803    let mut datoms = Vec::new();
804    for _ in 0..count {
805        let e = EntityId::from_raw(u64_be(&mut bytes)?);
806        let a = EntityId::from_raw(u64_be(&mut bytes)?);
807        let tx = EntityId::from_raw(u64_be(&mut bytes)?);
808        let added = take(&mut bytes, 1)?[0] != 0;
809        let len = usize::try_from(u64_be(&mut bytes)?).map_err(|_| LogError::Corrupt)?;
810        let raw = take(&mut bytes, len)?;
811        let (v, used) = decode_value(raw).map_err(|_| LogError::Corrupt)?;
812        if used != len {
813            return Err(LogError::Corrupt);
814        }
815        datoms.push(Datom { e, a, v, tx, added });
816    }
817    if !bytes.is_empty() {
818        return Err(LogError::Corrupt);
819    }
820    Ok(TxRecord {
821        t,
822        tx_instant,
823        datoms,
824    })
825}
826
827fn frame_header(payload_len: usize) -> Result<[u8; 8], LogError> {
828    let payload_len = u64::try_from(payload_len).map_err(|_| LogError::Corrupt)?;
829    if payload_len & CHECKSUMMED_FRAME != 0 {
830        return Err(LogError::Corrupt);
831    }
832    Ok((payload_len | CHECKSUMMED_FRAME).to_be_bytes())
833}
834
835fn frame_payload_len(header: [u8; 8]) -> Result<(usize, bool), LogError> {
836    let encoded = u64::from_be_bytes(header);
837    let checksummed = encoded & CHECKSUMMED_FRAME != 0;
838    let payload_len = encoded & !CHECKSUMMED_FRAME;
839    Ok((
840        usize::try_from(payload_len).map_err(|_| LogError::Corrupt)?,
841        checksummed,
842    ))
843}
844
845fn frame_checksum(header: [u8; 8], payload: &[u8]) -> u32 {
846    crc32c::crc32c_append(crc32c::crc32c(&header), payload)
847}
848
849/// Reads fully written records plus the byte length of that durable prefix.
850///
851/// A record cut short by a crash mid-append (truncated length prefix or
852/// payload/checksum) ends the scan; a fully present record with a checksum
853/// mismatch or invalid payload is genuine corruption and errors. Legacy
854/// length-only records remain readable, while newly written records set the
855/// high bit of the length word and carry a trailing CRC32C.
856fn read_records(path: &Path) -> Result<(Vec<TxRecord>, u64), LogError> {
857    let mut file = File::open(path)?;
858    let file_len = file.metadata()?.len();
859    let mut records = Vec::new();
860    let mut durable_len = 0_u64;
861    loop {
862        let mut len = [0; 8];
863        match file.read_exact(&mut len) {
864            Ok(()) => {}
865            Err(e) if e.kind() == io::ErrorKind::UnexpectedEof => break,
866            Err(e) => return Err(e.into()),
867        }
868        let (payload_len, checksummed) = frame_payload_len(len)?;
869        let frame_len = 8_u64
870            .checked_add(u64::try_from(payload_len).map_err(|_| LogError::Corrupt)?)
871            .and_then(|len| {
872                len.checked_add(if checksummed {
873                    u64::try_from(FRAME_CHECKSUM_LEN).expect("checksum length fits u64")
874                } else {
875                    0
876                })
877            })
878            .ok_or(LogError::Corrupt)?;
879        // Check the bytes remaining before allocating from an untrusted length
880        // word. A short final frame is the recoverable crash-tail case.
881        if file_len.saturating_sub(durable_len) < frame_len {
882            break;
883        }
884        let mut payload = vec![0; payload_len];
885        match file.read_exact(&mut payload) {
886            Ok(()) => {}
887            Err(e) if e.kind() == io::ErrorKind::UnexpectedEof => break,
888            Err(e) => return Err(e.into()),
889        }
890        if checksummed {
891            let mut stored_checksum = [0; FRAME_CHECKSUM_LEN];
892            match file.read_exact(&mut stored_checksum) {
893                Ok(()) => {}
894                Err(e) if e.kind() == io::ErrorKind::UnexpectedEof => break,
895                Err(e) => return Err(e.into()),
896            }
897            if u32::from_be_bytes(stored_checksum) != frame_checksum(len, &payload) {
898                return Err(LogError::Corrupt);
899            }
900        }
901        records.push(decode_record(&payload)?);
902        durable_len = durable_len
903            .checked_add(frame_len)
904            .ok_or(LogError::Corrupt)?;
905    }
906    Ok((records, durable_len))
907}
908
909/// Appends one checksummed, length-prefixed encoded record to `out`.
910///
911/// The high bit of the length word identifies the checksummed frame format;
912/// the remaining 63 bits are the payload length. A big-endian CRC32C over the
913/// encoded length word and payload follows the payload.
914///
915/// # Errors
916/// Returns an error if the record payload length is not representable.
917pub fn append_framed_record(out: &mut Vec<u8>, record: &TxRecord) -> Result<(), LogError> {
918    let payload = encode_record(record);
919    let header = frame_header(payload.len())?;
920    out.extend_from_slice(&header);
921    out.extend_from_slice(&payload);
922    out.extend_from_slice(&frame_checksum(header, &payload).to_be_bytes());
923    Ok(())
924}
925
926/// Decodes all records from a framed byte slice.
927///
928/// Unlike filesystem crash recovery, native stores publish whole values
929/// atomically, so any trailing partial frame is treated as corruption.
930/// Both legacy length-only frames and checksummed frames are accepted.
931///
932/// # Errors
933/// Returns an error when any frame is truncated, has an invalid length, or
934/// checksum, or contains a corrupt encoded transaction record.
935pub fn decode_framed_records(mut bytes: &[u8]) -> Result<Vec<TxRecord>, LogError> {
936    let mut records = Vec::new();
937    while !bytes.is_empty() {
938        if bytes.len() < 8 {
939            return Err(LogError::Corrupt);
940        }
941        let header: [u8; 8] = bytes[..8].try_into().map_err(|_| LogError::Corrupt)?;
942        let (payload_len, checksummed) = frame_payload_len(header)?;
943        bytes = &bytes[8..];
944        let payload = bytes.get(..payload_len).ok_or(LogError::Corrupt)?;
945        bytes = &bytes[payload_len..];
946        if checksummed {
947            let stored_checksum = u32::from_be_bytes(
948                bytes
949                    .get(..FRAME_CHECKSUM_LEN)
950                    .ok_or(LogError::Corrupt)?
951                    .try_into()
952                    .map_err(|_| LogError::Corrupt)?,
953            );
954            if stored_checksum != frame_checksum(header, payload) {
955                return Err(LogError::Corrupt);
956            }
957            bytes = &bytes[FRAME_CHECKSUM_LEN..];
958        }
959        records.push(decode_record(payload)?);
960    }
961    Ok(records)
962}