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