1use async_trait::async_trait;
4use corium_core::{
5 Datom, EntityId,
6 encoding::{decode_value, encode_value},
7};
8use std::{
9 collections::HashMap,
10 fs::{self, File, OpenOptions},
11 io::{self, Write},
12 path::{Path, PathBuf},
13 sync::{Arc, Mutex, RwLock},
14};
15use thiserror::Error;
16
17const CHECKSUMMED_FRAME: u64 = 1 << 63;
18const FRAME_CHECKSUM_LEN: usize = size_of::<u32>();
19const RANGE_READ_CHUNK_BYTES: u64 = 4 * 1024 * 1024;
20const MAX_CACHED_READ_VERSION_FILES: usize = 8;
21
22#[derive(Clone, Debug, Eq, PartialEq)]
24pub struct TxRecord {
25 pub t: u64,
27 pub tx_instant: i64,
29 pub datoms: Vec<Datom>,
31}
32
33#[derive(Debug, Error)]
35pub enum LogError {
36 #[error("log I/O failed: {0}")]
38 Io(#[from] io::Error),
39 #[error("corrupt transaction log")]
41 Corrupt,
42 #[error("native transaction log store failed: {0}")]
44 Native(String),
45 #[error("this transaction log requires asynchronous access")]
47 AsyncOnly,
48}
49
50#[async_trait]
52pub trait TransactionLog: Send + Sync {
53 fn append(&self, record: &TxRecord) -> Result<(), LogError>;
58 async fn append_async(&self, record: &TxRecord) -> Result<(), LogError> {
65 self.append(record)
66 }
67 async fn append_batch_async(&self, records: &[TxRecord]) -> Result<(), LogError> {
77 for record in records {
78 self.append_async(record).await?;
79 }
80 Ok(())
81 }
82 fn tx_range(&self, start: u64, end: Option<u64>) -> Result<Vec<TxRecord>, LogError>;
87 async fn tx_range_async(
92 &self,
93 start: u64,
94 end: Option<u64>,
95 ) -> Result<Vec<TxRecord>, LogError> {
96 self.tx_range(start, end)
97 }
98 fn replay(&self) -> Result<Vec<TxRecord>, LogError> {
103 self.tx_range(0, None)
104 }
105 async fn replay_async(&self) -> Result<Vec<TxRecord>, LogError> {
110 self.tx_range_async(0, None).await
111 }
112}
113
114#[derive(Clone, Default)]
116pub struct MemoryLog(Arc<RwLock<Vec<TxRecord>>>);
117impl TransactionLog for MemoryLog {
118 fn append(&self, record: &TxRecord) -> Result<(), LogError> {
119 let mut records = self.0.write().expect("poisoned log lock");
120 if records.last().map_or(1, |r| r.t + 1) != record.t {
121 return Err(LogError::Corrupt);
122 }
123 records.push(record.clone());
124 Ok(())
125 }
126 fn tx_range(&self, start: u64, end: Option<u64>) -> Result<Vec<TxRecord>, LogError> {
127 Ok(self
128 .0
129 .read()
130 .expect("poisoned log lock")
131 .iter()
132 .filter(|r| r.t >= start && end.is_none_or(|e| r.t < e))
133 .cloned()
134 .collect())
135 }
136}
137
138pub struct FileLog {
148 state: RwLock<IndexedFile>,
149}
150
151impl FileLog {
152 pub fn open(path: impl AsRef<Path>) -> Result<Self, LogError> {
158 let path = path.as_ref().to_path_buf();
159 if let Some(parent) = path.parent() {
160 fs::create_dir_all(parent)?;
161 }
162 let file = IndexedFile::open(&path, true, true)?;
163 file.validate_contiguous_prefix(file.frames.len())?;
164 Ok(Self {
165 state: RwLock::new(file),
166 })
167 }
168}
169impl TransactionLog for FileLog {
170 fn append(&self, record: &TxRecord) -> Result<(), LogError> {
171 let mut state = self.state.write().expect("poisoned log lock");
172 state.refresh()?;
173 state.validate_contiguous_prefix(state.frames.len())?;
174 if next_t(&state.frames)? != record.t {
175 return Err(LogError::Corrupt);
176 }
177 state.append(record)
178 }
179 fn tx_range(&self, start: u64, end: Option<u64>) -> Result<Vec<TxRecord>, LogError> {
180 if end.is_some_and(|end| end <= start) {
181 return Ok(Vec::new());
182 }
183 let indexed = {
184 let state = self.state.read().expect("poisoned log lock");
185 range_is_indexed(&state.frames, end)?
186 };
187 if !indexed {
188 let mut state = self.state.write().expect("poisoned log lock");
189 state.refresh()?;
190 state.validate_contiguous_prefix(state.frames.len())?;
191 }
192 self.state
193 .read()
194 .expect("poisoned log lock")
195 .tx_range(start, end)
196 }
197}
198
199pub struct VersionedLog {
211 dir: PathBuf,
212 name: String,
213 state: RwLock<VersionedLogState>,
214}
215
216struct VersionedLogState {
217 files: Vec<VersionedFile>,
218 write_version: Option<u64>,
219 next_t: u64,
220}
221
222struct VersionedFile {
223 version: u64,
224 file: IndexedFile,
225}
226
227impl VersionedLog {
228 pub fn open(dir: impl AsRef<Path>, name: &str, write_version: u64) -> Result<Self, LogError> {
236 let dir = dir.as_ref().to_path_buf();
237 fs::create_dir_all(&dir)?;
238 let write_path = version_path(&dir, name, write_version);
239 let mut files = Vec::new();
240 for (version, path) in version_files(&dir, name) {
241 let writable = version == write_version;
242 files.push(VersionedFile {
243 version,
244 file: IndexedFile::open(&path, writable, writable)?,
245 });
246 close_cold_version_files(&mut files);
247 }
248 if !files.iter().any(|file| file.version == write_version) {
249 files.push(VersionedFile {
250 version: write_version,
251 file: IndexedFile::open(&write_path, true, true)?,
252 });
253 files.sort_by_key(|file| file.version);
254 }
255 close_cold_version_files(&mut files);
256 let cutoffs = validated_version_cutoffs(&files)?;
257 let next_t = merged_next_t(&files, &cutoffs)?;
258 Ok(Self {
259 dir,
260 name: name.to_owned(),
261 state: RwLock::new(VersionedLogState {
262 files,
263 write_version: Some(write_version),
264 next_t,
265 }),
266 })
267 }
268
269 pub fn open_read_only(dir: impl AsRef<Path>, name: &str) -> Result<Self, LogError> {
275 let dir = dir.as_ref().to_path_buf();
276 let mut files = open_version_files(&dir, name)?;
277 close_cold_version_files(&mut files);
278 let cutoffs = validated_version_cutoffs(&files)?;
279 Ok(Self {
280 name: name.to_owned(),
281 state: RwLock::new(VersionedLogState {
282 next_t: merged_next_t(&files, &cutoffs)?,
283 write_version: None,
284 files,
285 }),
286 dir,
287 })
288 }
289
290 #[must_use]
292 pub fn exists(dir: impl AsRef<Path>, name: &str) -> bool {
293 !version_files(dir.as_ref(), name).is_empty()
294 }
295
296 pub fn delete_all(dir: impl AsRef<Path>, name: &str) -> Result<(), LogError> {
301 for (_, path) in version_files(dir.as_ref(), name) {
302 match fs::remove_file(&path) {
303 Ok(()) => {}
304 Err(error) if error.kind() == io::ErrorKind::NotFound => {}
305 Err(error) => return Err(error.into()),
306 }
307 }
308 Ok(())
309 }
310}
311
312impl TransactionLog for VersionedLog {
313 fn append(&self, record: &TxRecord) -> Result<(), LogError> {
314 let mut state = self.state.write().expect("poisoned log lock");
315 if state.next_t != record.t {
316 return Err(LogError::Corrupt);
317 }
318 let write_version = state
319 .write_version
320 .ok_or_else(|| LogError::Native("transaction log is read-only".into()))?;
321 let write_index = state
322 .files
323 .iter()
324 .position(|file| file.version == write_version)
325 .ok_or(LogError::Corrupt)?;
326 let cutoffs = version_cutoffs(&state.files);
327 if cutoffs[write_index] == u64::MAX
331 && !state.files[write_index].file.frames.is_empty()
332 && next_t(&state.files[write_index].file.frames)? != record.t
333 {
334 return Err(LogError::Corrupt);
335 }
336 state.files[write_index].file.append(record)?;
337 state.next_t = state.next_t.checked_add(1).ok_or(LogError::Corrupt)?;
338 Ok(())
339 }
340
341 fn tx_range(&self, start: u64, end: Option<u64>) -> Result<Vec<TxRecord>, LogError> {
342 if end.is_some_and(|end| end <= start) {
343 return Ok(Vec::new());
344 }
345 {
346 let state = self.state.read().expect("poisoned log lock");
347 let cutoffs = validated_version_cutoffs(&state.files)?;
348 if range_is_merged_indexed(&state.files, &cutoffs, end)? {
349 return read_merged_range(&state.files, &cutoffs, start, end);
350 }
351 }
352 {
353 let mut state = self.state.write().expect("poisoned log lock");
354 refresh_version_files(&self.dir, &self.name, &mut state.files)?;
355 close_cold_version_files(&mut state.files);
356 validated_version_cutoffs(&state.files)?;
357 }
358 let state = self.state.read().expect("poisoned log lock");
359 let cutoffs = validated_version_cutoffs(&state.files)?;
360 read_merged_range(&state.files, &cutoffs, start, end)
361 }
362}
363
364fn merge_versions(mut per_version: Vec<Vec<TxRecord>>) -> Vec<TxRecord> {
368 let mut cutoff = u64::MAX;
369 for records in per_version.iter_mut().rev() {
370 let first = records.first().map(|r| r.t);
371 records.retain(|r| r.t < cutoff);
372 if let Some(first) = first {
373 cutoff = cutoff.min(first);
374 }
375 }
376 per_version.into_iter().flatten().collect()
377}
378
379#[async_trait]
395pub trait NativeLogStorage: Send + Sync {
396 async fn put_batch(
409 &self,
410 name: &str,
411 version: u64,
412 records: &[(u64, Vec<u8>)],
413 ) -> Result<bool, LogError>;
414 async fn read_record(
420 &self,
421 name: &str,
422 version: u64,
423 t: u64,
424 ) -> Result<Option<Vec<u8>>, LogError>;
425 async fn list_records(&self, name: &str) -> Result<Vec<(u64, u64)>, LogError>;
431 async fn read_legacy_chunk(
437 &self,
438 name: &str,
439 version: u64,
440 chunk: u64,
441 ) -> Result<Option<Vec<u8>>, LogError>;
442 async fn list_legacy_chunks(&self, name: &str) -> Result<Vec<(u64, u64)>, LogError>;
450 async fn delete_all(&self, name: &str) -> Result<(), LogError>;
455}
456
457pub struct NativeVersionedLog<S: ?Sized> {
466 storage: Arc<S>,
467 name: String,
468 write_version: u64,
469 read_only: bool,
470 next_t: tokio::sync::Mutex<u64>,
472}
473
474impl<S: NativeLogStorage + ?Sized + 'static> NativeVersionedLog<S> {
475 pub async fn open(storage: Arc<S>, name: &str, write_version: u64) -> Result<Self, LogError> {
480 let records = read_native_merged(storage.as_ref(), name).await?;
484 let next_t = records.last().map_or(1, |r| r.t + 1);
485 Ok(Self {
486 storage,
487 name: name.to_owned(),
488 write_version,
489 read_only: false,
490 next_t: tokio::sync::Mutex::new(next_t),
491 })
492 }
493
494 #[must_use]
497 pub fn open_read_only(storage: Arc<S>, name: &str) -> Self {
498 Self {
499 storage,
500 name: name.to_owned(),
501 write_version: 0,
502 read_only: true,
503 next_t: tokio::sync::Mutex::new(0),
504 }
505 }
506}
507
508#[async_trait]
509impl<S: NativeLogStorage + ?Sized + 'static> TransactionLog for NativeVersionedLog<S> {
510 fn append(&self, record: &TxRecord) -> Result<(), LogError> {
511 let _ = record;
512 Err(LogError::AsyncOnly)
513 }
514
515 async fn append_async(&self, record: &TxRecord) -> Result<(), LogError> {
516 self.append_batch_async(std::slice::from_ref(record)).await
517 }
518
519 async fn append_batch_async(&self, records: &[TxRecord]) -> Result<(), LogError> {
520 if self.read_only {
521 return Err(LogError::Native("transaction log is read-only".into()));
522 }
523 if records.is_empty() {
524 return Ok(());
525 }
526 let mut next_t = self.next_t.lock().await;
527 for (offset, record) in records.iter().enumerate() {
529 if record.t != *next_t + offset as u64 {
530 return Err(LogError::Corrupt);
531 }
532 }
533 let framed = records
534 .iter()
535 .map(|record| {
536 let mut bytes = Vec::new();
537 append_framed_record(&mut bytes, record)?;
538 Ok((record.t, bytes))
539 })
540 .collect::<Result<Vec<_>, LogError>>()?;
541 if !self
546 .storage
547 .put_batch(&self.name, self.write_version, &framed)
548 .await?
549 {
550 return Err(LogError::Corrupt);
551 }
552 *next_t += records.len() as u64;
553 Ok(())
554 }
555
556 fn tx_range(&self, start: u64, end: Option<u64>) -> Result<Vec<TxRecord>, LogError> {
557 let _ = (start, end);
558 Err(LogError::AsyncOnly)
559 }
560
561 async fn tx_range_async(
562 &self,
563 start: u64,
564 end: Option<u64>,
565 ) -> Result<Vec<TxRecord>, LogError> {
566 let _guard = self.next_t.lock().await;
569 Ok(read_native_merged(self.storage.as_ref(), &self.name)
570 .await?
571 .into_iter()
572 .filter(|r| r.t >= start && end.is_none_or(|e| r.t < e))
573 .collect())
574 }
575}
576
577async fn read_native_merged<S: NativeLogStorage + ?Sized>(
578 storage: &S,
579 name: &str,
580) -> Result<Vec<TxRecord>, LogError> {
581 use std::collections::BTreeMap;
582
583 let mut per_version: BTreeMap<u64, Vec<TxRecord>> = BTreeMap::new();
586
587 let mut chunks = storage.list_legacy_chunks(name).await?;
591 chunks.sort_unstable();
592 for (version, chunk) in chunks {
593 let bytes = storage
594 .read_legacy_chunk(name, version, chunk)
595 .await?
596 .unwrap_or_default();
597 per_version
598 .entry(version)
599 .or_default()
600 .extend(decode_framed_records(&bytes)?);
601 }
602
603 let mut records = storage.list_records(name).await?;
605 records.sort_unstable();
606 for (version, t) in records {
607 let bytes = storage
608 .read_record(name, version, t)
609 .await?
610 .unwrap_or_default();
611 per_version
612 .entry(version)
613 .or_default()
614 .extend(decode_framed_records(&bytes)?);
615 }
616
617 let per_version: Vec<Vec<TxRecord>> = per_version
622 .into_values()
623 .map(|mut records| {
624 records.sort_by_key(|record| record.t);
625 records
626 })
627 .collect();
628 let merged = merge_versions(per_version);
629 for pair in merged.windows(2) {
630 if pair[1].t != pair[0].t + 1 {
631 return Err(LogError::Corrupt);
632 }
633 }
634 Ok(merged)
635}
636
637type VersionedRecords = Arc<Mutex<Vec<(u64, TxRecord)>>>;
640
641#[derive(Clone, Default)]
647pub struct MemLogRegistry {
648 logs: Arc<Mutex<HashMap<String, VersionedRecords>>>,
649}
650
651impl MemLogRegistry {
652 #[must_use]
654 pub fn new() -> Self {
655 Self::default()
656 }
657
658 fn entry(&self, name: &str) -> VersionedRecords {
659 Arc::clone(
660 self.logs
661 .lock()
662 .unwrap_or_else(std::sync::PoisonError::into_inner)
663 .entry(name.to_owned())
664 .or_default(),
665 )
666 }
667
668 #[must_use]
671 pub fn open(&self, name: &str, write_version: u64) -> MemVersionedLog {
672 let records = self.entry(name);
673 let next_t = {
674 let guard = records
675 .lock()
676 .unwrap_or_else(std::sync::PoisonError::into_inner);
677 MemVersionedLog::merged(&guard)
678 .last()
679 .map_or(1, |r| r.t + 1)
680 };
681 MemVersionedLog {
682 records,
683 write_version,
684 next_t: Mutex::new(next_t),
685 }
686 }
687
688 #[must_use]
690 pub fn exists(&self, name: &str) -> bool {
691 self.logs
692 .lock()
693 .unwrap_or_else(std::sync::PoisonError::into_inner)
694 .get(name)
695 .is_some_and(|entry| {
696 !entry
697 .lock()
698 .unwrap_or_else(std::sync::PoisonError::into_inner)
699 .is_empty()
700 })
701 }
702
703 pub fn delete_all(&self, name: &str) {
705 self.logs
706 .lock()
707 .unwrap_or_else(std::sync::PoisonError::into_inner)
708 .remove(name);
709 }
710}
711
712pub struct MemVersionedLog {
716 records: VersionedRecords,
717 write_version: u64,
718 next_t: Mutex<u64>,
722}
723
724impl MemVersionedLog {
725 fn merged(records: &[(u64, TxRecord)]) -> Vec<TxRecord> {
726 let mut versions: Vec<u64> = records.iter().map(|(version, _)| *version).collect();
727 versions.sort_unstable();
728 versions.dedup();
729 let per_version = versions
730 .into_iter()
731 .map(|version| {
732 records
733 .iter()
734 .filter(|(record_version, _)| *record_version == version)
735 .map(|(_, record)| record.clone())
736 .collect::<Vec<_>>()
737 })
738 .collect();
739 merge_versions(per_version)
740 }
741}
742
743impl TransactionLog for MemVersionedLog {
744 fn append(&self, record: &TxRecord) -> Result<(), LogError> {
745 let mut next_t = self
746 .next_t
747 .lock()
748 .unwrap_or_else(std::sync::PoisonError::into_inner);
749 if *next_t != record.t {
750 return Err(LogError::Corrupt);
751 }
752 self.records
753 .lock()
754 .unwrap_or_else(std::sync::PoisonError::into_inner)
755 .push((self.write_version, record.clone()));
756 *next_t += 1;
757 Ok(())
758 }
759
760 fn tx_range(&self, start: u64, end: Option<u64>) -> Result<Vec<TxRecord>, LogError> {
761 let records = self
762 .records
763 .lock()
764 .unwrap_or_else(std::sync::PoisonError::into_inner);
765 Ok(Self::merged(&records)
766 .into_iter()
767 .filter(|r| r.t >= start && end.is_none_or(|e| r.t < e))
768 .collect())
769 }
770}
771
772#[derive(Clone, Copy)]
773struct FrameIndex {
774 t: u64,
775 offset: u64,
776 len: u64,
777}
778
779struct IndexedFile {
785 path: PathBuf,
786 file: Option<Arc<File>>,
787 writable: bool,
788 poisoned: bool,
789 frames: Vec<FrameIndex>,
790 durable_len: u64,
791 first_gap: Option<usize>,
792}
793
794impl IndexedFile {
795 fn open(path: &Path, writable: bool, truncate_torn: bool) -> Result<Self, LogError> {
796 let file = Arc::new(open_index_file(path, writable)?);
797 let (frames, durable_len) = scan_frames(file.as_ref(), 0)?;
798 validate_sorted_frames(&frames)?;
799 if truncate_torn && file.metadata()?.len() > durable_len {
800 file.set_len(durable_len)?;
801 file.sync_all()?;
802 }
803 Ok(Self {
804 path: path.to_path_buf(),
805 file: Some(file),
806 writable,
807 poisoned: false,
808 first_gap: first_gap_index(&frames),
809 frames,
810 durable_len,
811 })
812 }
813
814 fn refresh(&mut self) -> Result<(), LogError> {
815 self.ensure_healthy()?;
816 let file_len = self.physical_len()?;
817 if file_len < self.durable_len {
818 let file = self.open_for_read()?;
819 let (frames, durable_len) = scan_frames(file.as_ref(), 0)?;
820 validate_sorted_frames(&frames)?;
821 self.first_gap = first_gap_index(&frames);
822 self.frames = frames;
823 self.durable_len = durable_len;
824 } else if file_len > self.durable_len {
825 let file = self.open_for_read()?;
826 let (new_frames, durable_len) = scan_frames(file.as_ref(), self.durable_len)?;
827 validate_sorted_extension(&self.frames, &new_frames)?;
828 let existing_len = self.frames.len();
829 if self.first_gap.is_none() {
830 self.first_gap =
831 extension_first_gap(&self.frames, &new_frames).map(|gap| existing_len + gap);
832 }
833 self.frames.extend(new_frames);
834 self.durable_len = durable_len;
835 }
836 Ok(())
837 }
838
839 fn append(&mut self, record: &TxRecord) -> Result<(), LogError> {
840 self.ensure_healthy()?;
841 if !self.writable {
842 return Err(LogError::Native("transaction log is read-only".into()));
843 }
844 let file = self.open_for_read()?;
845 if file.metadata()?.len() != self.durable_len {
848 return Err(LogError::Corrupt);
849 }
850
851 let mut frame = Vec::new();
852 append_framed_record(&mut frame, record)?;
853 let frame_len = u64::try_from(frame.len()).map_err(|_| LogError::Corrupt)?;
854 let offset = self.durable_len;
855 let mut writer = file.as_ref();
856 if let Err(error) = writer.write_all(&frame) {
857 self.poisoned = true;
860 return Err(error.into());
861 }
862 if let Err(error) = file.sync_all() {
863 self.poisoned = true;
864 return Err(error.into());
865 }
866 self.durable_len = self
867 .durable_len
868 .checked_add(frame_len)
869 .ok_or(LogError::Corrupt)?;
870 if self.first_gap.is_none()
871 && self
872 .frames
873 .last()
874 .is_some_and(|previous| previous.t.checked_add(1) != Some(record.t))
875 {
876 self.first_gap = Some(self.frames.len());
877 }
878 self.frames.push(FrameIndex {
879 t: record.t,
880 offset,
881 len: frame_len,
882 });
883 Ok(())
884 }
885
886 fn tx_range(&self, start: u64, end: Option<u64>) -> Result<Vec<TxRecord>, LogError> {
887 self.ensure_healthy()?;
888 let first = self.frames.partition_point(|frame| frame.t < start);
889 let last = end.map_or(self.frames.len(), |end| {
890 self.frames.partition_point(|frame| frame.t < end)
891 });
892 if first >= last {
893 return Ok(Vec::new());
894 }
895
896 let file = self.open_for_read()?;
897 let mut records = Vec::with_capacity(last - first);
898 let mut chunk_first = first;
899 while chunk_first < last {
900 let offset = self.frames[chunk_first].offset;
901 let mut chunk_last = chunk_first + 1;
902 while chunk_last < last {
903 let candidate_end = frame_end(self.frames[chunk_last])?;
904 if candidate_end.checked_sub(offset).ok_or(LogError::Corrupt)?
905 > RANGE_READ_CHUNK_BYTES
906 {
907 break;
908 }
909 chunk_last += 1;
910 }
911 let byte_end = frame_end(self.frames[chunk_last - 1])?;
912 let byte_len = usize::try_from(byte_end.checked_sub(offset).ok_or(LogError::Corrupt)?)
913 .map_err(|_| LogError::Corrupt)?;
914 let mut bytes = vec![0; byte_len];
915 read_exact_at(file.as_ref(), &mut bytes, offset)?;
916 let chunk_records = decode_framed_records(&bytes)?;
917 if chunk_records.len() != chunk_last - chunk_first
918 || chunk_records
919 .iter()
920 .zip(&self.frames[chunk_first..chunk_last])
921 .any(|(record, frame)| record.t != frame.t)
922 {
923 return Err(LogError::Corrupt);
924 }
925 records.extend(chunk_records);
926 chunk_first = chunk_last;
927 }
928 Ok(records)
929 }
930
931 fn validate_contiguous_prefix(&self, retained: usize) -> Result<(), LogError> {
932 if self.first_gap.is_some_and(|gap| gap < retained) {
933 return Err(LogError::Corrupt);
934 }
935 Ok(())
936 }
937
938 fn ensure_healthy(&self) -> Result<(), LogError> {
939 if self.poisoned {
940 return Err(io::Error::other("transaction log handle is poisoned").into());
941 }
942 Ok(())
943 }
944
945 fn open_for_read(&self) -> Result<Arc<File>, LogError> {
946 self.file.as_ref().map_or_else(
947 || Ok(Arc::new(open_index_file(&self.path, false)?)),
948 |file| Ok(Arc::clone(file)),
949 )
950 }
951
952 fn physical_len(&self) -> Result<u64, LogError> {
953 Ok(self
954 .file
955 .as_ref()
956 .map_or_else(|| fs::metadata(&self.path), |file| file.metadata())?
957 .len())
958 }
959
960 fn close_cached_reader(&mut self) {
961 if !self.writable {
962 self.file = None;
963 }
964 }
965}
966
967fn open_index_file(path: &Path, writable: bool) -> Result<File, io::Error> {
968 let mut options = OpenOptions::new();
969 options.read(true);
970 if writable {
971 options.create(true).write(true).append(true);
972 }
973 options.open(path)
974}
975
976fn validate_sorted_frames(frames: &[FrameIndex]) -> Result<(), LogError> {
977 for pair in frames.windows(2) {
978 if pair[0].t >= pair[1].t {
979 return Err(LogError::Corrupt);
980 }
981 }
982 Ok(())
983}
984
985fn validate_sorted_extension(
986 existing: &[FrameIndex],
987 appended: &[FrameIndex],
988) -> Result<(), LogError> {
989 validate_sorted_frames(appended)?;
990 if let (Some(previous), Some(next)) = (existing.last(), appended.first())
991 && previous.t >= next.t
992 {
993 return Err(LogError::Corrupt);
994 }
995 Ok(())
996}
997
998fn first_gap_index(frames: &[FrameIndex]) -> Option<usize> {
999 frames
1000 .windows(2)
1001 .position(|pair| pair[0].t.checked_add(1) != Some(pair[1].t))
1002 .map(|index| index + 1)
1003}
1004
1005fn extension_first_gap(existing: &[FrameIndex], appended: &[FrameIndex]) -> Option<usize> {
1006 if let (Some(previous), Some(next)) = (existing.last(), appended.first())
1007 && previous.t.checked_add(1) != Some(next.t)
1008 {
1009 return Some(0);
1010 }
1011 first_gap_index(appended)
1012}
1013
1014fn frame_end(frame: FrameIndex) -> Result<u64, LogError> {
1015 frame.offset.checked_add(frame.len).ok_or(LogError::Corrupt)
1016}
1017
1018fn next_t(frames: &[FrameIndex]) -> Result<u64, LogError> {
1019 frames.last().map_or(Ok(1), |frame| {
1020 frame.t.checked_add(1).ok_or(LogError::Corrupt)
1021 })
1022}
1023
1024fn range_is_indexed(frames: &[FrameIndex], end: Option<u64>) -> Result<bool, LogError> {
1025 end.map_or(Ok(false), |end| Ok(end <= next_t(frames)?))
1026}
1027
1028fn version_path(dir: &Path, name: &str, version: u64) -> PathBuf {
1029 if version == 0 {
1030 dir.join(format!("{name}.log"))
1031 } else {
1032 dir.join(format!("{name}.v{version}.log"))
1033 }
1034}
1035
1036fn version_files(dir: &Path, name: &str) -> Vec<(u64, PathBuf)> {
1038 let mut files = Vec::new();
1039 let legacy = version_path(dir, name, 0);
1040 if legacy.is_file() {
1041 files.push((0, legacy));
1042 }
1043 let prefix = format!("{name}.v");
1044 if let Ok(entries) = fs::read_dir(dir) {
1045 for entry in entries.flatten() {
1046 let file_name = entry.file_name();
1047 let Some(text) = file_name.to_str() else {
1048 continue;
1049 };
1050 if let Some(version) = text
1051 .strip_prefix(&prefix)
1052 .and_then(|rest| rest.strip_suffix(".log"))
1053 .and_then(|v| v.parse::<u64>().ok())
1054 && version > 0
1055 {
1056 files.push((version, entry.path()));
1057 }
1058 }
1059 }
1060 files.sort_by_key(|(version, _)| *version);
1061 files
1062}
1063
1064fn open_version_files(dir: &Path, name: &str) -> Result<Vec<VersionedFile>, LogError> {
1065 let mut files = Vec::new();
1066 for (version, path) in version_files(dir, name) {
1067 files.push(VersionedFile {
1068 version,
1069 file: IndexedFile::open(&path, false, false)?,
1070 });
1071 close_cold_version_files(&mut files);
1072 }
1073 Ok(files)
1074}
1075
1076fn refresh_version_files(
1077 dir: &Path,
1078 name: &str,
1079 files: &mut Vec<VersionedFile>,
1080) -> Result<(), LogError> {
1081 for file in &mut *files {
1082 file.file.refresh()?;
1083 }
1084
1085 for (version, path) in version_files(dir, name) {
1086 if files.iter().all(|file| file.version != version) {
1087 files.push(VersionedFile {
1088 version,
1089 file: IndexedFile::open(&path, false, false)?,
1090 });
1091 close_cold_version_files(files);
1092 }
1093 }
1094 files.sort_by_key(|file| file.version);
1095 Ok(())
1096}
1097
1098fn close_cold_version_files(files: &mut [VersionedFile]) {
1099 let mut cached_readers = 0;
1100 for file in files.iter_mut().rev() {
1101 if file.file.writable {
1102 continue;
1103 }
1104 if file.file.file.is_some() {
1105 if cached_readers < MAX_CACHED_READ_VERSION_FILES {
1106 cached_readers += 1;
1107 } else {
1108 file.file.close_cached_reader();
1109 }
1110 }
1111 }
1112}
1113
1114fn version_cutoffs(files: &[VersionedFile]) -> Vec<u64> {
1117 let mut cutoffs = vec![u64::MAX; files.len()];
1118 let mut cutoff = u64::MAX;
1119 for (index, file) in files.iter().enumerate().rev() {
1120 cutoffs[index] = cutoff;
1121 if let Some(first) = file.file.frames.first() {
1122 cutoff = cutoff.min(first.t);
1123 }
1124 }
1125 cutoffs
1126}
1127
1128fn validated_version_cutoffs(files: &[VersionedFile]) -> Result<Vec<u64>, LogError> {
1129 let cutoffs = version_cutoffs(files);
1130 let mut previous_t: Option<u64> = None;
1131 for (file, cutoff) in files.iter().zip(&cutoffs) {
1132 let retained = file.file.frames.partition_point(|frame| frame.t < *cutoff);
1133 if retained == 0 {
1134 continue;
1135 }
1136 file.file.validate_contiguous_prefix(retained)?;
1137 let first_t = file.file.frames[0].t;
1138 if previous_t.is_some_and(|previous| previous.checked_add(1) != Some(first_t)) {
1139 return Err(LogError::Corrupt);
1140 }
1141 previous_t = Some(file.file.frames[retained - 1].t);
1142 }
1143 Ok(cutoffs)
1144}
1145
1146fn merged_next_t(files: &[VersionedFile], cutoffs: &[u64]) -> Result<u64, LogError> {
1147 for (file, cutoff) in files.iter().zip(cutoffs).rev() {
1148 let retained = file.file.frames.partition_point(|frame| frame.t < *cutoff);
1149 if retained > 0 {
1150 return file.file.frames[retained - 1]
1151 .t
1152 .checked_add(1)
1153 .ok_or(LogError::Corrupt);
1154 }
1155 }
1156 Ok(1)
1157}
1158
1159fn range_is_merged_indexed(
1160 files: &[VersionedFile],
1161 cutoffs: &[u64],
1162 end: Option<u64>,
1163) -> Result<bool, LogError> {
1164 end.map_or(Ok(false), |end| Ok(end <= merged_next_t(files, cutoffs)?))
1165}
1166
1167fn read_merged_range(
1168 files: &[VersionedFile],
1169 cutoffs: &[u64],
1170 start: u64,
1171 end: Option<u64>,
1172) -> Result<Vec<TxRecord>, LogError> {
1173 let mut records = Vec::new();
1174 for (file, cutoff) in files.iter().zip(cutoffs) {
1175 let end = Some(end.map_or(*cutoff, |end| end.min(*cutoff)));
1176 records.extend(file.file.tx_range(start, end)?);
1177 }
1178 Ok(records)
1179}
1180
1181fn encode_record(record: &TxRecord) -> Vec<u8> {
1182 let mut out = Vec::new();
1183 out.extend_from_slice(&record.t.to_be_bytes());
1184 out.extend_from_slice(&record.tx_instant.to_be_bytes());
1185 out.extend_from_slice(&(record.datoms.len() as u64).to_be_bytes());
1186 for d in &record.datoms {
1187 out.extend_from_slice(&d.e.raw().to_be_bytes());
1188 out.extend_from_slice(&d.a.raw().to_be_bytes());
1189 out.extend_from_slice(&d.tx.raw().to_be_bytes());
1190 out.push(u8::from(d.added));
1191 let v = encode_value(&d.v);
1192 out.extend_from_slice(&(v.len() as u64).to_be_bytes());
1193 out.extend_from_slice(&v);
1194 }
1195 out
1196}
1197fn decode_record(mut bytes: &[u8]) -> Result<TxRecord, LogError> {
1198 fn take<'a>(bytes: &mut &'a [u8], n: usize) -> Result<&'a [u8], LogError> {
1199 let value = bytes.get(..n).ok_or(LogError::Corrupt)?;
1200 *bytes = &bytes[n..];
1201 Ok(value)
1202 }
1203 fn u64_be(bytes: &mut &[u8]) -> Result<u64, LogError> {
1204 Ok(u64::from_be_bytes(
1205 take(bytes, 8)?.try_into().map_err(|_| LogError::Corrupt)?,
1206 ))
1207 }
1208 let t = u64_be(&mut bytes)?;
1209 let tx_instant = i64::from_be_bytes(
1210 take(&mut bytes, 8)?
1211 .try_into()
1212 .map_err(|_| LogError::Corrupt)?,
1213 );
1214 let count = u64_be(&mut bytes)?;
1215 let mut datoms = Vec::new();
1216 for _ in 0..count {
1217 let e = EntityId::from_raw(u64_be(&mut bytes)?);
1218 let a = EntityId::from_raw(u64_be(&mut bytes)?);
1219 let tx = EntityId::from_raw(u64_be(&mut bytes)?);
1220 let added = take(&mut bytes, 1)?[0] != 0;
1221 let len = usize::try_from(u64_be(&mut bytes)?).map_err(|_| LogError::Corrupt)?;
1222 let raw = take(&mut bytes, len)?;
1223 let (v, used) = decode_value(raw).map_err(|_| LogError::Corrupt)?;
1224 if used != len {
1225 return Err(LogError::Corrupt);
1226 }
1227 datoms.push(Datom { e, a, v, tx, added });
1228 }
1229 if !bytes.is_empty() {
1230 return Err(LogError::Corrupt);
1231 }
1232 Ok(TxRecord {
1233 t,
1234 tx_instant,
1235 datoms,
1236 })
1237}
1238
1239fn frame_header(payload_len: usize) -> Result<[u8; 8], LogError> {
1240 let payload_len = u64::try_from(payload_len).map_err(|_| LogError::Corrupt)?;
1241 if payload_len & CHECKSUMMED_FRAME != 0 {
1242 return Err(LogError::Corrupt);
1243 }
1244 Ok((payload_len | CHECKSUMMED_FRAME).to_be_bytes())
1245}
1246
1247fn frame_payload_len(header: [u8; 8]) -> Result<(usize, bool), LogError> {
1248 let encoded = u64::from_be_bytes(header);
1249 let checksummed = encoded & CHECKSUMMED_FRAME != 0;
1250 let payload_len = encoded & !CHECKSUMMED_FRAME;
1251 Ok((
1252 usize::try_from(payload_len).map_err(|_| LogError::Corrupt)?,
1253 checksummed,
1254 ))
1255}
1256
1257fn frame_checksum(header: [u8; 8], payload: &[u8]) -> u32 {
1258 crc32c::crc32c_append(crc32c::crc32c(&header), payload)
1259}
1260
1261#[cfg(unix)]
1262fn read_exact_at(file: &File, mut bytes: &mut [u8], mut offset: u64) -> io::Result<()> {
1263 use std::os::unix::fs::FileExt;
1264 while !bytes.is_empty() {
1265 match file.read_at(bytes, offset) {
1266 Ok(0) => return Err(io::ErrorKind::UnexpectedEof.into()),
1267 Ok(read) => {
1268 offset = offset
1269 .checked_add(u64::try_from(read).expect("read length fits u64"))
1270 .ok_or_else(|| io::Error::other("file offset overflow"))?;
1271 bytes = &mut bytes[read..];
1272 }
1273 Err(error) if error.kind() == io::ErrorKind::Interrupted => {}
1274 Err(error) => return Err(error),
1275 }
1276 }
1277 Ok(())
1278}
1279
1280#[cfg(windows)]
1281fn read_exact_at(file: &File, mut bytes: &mut [u8], mut offset: u64) -> io::Result<()> {
1282 use std::os::windows::fs::FileExt;
1283 while !bytes.is_empty() {
1284 match file.seek_read(bytes, offset) {
1285 Ok(0) => return Err(io::ErrorKind::UnexpectedEof.into()),
1286 Ok(read) => {
1287 offset = offset
1288 .checked_add(u64::try_from(read).expect("read length fits u64"))
1289 .ok_or_else(|| io::Error::other("file offset overflow"))?;
1290 bytes = &mut bytes[read..];
1291 }
1292 Err(error) if error.kind() == io::ErrorKind::Interrupted => {}
1293 Err(error) => return Err(error),
1294 }
1295 }
1296 Ok(())
1297}
1298
1299#[cfg(not(any(unix, windows)))]
1300fn read_exact_at(file: &File, bytes: &mut [u8], offset: u64) -> io::Result<()> {
1301 use std::io::{Read, Seek, SeekFrom};
1302 let mut file = file.try_clone()?;
1303 file.seek(SeekFrom::Start(offset))?;
1304 file.read_exact(bytes)
1305}
1306
1307fn scan_frames(file: &File, offset: u64) -> Result<(Vec<FrameIndex>, u64), LogError> {
1316 let file_len = file.metadata()?.len();
1317 let mut frames = Vec::new();
1318 let mut durable_len = offset;
1319 loop {
1320 if file_len.saturating_sub(durable_len) < 8 {
1321 break;
1322 }
1323 let mut len = [0; 8];
1324 read_exact_at(file, &mut len, durable_len)?;
1325 let (payload_len, checksummed) = frame_payload_len(len)?;
1326 let frame_len = 8_u64
1327 .checked_add(u64::try_from(payload_len).map_err(|_| LogError::Corrupt)?)
1328 .and_then(|len| {
1329 len.checked_add(if checksummed {
1330 u64::try_from(FRAME_CHECKSUM_LEN).expect("checksum length fits u64")
1331 } else {
1332 0
1333 })
1334 })
1335 .ok_or(LogError::Corrupt)?;
1336 if file_len.saturating_sub(durable_len) < frame_len {
1339 break;
1340 }
1341 let mut payload = vec![0; payload_len];
1342 let payload_offset = durable_len.checked_add(8).ok_or(LogError::Corrupt)?;
1343 read_exact_at(file, &mut payload, payload_offset)?;
1344 if checksummed {
1345 let mut stored_checksum = [0; FRAME_CHECKSUM_LEN];
1346 let checksum_offset = payload_offset
1347 .checked_add(u64::try_from(payload_len).map_err(|_| LogError::Corrupt)?)
1348 .ok_or(LogError::Corrupt)?;
1349 read_exact_at(file, &mut stored_checksum, checksum_offset)?;
1350 if u32::from_be_bytes(stored_checksum) != frame_checksum(len, &payload) {
1351 return Err(LogError::Corrupt);
1352 }
1353 }
1354 let record = decode_record(&payload)?;
1355 frames.push(FrameIndex {
1356 t: record.t,
1357 offset: durable_len,
1358 len: frame_len,
1359 });
1360 durable_len = durable_len
1361 .checked_add(frame_len)
1362 .ok_or(LogError::Corrupt)?;
1363 }
1364 Ok((frames, durable_len))
1365}
1366
1367pub fn append_framed_record(out: &mut Vec<u8>, record: &TxRecord) -> Result<(), LogError> {
1376 let payload = encode_record(record);
1377 let header = frame_header(payload.len())?;
1378 out.extend_from_slice(&header);
1379 out.extend_from_slice(&payload);
1380 out.extend_from_slice(&frame_checksum(header, &payload).to_be_bytes());
1381 Ok(())
1382}
1383
1384pub fn decode_framed_records(mut bytes: &[u8]) -> Result<Vec<TxRecord>, LogError> {
1394 let mut records = Vec::new();
1395 while !bytes.is_empty() {
1396 if bytes.len() < 8 {
1397 return Err(LogError::Corrupt);
1398 }
1399 let header: [u8; 8] = bytes[..8].try_into().map_err(|_| LogError::Corrupt)?;
1400 let (payload_len, checksummed) = frame_payload_len(header)?;
1401 bytes = &bytes[8..];
1402 let payload = bytes.get(..payload_len).ok_or(LogError::Corrupt)?;
1403 bytes = &bytes[payload_len..];
1404 if checksummed {
1405 let stored_checksum = u32::from_be_bytes(
1406 bytes
1407 .get(..FRAME_CHECKSUM_LEN)
1408 .ok_or(LogError::Corrupt)?
1409 .try_into()
1410 .map_err(|_| LogError::Corrupt)?,
1411 );
1412 if stored_checksum != frame_checksum(header, payload) {
1413 return Err(LogError::Corrupt);
1414 }
1415 bytes = &bytes[FRAME_CHECKSUM_LEN..];
1416 }
1417 records.push(decode_record(payload)?);
1418 }
1419 Ok(records)
1420}
1421
1422#[cfg(test)]
1423mod tests {
1424 use super::*;
1425
1426 #[test]
1427 fn versioned_log_bounds_cached_read_descriptors() {
1428 let dir = tempfile::tempdir().expect("tempdir");
1429 let segment_count = MAX_CACHED_READ_VERSION_FILES + 5;
1430 for version in 1..=u64::try_from(segment_count).expect("segment count fits u64") {
1431 File::create(version_path(dir.path(), "db", version)).expect("create segment");
1432 }
1433
1434 let log = VersionedLog::open_read_only(dir.path(), "db").expect("open log");
1435 let state = log.state.read().expect("log lock");
1436 assert_eq!(state.files.len(), segment_count);
1437 assert!(
1438 state
1439 .files
1440 .iter()
1441 .filter(|file| file.file.file.is_some())
1442 .count()
1443 <= MAX_CACHED_READ_VERSION_FILES
1444 );
1445 }
1446}