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, 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>();
19const RANGE_READ_CHUNK_BYTES: u64 = 4 * 1024 * 1024;
20const MAX_CACHED_READ_VERSION_FILES: usize = 8;
21
22/// One committed transaction record.
23#[derive(Clone, Debug, Eq, PartialEq)]
24pub struct TxRecord {
25    /// Monotonic transaction number.
26    pub t: u64,
27    /// Monotonic UTC millisecond timestamp.
28    pub tx_instant: i64,
29    /// Facts asserted/retracted by the transaction.
30    pub datoms: Vec<Datom>,
31}
32
33/// Log errors.
34#[derive(Debug, Error)]
35pub enum LogError {
36    /// Filesystem error.
37    #[error("log I/O failed: {0}")]
38    Io(#[from] io::Error),
39    /// Malformed or incomplete log data.
40    #[error("corrupt transaction log")]
41    Corrupt,
42    /// Native store backend failure.
43    #[error("native transaction log store failed: {0}")]
44    Native(String),
45    /// The operation requires the asynchronous log interface.
46    #[error("this transaction log requires asynchronous access")]
47    AsyncOnly,
48}
49
50/// Common transaction log interface.
51#[async_trait]
52pub trait TransactionLog: Send + Sync {
53    /// Durably appends exactly the next transaction.
54    ///
55    /// # Errors
56    /// Returns an error for I/O failure, corruption, or a non-contiguous `t`.
57    fn append(&self, record: &TxRecord) -> Result<(), LogError>;
58    /// Durably appends exactly the next transaction without blocking an async
59    /// runtime worker. Synchronous logs use [`Self::append`] by default;
60    /// storage-backed logs override this method and await their backend.
61    ///
62    /// # Errors
63    /// Returns an error for I/O failure, corruption, or a non-contiguous `t`.
64    async fn append_async(&self, record: &TxRecord) -> Result<(), LogError> {
65        self.append(record)
66    }
67    /// Durably appends a contiguous run of transactions under a single
68    /// durability boundary where the backend supports one (one `fsync`, one
69    /// object, one database transaction), so group commit amortizes the
70    /// per-append cost across the batch. `records` must be contiguous in `t`
71    /// starting at the log's next expected `t`; an empty slice is a no-op. The
72    /// default appends them one at a time; batching backends override this.
73    ///
74    /// # Errors
75    /// Returns an error for I/O failure, corruption, or a non-contiguous `t`.
76    async fn append_batch_async(&self, records: &[TxRecord]) -> Result<(), LogError> {
77        for record in records {
78            self.append_async(record).await?;
79        }
80        Ok(())
81    }
82    /// Returns records in the half-open transaction range `[start, end)`.
83    ///
84    /// # Errors
85    /// Returns an error when stored records cannot be read or decoded.
86    fn tx_range(&self, start: u64, end: Option<u64>) -> Result<Vec<TxRecord>, LogError>;
87    /// Asynchronous form of [`Self::tx_range`].
88    ///
89    /// # Errors
90    /// Returns an error when stored records cannot be read or decoded.
91    async fn tx_range_async(
92        &self,
93        start: u64,
94        end: Option<u64>,
95    ) -> Result<Vec<TxRecord>, LogError> {
96        self.tx_range(start, end)
97    }
98    /// Replays every committed record.
99    ///
100    /// # Errors
101    /// Returns an error when stored records cannot be read or decoded.
102    fn replay(&self) -> Result<Vec<TxRecord>, LogError> {
103        self.tx_range(0, None)
104    }
105    /// Asynchronously replays every committed record.
106    ///
107    /// # Errors
108    /// Returns an error when stored records cannot be read or decoded.
109    async fn replay_async(&self) -> Result<Vec<TxRecord>, LogError> {
110        self.tx_range_async(0, None).await
111    }
112}
113
114/// In-memory log implementation.
115#[derive(Clone, Default)]
116pub struct MemoryLog(Arc<RwLock<Vec<TxRecord>>>);
117impl TransactionLog for MemoryLog {
118    fn append(&self, record: &TxRecord) -> Result<(), LogError> {
119        let mut records = self.0.write().expect("poisoned log lock");
120        if records.last().map_or(1, |r| r.t + 1) != record.t {
121            return Err(LogError::Corrupt);
122        }
123        records.push(record.clone());
124        Ok(())
125    }
126    fn tx_range(&self, start: u64, end: Option<u64>) -> Result<Vec<TxRecord>, LogError> {
127        Ok(self
128            .0
129            .read()
130            .expect("poisoned log lock")
131            .iter()
132            .filter(|r| r.t >= start && end.is_none_or(|e| r.t < e))
133            .cloned()
134            .collect())
135    }
136}
137
138/// Filesystem append log. Each append is flushed and `fsync`ed before returning.
139///
140/// The file descriptor remains open for the lifetime of the log. An in-memory
141/// `(t, byte offset, frame length)` index is built during recovery, extended
142/// on append, and used to read only the frames selected by a range scan.
143///
144/// A crash mid-append leaves a torn, never-acked record at the tail; `open`
145/// truncates it away so replay stops at the durability point of the last
146/// acked transaction and later appends extend a clean tail.
147pub struct FileLog {
148    state: RwLock<IndexedFile>,
149}
150
151impl FileLog {
152    /// Opens or creates a log file, dropping any torn tail left by a crash.
153    ///
154    /// # Errors
155    /// Returns an error if the file cannot be created or a fully written
156    /// record is corrupt.
157    pub fn open(path: impl AsRef<Path>) -> Result<Self, LogError> {
158        let path = path.as_ref().to_path_buf();
159        if let Some(parent) = path.parent() {
160            fs::create_dir_all(parent)?;
161        }
162        let file = IndexedFile::open(&path, true, true)?;
163        file.validate_contiguous_prefix(file.frames.len())?;
164        Ok(Self {
165            state: RwLock::new(file),
166        })
167    }
168}
169impl TransactionLog for FileLog {
170    fn append(&self, record: &TxRecord) -> Result<(), LogError> {
171        let mut state = self.state.write().expect("poisoned log lock");
172        state.refresh()?;
173        state.validate_contiguous_prefix(state.frames.len())?;
174        if next_t(&state.frames)? != record.t {
175            return Err(LogError::Corrupt);
176        }
177        state.append(record)
178    }
179    fn tx_range(&self, start: u64, end: Option<u64>) -> Result<Vec<TxRecord>, LogError> {
180        if end.is_some_and(|end| end <= start) {
181            return Ok(Vec::new());
182        }
183        let indexed = {
184            let state = self.state.read().expect("poisoned log lock");
185            range_is_indexed(&state.frames, end)?
186        };
187        if !indexed {
188            let mut state = self.state.write().expect("poisoned log lock");
189            state.refresh()?;
190            state.validate_contiguous_prefix(state.frames.len())?;
191        }
192        self.state
193            .read()
194            .expect("poisoned log lock")
195            .tx_range(start, end)
196    }
197}
198
199/// A transaction log split into per-lease-version files for HA append
200/// isolation (see `docs/design/log-and-transactor.md`).
201///
202/// The active writer under lease version `V` appends only to
203/// `{name}.v{V}.log` (the pre-HA `{name}.log` reads as version 0). Readers
204/// merge the files in version order and drop any record in an older file
205/// whose `t` is at or past the first record of a later file: such records
206/// were appended by a deposed writer after a takeover and were never
207/// acknowledged, because acknowledgement re-verifies lease ownership after
208/// the durable append. A deposed writer therefore cannot corrupt or fork
209/// the log — its stale appends land in a file nobody considers current.
210pub struct VersionedLog {
211    dir: PathBuf,
212    name: String,
213    state: RwLock<VersionedLogState>,
214}
215
216struct VersionedLogState {
217    files: Vec<VersionedFile>,
218    write_version: Option<u64>,
219    next_t: u64,
220}
221
222struct VersionedFile {
223    version: u64,
224    file: IndexedFile,
225}
226
227impl VersionedLog {
228    /// Opens the log for writing under `write_version`, creating the
229    /// version file if needed and dropping any torn tail it carries.
230    /// Files of other versions are never modified.
231    ///
232    /// # Errors
233    /// Returns an error if files cannot be read/created or a fully written
234    /// record is corrupt.
235    pub fn open(dir: impl AsRef<Path>, name: &str, write_version: u64) -> Result<Self, LogError> {
236        let dir = dir.as_ref().to_path_buf();
237        fs::create_dir_all(&dir)?;
238        let write_path = version_path(&dir, name, write_version);
239        let mut files = Vec::new();
240        for (version, path) in version_files(&dir, name) {
241            let writable = version == write_version;
242            files.push(VersionedFile {
243                version,
244                file: IndexedFile::open(&path, writable, writable)?,
245            });
246            close_cold_version_files(&mut files);
247        }
248        if !files.iter().any(|file| file.version == write_version) {
249            files.push(VersionedFile {
250                version: write_version,
251                file: IndexedFile::open(&write_path, true, true)?,
252            });
253            files.sort_by_key(|file| file.version);
254        }
255        close_cold_version_files(&mut files);
256        let cutoffs = validated_version_cutoffs(&files)?;
257        let next_t = merged_next_t(&files, &cutoffs)?;
258        Ok(Self {
259            dir,
260            name: name.to_owned(),
261            state: RwLock::new(VersionedLogState {
262                files,
263                write_version: Some(write_version),
264                next_t,
265            }),
266        })
267    }
268
269    /// Opens the log read-only (independent inspection or backup); appends fail.
270    ///
271    /// # Errors
272    /// Returns an error when the directory cannot be read or a fully
273    /// written record is corrupt.
274    pub fn open_read_only(dir: impl AsRef<Path>, name: &str) -> Result<Self, LogError> {
275        let dir = dir.as_ref().to_path_buf();
276        let mut files = open_version_files(&dir, name)?;
277        close_cold_version_files(&mut files);
278        let cutoffs = validated_version_cutoffs(&files)?;
279        Ok(Self {
280            name: name.to_owned(),
281            state: RwLock::new(VersionedLogState {
282                next_t: merged_next_t(&files, &cutoffs)?,
283                write_version: None,
284                files,
285            }),
286            dir,
287        })
288    }
289
290    /// Reports whether any log file exists for this database.
291    #[must_use]
292    pub fn exists(dir: impl AsRef<Path>, name: &str) -> bool {
293        !version_files(dir.as_ref(), name).is_empty()
294    }
295
296    /// Deletes every version file for this database.
297    ///
298    /// # Errors
299    /// Returns an error when a file cannot be removed.
300    pub fn delete_all(dir: impl AsRef<Path>, name: &str) -> Result<(), LogError> {
301        for (_, path) in version_files(dir.as_ref(), name) {
302            match fs::remove_file(&path) {
303                Ok(()) => {}
304                Err(error) if error.kind() == io::ErrorKind::NotFound => {}
305                Err(error) => return Err(error.into()),
306            }
307        }
308        Ok(())
309    }
310}
311
312impl TransactionLog for VersionedLog {
313    fn append(&self, record: &TxRecord) -> Result<(), LogError> {
314        let mut state = self.state.write().expect("poisoned log lock");
315        if state.next_t != record.t {
316            return Err(LogError::Corrupt);
317        }
318        let write_version = state
319            .write_version
320            .ok_or_else(|| LogError::Native("transaction log is read-only".into()))?;
321        let write_index = state
322            .files
323            .iter()
324            .position(|file| file.version == write_version)
325            .ok_or(LogError::Corrupt)?;
326        let cutoffs = version_cutoffs(&state.files);
327        // A current-version writer must extend its own local tail. A deposed
328        // writer may append beyond a later version's cutoff; that harmless gap
329        // lives entirely in the dead suffix and is discarded during merge.
330        if cutoffs[write_index] == u64::MAX
331            && !state.files[write_index].file.frames.is_empty()
332            && next_t(&state.files[write_index].file.frames)? != record.t
333        {
334            return Err(LogError::Corrupt);
335        }
336        state.files[write_index].file.append(record)?;
337        state.next_t = state.next_t.checked_add(1).ok_or(LogError::Corrupt)?;
338        Ok(())
339    }
340
341    fn tx_range(&self, start: u64, end: Option<u64>) -> Result<Vec<TxRecord>, LogError> {
342        if end.is_some_and(|end| end <= start) {
343            return Ok(Vec::new());
344        }
345        {
346            let state = self.state.read().expect("poisoned log lock");
347            let cutoffs = validated_version_cutoffs(&state.files)?;
348            if range_is_merged_indexed(&state.files, &cutoffs, end)? {
349                return read_merged_range(&state.files, &cutoffs, start, end);
350            }
351        }
352        {
353            let mut state = self.state.write().expect("poisoned log lock");
354            refresh_version_files(&self.dir, &self.name, &mut state.files)?;
355            close_cold_version_files(&mut state.files);
356            validated_version_cutoffs(&state.files)?;
357        }
358        let state = self.state.read().expect("poisoned log lock");
359        let cutoffs = validated_version_cutoffs(&state.files)?;
360        read_merged_range(&state.files, &cutoffs, start, end)
361    }
362}
363
364/// Applies the takeover cutoff rule to per-version record lists: a record in
365/// an older version dies once any later version begins at or below its `t`,
366/// dropping only the never-acked stale appends of a deposed writer.
367fn merge_versions(mut per_version: Vec<Vec<TxRecord>>) -> Vec<TxRecord> {
368    let mut cutoff = u64::MAX;
369    for records in per_version.iter_mut().rev() {
370        let first = records.first().map(|r| r.t);
371        records.retain(|r| r.t < cutoff);
372        if let Some(first) = first {
373            cutoff = cutoff.min(first);
374        }
375    }
376    per_version.into_iter().flatten().collect()
377}
378
379/// Asynchronous object store for transaction-log records.
380///
381/// Implementations adapt the same native storage system used for blobs and
382/// roots. The live log is written **one object per transaction** — a
383/// create-only write keyed `(name, version, t)` whose success is the
384/// durability point — so an append is O(1) (a small insert) instead of a
385/// read-modify-write of a growing chunk. On a SQL backend that is a
386/// row-per-commit insert; on an object store, a create-only `PUT`.
387///
388/// Earlier releases wrote a different layout: a sequence of *chunk* objects
389/// `(name, version, chunk)`, each a run of framed records, appended in place
390/// and rolled at a size cap. Those objects are still read, read-only, through
391/// [`Self::list_legacy_chunks`] / [`Self::read_legacy_chunk`], so a log
392/// written by an older binary keeps replaying after an upgrade; new records
393/// are always written in the per-transaction layout.
394#[async_trait]
395pub trait NativeLogStorage: Send + Sync {
396    /// Create-only, atomic write of one contiguous batch of transactions as a
397    /// single object, keyed by the batch's last `t`. Each element is
398    /// `(t, framed_bytes)` in ascending `t`; the object holds their framed
399    /// bytes concatenated (the same encoding a multi-record chunk uses).
400    /// Returns `Ok(true)` when written and `Ok(false)` when an object already
401    /// exists for that last-`t` — a lost create race or a retry of an
402    /// already-durable batch. The create-only condition is the log's fence: a
403    /// given `(version, last t)` is written at most once, and the batch is
404    /// durable in full or not at all.
405    ///
406    /// # Errors
407    /// Returns an error when the native backend cannot publish the object.
408    async fn put_batch(
409        &self,
410        name: &str,
411        version: u64,
412        records: &[(u64, Vec<u8>)],
413    ) -> Result<bool, LogError>;
414    /// Reads the bytes of the batch object keyed by last-`t` `t` for
415    /// `(name, version)`.
416    ///
417    /// # Errors
418    /// Returns an error when the native backend cannot read the object.
419    async fn read_record(
420        &self,
421        name: &str,
422        version: u64,
423        t: u64,
424    ) -> Result<Option<Vec<u8>>, LogError>;
425    /// Lists every `(version, last-t)` batch object present for `name`.
426    ///
427    /// # Errors
428    /// Returns an error when the native backend cannot enumerate log objects
429    /// or returns an invalid identifier.
430    async fn list_records(&self, name: &str) -> Result<Vec<(u64, u64)>, LogError>;
431    /// Reads one legacy chunk object (pre-per-record layout), for read-only
432    /// replay of logs written by older binaries.
433    ///
434    /// # Errors
435    /// Returns an error when the native backend cannot read the chunk object.
436    async fn read_legacy_chunk(
437        &self,
438        name: &str,
439        version: u64,
440        chunk: u64,
441    ) -> Result<Option<Vec<u8>>, LogError>;
442    /// Lists every legacy `(version, chunk)` object present for `name`, for
443    /// read-only replay of older logs. Returns an empty list on a store that
444    /// only ever wrote the per-record layout.
445    ///
446    /// # Errors
447    /// Returns an error when the native backend cannot enumerate log objects
448    /// or returns an invalid identifier.
449    async fn list_legacy_chunks(&self, name: &str) -> Result<Vec<(u64, u64)>, LogError>;
450    /// Deletes every log object for `name`, in both layouts.
451    ///
452    /// # Errors
453    /// Returns an error when the native backend cannot remove an object.
454    async fn delete_all(&self, name: &str) -> Result<(), LogError>;
455}
456
457/// Versioned transaction log backed by a native key/value-style store.
458///
459/// Each transaction is written as its own create-only object keyed
460/// `(name, write_version, t)`, so an append is a single small insert with no
461/// read and no growing buffer. The writer is the sole appender under its lease
462/// version (the fence gives each active owner its own version; a deposed
463/// writer's stale appends land in a version the takeover cutoff discards), so
464/// tracking `next_t` in memory is all the append state required.
465pub struct NativeVersionedLog<S: ?Sized> {
466    storage: Arc<S>,
467    name: String,
468    write_version: u64,
469    read_only: bool,
470    /// Next `t` this writer will accept; also serializes concurrent appends.
471    next_t: tokio::sync::Mutex<u64>,
472}
473
474impl<S: NativeLogStorage + ?Sized + 'static> NativeVersionedLog<S> {
475    /// Opens the log for writing under `write_version`.
476    ///
477    /// # Errors
478    /// Returns an error when stored records cannot be read or decoded.
479    pub async fn open(storage: Arc<S>, name: &str, write_version: u64) -> Result<Self, LogError> {
480        // The merged view across every version and both layouts establishes the
481        // next `t` — the takeover cutoff may place it past this writer's own
482        // last record.
483        let records = read_native_merged(storage.as_ref(), name).await?;
484        let next_t = records.last().map_or(1, |r| r.t + 1);
485        Ok(Self {
486            storage,
487            name: name.to_owned(),
488            write_version,
489            read_only: false,
490            next_t: tokio::sync::Mutex::new(next_t),
491        })
492    }
493
494    /// Opens the log for read-only range replay without first scanning it to
495    /// initialize writer state.
496    #[must_use]
497    pub fn open_read_only(storage: Arc<S>, name: &str) -> Self {
498        Self {
499            storage,
500            name: name.to_owned(),
501            write_version: 0,
502            read_only: true,
503            next_t: tokio::sync::Mutex::new(0),
504        }
505    }
506}
507
508#[async_trait]
509impl<S: NativeLogStorage + ?Sized + 'static> TransactionLog for NativeVersionedLog<S> {
510    fn append(&self, record: &TxRecord) -> Result<(), LogError> {
511        let _ = record;
512        Err(LogError::AsyncOnly)
513    }
514
515    async fn append_async(&self, record: &TxRecord) -> Result<(), LogError> {
516        self.append_batch_async(std::slice::from_ref(record)).await
517    }
518
519    async fn append_batch_async(&self, records: &[TxRecord]) -> Result<(), LogError> {
520        if self.read_only {
521            return Err(LogError::Native("transaction log is read-only".into()));
522        }
523        if records.is_empty() {
524            return Ok(());
525        }
526        let mut next_t = self.next_t.lock().await;
527        // The batch must be exactly the next contiguous run of transactions.
528        for (offset, record) in records.iter().enumerate() {
529            if record.t != *next_t + offset as u64 {
530                return Err(LogError::Corrupt);
531            }
532        }
533        let framed = records
534            .iter()
535            .map(|record| {
536                let mut bytes = Vec::new();
537                append_framed_record(&mut bytes, record)?;
538                Ok((record.t, bytes))
539            })
540            .collect::<Result<Vec<_>, LogError>>()?;
541        // Create-only write of the whole batch as one object. As the sole
542        // appender under this lease version, an object that already exists for
543        // this last-`t` is a duplicate or a racing writer under our version —
544        // never a legitimate append — so reject it rather than overwrite.
545        if !self
546            .storage
547            .put_batch(&self.name, self.write_version, &framed)
548            .await?
549        {
550            return Err(LogError::Corrupt);
551        }
552        *next_t += records.len() as u64;
553        Ok(())
554    }
555
556    fn tx_range(&self, start: u64, end: Option<u64>) -> Result<Vec<TxRecord>, LogError> {
557        let _ = (start, end);
558        Err(LogError::AsyncOnly)
559    }
560
561    async fn tx_range_async(
562        &self,
563        start: u64,
564        end: Option<u64>,
565    ) -> Result<Vec<TxRecord>, LogError> {
566        // Range/replay must merge every version (for the takeover cutoff), so
567        // they read the store; the lock only serializes them with appends.
568        let _guard = self.next_t.lock().await;
569        Ok(read_native_merged(self.storage.as_ref(), &self.name)
570            .await?
571            .into_iter()
572            .filter(|r| r.t >= start && end.is_none_or(|e| r.t < e))
573            .collect())
574    }
575}
576
577async fn read_native_merged<S: NativeLogStorage + ?Sized>(
578    storage: &S,
579    name: &str,
580) -> Result<Vec<TxRecord>, LogError> {
581    use std::collections::BTreeMap;
582
583    // Gather every version's records from both layouts, keyed by version so the
584    // cross-version takeover cutoff below sees them in ascending version order.
585    let mut per_version: BTreeMap<u64, Vec<TxRecord>> = BTreeMap::new();
586
587    // Legacy chunk objects (read-only): a version's chunks concatenate in chunk
588    // order, which is the order they were filled — i.e. transaction order.
589    // Empty on a store that only ever wrote the per-record layout.
590    let mut chunks = storage.list_legacy_chunks(name).await?;
591    chunks.sort_unstable();
592    for (version, chunk) in chunks {
593        let bytes = storage
594            .read_legacy_chunk(name, version, chunk)
595            .await?
596            .unwrap_or_default();
597        per_version
598            .entry(version)
599            .or_default()
600            .extend(decode_framed_records(&bytes)?);
601    }
602
603    // Per-record objects, one framed record each.
604    let mut records = storage.list_records(name).await?;
605    records.sort_unstable();
606    for (version, t) in records {
607        let bytes = storage
608            .read_record(name, version, t)
609            .await?
610            .unwrap_or_default();
611        per_version
612            .entry(version)
613            .or_default()
614            .extend(decode_framed_records(&bytes)?);
615    }
616
617    // Order each version's records by `t` (a version is written in a single
618    // layout in practice; sorting keeps even a version that carries both — a
619    // legacy tail then per-record appends — correct), then apply the takeover
620    // cutoff across versions.
621    let per_version: Vec<Vec<TxRecord>> = per_version
622        .into_values()
623        .map(|mut records| {
624            records.sort_by_key(|record| record.t);
625            records
626        })
627        .collect();
628    let merged = merge_versions(per_version);
629    for pair in merged.windows(2) {
630        if pair[1].t != pair[0].t + 1 {
631            return Err(LogError::Corrupt);
632        }
633    }
634    Ok(merged)
635}
636
637/// Shared store of one log's records, each tagged with the lease version it
638/// was appended under.
639type VersionedRecords = Arc<Mutex<Vec<(u64, TxRecord)>>>;
640
641/// Process-shared registry of in-memory transaction logs, keyed by database
642/// name. It plays the role the log directory plays for [`VersionedLog`]:
643/// opening the same name (under any lease version) reaches the same records,
644/// so a mem-backed transactor recovers state across `open`/`create` calls
645/// within one process. Cloning a registry shares its storage.
646#[derive(Clone, Default)]
647pub struct MemLogRegistry {
648    logs: Arc<Mutex<HashMap<String, VersionedRecords>>>,
649}
650
651impl MemLogRegistry {
652    /// Creates an empty registry.
653    #[must_use]
654    pub fn new() -> Self {
655        Self::default()
656    }
657
658    fn entry(&self, name: &str) -> VersionedRecords {
659        Arc::clone(
660            self.logs
661                .lock()
662                .unwrap_or_else(std::sync::PoisonError::into_inner)
663                .entry(name.to_owned())
664                .or_default(),
665        )
666    }
667
668    /// Opens the named log for writing under `write_version`, mirroring
669    /// [`VersionedLog::open`] with in-memory storage.
670    #[must_use]
671    pub fn open(&self, name: &str, write_version: u64) -> MemVersionedLog {
672        let records = self.entry(name);
673        let next_t = {
674            let guard = records
675                .lock()
676                .unwrap_or_else(std::sync::PoisonError::into_inner);
677            MemVersionedLog::merged(&guard)
678                .last()
679                .map_or(1, |r| r.t + 1)
680        };
681        MemVersionedLog {
682            records,
683            write_version,
684            next_t: Mutex::new(next_t),
685        }
686    }
687
688    /// Reports whether any records exist for the named log.
689    #[must_use]
690    pub fn exists(&self, name: &str) -> bool {
691        self.logs
692            .lock()
693            .unwrap_or_else(std::sync::PoisonError::into_inner)
694            .get(name)
695            .is_some_and(|entry| {
696                !entry
697                    .lock()
698                    .unwrap_or_else(std::sync::PoisonError::into_inner)
699                    .is_empty()
700            })
701    }
702
703    /// Discards every record for the named log.
704    pub fn delete_all(&self, name: &str) {
705        self.logs
706            .lock()
707            .unwrap_or_else(std::sync::PoisonError::into_inner)
708            .remove(name);
709    }
710}
711
712/// An in-memory transaction log with the same per-lease-version merge
713/// semantics as [`VersionedLog`], obtained from a [`MemLogRegistry`]. Used by
714/// the mem-backed transactor: fully ephemeral, confined to one process.
715pub struct MemVersionedLog {
716    records: VersionedRecords,
717    write_version: u64,
718    /// The next `t` this writer will accept, tracked per opened instance
719    /// exactly as [`VersionedLog`] does — a deposed writer keeps appending
720    /// under its own stale count, and the merge cutoff discards those records.
721    next_t: Mutex<u64>,
722}
723
724impl MemVersionedLog {
725    fn merged(records: &[(u64, TxRecord)]) -> Vec<TxRecord> {
726        let mut versions: Vec<u64> = records.iter().map(|(version, _)| *version).collect();
727        versions.sort_unstable();
728        versions.dedup();
729        let per_version = versions
730            .into_iter()
731            .map(|version| {
732                records
733                    .iter()
734                    .filter(|(record_version, _)| *record_version == version)
735                    .map(|(_, record)| record.clone())
736                    .collect::<Vec<_>>()
737            })
738            .collect();
739        merge_versions(per_version)
740    }
741}
742
743impl TransactionLog for MemVersionedLog {
744    fn append(&self, record: &TxRecord) -> Result<(), LogError> {
745        let mut next_t = self
746            .next_t
747            .lock()
748            .unwrap_or_else(std::sync::PoisonError::into_inner);
749        if *next_t != record.t {
750            return Err(LogError::Corrupt);
751        }
752        self.records
753            .lock()
754            .unwrap_or_else(std::sync::PoisonError::into_inner)
755            .push((self.write_version, record.clone()));
756        *next_t += 1;
757        Ok(())
758    }
759
760    fn tx_range(&self, start: u64, end: Option<u64>) -> Result<Vec<TxRecord>, LogError> {
761        let records = self
762            .records
763            .lock()
764            .unwrap_or_else(std::sync::PoisonError::into_inner);
765        Ok(Self::merged(&records)
766            .into_iter()
767            .filter(|r| r.t >= start && end.is_none_or(|e| r.t < e))
768            .collect())
769    }
770}
771
772#[derive(Clone, Copy)]
773struct FrameIndex {
774    t: u64,
775    offset: u64,
776    len: u64,
777}
778
779/// One cached file descriptor and its in-memory transaction-to-byte index.
780///
781/// Opening scans and validates the durable prefix once. Later refreshes scan
782/// only bytes appended after that prefix, and concurrent range readers use
783/// positional I/O for the selected frames instead of decoding from byte zero.
784struct IndexedFile {
785    path: PathBuf,
786    file: Option<Arc<File>>,
787    writable: bool,
788    poisoned: bool,
789    frames: Vec<FrameIndex>,
790    durable_len: u64,
791    first_gap: Option<usize>,
792}
793
794impl IndexedFile {
795    fn open(path: &Path, writable: bool, truncate_torn: bool) -> Result<Self, LogError> {
796        let file = Arc::new(open_index_file(path, writable)?);
797        let (frames, durable_len) = scan_frames(file.as_ref(), 0)?;
798        validate_sorted_frames(&frames)?;
799        if truncate_torn && file.metadata()?.len() > durable_len {
800            file.set_len(durable_len)?;
801            file.sync_all()?;
802        }
803        Ok(Self {
804            path: path.to_path_buf(),
805            file: Some(file),
806            writable,
807            poisoned: false,
808            first_gap: first_gap_index(&frames),
809            frames,
810            durable_len,
811        })
812    }
813
814    fn refresh(&mut self) -> Result<(), LogError> {
815        self.ensure_healthy()?;
816        let file_len = self.physical_len()?;
817        if file_len < self.durable_len {
818            let file = self.open_for_read()?;
819            let (frames, durable_len) = scan_frames(file.as_ref(), 0)?;
820            validate_sorted_frames(&frames)?;
821            self.first_gap = first_gap_index(&frames);
822            self.frames = frames;
823            self.durable_len = durable_len;
824        } else if file_len > self.durable_len {
825            let file = self.open_for_read()?;
826            let (new_frames, durable_len) = scan_frames(file.as_ref(), self.durable_len)?;
827            validate_sorted_extension(&self.frames, &new_frames)?;
828            let existing_len = self.frames.len();
829            if self.first_gap.is_none() {
830                self.first_gap =
831                    extension_first_gap(&self.frames, &new_frames).map(|gap| existing_len + gap);
832            }
833            self.frames.extend(new_frames);
834            self.durable_len = durable_len;
835        }
836        Ok(())
837    }
838
839    fn append(&mut self, record: &TxRecord) -> Result<(), LogError> {
840        self.ensure_healthy()?;
841        if !self.writable {
842            return Err(LogError::Native("transaction log is read-only".into()));
843        }
844        let file = self.open_for_read()?;
845        // Never truncate here: durable_len belongs to this handle, and bytes
846        // beyond it may be an acknowledged append made by another handle.
847        if file.metadata()?.len() != self.durable_len {
848            return Err(LogError::Corrupt);
849        }
850
851        let mut frame = Vec::new();
852        append_framed_record(&mut frame, record)?;
853        let frame_len = u64::try_from(frame.len()).map_err(|_| LogError::Corrupt)?;
854        let offset = self.durable_len;
855        let mut writer = file.as_ref();
856        if let Err(error) = writer.write_all(&frame) {
857            // A short write makes the tail uncertain. Poison this handle so a
858            // retry cannot adopt or append beyond a never-acknowledged frame.
859            self.poisoned = true;
860            return Err(error.into());
861        }
862        if let Err(error) = file.sync_all() {
863            self.poisoned = true;
864            return Err(error.into());
865        }
866        self.durable_len = self
867            .durable_len
868            .checked_add(frame_len)
869            .ok_or(LogError::Corrupt)?;
870        if self.first_gap.is_none()
871            && self
872                .frames
873                .last()
874                .is_some_and(|previous| previous.t.checked_add(1) != Some(record.t))
875        {
876            self.first_gap = Some(self.frames.len());
877        }
878        self.frames.push(FrameIndex {
879            t: record.t,
880            offset,
881            len: frame_len,
882        });
883        Ok(())
884    }
885
886    fn tx_range(&self, start: u64, end: Option<u64>) -> Result<Vec<TxRecord>, LogError> {
887        self.ensure_healthy()?;
888        let first = self.frames.partition_point(|frame| frame.t < start);
889        let last = end.map_or(self.frames.len(), |end| {
890            self.frames.partition_point(|frame| frame.t < end)
891        });
892        if first >= last {
893            return Ok(Vec::new());
894        }
895
896        let file = self.open_for_read()?;
897        let mut records = Vec::with_capacity(last - first);
898        let mut chunk_first = first;
899        while chunk_first < last {
900            let offset = self.frames[chunk_first].offset;
901            let mut chunk_last = chunk_first + 1;
902            while chunk_last < last {
903                let candidate_end = frame_end(self.frames[chunk_last])?;
904                if candidate_end.checked_sub(offset).ok_or(LogError::Corrupt)?
905                    > RANGE_READ_CHUNK_BYTES
906                {
907                    break;
908                }
909                chunk_last += 1;
910            }
911            let byte_end = frame_end(self.frames[chunk_last - 1])?;
912            let byte_len = usize::try_from(byte_end.checked_sub(offset).ok_or(LogError::Corrupt)?)
913                .map_err(|_| LogError::Corrupt)?;
914            let mut bytes = vec![0; byte_len];
915            read_exact_at(file.as_ref(), &mut bytes, offset)?;
916            let chunk_records = decode_framed_records(&bytes)?;
917            if chunk_records.len() != chunk_last - chunk_first
918                || chunk_records
919                    .iter()
920                    .zip(&self.frames[chunk_first..chunk_last])
921                    .any(|(record, frame)| record.t != frame.t)
922            {
923                return Err(LogError::Corrupt);
924            }
925            records.extend(chunk_records);
926            chunk_first = chunk_last;
927        }
928        Ok(records)
929    }
930
931    fn validate_contiguous_prefix(&self, retained: usize) -> Result<(), LogError> {
932        if self.first_gap.is_some_and(|gap| gap < retained) {
933            return Err(LogError::Corrupt);
934        }
935        Ok(())
936    }
937
938    fn ensure_healthy(&self) -> Result<(), LogError> {
939        if self.poisoned {
940            return Err(io::Error::other("transaction log handle is poisoned").into());
941        }
942        Ok(())
943    }
944
945    fn open_for_read(&self) -> Result<Arc<File>, LogError> {
946        self.file.as_ref().map_or_else(
947            || Ok(Arc::new(open_index_file(&self.path, false)?)),
948            |file| Ok(Arc::clone(file)),
949        )
950    }
951
952    fn physical_len(&self) -> Result<u64, LogError> {
953        Ok(self
954            .file
955            .as_ref()
956            .map_or_else(|| fs::metadata(&self.path), |file| file.metadata())?
957            .len())
958    }
959
960    fn close_cached_reader(&mut self) {
961        if !self.writable {
962            self.file = None;
963        }
964    }
965}
966
967fn open_index_file(path: &Path, writable: bool) -> Result<File, io::Error> {
968    let mut options = OpenOptions::new();
969    options.read(true);
970    if writable {
971        options.create(true).write(true).append(true);
972    }
973    options.open(path)
974}
975
976fn validate_sorted_frames(frames: &[FrameIndex]) -> Result<(), LogError> {
977    for pair in frames.windows(2) {
978        if pair[0].t >= pair[1].t {
979            return Err(LogError::Corrupt);
980        }
981    }
982    Ok(())
983}
984
985fn validate_sorted_extension(
986    existing: &[FrameIndex],
987    appended: &[FrameIndex],
988) -> Result<(), LogError> {
989    validate_sorted_frames(appended)?;
990    if let (Some(previous), Some(next)) = (existing.last(), appended.first())
991        && previous.t >= next.t
992    {
993        return Err(LogError::Corrupt);
994    }
995    Ok(())
996}
997
998fn first_gap_index(frames: &[FrameIndex]) -> Option<usize> {
999    frames
1000        .windows(2)
1001        .position(|pair| pair[0].t.checked_add(1) != Some(pair[1].t))
1002        .map(|index| index + 1)
1003}
1004
1005fn extension_first_gap(existing: &[FrameIndex], appended: &[FrameIndex]) -> Option<usize> {
1006    if let (Some(previous), Some(next)) = (existing.last(), appended.first())
1007        && previous.t.checked_add(1) != Some(next.t)
1008    {
1009        return Some(0);
1010    }
1011    first_gap_index(appended)
1012}
1013
1014fn frame_end(frame: FrameIndex) -> Result<u64, LogError> {
1015    frame.offset.checked_add(frame.len).ok_or(LogError::Corrupt)
1016}
1017
1018fn next_t(frames: &[FrameIndex]) -> Result<u64, LogError> {
1019    frames.last().map_or(Ok(1), |frame| {
1020        frame.t.checked_add(1).ok_or(LogError::Corrupt)
1021    })
1022}
1023
1024fn range_is_indexed(frames: &[FrameIndex], end: Option<u64>) -> Result<bool, LogError> {
1025    end.map_or(Ok(false), |end| Ok(end <= next_t(frames)?))
1026}
1027
1028fn version_path(dir: &Path, name: &str, version: u64) -> PathBuf {
1029    if version == 0 {
1030        dir.join(format!("{name}.log"))
1031    } else {
1032        dir.join(format!("{name}.v{version}.log"))
1033    }
1034}
1035
1036/// Existing version files for `name`, sorted by version.
1037fn version_files(dir: &Path, name: &str) -> Vec<(u64, PathBuf)> {
1038    let mut files = Vec::new();
1039    let legacy = version_path(dir, name, 0);
1040    if legacy.is_file() {
1041        files.push((0, legacy));
1042    }
1043    let prefix = format!("{name}.v");
1044    if let Ok(entries) = fs::read_dir(dir) {
1045        for entry in entries.flatten() {
1046            let file_name = entry.file_name();
1047            let Some(text) = file_name.to_str() else {
1048                continue;
1049            };
1050            if let Some(version) = text
1051                .strip_prefix(&prefix)
1052                .and_then(|rest| rest.strip_suffix(".log"))
1053                .and_then(|v| v.parse::<u64>().ok())
1054                && version > 0
1055            {
1056                files.push((version, entry.path()));
1057            }
1058        }
1059    }
1060    files.sort_by_key(|(version, _)| *version);
1061    files
1062}
1063
1064fn open_version_files(dir: &Path, name: &str) -> Result<Vec<VersionedFile>, LogError> {
1065    let mut files = Vec::new();
1066    for (version, path) in version_files(dir, name) {
1067        files.push(VersionedFile {
1068            version,
1069            file: IndexedFile::open(&path, false, false)?,
1070        });
1071        close_cold_version_files(&mut files);
1072    }
1073    Ok(files)
1074}
1075
1076fn refresh_version_files(
1077    dir: &Path,
1078    name: &str,
1079    files: &mut Vec<VersionedFile>,
1080) -> Result<(), LogError> {
1081    for file in &mut *files {
1082        file.file.refresh()?;
1083    }
1084
1085    for (version, path) in version_files(dir, name) {
1086        if files.iter().all(|file| file.version != version) {
1087            files.push(VersionedFile {
1088                version,
1089                file: IndexedFile::open(&path, false, false)?,
1090            });
1091            close_cold_version_files(files);
1092        }
1093    }
1094    files.sort_by_key(|file| file.version);
1095    Ok(())
1096}
1097
1098fn close_cold_version_files(files: &mut [VersionedFile]) {
1099    let mut cached_readers = 0;
1100    for file in files.iter_mut().rev() {
1101        if file.file.writable {
1102            continue;
1103        }
1104        if file.file.file.is_some() {
1105            if cached_readers < MAX_CACHED_READ_VERSION_FILES {
1106                cached_readers += 1;
1107            } else {
1108                file.file.close_cached_reader();
1109            }
1110        }
1111    }
1112}
1113
1114/// For each version, the first transaction in any later version. Frames at
1115/// or beyond this cutoff are stale appends from a deposed writer.
1116fn version_cutoffs(files: &[VersionedFile]) -> Vec<u64> {
1117    let mut cutoffs = vec![u64::MAX; files.len()];
1118    let mut cutoff = u64::MAX;
1119    for (index, file) in files.iter().enumerate().rev() {
1120        cutoffs[index] = cutoff;
1121        if let Some(first) = file.file.frames.first() {
1122            cutoff = cutoff.min(first.t);
1123        }
1124    }
1125    cutoffs
1126}
1127
1128fn validated_version_cutoffs(files: &[VersionedFile]) -> Result<Vec<u64>, LogError> {
1129    let cutoffs = version_cutoffs(files);
1130    let mut previous_t: Option<u64> = None;
1131    for (file, cutoff) in files.iter().zip(&cutoffs) {
1132        let retained = file.file.frames.partition_point(|frame| frame.t < *cutoff);
1133        if retained == 0 {
1134            continue;
1135        }
1136        file.file.validate_contiguous_prefix(retained)?;
1137        let first_t = file.file.frames[0].t;
1138        if previous_t.is_some_and(|previous| previous.checked_add(1) != Some(first_t)) {
1139            return Err(LogError::Corrupt);
1140        }
1141        previous_t = Some(file.file.frames[retained - 1].t);
1142    }
1143    Ok(cutoffs)
1144}
1145
1146fn merged_next_t(files: &[VersionedFile], cutoffs: &[u64]) -> Result<u64, LogError> {
1147    for (file, cutoff) in files.iter().zip(cutoffs).rev() {
1148        let retained = file.file.frames.partition_point(|frame| frame.t < *cutoff);
1149        if retained > 0 {
1150            return file.file.frames[retained - 1]
1151                .t
1152                .checked_add(1)
1153                .ok_or(LogError::Corrupt);
1154        }
1155    }
1156    Ok(1)
1157}
1158
1159fn range_is_merged_indexed(
1160    files: &[VersionedFile],
1161    cutoffs: &[u64],
1162    end: Option<u64>,
1163) -> Result<bool, LogError> {
1164    end.map_or(Ok(false), |end| Ok(end <= merged_next_t(files, cutoffs)?))
1165}
1166
1167fn read_merged_range(
1168    files: &[VersionedFile],
1169    cutoffs: &[u64],
1170    start: u64,
1171    end: Option<u64>,
1172) -> Result<Vec<TxRecord>, LogError> {
1173    let mut records = Vec::new();
1174    for (file, cutoff) in files.iter().zip(cutoffs) {
1175        let end = Some(end.map_or(*cutoff, |end| end.min(*cutoff)));
1176        records.extend(file.file.tx_range(start, end)?);
1177    }
1178    Ok(records)
1179}
1180
1181fn encode_record(record: &TxRecord) -> Vec<u8> {
1182    let mut out = Vec::new();
1183    out.extend_from_slice(&record.t.to_be_bytes());
1184    out.extend_from_slice(&record.tx_instant.to_be_bytes());
1185    out.extend_from_slice(&(record.datoms.len() as u64).to_be_bytes());
1186    for d in &record.datoms {
1187        out.extend_from_slice(&d.e.raw().to_be_bytes());
1188        out.extend_from_slice(&d.a.raw().to_be_bytes());
1189        out.extend_from_slice(&d.tx.raw().to_be_bytes());
1190        out.push(u8::from(d.added));
1191        let v = encode_value(&d.v);
1192        out.extend_from_slice(&(v.len() as u64).to_be_bytes());
1193        out.extend_from_slice(&v);
1194    }
1195    out
1196}
1197fn decode_record(mut bytes: &[u8]) -> Result<TxRecord, LogError> {
1198    fn take<'a>(bytes: &mut &'a [u8], n: usize) -> Result<&'a [u8], LogError> {
1199        let value = bytes.get(..n).ok_or(LogError::Corrupt)?;
1200        *bytes = &bytes[n..];
1201        Ok(value)
1202    }
1203    fn u64_be(bytes: &mut &[u8]) -> Result<u64, LogError> {
1204        Ok(u64::from_be_bytes(
1205            take(bytes, 8)?.try_into().map_err(|_| LogError::Corrupt)?,
1206        ))
1207    }
1208    let t = u64_be(&mut bytes)?;
1209    let tx_instant = i64::from_be_bytes(
1210        take(&mut bytes, 8)?
1211            .try_into()
1212            .map_err(|_| LogError::Corrupt)?,
1213    );
1214    let count = u64_be(&mut bytes)?;
1215    let mut datoms = Vec::new();
1216    for _ in 0..count {
1217        let e = EntityId::from_raw(u64_be(&mut bytes)?);
1218        let a = EntityId::from_raw(u64_be(&mut bytes)?);
1219        let tx = EntityId::from_raw(u64_be(&mut bytes)?);
1220        let added = take(&mut bytes, 1)?[0] != 0;
1221        let len = usize::try_from(u64_be(&mut bytes)?).map_err(|_| LogError::Corrupt)?;
1222        let raw = take(&mut bytes, len)?;
1223        let (v, used) = decode_value(raw).map_err(|_| LogError::Corrupt)?;
1224        if used != len {
1225            return Err(LogError::Corrupt);
1226        }
1227        datoms.push(Datom { e, a, v, tx, added });
1228    }
1229    if !bytes.is_empty() {
1230        return Err(LogError::Corrupt);
1231    }
1232    Ok(TxRecord {
1233        t,
1234        tx_instant,
1235        datoms,
1236    })
1237}
1238
1239fn frame_header(payload_len: usize) -> Result<[u8; 8], LogError> {
1240    let payload_len = u64::try_from(payload_len).map_err(|_| LogError::Corrupt)?;
1241    if payload_len & CHECKSUMMED_FRAME != 0 {
1242        return Err(LogError::Corrupt);
1243    }
1244    Ok((payload_len | CHECKSUMMED_FRAME).to_be_bytes())
1245}
1246
1247fn frame_payload_len(header: [u8; 8]) -> Result<(usize, bool), LogError> {
1248    let encoded = u64::from_be_bytes(header);
1249    let checksummed = encoded & CHECKSUMMED_FRAME != 0;
1250    let payload_len = encoded & !CHECKSUMMED_FRAME;
1251    Ok((
1252        usize::try_from(payload_len).map_err(|_| LogError::Corrupt)?,
1253        checksummed,
1254    ))
1255}
1256
1257fn frame_checksum(header: [u8; 8], payload: &[u8]) -> u32 {
1258    crc32c::crc32c_append(crc32c::crc32c(&header), payload)
1259}
1260
1261#[cfg(unix)]
1262fn read_exact_at(file: &File, mut bytes: &mut [u8], mut offset: u64) -> io::Result<()> {
1263    use std::os::unix::fs::FileExt;
1264    while !bytes.is_empty() {
1265        match file.read_at(bytes, offset) {
1266            Ok(0) => return Err(io::ErrorKind::UnexpectedEof.into()),
1267            Ok(read) => {
1268                offset = offset
1269                    .checked_add(u64::try_from(read).expect("read length fits u64"))
1270                    .ok_or_else(|| io::Error::other("file offset overflow"))?;
1271                bytes = &mut bytes[read..];
1272            }
1273            Err(error) if error.kind() == io::ErrorKind::Interrupted => {}
1274            Err(error) => return Err(error),
1275        }
1276    }
1277    Ok(())
1278}
1279
1280#[cfg(windows)]
1281fn read_exact_at(file: &File, mut bytes: &mut [u8], mut offset: u64) -> io::Result<()> {
1282    use std::os::windows::fs::FileExt;
1283    while !bytes.is_empty() {
1284        match file.seek_read(bytes, offset) {
1285            Ok(0) => return Err(io::ErrorKind::UnexpectedEof.into()),
1286            Ok(read) => {
1287                offset = offset
1288                    .checked_add(u64::try_from(read).expect("read length fits u64"))
1289                    .ok_or_else(|| io::Error::other("file offset overflow"))?;
1290                bytes = &mut bytes[read..];
1291            }
1292            Err(error) if error.kind() == io::ErrorKind::Interrupted => {}
1293            Err(error) => return Err(error),
1294        }
1295    }
1296    Ok(())
1297}
1298
1299#[cfg(not(any(unix, windows)))]
1300fn read_exact_at(file: &File, bytes: &mut [u8], offset: u64) -> io::Result<()> {
1301    use std::io::{Read, Seek, SeekFrom};
1302    let mut file = file.try_clone()?;
1303    file.seek(SeekFrom::Start(offset))?;
1304    file.read_exact(bytes)
1305}
1306
1307/// Indexes fully written records starting at `offset`, returning the byte end
1308/// of that durable prefix.
1309///
1310/// A record cut short by a crash mid-append (truncated length prefix or
1311/// payload/checksum) ends the scan; a fully present record with a checksum
1312/// mismatch or invalid payload is genuine corruption and errors. Legacy
1313/// length-only records remain readable, while newly written records set the
1314/// high bit of the length word and carry a trailing CRC32C.
1315fn scan_frames(file: &File, offset: u64) -> Result<(Vec<FrameIndex>, u64), LogError> {
1316    let file_len = file.metadata()?.len();
1317    let mut frames = Vec::new();
1318    let mut durable_len = offset;
1319    loop {
1320        if file_len.saturating_sub(durable_len) < 8 {
1321            break;
1322        }
1323        let mut len = [0; 8];
1324        read_exact_at(file, &mut len, durable_len)?;
1325        let (payload_len, checksummed) = frame_payload_len(len)?;
1326        let frame_len = 8_u64
1327            .checked_add(u64::try_from(payload_len).map_err(|_| LogError::Corrupt)?)
1328            .and_then(|len| {
1329                len.checked_add(if checksummed {
1330                    u64::try_from(FRAME_CHECKSUM_LEN).expect("checksum length fits u64")
1331                } else {
1332                    0
1333                })
1334            })
1335            .ok_or(LogError::Corrupt)?;
1336        // Check the bytes remaining before allocating from an untrusted length
1337        // word. A short final frame is the recoverable crash-tail case.
1338        if file_len.saturating_sub(durable_len) < frame_len {
1339            break;
1340        }
1341        let mut payload = vec![0; payload_len];
1342        let payload_offset = durable_len.checked_add(8).ok_or(LogError::Corrupt)?;
1343        read_exact_at(file, &mut payload, payload_offset)?;
1344        if checksummed {
1345            let mut stored_checksum = [0; FRAME_CHECKSUM_LEN];
1346            let checksum_offset = payload_offset
1347                .checked_add(u64::try_from(payload_len).map_err(|_| LogError::Corrupt)?)
1348                .ok_or(LogError::Corrupt)?;
1349            read_exact_at(file, &mut stored_checksum, checksum_offset)?;
1350            if u32::from_be_bytes(stored_checksum) != frame_checksum(len, &payload) {
1351                return Err(LogError::Corrupt);
1352            }
1353        }
1354        let record = decode_record(&payload)?;
1355        frames.push(FrameIndex {
1356            t: record.t,
1357            offset: durable_len,
1358            len: frame_len,
1359        });
1360        durable_len = durable_len
1361            .checked_add(frame_len)
1362            .ok_or(LogError::Corrupt)?;
1363    }
1364    Ok((frames, durable_len))
1365}
1366
1367/// Appends one checksummed, length-prefixed encoded record to `out`.
1368///
1369/// The high bit of the length word identifies the checksummed frame format;
1370/// the remaining 63 bits are the payload length. A big-endian CRC32C over the
1371/// encoded length word and payload follows the payload.
1372///
1373/// # Errors
1374/// Returns an error if the record payload length is not representable.
1375pub fn append_framed_record(out: &mut Vec<u8>, record: &TxRecord) -> Result<(), LogError> {
1376    let payload = encode_record(record);
1377    let header = frame_header(payload.len())?;
1378    out.extend_from_slice(&header);
1379    out.extend_from_slice(&payload);
1380    out.extend_from_slice(&frame_checksum(header, &payload).to_be_bytes());
1381    Ok(())
1382}
1383
1384/// Decodes all records from a framed byte slice.
1385///
1386/// Unlike filesystem crash recovery, native stores publish whole values
1387/// atomically, so any trailing partial frame is treated as corruption.
1388/// Both legacy length-only frames and checksummed frames are accepted.
1389///
1390/// # Errors
1391/// Returns an error when any frame is truncated, has an invalid length, or
1392/// checksum, or contains a corrupt encoded transaction record.
1393pub fn decode_framed_records(mut bytes: &[u8]) -> Result<Vec<TxRecord>, LogError> {
1394    let mut records = Vec::new();
1395    while !bytes.is_empty() {
1396        if bytes.len() < 8 {
1397            return Err(LogError::Corrupt);
1398        }
1399        let header: [u8; 8] = bytes[..8].try_into().map_err(|_| LogError::Corrupt)?;
1400        let (payload_len, checksummed) = frame_payload_len(header)?;
1401        bytes = &bytes[8..];
1402        let payload = bytes.get(..payload_len).ok_or(LogError::Corrupt)?;
1403        bytes = &bytes[payload_len..];
1404        if checksummed {
1405            let stored_checksum = u32::from_be_bytes(
1406                bytes
1407                    .get(..FRAME_CHECKSUM_LEN)
1408                    .ok_or(LogError::Corrupt)?
1409                    .try_into()
1410                    .map_err(|_| LogError::Corrupt)?,
1411            );
1412            if stored_checksum != frame_checksum(header, payload) {
1413                return Err(LogError::Corrupt);
1414            }
1415            bytes = &bytes[FRAME_CHECKSUM_LEN..];
1416        }
1417        records.push(decode_record(payload)?);
1418    }
1419    Ok(records)
1420}
1421
1422#[cfg(test)]
1423mod tests {
1424    use super::*;
1425
1426    #[test]
1427    fn versioned_log_bounds_cached_read_descriptors() {
1428        let dir = tempfile::tempdir().expect("tempdir");
1429        let segment_count = MAX_CACHED_READ_VERSION_FILES + 5;
1430        for version in 1..=u64::try_from(segment_count).expect("segment count fits u64") {
1431            File::create(version_path(dir.path(), "db", version)).expect("create segment");
1432        }
1433
1434        let log = VersionedLog::open_read_only(dir.path(), "db").expect("open log");
1435        let state = log.state.read().expect("log lock");
1436        assert_eq!(state.files.len(), segment_count);
1437        assert!(
1438            state
1439                .files
1440                .iter()
1441                .filter(|file| file.file.file.is_some())
1442                .count()
1443                <= MAX_CACHED_READ_VERSION_FILES
1444        );
1445    }
1446}