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 (offline inspection, 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    /// Next `t` this writer will accept; also serializes concurrent appends.
415    next_t: tokio::sync::Mutex<u64>,
416}
417
418impl<S: NativeLogStorage + ?Sized + 'static> NativeVersionedLog<S> {
419    /// Opens the log for writing under `write_version`.
420    ///
421    /// # Errors
422    /// Returns an error when stored records cannot be read or decoded.
423    pub async fn open(storage: Arc<S>, name: &str, write_version: u64) -> Result<Self, LogError> {
424        // The merged view across every version and both layouts establishes the
425        // next `t` — the takeover cutoff may place it past this writer's own
426        // last record.
427        let records = read_native_merged(storage.as_ref(), name).await?;
428        let next_t = records.last().map_or(1, |r| r.t + 1);
429        Ok(Self {
430            storage,
431            name: name.to_owned(),
432            write_version,
433            next_t: tokio::sync::Mutex::new(next_t),
434        })
435    }
436}
437
438#[async_trait]
439impl<S: NativeLogStorage + ?Sized + 'static> TransactionLog for NativeVersionedLog<S> {
440    fn append(&self, record: &TxRecord) -> Result<(), LogError> {
441        let _ = record;
442        Err(LogError::AsyncOnly)
443    }
444
445    async fn append_async(&self, record: &TxRecord) -> Result<(), LogError> {
446        self.append_batch_async(std::slice::from_ref(record)).await
447    }
448
449    async fn append_batch_async(&self, records: &[TxRecord]) -> Result<(), LogError> {
450        if records.is_empty() {
451            return Ok(());
452        }
453        let mut next_t = self.next_t.lock().await;
454        // The batch must be exactly the next contiguous run of transactions.
455        for (offset, record) in records.iter().enumerate() {
456            if record.t != *next_t + offset as u64 {
457                return Err(LogError::Corrupt);
458            }
459        }
460        let framed = records
461            .iter()
462            .map(|record| {
463                let mut bytes = Vec::new();
464                append_framed_record(&mut bytes, record)?;
465                Ok((record.t, bytes))
466            })
467            .collect::<Result<Vec<_>, LogError>>()?;
468        // Create-only write of the whole batch as one object. As the sole
469        // appender under this lease version, an object that already exists for
470        // this last-`t` is a duplicate or a racing writer under our version —
471        // never a legitimate append — so reject it rather than overwrite.
472        if !self
473            .storage
474            .put_batch(&self.name, self.write_version, &framed)
475            .await?
476        {
477            return Err(LogError::Corrupt);
478        }
479        *next_t += records.len() as u64;
480        Ok(())
481    }
482
483    fn tx_range(&self, start: u64, end: Option<u64>) -> Result<Vec<TxRecord>, LogError> {
484        let _ = (start, end);
485        Err(LogError::AsyncOnly)
486    }
487
488    async fn tx_range_async(
489        &self,
490        start: u64,
491        end: Option<u64>,
492    ) -> Result<Vec<TxRecord>, LogError> {
493        // Range/replay must merge every version (for the takeover cutoff), so
494        // they read the store; the lock only serializes them with appends.
495        let _guard = self.next_t.lock().await;
496        Ok(read_native_merged(self.storage.as_ref(), &self.name)
497            .await?
498            .into_iter()
499            .filter(|r| r.t >= start && end.is_none_or(|e| r.t < e))
500            .collect())
501    }
502}
503
504async fn read_native_merged<S: NativeLogStorage + ?Sized>(
505    storage: &S,
506    name: &str,
507) -> Result<Vec<TxRecord>, LogError> {
508    use std::collections::BTreeMap;
509
510    // Gather every version's records from both layouts, keyed by version so the
511    // cross-version takeover cutoff below sees them in ascending version order.
512    let mut per_version: BTreeMap<u64, Vec<TxRecord>> = BTreeMap::new();
513
514    // Legacy chunk objects (read-only): a version's chunks concatenate in chunk
515    // order, which is the order they were filled — i.e. transaction order.
516    // Empty on a store that only ever wrote the per-record layout.
517    let mut chunks = storage.list_legacy_chunks(name).await?;
518    chunks.sort_unstable();
519    for (version, chunk) in chunks {
520        let bytes = storage
521            .read_legacy_chunk(name, version, chunk)
522            .await?
523            .unwrap_or_default();
524        per_version
525            .entry(version)
526            .or_default()
527            .extend(decode_framed_records(&bytes)?);
528    }
529
530    // Per-record objects, one framed record each.
531    let mut records = storage.list_records(name).await?;
532    records.sort_unstable();
533    for (version, t) in records {
534        let bytes = storage
535            .read_record(name, version, t)
536            .await?
537            .unwrap_or_default();
538        per_version
539            .entry(version)
540            .or_default()
541            .extend(decode_framed_records(&bytes)?);
542    }
543
544    // Order each version's records by `t` (a version is written in a single
545    // layout in practice; sorting keeps even a version that carries both — a
546    // legacy tail then per-record appends — correct), then apply the takeover
547    // cutoff across versions.
548    let per_version: Vec<Vec<TxRecord>> = per_version
549        .into_values()
550        .map(|mut records| {
551            records.sort_by_key(|record| record.t);
552            records
553        })
554        .collect();
555    let merged = merge_versions(per_version);
556    for pair in merged.windows(2) {
557        if pair[1].t != pair[0].t + 1 {
558            return Err(LogError::Corrupt);
559        }
560    }
561    Ok(merged)
562}
563
564/// Shared store of one log's records, each tagged with the lease version it
565/// was appended under.
566type VersionedRecords = Arc<Mutex<Vec<(u64, TxRecord)>>>;
567
568/// Process-shared registry of in-memory transaction logs, keyed by database
569/// name. It plays the role the log directory plays for [`VersionedLog`]:
570/// opening the same name (under any lease version) reaches the same records,
571/// so a mem-backed transactor recovers state across `open`/`create` calls
572/// within one process. Cloning a registry shares its storage.
573#[derive(Clone, Default)]
574pub struct MemLogRegistry {
575    logs: Arc<Mutex<HashMap<String, VersionedRecords>>>,
576}
577
578impl MemLogRegistry {
579    /// Creates an empty registry.
580    #[must_use]
581    pub fn new() -> Self {
582        Self::default()
583    }
584
585    fn entry(&self, name: &str) -> VersionedRecords {
586        Arc::clone(
587            self.logs
588                .lock()
589                .unwrap_or_else(std::sync::PoisonError::into_inner)
590                .entry(name.to_owned())
591                .or_default(),
592        )
593    }
594
595    /// Opens the named log for writing under `write_version`, mirroring
596    /// [`VersionedLog::open`] with in-memory storage.
597    #[must_use]
598    pub fn open(&self, name: &str, write_version: u64) -> MemVersionedLog {
599        let records = self.entry(name);
600        let next_t = {
601            let guard = records
602                .lock()
603                .unwrap_or_else(std::sync::PoisonError::into_inner);
604            MemVersionedLog::merged(&guard)
605                .last()
606                .map_or(1, |r| r.t + 1)
607        };
608        MemVersionedLog {
609            records,
610            write_version,
611            next_t: Mutex::new(next_t),
612        }
613    }
614
615    /// Reports whether any records exist for the named log.
616    #[must_use]
617    pub fn exists(&self, name: &str) -> bool {
618        self.logs
619            .lock()
620            .unwrap_or_else(std::sync::PoisonError::into_inner)
621            .get(name)
622            .is_some_and(|entry| {
623                !entry
624                    .lock()
625                    .unwrap_or_else(std::sync::PoisonError::into_inner)
626                    .is_empty()
627            })
628    }
629
630    /// Discards every record for the named log.
631    pub fn delete_all(&self, name: &str) {
632        self.logs
633            .lock()
634            .unwrap_or_else(std::sync::PoisonError::into_inner)
635            .remove(name);
636    }
637}
638
639/// An in-memory transaction log with the same per-lease-version merge
640/// semantics as [`VersionedLog`], obtained from a [`MemLogRegistry`]. Used by
641/// the mem-backed transactor: fully ephemeral, confined to one process.
642pub struct MemVersionedLog {
643    records: VersionedRecords,
644    write_version: u64,
645    /// The next `t` this writer will accept, tracked per opened instance
646    /// exactly as [`VersionedLog`] does — a deposed writer keeps appending
647    /// under its own stale count, and the merge cutoff discards those records.
648    next_t: Mutex<u64>,
649}
650
651impl MemVersionedLog {
652    fn merged(records: &[(u64, TxRecord)]) -> Vec<TxRecord> {
653        let mut versions: Vec<u64> = records.iter().map(|(version, _)| *version).collect();
654        versions.sort_unstable();
655        versions.dedup();
656        let per_version = versions
657            .into_iter()
658            .map(|version| {
659                records
660                    .iter()
661                    .filter(|(record_version, _)| *record_version == version)
662                    .map(|(_, record)| record.clone())
663                    .collect::<Vec<_>>()
664            })
665            .collect();
666        merge_versions(per_version)
667    }
668}
669
670impl TransactionLog for MemVersionedLog {
671    fn append(&self, record: &TxRecord) -> Result<(), LogError> {
672        let mut next_t = self
673            .next_t
674            .lock()
675            .unwrap_or_else(std::sync::PoisonError::into_inner);
676        if *next_t != record.t {
677            return Err(LogError::Corrupt);
678        }
679        self.records
680            .lock()
681            .unwrap_or_else(std::sync::PoisonError::into_inner)
682            .push((self.write_version, record.clone()));
683        *next_t += 1;
684        Ok(())
685    }
686
687    fn tx_range(&self, start: u64, end: Option<u64>) -> Result<Vec<TxRecord>, LogError> {
688        let records = self
689            .records
690            .lock()
691            .unwrap_or_else(std::sync::PoisonError::into_inner);
692        Ok(Self::merged(&records)
693            .into_iter()
694            .filter(|r| r.t >= start && end.is_none_or(|e| r.t < e))
695            .collect())
696    }
697}
698
699fn version_path(dir: &Path, name: &str, version: u64) -> PathBuf {
700    if version == 0 {
701        dir.join(format!("{name}.log"))
702    } else {
703        dir.join(format!("{name}.v{version}.log"))
704    }
705}
706
707/// Existing version files for `name`, sorted by version.
708fn version_files(dir: &Path, name: &str) -> Vec<(u64, PathBuf)> {
709    let mut files = Vec::new();
710    let legacy = version_path(dir, name, 0);
711    if legacy.is_file() {
712        files.push((0, legacy));
713    }
714    let prefix = format!("{name}.v");
715    if let Ok(entries) = fs::read_dir(dir) {
716        for entry in entries.flatten() {
717            let file_name = entry.file_name();
718            let Some(text) = file_name.to_str() else {
719                continue;
720            };
721            if let Some(version) = text
722                .strip_prefix(&prefix)
723                .and_then(|rest| rest.strip_suffix(".log"))
724                .and_then(|v| v.parse::<u64>().ok())
725                && version > 0
726            {
727                files.push((version, entry.path()));
728            }
729        }
730    }
731    files.sort_by_key(|(version, _)| *version);
732    files
733}
734
735/// Merges every version file, applying the takeover cutoff rule, and
736/// verifies the surviving sequence is contiguous.
737fn read_merged(dir: &Path, name: &str) -> Result<Vec<TxRecord>, LogError> {
738    let files = version_files(dir, name);
739    let mut per_file: Vec<Vec<TxRecord>> = Vec::with_capacity(files.len());
740    for (_, path) in &files {
741        per_file.push(read_records(path)?.0);
742    }
743    // A record in an older file is dead once any later file starts at or
744    // below its t: every record acked under version v precedes the first
745    // record of every later version (the successor replayed it before
746    // choosing its own first t), so only never-acked stale appends die.
747    let merged = merge_versions(per_file);
748    for pair in merged.windows(2) {
749        if pair[1].t != pair[0].t + 1 {
750            return Err(LogError::Corrupt);
751        }
752    }
753    Ok(merged)
754}
755
756fn encode_record(record: &TxRecord) -> Vec<u8> {
757    let mut out = Vec::new();
758    out.extend_from_slice(&record.t.to_be_bytes());
759    out.extend_from_slice(&record.tx_instant.to_be_bytes());
760    out.extend_from_slice(&(record.datoms.len() as u64).to_be_bytes());
761    for d in &record.datoms {
762        out.extend_from_slice(&d.e.raw().to_be_bytes());
763        out.extend_from_slice(&d.a.raw().to_be_bytes());
764        out.extend_from_slice(&d.tx.raw().to_be_bytes());
765        out.push(u8::from(d.added));
766        let v = encode_value(&d.v);
767        out.extend_from_slice(&(v.len() as u64).to_be_bytes());
768        out.extend_from_slice(&v);
769    }
770    out
771}
772fn decode_record(mut bytes: &[u8]) -> Result<TxRecord, LogError> {
773    fn take<'a>(bytes: &mut &'a [u8], n: usize) -> Result<&'a [u8], LogError> {
774        let value = bytes.get(..n).ok_or(LogError::Corrupt)?;
775        *bytes = &bytes[n..];
776        Ok(value)
777    }
778    fn u64_be(bytes: &mut &[u8]) -> Result<u64, LogError> {
779        Ok(u64::from_be_bytes(
780            take(bytes, 8)?.try_into().map_err(|_| LogError::Corrupt)?,
781        ))
782    }
783    let t = u64_be(&mut bytes)?;
784    let tx_instant = i64::from_be_bytes(
785        take(&mut bytes, 8)?
786            .try_into()
787            .map_err(|_| LogError::Corrupt)?,
788    );
789    let count = u64_be(&mut bytes)?;
790    let mut datoms = Vec::new();
791    for _ in 0..count {
792        let e = EntityId::from_raw(u64_be(&mut bytes)?);
793        let a = EntityId::from_raw(u64_be(&mut bytes)?);
794        let tx = EntityId::from_raw(u64_be(&mut bytes)?);
795        let added = take(&mut bytes, 1)?[0] != 0;
796        let len = usize::try_from(u64_be(&mut bytes)?).map_err(|_| LogError::Corrupt)?;
797        let raw = take(&mut bytes, len)?;
798        let (v, used) = decode_value(raw).map_err(|_| LogError::Corrupt)?;
799        if used != len {
800            return Err(LogError::Corrupt);
801        }
802        datoms.push(Datom { e, a, v, tx, added });
803    }
804    if !bytes.is_empty() {
805        return Err(LogError::Corrupt);
806    }
807    Ok(TxRecord {
808        t,
809        tx_instant,
810        datoms,
811    })
812}
813/// Reads fully written records plus the byte length of that durable prefix.
814///
815/// A record cut short by a crash mid-append (truncated length prefix or
816/// payload) ends the scan; a fully present record that fails to decode is
817/// genuine corruption and errors.
818fn read_records(path: &Path) -> Result<(Vec<TxRecord>, u64), LogError> {
819    let mut file = File::open(path)?;
820    let mut records = Vec::new();
821    let mut durable_len = 0_u64;
822    loop {
823        let mut len = [0; 8];
824        match file.read_exact(&mut len) {
825            Ok(()) => {}
826            Err(e) if e.kind() == io::ErrorKind::UnexpectedEof => break,
827            Err(e) => return Err(e.into()),
828        }
829        let len = usize::try_from(u64::from_be_bytes(len)).map_err(|_| LogError::Corrupt)?;
830        let mut payload = vec![0; len];
831        match file.read_exact(&mut payload) {
832            Ok(()) => {}
833            Err(e) if e.kind() == io::ErrorKind::UnexpectedEof => break,
834            Err(e) => return Err(e.into()),
835        }
836        records.push(decode_record(&payload)?);
837        durable_len += 8 + len as u64;
838    }
839    Ok((records, durable_len))
840}
841
842/// Appends one length-prefixed encoded record to `out`.
843///
844/// # Errors
845/// Returns an error if the record payload length is not representable.
846pub fn append_framed_record(out: &mut Vec<u8>, record: &TxRecord) -> Result<(), LogError> {
847    let payload = encode_record(record);
848    out.extend_from_slice(
849        &u64::try_from(payload.len())
850            .map_err(|_| LogError::Corrupt)?
851            .to_be_bytes(),
852    );
853    out.extend_from_slice(&payload);
854    Ok(())
855}
856
857/// Decodes all records from a length-prefixed byte slice.
858///
859/// Unlike filesystem crash recovery, native stores publish whole values
860/// atomically, so any trailing partial frame is treated as corruption.
861///
862/// # Errors
863/// Returns an error when any frame is truncated, has an invalid length, or
864/// contains a corrupt encoded transaction record.
865pub fn decode_framed_records(mut bytes: &[u8]) -> Result<Vec<TxRecord>, LogError> {
866    let mut records = Vec::new();
867    while !bytes.is_empty() {
868        if bytes.len() < 8 {
869            return Err(LogError::Corrupt);
870        }
871        let len = usize::try_from(u64::from_be_bytes(
872            bytes[..8].try_into().map_err(|_| LogError::Corrupt)?,
873        ))
874        .map_err(|_| LogError::Corrupt)?;
875        bytes = &bytes[8..];
876        let payload = bytes.get(..len).ok_or(LogError::Corrupt)?;
877        records.push(decode_record(payload)?);
878        bytes = &bytes[len..];
879    }
880    Ok(records)
881}