1use 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#[derive(Clone, Debug, Eq, PartialEq)]
28pub struct TxRecord {
29 pub t: u64,
31 pub tx_instant: i64,
33 pub datoms: Vec<Datom>,
35}
36
37#[derive(Debug, Error)]
39pub enum LogError {
40 #[error("log I/O failed: {0}")]
42 Io(#[from] io::Error),
43 #[error("corrupt transaction log")]
45 Corrupt,
46 #[error("native transaction log store failed: {0}")]
48 Native(String),
49 #[error("this transaction log requires asynchronous access")]
51 AsyncOnly,
52 #[error("transaction log is encrypted; no storage key is configured")]
54 Encrypted,
55 #[error("transaction log is not encrypted, but a storage key is configured")]
59 Unencrypted,
60 #[error("transaction log record uses storage key epoch {0}, which is unavailable")]
62 MissingKeyEpoch(u32),
63 #[error("transaction log record encryption failed: {0}")]
65 Crypt(#[from] CryptError),
66}
67
68pub struct LogCipher {
79 lineage: Vec<u8>,
80 keys: RwLock<Arc<KeySnapshot>>,
81}
82
83struct 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(¤t_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 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 #[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 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 #[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 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#[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 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#[async_trait]
268pub trait TransactionLog: Send + Sync {
269 fn append(&self, record: &TxRecord) -> Result<(), LogError>;
274 async fn append_async(&self, record: &TxRecord) -> Result<(), LogError> {
281 self.append(record)
282 }
283 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 fn tx_range(&self, start: u64, end: Option<u64>) -> Result<Vec<TxRecord>, LogError>;
303 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 fn replay(&self) -> Result<Vec<TxRecord>, LogError> {
319 self.tx_range(0, None)
320 }
321 async fn replay_async(&self) -> Result<Vec<TxRecord>, LogError> {
326 self.tx_range_async(0, None).await
327 }
328}
329
330#[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
354pub struct FileLog {
364 state: RwLock<IndexedFile>,
365}
366
367impl FileLog {
368 pub fn open(path: impl AsRef<Path>) -> Result<Self, LogError> {
374 Self::open_with(path, None)
375 }
376
377 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
431pub 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 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 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 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 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 #[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 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 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
661fn 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#[async_trait]
692pub trait NativeLogStorage: Send + Sync {
693 async fn put_batch(
706 &self,
707 name: &str,
708 version: u64,
709 records: &[(u64, Vec<u8>)],
710 ) -> Result<bool, LogError>;
711 async fn read_record(
717 &self,
718 name: &str,
719 version: u64,
720 t: u64,
721 ) -> Result<Option<Vec<u8>>, LogError>;
722 async fn list_records(&self, name: &str) -> Result<Vec<(u64, u64)>, LogError>;
728 async fn read_legacy_chunk(
734 &self,
735 name: &str,
736 version: u64,
737 chunk: u64,
738 ) -> Result<Option<Vec<u8>>, LogError>;
739 async fn list_legacy_chunks(&self, name: &str) -> Result<Vec<(u64, u64)>, LogError>;
747 async fn delete_all(&self, name: &str) -> Result<(), LogError>;
752}
753
754pub 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: tokio::sync::Mutex<u64>,
770}
771
772impl<S: NativeLogStorage + ?Sized + 'static> NativeVersionedLog<S> {
773 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 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 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 #[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 #[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 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 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 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 let mut per_version: BTreeMap<u64, Vec<TxRecord>> = BTreeMap::new();
923
924 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 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 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
980type VersionedRecords = Arc<Mutex<Vec<(u64, TxRecord)>>>;
983
984#[derive(Clone, Default)]
990pub struct MemLogRegistry {
991 logs: Arc<Mutex<HashMap<String, VersionedRecords>>>,
992}
993
994impl MemLogRegistry {
995 #[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 #[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 #[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 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
1055pub struct MemVersionedLog {
1059 records: VersionedRecords,
1060 write_version: u64,
1061 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
1122struct 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 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 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
1387fn 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
1480fn 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
1673fn 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 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
1737fn 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
1752pub fn append_framed_record(out: &mut Vec<u8>, record: &TxRecord) -> Result<(), LogError> {
1757 append_framed_payload(out, &encode_record(record))
1758}
1759
1760pub 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
1779pub fn decode_framed_records(bytes: &[u8]) -> Result<Vec<TxRecord>, LogError> {
1785 decode_framed_payloads(bytes, &RecordCodec::plaintext())
1786}
1787
1788pub 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
1806fn 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 #[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 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}