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