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 current_epoch: u32,
81 keys: BTreeMap<u32, SecretKey>,
82}
83
84impl LogCipher {
85 pub fn new(
95 lineage: impl Into<Vec<u8>>,
96 current_epoch: u32,
97 keys: impl IntoIterator<Item = (u32, SecretKey)>,
98 ) -> Result<Self, LogError> {
99 let keys = keys.into_iter().collect::<BTreeMap<_, _>>();
100 if !keys.contains_key(¤t_epoch) {
101 return Err(LogError::MissingKeyEpoch(current_epoch));
102 }
103 Ok(Self {
104 lineage: lineage.into(),
105 current_epoch,
106 keys,
107 })
108 }
109
110 #[must_use]
112 pub fn with_key(lineage: impl Into<Vec<u8>>, epoch: u32, key: SecretKey) -> Self {
113 Self {
114 lineage: lineage.into(),
115 current_epoch: epoch,
116 keys: BTreeMap::from([(epoch, key)]),
117 }
118 }
119
120 #[must_use]
122 pub fn current_epoch(&self) -> u32 {
123 self.current_epoch
124 }
125
126 fn key(&self, epoch: u32) -> Result<&SecretKey, LogError> {
127 self.keys
128 .get(&epoch)
129 .ok_or(LogError::MissingKeyEpoch(epoch))
130 }
131
132 fn seal(&self, log_version: u64, t: u64, plaintext: &[u8]) -> Result<Vec<u8>, LogError> {
133 Ok(encrypt_log_record(
134 self.key(self.current_epoch)?,
135 self.current_epoch,
136 &self.lineage,
137 log_version,
138 t,
139 plaintext,
140 )?)
141 }
142
143 fn open(
146 &self,
147 log_version: u64,
148 header: LogHeader,
149 payload: &[u8],
150 ) -> Result<Vec<u8>, LogError> {
151 Ok(decrypt_log_record(
152 self.key(header.epoch)?,
153 &self.lineage,
154 log_version,
155 payload,
156 )?)
157 }
158}
159
160#[derive(Clone, Default)]
163struct RecordCodec {
164 cipher: Option<Arc<LogCipher>>,
165 log_version: u64,
166}
167
168impl RecordCodec {
169 fn plaintext() -> Self {
170 Self::default()
171 }
172
173 fn new(cipher: Option<Arc<LogCipher>>, log_version: u64) -> Self {
174 Self {
175 cipher,
176 log_version,
177 }
178 }
179
180 fn encode(&self, record: &TxRecord) -> Result<Vec<u8>, LogError> {
181 let encoded = encode_record(record);
182 match &self.cipher {
183 Some(cipher) => cipher.seal(self.log_version, record.t, &encoded),
184 None => Ok(encoded),
185 }
186 }
187
188 fn decode(&self, payload: &[u8]) -> Result<TxRecord, LogError> {
189 match (&self.cipher, is_encrypted_log_record(payload)) {
190 (Some(cipher), true) => {
191 let header = parse_log_header(payload)?;
192 let record = decode_record(&cipher.open(self.log_version, header, payload)?)?;
193 if record.t != header.t {
200 return Err(LogError::Corrupt);
201 }
202 Ok(record)
203 }
204 (None, true) => Err(LogError::Encrypted),
205 (Some(_), false) => Err(LogError::Unencrypted),
206 (None, false) => decode_record(payload),
207 }
208 }
209}
210
211#[async_trait]
213pub trait TransactionLog: Send + Sync {
214 fn append(&self, record: &TxRecord) -> Result<(), LogError>;
219 async fn append_async(&self, record: &TxRecord) -> Result<(), LogError> {
226 self.append(record)
227 }
228 async fn append_batch_async(&self, records: &[TxRecord]) -> Result<(), LogError> {
238 for record in records {
239 self.append_async(record).await?;
240 }
241 Ok(())
242 }
243 fn tx_range(&self, start: u64, end: Option<u64>) -> Result<Vec<TxRecord>, LogError>;
248 async fn tx_range_async(
253 &self,
254 start: u64,
255 end: Option<u64>,
256 ) -> Result<Vec<TxRecord>, LogError> {
257 self.tx_range(start, end)
258 }
259 fn replay(&self) -> Result<Vec<TxRecord>, LogError> {
264 self.tx_range(0, None)
265 }
266 async fn replay_async(&self) -> Result<Vec<TxRecord>, LogError> {
271 self.tx_range_async(0, None).await
272 }
273}
274
275#[derive(Clone, Default)]
277pub struct MemoryLog(Arc<RwLock<Vec<TxRecord>>>);
278impl TransactionLog for MemoryLog {
279 fn append(&self, record: &TxRecord) -> Result<(), LogError> {
280 let mut records = self.0.write().expect("poisoned log lock");
281 if records.last().map_or(1, |r| r.t + 1) != record.t {
282 return Err(LogError::Corrupt);
283 }
284 records.push(record.clone());
285 Ok(())
286 }
287 fn tx_range(&self, start: u64, end: Option<u64>) -> Result<Vec<TxRecord>, LogError> {
288 Ok(self
289 .0
290 .read()
291 .expect("poisoned log lock")
292 .iter()
293 .filter(|r| r.t >= start && end.is_none_or(|e| r.t < e))
294 .cloned()
295 .collect())
296 }
297}
298
299pub struct FileLog {
309 state: RwLock<IndexedFile>,
310}
311
312impl FileLog {
313 pub fn open(path: impl AsRef<Path>) -> Result<Self, LogError> {
319 Self::open_with(path, None)
320 }
321
322 pub fn open_sealed(path: impl AsRef<Path>, cipher: Arc<LogCipher>) -> Result<Self, LogError> {
331 Self::open_with(path, Some(cipher))
332 }
333
334 fn open_with(path: impl AsRef<Path>, cipher: Option<Arc<LogCipher>>) -> Result<Self, LogError> {
335 let path = path.as_ref().to_path_buf();
336 if let Some(parent) = path.parent() {
337 fs::create_dir_all(parent)?;
338 }
339 let file = IndexedFile::open(&path, true, true, RecordCodec::new(cipher, 0))?;
340 file.validate_contiguous_prefix(file.frames.len())?;
341 Ok(Self {
342 state: RwLock::new(file),
343 })
344 }
345}
346impl TransactionLog for FileLog {
347 fn append(&self, record: &TxRecord) -> Result<(), LogError> {
348 let mut state = self.state.write().expect("poisoned log lock");
349 state.refresh()?;
350 state.validate_contiguous_prefix(state.frames.len())?;
351 if next_t(&state.frames)? != record.t {
352 return Err(LogError::Corrupt);
353 }
354 state.append(record)
355 }
356 fn tx_range(&self, start: u64, end: Option<u64>) -> Result<Vec<TxRecord>, LogError> {
357 if end.is_some_and(|end| end <= start) {
358 return Ok(Vec::new());
359 }
360 let indexed = {
361 let state = self.state.read().expect("poisoned log lock");
362 range_is_indexed(&state.frames, end)?
363 };
364 if !indexed {
365 let mut state = self.state.write().expect("poisoned log lock");
366 state.refresh()?;
367 state.validate_contiguous_prefix(state.frames.len())?;
368 }
369 self.state
370 .read()
371 .expect("poisoned log lock")
372 .tx_range(start, end)
373 }
374}
375
376pub struct VersionedLog {
388 dir: PathBuf,
389 name: String,
390 cipher: Option<Arc<LogCipher>>,
391 state: RwLock<VersionedLogState>,
392}
393
394struct VersionedLogState {
395 files: Vec<VersionedFile>,
396 write_version: Option<u64>,
397 next_t: u64,
398}
399
400struct VersionedFile {
401 version: u64,
402 file: IndexedFile,
403}
404
405impl VersionedLog {
406 pub fn open(dir: impl AsRef<Path>, name: &str, write_version: u64) -> Result<Self, LogError> {
414 Self::open_with(dir, name, write_version, None)
415 }
416
417 pub fn open_sealed(
426 dir: impl AsRef<Path>,
427 name: &str,
428 write_version: u64,
429 cipher: Arc<LogCipher>,
430 ) -> Result<Self, LogError> {
431 Self::open_with(dir, name, write_version, Some(cipher))
432 }
433
434 fn open_with(
435 dir: impl AsRef<Path>,
436 name: &str,
437 write_version: u64,
438 cipher: Option<Arc<LogCipher>>,
439 ) -> Result<Self, LogError> {
440 let dir = dir.as_ref().to_path_buf();
441 fs::create_dir_all(&dir)?;
442 let write_path = version_path(&dir, name, write_version);
443 let mut files = Vec::new();
444 for (version, path) in version_files(&dir, name) {
445 let writable = version == write_version;
446 files.push(VersionedFile {
447 version,
448 file: IndexedFile::open(
449 &path,
450 writable,
451 writable,
452 RecordCodec::new(cipher.clone(), version),
453 )?,
454 });
455 close_cold_version_files(&mut files);
456 }
457 if !files.iter().any(|file| file.version == write_version) {
458 files.push(VersionedFile {
459 version: write_version,
460 file: IndexedFile::open(
461 &write_path,
462 true,
463 true,
464 RecordCodec::new(cipher.clone(), write_version),
465 )?,
466 });
467 files.sort_by_key(|file| file.version);
468 }
469 close_cold_version_files(&mut files);
470 let cutoffs = validated_version_cutoffs(&files)?;
471 let next_t = merged_next_t(&files, &cutoffs)?;
472 Ok(Self {
473 dir,
474 name: name.to_owned(),
475 cipher,
476 state: RwLock::new(VersionedLogState {
477 files,
478 write_version: Some(write_version),
479 next_t,
480 }),
481 })
482 }
483
484 pub fn open_read_only(dir: impl AsRef<Path>, name: &str) -> Result<Self, LogError> {
490 Self::open_read_only_with(dir, name, None)
491 }
492
493 pub fn open_read_only_sealed(
499 dir: impl AsRef<Path>,
500 name: &str,
501 cipher: Arc<LogCipher>,
502 ) -> Result<Self, LogError> {
503 Self::open_read_only_with(dir, name, Some(cipher))
504 }
505
506 fn open_read_only_with(
507 dir: impl AsRef<Path>,
508 name: &str,
509 cipher: Option<Arc<LogCipher>>,
510 ) -> Result<Self, LogError> {
511 let dir = dir.as_ref().to_path_buf();
512 let mut files = open_version_files(&dir, name, cipher.as_ref())?;
513 close_cold_version_files(&mut files);
514 let cutoffs = validated_version_cutoffs(&files)?;
515 Ok(Self {
516 name: name.to_owned(),
517 cipher,
518 state: RwLock::new(VersionedLogState {
519 next_t: merged_next_t(&files, &cutoffs)?,
520 write_version: None,
521 files,
522 }),
523 dir,
524 })
525 }
526
527 #[must_use]
529 pub fn exists(dir: impl AsRef<Path>, name: &str) -> bool {
530 !version_files(dir.as_ref(), name).is_empty()
531 }
532
533 pub fn delete_all(dir: impl AsRef<Path>, name: &str) -> Result<(), LogError> {
538 for (_, path) in version_files(dir.as_ref(), name) {
539 match fs::remove_file(&path) {
540 Ok(()) => {}
541 Err(error) if error.kind() == io::ErrorKind::NotFound => {}
542 Err(error) => return Err(error.into()),
543 }
544 }
545 Ok(())
546 }
547}
548
549impl TransactionLog for VersionedLog {
550 fn append(&self, record: &TxRecord) -> Result<(), LogError> {
551 let mut state = self.state.write().expect("poisoned log lock");
552 if state.next_t != record.t {
553 return Err(LogError::Corrupt);
554 }
555 let write_version = state
556 .write_version
557 .ok_or_else(|| LogError::Native("transaction log is read-only".into()))?;
558 let write_index = state
559 .files
560 .iter()
561 .position(|file| file.version == write_version)
562 .ok_or(LogError::Corrupt)?;
563 let cutoffs = version_cutoffs(&state.files);
564 if cutoffs[write_index] == u64::MAX
568 && !state.files[write_index].file.frames.is_empty()
569 && next_t(&state.files[write_index].file.frames)? != record.t
570 {
571 return Err(LogError::Corrupt);
572 }
573 state.files[write_index].file.append(record)?;
574 state.next_t = state.next_t.checked_add(1).ok_or(LogError::Corrupt)?;
575 Ok(())
576 }
577
578 fn tx_range(&self, start: u64, end: Option<u64>) -> Result<Vec<TxRecord>, LogError> {
579 if end.is_some_and(|end| end <= start) {
580 return Ok(Vec::new());
581 }
582 {
583 let state = self.state.read().expect("poisoned log lock");
584 let cutoffs = validated_version_cutoffs(&state.files)?;
585 if range_is_merged_indexed(&state.files, &cutoffs, end)? {
586 return read_merged_range(&state.files, &cutoffs, start, end);
587 }
588 }
589 {
590 let mut state = self.state.write().expect("poisoned log lock");
591 refresh_version_files(
592 &self.dir,
593 &self.name,
594 self.cipher.as_ref(),
595 &mut state.files,
596 )?;
597 close_cold_version_files(&mut state.files);
598 validated_version_cutoffs(&state.files)?;
599 }
600 let state = self.state.read().expect("poisoned log lock");
601 let cutoffs = validated_version_cutoffs(&state.files)?;
602 read_merged_range(&state.files, &cutoffs, start, end)
603 }
604}
605
606fn merge_versions(mut per_version: Vec<Vec<TxRecord>>) -> Vec<TxRecord> {
610 let mut cutoff = u64::MAX;
611 for records in per_version.iter_mut().rev() {
612 let first = records.first().map(|r| r.t);
613 records.retain(|r| r.t < cutoff);
614 if let Some(first) = first {
615 cutoff = cutoff.min(first);
616 }
617 }
618 per_version.into_iter().flatten().collect()
619}
620
621#[async_trait]
637pub trait NativeLogStorage: Send + Sync {
638 async fn put_batch(
651 &self,
652 name: &str,
653 version: u64,
654 records: &[(u64, Vec<u8>)],
655 ) -> Result<bool, LogError>;
656 async fn read_record(
662 &self,
663 name: &str,
664 version: u64,
665 t: u64,
666 ) -> Result<Option<Vec<u8>>, LogError>;
667 async fn list_records(&self, name: &str) -> Result<Vec<(u64, u64)>, LogError>;
673 async fn read_legacy_chunk(
679 &self,
680 name: &str,
681 version: u64,
682 chunk: u64,
683 ) -> Result<Option<Vec<u8>>, LogError>;
684 async fn list_legacy_chunks(&self, name: &str) -> Result<Vec<(u64, u64)>, LogError>;
692 async fn delete_all(&self, name: &str) -> Result<(), LogError>;
697}
698
699pub struct NativeVersionedLog<S: ?Sized> {
708 storage: Arc<S>,
709 name: String,
710 write_version: u64,
711 read_only: bool,
712 cipher: Option<Arc<LogCipher>>,
713 next_t: tokio::sync::Mutex<u64>,
715}
716
717impl<S: NativeLogStorage + ?Sized + 'static> NativeVersionedLog<S> {
718 pub async fn open(storage: Arc<S>, name: &str, write_version: u64) -> Result<Self, LogError> {
723 Self::open_with(storage, name, write_version, None).await
724 }
725
726 pub async fn open_sealed(
732 storage: Arc<S>,
733 name: &str,
734 write_version: u64,
735 cipher: Arc<LogCipher>,
736 ) -> Result<Self, LogError> {
737 Self::open_with(storage, name, write_version, Some(cipher)).await
738 }
739
740 async fn open_with(
741 storage: Arc<S>,
742 name: &str,
743 write_version: u64,
744 cipher: Option<Arc<LogCipher>>,
745 ) -> Result<Self, LogError> {
746 let records = read_native_merged(storage.as_ref(), name, cipher.as_ref()).await?;
750 let next_t = records.last().map_or(1, |r| r.t + 1);
751 Ok(Self {
752 storage,
753 name: name.to_owned(),
754 write_version,
755 read_only: false,
756 cipher,
757 next_t: tokio::sync::Mutex::new(next_t),
758 })
759 }
760
761 #[must_use]
764 pub fn open_read_only(storage: Arc<S>, name: &str) -> Self {
765 Self::open_read_only_with(storage, name, None)
766 }
767
768 #[must_use]
770 pub fn open_read_only_sealed(storage: Arc<S>, name: &str, cipher: Arc<LogCipher>) -> Self {
771 Self::open_read_only_with(storage, name, Some(cipher))
772 }
773
774 fn open_read_only_with(storage: Arc<S>, name: &str, cipher: Option<Arc<LogCipher>>) -> Self {
775 Self {
776 storage,
777 name: name.to_owned(),
778 write_version: 0,
779 read_only: true,
780 cipher,
781 next_t: tokio::sync::Mutex::new(0),
782 }
783 }
784}
785
786#[async_trait]
787impl<S: NativeLogStorage + ?Sized + 'static> TransactionLog for NativeVersionedLog<S> {
788 fn append(&self, record: &TxRecord) -> Result<(), LogError> {
789 let _ = record;
790 Err(LogError::AsyncOnly)
791 }
792
793 async fn append_async(&self, record: &TxRecord) -> Result<(), LogError> {
794 self.append_batch_async(std::slice::from_ref(record)).await
795 }
796
797 async fn append_batch_async(&self, records: &[TxRecord]) -> Result<(), LogError> {
798 if self.read_only {
799 return Err(LogError::Native("transaction log is read-only".into()));
800 }
801 if records.is_empty() {
802 return Ok(());
803 }
804 let mut next_t = self.next_t.lock().await;
805 for (offset, record) in records.iter().enumerate() {
807 if record.t != *next_t + offset as u64 {
808 return Err(LogError::Corrupt);
809 }
810 }
811 let codec = RecordCodec::new(self.cipher.clone(), self.write_version);
812 let framed = records
813 .iter()
814 .map(|record| {
815 let mut bytes = Vec::new();
816 append_framed_payload(&mut bytes, &codec.encode(record)?)?;
817 Ok((record.t, bytes))
818 })
819 .collect::<Result<Vec<_>, LogError>>()?;
820 if !self
825 .storage
826 .put_batch(&self.name, self.write_version, &framed)
827 .await?
828 {
829 return Err(LogError::Corrupt);
830 }
831 *next_t += records.len() as u64;
832 Ok(())
833 }
834
835 fn tx_range(&self, start: u64, end: Option<u64>) -> Result<Vec<TxRecord>, LogError> {
836 let _ = (start, end);
837 Err(LogError::AsyncOnly)
838 }
839
840 async fn tx_range_async(
841 &self,
842 start: u64,
843 end: Option<u64>,
844 ) -> Result<Vec<TxRecord>, LogError> {
845 let _guard = self.next_t.lock().await;
848 Ok(
849 read_native_merged(self.storage.as_ref(), &self.name, self.cipher.as_ref())
850 .await?
851 .into_iter()
852 .filter(|r| r.t >= start && end.is_none_or(|e| r.t < e))
853 .collect(),
854 )
855 }
856}
857
858async fn read_native_merged<S: NativeLogStorage + ?Sized>(
859 storage: &S,
860 name: &str,
861 cipher: Option<&Arc<LogCipher>>,
862) -> Result<Vec<TxRecord>, LogError> {
863 use std::collections::BTreeMap;
864
865 let mut per_version: BTreeMap<u64, Vec<TxRecord>> = BTreeMap::new();
868
869 let mut chunks = storage.list_legacy_chunks(name).await?;
873 chunks.sort_unstable();
874 for (version, chunk) in chunks {
875 let bytes = storage
876 .read_legacy_chunk(name, version, chunk)
877 .await?
878 .unwrap_or_default();
879 per_version
880 .entry(version)
881 .or_default()
882 .extend(decode_framed_payloads(
883 &bytes,
884 &RecordCodec::new(cipher.map(Arc::clone), version),
885 )?);
886 }
887
888 let mut records = storage.list_records(name).await?;
890 records.sort_unstable();
891 for (version, t) in records {
892 let bytes = storage
893 .read_record(name, version, t)
894 .await?
895 .unwrap_or_default();
896 per_version
897 .entry(version)
898 .or_default()
899 .extend(decode_framed_payloads(
900 &bytes,
901 &RecordCodec::new(cipher.map(Arc::clone), version),
902 )?);
903 }
904
905 let per_version: Vec<Vec<TxRecord>> = per_version
910 .into_values()
911 .map(|mut records| {
912 records.sort_by_key(|record| record.t);
913 records
914 })
915 .collect();
916 let merged = merge_versions(per_version);
917 for pair in merged.windows(2) {
918 if pair[1].t != pair[0].t + 1 {
919 return Err(LogError::Corrupt);
920 }
921 }
922 Ok(merged)
923}
924
925type VersionedRecords = Arc<Mutex<Vec<(u64, TxRecord)>>>;
928
929#[derive(Clone, Default)]
935pub struct MemLogRegistry {
936 logs: Arc<Mutex<HashMap<String, VersionedRecords>>>,
937}
938
939impl MemLogRegistry {
940 #[must_use]
942 pub fn new() -> Self {
943 Self::default()
944 }
945
946 fn entry(&self, name: &str) -> VersionedRecords {
947 Arc::clone(
948 self.logs
949 .lock()
950 .unwrap_or_else(std::sync::PoisonError::into_inner)
951 .entry(name.to_owned())
952 .or_default(),
953 )
954 }
955
956 #[must_use]
959 pub fn open(&self, name: &str, write_version: u64) -> MemVersionedLog {
960 let records = self.entry(name);
961 let next_t = {
962 let guard = records
963 .lock()
964 .unwrap_or_else(std::sync::PoisonError::into_inner);
965 MemVersionedLog::merged(&guard)
966 .last()
967 .map_or(1, |r| r.t + 1)
968 };
969 MemVersionedLog {
970 records,
971 write_version,
972 next_t: Mutex::new(next_t),
973 }
974 }
975
976 #[must_use]
978 pub fn exists(&self, name: &str) -> bool {
979 self.logs
980 .lock()
981 .unwrap_or_else(std::sync::PoisonError::into_inner)
982 .get(name)
983 .is_some_and(|entry| {
984 !entry
985 .lock()
986 .unwrap_or_else(std::sync::PoisonError::into_inner)
987 .is_empty()
988 })
989 }
990
991 pub fn delete_all(&self, name: &str) {
993 self.logs
994 .lock()
995 .unwrap_or_else(std::sync::PoisonError::into_inner)
996 .remove(name);
997 }
998}
999
1000pub struct MemVersionedLog {
1004 records: VersionedRecords,
1005 write_version: u64,
1006 next_t: Mutex<u64>,
1010}
1011
1012impl MemVersionedLog {
1013 fn merged(records: &[(u64, TxRecord)]) -> Vec<TxRecord> {
1014 let mut versions: Vec<u64> = records.iter().map(|(version, _)| *version).collect();
1015 versions.sort_unstable();
1016 versions.dedup();
1017 let per_version = versions
1018 .into_iter()
1019 .map(|version| {
1020 records
1021 .iter()
1022 .filter(|(record_version, _)| *record_version == version)
1023 .map(|(_, record)| record.clone())
1024 .collect::<Vec<_>>()
1025 })
1026 .collect();
1027 merge_versions(per_version)
1028 }
1029}
1030
1031impl TransactionLog for MemVersionedLog {
1032 fn append(&self, record: &TxRecord) -> Result<(), LogError> {
1033 let mut next_t = self
1034 .next_t
1035 .lock()
1036 .unwrap_or_else(std::sync::PoisonError::into_inner);
1037 if *next_t != record.t {
1038 return Err(LogError::Corrupt);
1039 }
1040 self.records
1041 .lock()
1042 .unwrap_or_else(std::sync::PoisonError::into_inner)
1043 .push((self.write_version, record.clone()));
1044 *next_t += 1;
1045 Ok(())
1046 }
1047
1048 fn tx_range(&self, start: u64, end: Option<u64>) -> Result<Vec<TxRecord>, LogError> {
1049 let records = self
1050 .records
1051 .lock()
1052 .unwrap_or_else(std::sync::PoisonError::into_inner);
1053 Ok(Self::merged(&records)
1054 .into_iter()
1055 .filter(|r| r.t >= start && end.is_none_or(|e| r.t < e))
1056 .collect())
1057 }
1058}
1059
1060#[derive(Clone, Copy)]
1061struct FrameIndex {
1062 t: u64,
1063 offset: u64,
1064 len: u64,
1065}
1066
1067struct IndexedFile {
1073 path: PathBuf,
1074 file: Option<Arc<File>>,
1075 writable: bool,
1076 poisoned: bool,
1077 frames: Vec<FrameIndex>,
1078 durable_len: u64,
1079 first_gap: Option<usize>,
1080 codec: RecordCodec,
1081}
1082
1083impl IndexedFile {
1084 fn open(
1085 path: &Path,
1086 writable: bool,
1087 truncate_torn: bool,
1088 codec: RecordCodec,
1089 ) -> Result<Self, LogError> {
1090 let file = Arc::new(open_index_file(path, writable)?);
1091 let (frames, durable_len) = scan_frames(file.as_ref(), 0, &codec)?;
1092 validate_sorted_frames(&frames)?;
1093 if truncate_torn && file.metadata()?.len() > durable_len {
1094 file.set_len(durable_len)?;
1095 file.sync_all()?;
1096 }
1097 Ok(Self {
1098 path: path.to_path_buf(),
1099 file: Some(file),
1100 writable,
1101 poisoned: false,
1102 first_gap: first_gap_index(&frames),
1103 frames,
1104 durable_len,
1105 codec,
1106 })
1107 }
1108
1109 fn refresh(&mut self) -> Result<(), LogError> {
1110 self.ensure_healthy()?;
1111 let file_len = self.physical_len()?;
1112 if file_len < self.durable_len {
1113 let file = self.open_for_read()?;
1114 let (frames, durable_len) = scan_frames(file.as_ref(), 0, &self.codec)?;
1115 validate_sorted_frames(&frames)?;
1116 self.first_gap = first_gap_index(&frames);
1117 self.frames = frames;
1118 self.durable_len = durable_len;
1119 } else if file_len > self.durable_len {
1120 let file = self.open_for_read()?;
1121 let (new_frames, durable_len) =
1122 scan_frames(file.as_ref(), self.durable_len, &self.codec)?;
1123 validate_sorted_extension(&self.frames, &new_frames)?;
1124 let existing_len = self.frames.len();
1125 if self.first_gap.is_none() {
1126 self.first_gap =
1127 extension_first_gap(&self.frames, &new_frames).map(|gap| existing_len + gap);
1128 }
1129 self.frames.extend(new_frames);
1130 self.durable_len = durable_len;
1131 }
1132 Ok(())
1133 }
1134
1135 fn append(&mut self, record: &TxRecord) -> Result<(), LogError> {
1136 self.ensure_healthy()?;
1137 if !self.writable {
1138 return Err(LogError::Native("transaction log is read-only".into()));
1139 }
1140 let file = self.open_for_read()?;
1141 if file.metadata()?.len() != self.durable_len {
1144 return Err(LogError::Corrupt);
1145 }
1146
1147 let mut frame = Vec::new();
1148 append_framed_payload(&mut frame, &self.codec.encode(record)?)?;
1149 let frame_len = u64::try_from(frame.len()).map_err(|_| LogError::Corrupt)?;
1150 let offset = self.durable_len;
1151 let mut writer = file.as_ref();
1152 if let Err(error) = writer.write_all(&frame) {
1153 self.poisoned = true;
1156 return Err(error.into());
1157 }
1158 if let Err(error) = file.sync_all() {
1159 self.poisoned = true;
1160 return Err(error.into());
1161 }
1162 self.durable_len = self
1163 .durable_len
1164 .checked_add(frame_len)
1165 .ok_or(LogError::Corrupt)?;
1166 if self.first_gap.is_none()
1167 && self
1168 .frames
1169 .last()
1170 .is_some_and(|previous| previous.t.checked_add(1) != Some(record.t))
1171 {
1172 self.first_gap = Some(self.frames.len());
1173 }
1174 self.frames.push(FrameIndex {
1175 t: record.t,
1176 offset,
1177 len: frame_len,
1178 });
1179 Ok(())
1180 }
1181
1182 fn tx_range(&self, start: u64, end: Option<u64>) -> Result<Vec<TxRecord>, LogError> {
1183 self.ensure_healthy()?;
1184 let first = self.frames.partition_point(|frame| frame.t < start);
1185 let last = end.map_or(self.frames.len(), |end| {
1186 self.frames.partition_point(|frame| frame.t < end)
1187 });
1188 if first >= last {
1189 return Ok(Vec::new());
1190 }
1191
1192 let file = self.open_for_read()?;
1193 let mut records = Vec::with_capacity(last - first);
1194 let mut chunk_first = first;
1195 while chunk_first < last {
1196 let offset = self.frames[chunk_first].offset;
1197 let mut chunk_last = chunk_first + 1;
1198 while chunk_last < last {
1199 let candidate_end = frame_end(self.frames[chunk_last])?;
1200 if candidate_end.checked_sub(offset).ok_or(LogError::Corrupt)?
1201 > RANGE_READ_CHUNK_BYTES
1202 {
1203 break;
1204 }
1205 chunk_last += 1;
1206 }
1207 let byte_end = frame_end(self.frames[chunk_last - 1])?;
1208 let byte_len = usize::try_from(byte_end.checked_sub(offset).ok_or(LogError::Corrupt)?)
1209 .map_err(|_| LogError::Corrupt)?;
1210 let mut bytes = vec![0; byte_len];
1211 read_exact_at(file.as_ref(), &mut bytes, offset)?;
1212 let chunk_records = decode_framed_payloads(&bytes, &self.codec)?;
1213 if chunk_records.len() != chunk_last - chunk_first
1214 || chunk_records
1215 .iter()
1216 .zip(&self.frames[chunk_first..chunk_last])
1217 .any(|(record, frame)| record.t != frame.t)
1218 {
1219 return Err(LogError::Corrupt);
1220 }
1221 records.extend(chunk_records);
1222 chunk_first = chunk_last;
1223 }
1224 Ok(records)
1225 }
1226
1227 fn validate_contiguous_prefix(&self, retained: usize) -> Result<(), LogError> {
1228 if self.first_gap.is_some_and(|gap| gap < retained) {
1229 return Err(LogError::Corrupt);
1230 }
1231 Ok(())
1232 }
1233
1234 fn ensure_healthy(&self) -> Result<(), LogError> {
1235 if self.poisoned {
1236 return Err(io::Error::other("transaction log handle is poisoned").into());
1237 }
1238 Ok(())
1239 }
1240
1241 fn open_for_read(&self) -> Result<Arc<File>, LogError> {
1242 self.file.as_ref().map_or_else(
1243 || Ok(Arc::new(open_index_file(&self.path, false)?)),
1244 |file| Ok(Arc::clone(file)),
1245 )
1246 }
1247
1248 fn physical_len(&self) -> Result<u64, LogError> {
1249 Ok(self
1250 .file
1251 .as_ref()
1252 .map_or_else(|| fs::metadata(&self.path), |file| file.metadata())?
1253 .len())
1254 }
1255
1256 fn close_cached_reader(&mut self) {
1257 if !self.writable {
1258 self.file = None;
1259 }
1260 }
1261}
1262
1263fn open_index_file(path: &Path, writable: bool) -> Result<File, io::Error> {
1264 let mut options = OpenOptions::new();
1265 options.read(true);
1266 if writable {
1267 options.create(true).write(true).append(true);
1268 }
1269 options.open(path)
1270}
1271
1272fn validate_sorted_frames(frames: &[FrameIndex]) -> Result<(), LogError> {
1273 for pair in frames.windows(2) {
1274 if pair[0].t >= pair[1].t {
1275 return Err(LogError::Corrupt);
1276 }
1277 }
1278 Ok(())
1279}
1280
1281fn validate_sorted_extension(
1282 existing: &[FrameIndex],
1283 appended: &[FrameIndex],
1284) -> Result<(), LogError> {
1285 validate_sorted_frames(appended)?;
1286 if let (Some(previous), Some(next)) = (existing.last(), appended.first())
1287 && previous.t >= next.t
1288 {
1289 return Err(LogError::Corrupt);
1290 }
1291 Ok(())
1292}
1293
1294fn first_gap_index(frames: &[FrameIndex]) -> Option<usize> {
1295 frames
1296 .windows(2)
1297 .position(|pair| pair[0].t.checked_add(1) != Some(pair[1].t))
1298 .map(|index| index + 1)
1299}
1300
1301fn extension_first_gap(existing: &[FrameIndex], appended: &[FrameIndex]) -> Option<usize> {
1302 if let (Some(previous), Some(next)) = (existing.last(), appended.first())
1303 && previous.t.checked_add(1) != Some(next.t)
1304 {
1305 return Some(0);
1306 }
1307 first_gap_index(appended)
1308}
1309
1310fn frame_end(frame: FrameIndex) -> Result<u64, LogError> {
1311 frame.offset.checked_add(frame.len).ok_or(LogError::Corrupt)
1312}
1313
1314fn next_t(frames: &[FrameIndex]) -> Result<u64, LogError> {
1315 frames.last().map_or(Ok(1), |frame| {
1316 frame.t.checked_add(1).ok_or(LogError::Corrupt)
1317 })
1318}
1319
1320fn range_is_indexed(frames: &[FrameIndex], end: Option<u64>) -> Result<bool, LogError> {
1321 end.map_or(Ok(false), |end| Ok(end <= next_t(frames)?))
1322}
1323
1324fn version_path(dir: &Path, name: &str, version: u64) -> PathBuf {
1325 if version == 0 {
1326 dir.join(format!("{name}.log"))
1327 } else {
1328 dir.join(format!("{name}.v{version}.log"))
1329 }
1330}
1331
1332fn version_files(dir: &Path, name: &str) -> Vec<(u64, PathBuf)> {
1334 let mut files = Vec::new();
1335 let legacy = version_path(dir, name, 0);
1336 if legacy.is_file() {
1337 files.push((0, legacy));
1338 }
1339 let prefix = format!("{name}.v");
1340 if let Ok(entries) = fs::read_dir(dir) {
1341 for entry in entries.flatten() {
1342 let file_name = entry.file_name();
1343 let Some(text) = file_name.to_str() else {
1344 continue;
1345 };
1346 if let Some(version) = text
1347 .strip_prefix(&prefix)
1348 .and_then(|rest| rest.strip_suffix(".log"))
1349 .and_then(|v| v.parse::<u64>().ok())
1350 && version > 0
1351 {
1352 files.push((version, entry.path()));
1353 }
1354 }
1355 }
1356 files.sort_by_key(|(version, _)| *version);
1357 files
1358}
1359
1360fn open_version_files(
1361 dir: &Path,
1362 name: &str,
1363 cipher: Option<&Arc<LogCipher>>,
1364) -> Result<Vec<VersionedFile>, LogError> {
1365 let mut files = Vec::new();
1366 for (version, path) in version_files(dir, name) {
1367 files.push(VersionedFile {
1368 version,
1369 file: IndexedFile::open(
1370 &path,
1371 false,
1372 false,
1373 RecordCodec::new(cipher.map(Arc::clone), version),
1374 )?,
1375 });
1376 close_cold_version_files(&mut files);
1377 }
1378 Ok(files)
1379}
1380
1381fn refresh_version_files(
1382 dir: &Path,
1383 name: &str,
1384 cipher: Option<&Arc<LogCipher>>,
1385 files: &mut Vec<VersionedFile>,
1386) -> Result<(), LogError> {
1387 for file in &mut *files {
1388 file.file.refresh()?;
1389 }
1390
1391 for (version, path) in version_files(dir, name) {
1392 if files.iter().all(|file| file.version != version) {
1393 files.push(VersionedFile {
1394 version,
1395 file: IndexedFile::open(
1396 &path,
1397 false,
1398 false,
1399 RecordCodec::new(cipher.map(Arc::clone), version),
1400 )?,
1401 });
1402 close_cold_version_files(files);
1403 }
1404 }
1405 files.sort_by_key(|file| file.version);
1406 Ok(())
1407}
1408
1409fn close_cold_version_files(files: &mut [VersionedFile]) {
1410 let mut cached_readers = 0;
1411 for file in files.iter_mut().rev() {
1412 if file.file.writable {
1413 continue;
1414 }
1415 if file.file.file.is_some() {
1416 if cached_readers < MAX_CACHED_READ_VERSION_FILES {
1417 cached_readers += 1;
1418 } else {
1419 file.file.close_cached_reader();
1420 }
1421 }
1422 }
1423}
1424
1425fn version_cutoffs(files: &[VersionedFile]) -> Vec<u64> {
1428 let mut cutoffs = vec![u64::MAX; files.len()];
1429 let mut cutoff = u64::MAX;
1430 for (index, file) in files.iter().enumerate().rev() {
1431 cutoffs[index] = cutoff;
1432 if let Some(first) = file.file.frames.first() {
1433 cutoff = cutoff.min(first.t);
1434 }
1435 }
1436 cutoffs
1437}
1438
1439fn validated_version_cutoffs(files: &[VersionedFile]) -> Result<Vec<u64>, LogError> {
1440 let cutoffs = version_cutoffs(files);
1441 let mut previous_t: Option<u64> = None;
1442 for (file, cutoff) in files.iter().zip(&cutoffs) {
1443 let retained = file.file.frames.partition_point(|frame| frame.t < *cutoff);
1444 if retained == 0 {
1445 continue;
1446 }
1447 file.file.validate_contiguous_prefix(retained)?;
1448 let first_t = file.file.frames[0].t;
1449 if previous_t.is_some_and(|previous| previous.checked_add(1) != Some(first_t)) {
1450 return Err(LogError::Corrupt);
1451 }
1452 previous_t = Some(file.file.frames[retained - 1].t);
1453 }
1454 Ok(cutoffs)
1455}
1456
1457fn merged_next_t(files: &[VersionedFile], cutoffs: &[u64]) -> Result<u64, LogError> {
1458 for (file, cutoff) in files.iter().zip(cutoffs).rev() {
1459 let retained = file.file.frames.partition_point(|frame| frame.t < *cutoff);
1460 if retained > 0 {
1461 return file.file.frames[retained - 1]
1462 .t
1463 .checked_add(1)
1464 .ok_or(LogError::Corrupt);
1465 }
1466 }
1467 Ok(1)
1468}
1469
1470fn range_is_merged_indexed(
1471 files: &[VersionedFile],
1472 cutoffs: &[u64],
1473 end: Option<u64>,
1474) -> Result<bool, LogError> {
1475 end.map_or(Ok(false), |end| Ok(end <= merged_next_t(files, cutoffs)?))
1476}
1477
1478fn read_merged_range(
1479 files: &[VersionedFile],
1480 cutoffs: &[u64],
1481 start: u64,
1482 end: Option<u64>,
1483) -> Result<Vec<TxRecord>, LogError> {
1484 let mut records = Vec::new();
1485 for (file, cutoff) in files.iter().zip(cutoffs) {
1486 let end = Some(end.map_or(*cutoff, |end| end.min(*cutoff)));
1487 records.extend(file.file.tx_range(start, end)?);
1488 }
1489 Ok(records)
1490}
1491
1492fn encode_record(record: &TxRecord) -> Vec<u8> {
1493 let mut out = Vec::new();
1494 out.extend_from_slice(&record.t.to_be_bytes());
1495 out.extend_from_slice(&record.tx_instant.to_be_bytes());
1496 out.extend_from_slice(&(record.datoms.len() as u64).to_be_bytes());
1497 for d in &record.datoms {
1498 out.extend_from_slice(&d.e.raw().to_be_bytes());
1499 out.extend_from_slice(&d.a.raw().to_be_bytes());
1500 out.extend_from_slice(&d.tx.raw().to_be_bytes());
1501 out.push(u8::from(d.added));
1502 let v = encode_value(&d.v);
1503 out.extend_from_slice(&(v.len() as u64).to_be_bytes());
1504 out.extend_from_slice(&v);
1505 }
1506 out
1507}
1508fn decode_record(mut bytes: &[u8]) -> Result<TxRecord, LogError> {
1509 fn take<'a>(bytes: &mut &'a [u8], n: usize) -> Result<&'a [u8], LogError> {
1510 let value = bytes.get(..n).ok_or(LogError::Corrupt)?;
1511 *bytes = &bytes[n..];
1512 Ok(value)
1513 }
1514 fn u64_be(bytes: &mut &[u8]) -> Result<u64, LogError> {
1515 Ok(u64::from_be_bytes(
1516 take(bytes, 8)?.try_into().map_err(|_| LogError::Corrupt)?,
1517 ))
1518 }
1519 let t = u64_be(&mut bytes)?;
1520 let tx_instant = i64::from_be_bytes(
1521 take(&mut bytes, 8)?
1522 .try_into()
1523 .map_err(|_| LogError::Corrupt)?,
1524 );
1525 let count = u64_be(&mut bytes)?;
1526 let mut datoms = Vec::new();
1527 for _ in 0..count {
1528 let e = EntityId::from_raw(u64_be(&mut bytes)?);
1529 let a = EntityId::from_raw(u64_be(&mut bytes)?);
1530 let tx = EntityId::from_raw(u64_be(&mut bytes)?);
1531 let added = take(&mut bytes, 1)?[0] != 0;
1532 let len = usize::try_from(u64_be(&mut bytes)?).map_err(|_| LogError::Corrupt)?;
1533 let raw = take(&mut bytes, len)?;
1534 let (v, used) = decode_value(raw).map_err(|_| LogError::Corrupt)?;
1535 if used != len {
1536 return Err(LogError::Corrupt);
1537 }
1538 datoms.push(Datom { e, a, v, tx, added });
1539 }
1540 if !bytes.is_empty() {
1541 return Err(LogError::Corrupt);
1542 }
1543 Ok(TxRecord {
1544 t,
1545 tx_instant,
1546 datoms,
1547 })
1548}
1549
1550fn frame_header(payload_len: usize) -> Result<[u8; 8], LogError> {
1551 let payload_len = u64::try_from(payload_len).map_err(|_| LogError::Corrupt)?;
1552 if payload_len & CHECKSUMMED_FRAME != 0 {
1553 return Err(LogError::Corrupt);
1554 }
1555 Ok((payload_len | CHECKSUMMED_FRAME).to_be_bytes())
1556}
1557
1558fn frame_payload_len(header: [u8; 8]) -> Result<(usize, bool), LogError> {
1559 let encoded = u64::from_be_bytes(header);
1560 let checksummed = encoded & CHECKSUMMED_FRAME != 0;
1561 let payload_len = encoded & !CHECKSUMMED_FRAME;
1562 Ok((
1563 usize::try_from(payload_len).map_err(|_| LogError::Corrupt)?,
1564 checksummed,
1565 ))
1566}
1567
1568fn frame_checksum(header: [u8; 8], payload: &[u8]) -> u32 {
1569 crc32c::crc32c_append(crc32c::crc32c(&header), payload)
1570}
1571
1572#[cfg(unix)]
1573fn read_exact_at(file: &File, mut bytes: &mut [u8], mut offset: u64) -> io::Result<()> {
1574 use std::os::unix::fs::FileExt;
1575 while !bytes.is_empty() {
1576 match file.read_at(bytes, offset) {
1577 Ok(0) => return Err(io::ErrorKind::UnexpectedEof.into()),
1578 Ok(read) => {
1579 offset = offset
1580 .checked_add(u64::try_from(read).expect("read length fits u64"))
1581 .ok_or_else(|| io::Error::other("file offset overflow"))?;
1582 bytes = &mut bytes[read..];
1583 }
1584 Err(error) if error.kind() == io::ErrorKind::Interrupted => {}
1585 Err(error) => return Err(error),
1586 }
1587 }
1588 Ok(())
1589}
1590
1591#[cfg(windows)]
1592fn read_exact_at(file: &File, mut bytes: &mut [u8], mut offset: u64) -> io::Result<()> {
1593 use std::os::windows::fs::FileExt;
1594 while !bytes.is_empty() {
1595 match file.seek_read(bytes, offset) {
1596 Ok(0) => return Err(io::ErrorKind::UnexpectedEof.into()),
1597 Ok(read) => {
1598 offset = offset
1599 .checked_add(u64::try_from(read).expect("read length fits u64"))
1600 .ok_or_else(|| io::Error::other("file offset overflow"))?;
1601 bytes = &mut bytes[read..];
1602 }
1603 Err(error) if error.kind() == io::ErrorKind::Interrupted => {}
1604 Err(error) => return Err(error),
1605 }
1606 }
1607 Ok(())
1608}
1609
1610#[cfg(not(any(unix, windows)))]
1611fn read_exact_at(file: &File, bytes: &mut [u8], offset: u64) -> io::Result<()> {
1612 use std::io::{Read, Seek, SeekFrom};
1613 let mut file = file.try_clone()?;
1614 file.seek(SeekFrom::Start(offset))?;
1615 file.read_exact(bytes)
1616}
1617
1618fn scan_frames(
1627 file: &File,
1628 offset: u64,
1629 codec: &RecordCodec,
1630) -> Result<(Vec<FrameIndex>, u64), LogError> {
1631 let file_len = file.metadata()?.len();
1632 let mut frames = Vec::new();
1633 let mut durable_len = offset;
1634 loop {
1635 if file_len.saturating_sub(durable_len) < 8 {
1636 break;
1637 }
1638 let mut len = [0; 8];
1639 read_exact_at(file, &mut len, durable_len)?;
1640 let (payload_len, checksummed) = frame_payload_len(len)?;
1641 let frame_len = 8_u64
1642 .checked_add(u64::try_from(payload_len).map_err(|_| LogError::Corrupt)?)
1643 .and_then(|len| {
1644 len.checked_add(if checksummed {
1645 u64::try_from(FRAME_CHECKSUM_LEN).expect("checksum length fits u64")
1646 } else {
1647 0
1648 })
1649 })
1650 .ok_or(LogError::Corrupt)?;
1651 if file_len.saturating_sub(durable_len) < frame_len {
1654 break;
1655 }
1656 let mut payload = vec![0; payload_len];
1657 let payload_offset = durable_len.checked_add(8).ok_or(LogError::Corrupt)?;
1658 read_exact_at(file, &mut payload, payload_offset)?;
1659 if checksummed {
1660 let mut stored_checksum = [0; FRAME_CHECKSUM_LEN];
1661 let checksum_offset = payload_offset
1662 .checked_add(u64::try_from(payload_len).map_err(|_| LogError::Corrupt)?)
1663 .ok_or(LogError::Corrupt)?;
1664 read_exact_at(file, &mut stored_checksum, checksum_offset)?;
1665 if u32::from_be_bytes(stored_checksum) != frame_checksum(len, &payload) {
1666 return Err(LogError::Corrupt);
1667 }
1668 }
1669 let record = codec.decode(&payload)?;
1670 frames.push(FrameIndex {
1671 t: record.t,
1672 offset: durable_len,
1673 len: frame_len,
1674 });
1675 durable_len = durable_len
1676 .checked_add(frame_len)
1677 .ok_or(LogError::Corrupt)?;
1678 }
1679 Ok((frames, durable_len))
1680}
1681
1682fn append_framed_payload(out: &mut Vec<u8>, payload: &[u8]) -> Result<(), LogError> {
1690 let header = frame_header(payload.len())?;
1691 out.extend_from_slice(&header);
1692 out.extend_from_slice(payload);
1693 out.extend_from_slice(&frame_checksum(header, payload).to_be_bytes());
1694 Ok(())
1695}
1696
1697pub fn append_framed_record(out: &mut Vec<u8>, record: &TxRecord) -> Result<(), LogError> {
1702 append_framed_payload(out, &encode_record(record))
1703}
1704
1705pub fn append_framed_record_sealed(
1715 out: &mut Vec<u8>,
1716 record: &TxRecord,
1717 cipher: Option<&Arc<LogCipher>>,
1718 log_version: u64,
1719) -> Result<(), LogError> {
1720 let codec = RecordCodec::new(cipher.map(Arc::clone), log_version);
1721 append_framed_payload(out, &codec.encode(record)?)
1722}
1723
1724pub fn decode_framed_records(bytes: &[u8]) -> Result<Vec<TxRecord>, LogError> {
1730 decode_framed_payloads(bytes, &RecordCodec::plaintext())
1731}
1732
1733pub fn decode_framed_records_sealed(
1741 bytes: &[u8],
1742 cipher: Option<&Arc<LogCipher>>,
1743 log_version: u64,
1744) -> Result<Vec<TxRecord>, LogError> {
1745 decode_framed_payloads(
1746 bytes,
1747 &RecordCodec::new(cipher.map(Arc::clone), log_version),
1748 )
1749}
1750
1751fn decode_framed_payloads(
1757 mut bytes: &[u8],
1758 codec: &RecordCodec,
1759) -> Result<Vec<TxRecord>, LogError> {
1760 let mut records = Vec::new();
1761 while !bytes.is_empty() {
1762 if bytes.len() < 8 {
1763 return Err(LogError::Corrupt);
1764 }
1765 let header: [u8; 8] = bytes[..8].try_into().map_err(|_| LogError::Corrupt)?;
1766 let (payload_len, checksummed) = frame_payload_len(header)?;
1767 bytes = &bytes[8..];
1768 let payload = bytes.get(..payload_len).ok_or(LogError::Corrupt)?;
1769 bytes = &bytes[payload_len..];
1770 if checksummed {
1771 let stored_checksum = u32::from_be_bytes(
1772 bytes
1773 .get(..FRAME_CHECKSUM_LEN)
1774 .ok_or(LogError::Corrupt)?
1775 .try_into()
1776 .map_err(|_| LogError::Corrupt)?,
1777 );
1778 if stored_checksum != frame_checksum(header, payload) {
1779 return Err(LogError::Corrupt);
1780 }
1781 bytes = &bytes[FRAME_CHECKSUM_LEN..];
1782 }
1783 records.push(codec.decode(payload)?);
1784 }
1785 Ok(records)
1786}
1787
1788#[cfg(test)]
1789mod tests {
1790 use super::*;
1791
1792 #[test]
1797 fn a_header_t_disagreeing_with_its_payload_is_corrupt() {
1798 let key = SecretKey::new([7; 32]);
1799 let codec = RecordCodec::new(Some(Arc::new(LogCipher::with_key("db", 1, key.clone()))), 0);
1800 let record = TxRecord {
1801 t: 1,
1802 tx_instant: 5,
1803 datoms: Vec::new(),
1804 };
1805
1806 let honest = codec.encode(&record).expect("seal");
1807 assert_eq!(codec.decode(&honest).expect("decode"), record);
1808
1809 let forged =
1812 encrypt_log_record(&key, 1, b"db", 0, 2, &encode_record(&record)).expect("seal");
1813 assert_eq!(parse_log_header(&forged).expect("header").t, 2);
1814 assert!(matches!(codec.decode(&forged), Err(LogError::Corrupt)));
1815 }
1816
1817 #[test]
1818 fn versioned_log_bounds_cached_read_descriptors() {
1819 let dir = tempfile::tempdir().expect("tempdir");
1820 let segment_count = MAX_CACHED_READ_VERSION_FILES + 5;
1821 for version in 1..=u64::try_from(segment_count).expect("segment count fits u64") {
1822 File::create(version_path(dir.path(), "db", version)).expect("create segment");
1823 }
1824
1825 let log = VersionedLog::open_read_only(dir.path(), "db").expect("open log");
1826 let state = log.state.read().expect("log lock");
1827 assert_eq!(state.files.len(), segment_count);
1828 assert!(
1829 state
1830 .files
1831 .iter()
1832 .filter(|file| file.file.file.is_some())
1833 .count()
1834 <= MAX_CACHED_READ_VERSION_FILES
1835 );
1836 }
1837}