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