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, Read, 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>();
19
20#[derive(Clone, Debug, Eq, PartialEq)]
22pub struct TxRecord {
23 pub t: u64,
25 pub tx_instant: i64,
27 pub datoms: Vec<Datom>,
29}
30
31#[derive(Debug, Error)]
33pub enum LogError {
34 #[error("log I/O failed: {0}")]
36 Io(#[from] io::Error),
37 #[error("corrupt transaction log")]
39 Corrupt,
40 #[error("native transaction log store failed: {0}")]
42 Native(String),
43 #[error("this transaction log requires asynchronous access")]
45 AsyncOnly,
46}
47
48#[async_trait]
50pub trait TransactionLog: Send + Sync {
51 fn append(&self, record: &TxRecord) -> Result<(), LogError>;
56 async fn append_async(&self, record: &TxRecord) -> Result<(), LogError> {
63 self.append(record)
64 }
65 async fn append_batch_async(&self, records: &[TxRecord]) -> Result<(), LogError> {
75 for record in records {
76 self.append_async(record).await?;
77 }
78 Ok(())
79 }
80 fn tx_range(&self, start: u64, end: Option<u64>) -> Result<Vec<TxRecord>, LogError>;
85 async fn tx_range_async(
90 &self,
91 start: u64,
92 end: Option<u64>,
93 ) -> Result<Vec<TxRecord>, LogError> {
94 self.tx_range(start, end)
95 }
96 fn replay(&self) -> Result<Vec<TxRecord>, LogError> {
101 self.tx_range(0, None)
102 }
103 async fn replay_async(&self) -> Result<Vec<TxRecord>, LogError> {
108 self.tx_range_async(0, None).await
109 }
110}
111
112#[derive(Clone, Default)]
114pub struct MemoryLog(Arc<RwLock<Vec<TxRecord>>>);
115impl TransactionLog for MemoryLog {
116 fn append(&self, record: &TxRecord) -> Result<(), LogError> {
117 let mut records = self.0.write().expect("poisoned log lock");
118 if records.last().map_or(1, |r| r.t + 1) != record.t {
119 return Err(LogError::Corrupt);
120 }
121 records.push(record.clone());
122 Ok(())
123 }
124 fn tx_range(&self, start: u64, end: Option<u64>) -> Result<Vec<TxRecord>, LogError> {
125 Ok(self
126 .0
127 .read()
128 .expect("poisoned log lock")
129 .iter()
130 .filter(|r| r.t >= start && end.is_none_or(|e| r.t < e))
131 .cloned()
132 .collect())
133 }
134}
135
136pub struct FileLog {
142 path: PathBuf,
143 next_t: RwLock<u64>,
144}
145impl FileLog {
146 pub fn open(path: impl AsRef<Path>) -> Result<Self, LogError> {
152 let path = path.as_ref().to_path_buf();
153 if let Some(parent) = path.parent() {
154 fs::create_dir_all(parent)?;
155 }
156 OpenOptions::new().create(true).append(true).open(&path)?;
157 let (records, durable_len) = read_records(&path)?;
158 if fs::metadata(&path)?.len() > durable_len {
159 let file = OpenOptions::new().write(true).open(&path)?;
160 file.set_len(durable_len)?;
161 file.sync_all()?;
162 }
163 Ok(Self {
164 path,
165 next_t: RwLock::new(records.last().map_or(1, |r| r.t + 1)),
166 })
167 }
168}
169impl TransactionLog for FileLog {
170 fn append(&self, record: &TxRecord) -> Result<(), LogError> {
171 let mut next_t = self.next_t.write().expect("poisoned log lock");
172 if *next_t != record.t {
173 return Err(LogError::Corrupt);
174 }
175 let mut frame = Vec::new();
176 append_framed_record(&mut frame, record)?;
177 let mut file = OpenOptions::new().append(true).open(&self.path)?;
178 file.write_all(&frame)?;
179 file.sync_all()?;
180 *next_t += 1;
181 Ok(())
182 }
183 fn tx_range(&self, start: u64, end: Option<u64>) -> Result<Vec<TxRecord>, LogError> {
184 let _guard = self.next_t.read().expect("poisoned log lock");
185 Ok(read_records(&self.path)?
186 .0
187 .into_iter()
188 .filter(|r| r.t >= start && end.is_none_or(|e| r.t < e))
189 .collect())
190 }
191}
192
193pub struct VersionedLog {
205 dir: PathBuf,
206 name: String,
207 write_path: PathBuf,
208 next_t: RwLock<u64>,
209}
210
211impl VersionedLog {
212 pub fn open(dir: impl AsRef<Path>, name: &str, write_version: u64) -> Result<Self, LogError> {
220 let dir = dir.as_ref().to_path_buf();
221 fs::create_dir_all(&dir)?;
222 let write_path = version_path(&dir, name, write_version);
223 OpenOptions::new()
224 .create(true)
225 .append(true)
226 .open(&write_path)?;
227 let (_, durable_len) = read_records(&write_path)?;
228 if fs::metadata(&write_path)?.len() > durable_len {
229 let file = OpenOptions::new().write(true).open(&write_path)?;
230 file.set_len(durable_len)?;
231 file.sync_all()?;
232 }
233 let records = read_merged(&dir, name)?;
234 Ok(Self {
235 dir,
236 name: name.to_owned(),
237 write_path,
238 next_t: RwLock::new(records.last().map_or(1, |r| r.t + 1)),
239 })
240 }
241
242 pub fn open_read_only(dir: impl AsRef<Path>, name: &str) -> Result<Self, LogError> {
248 let dir = dir.as_ref().to_path_buf();
249 Ok(Self {
250 write_path: PathBuf::new(),
251 name: name.to_owned(),
252 next_t: RwLock::new(u64::MAX),
253 dir,
254 })
255 }
256
257 #[must_use]
259 pub fn exists(dir: impl AsRef<Path>, name: &str) -> bool {
260 !version_files(dir.as_ref(), name).is_empty()
261 }
262
263 pub fn delete_all(dir: impl AsRef<Path>, name: &str) -> Result<(), LogError> {
268 for (_, path) in version_files(dir.as_ref(), name) {
269 match fs::remove_file(&path) {
270 Ok(()) => {}
271 Err(error) if error.kind() == io::ErrorKind::NotFound => {}
272 Err(error) => return Err(error.into()),
273 }
274 }
275 Ok(())
276 }
277}
278
279impl TransactionLog for VersionedLog {
280 fn append(&self, record: &TxRecord) -> Result<(), LogError> {
281 let mut next_t = self.next_t.write().expect("poisoned log lock");
282 if *next_t != record.t {
283 return Err(LogError::Corrupt);
284 }
285 let mut frame = Vec::new();
286 append_framed_record(&mut frame, record)?;
287 let mut file = OpenOptions::new().append(true).open(&self.write_path)?;
288 file.write_all(&frame)?;
289 file.sync_all()?;
290 *next_t += 1;
291 Ok(())
292 }
293
294 fn tx_range(&self, start: u64, end: Option<u64>) -> Result<Vec<TxRecord>, LogError> {
295 let _guard = self.next_t.read().expect("poisoned log lock");
296 Ok(read_merged(&self.dir, &self.name)?
297 .into_iter()
298 .filter(|r| r.t >= start && end.is_none_or(|e| r.t < e))
299 .collect())
300 }
301}
302
303fn merge_versions(mut per_version: Vec<Vec<TxRecord>>) -> Vec<TxRecord> {
308 let mut cutoff = u64::MAX;
309 for records in per_version.iter_mut().rev() {
310 let first = records.first().map(|r| r.t);
311 records.retain(|r| r.t < cutoff);
312 if let Some(first) = first {
313 cutoff = cutoff.min(first);
314 }
315 }
316 per_version.into_iter().flatten().collect()
317}
318
319#[async_trait]
335pub trait NativeLogStorage: Send + Sync {
336 async fn put_batch(
349 &self,
350 name: &str,
351 version: u64,
352 records: &[(u64, Vec<u8>)],
353 ) -> Result<bool, LogError>;
354 async fn read_record(
360 &self,
361 name: &str,
362 version: u64,
363 t: u64,
364 ) -> Result<Option<Vec<u8>>, LogError>;
365 async fn list_records(&self, name: &str) -> Result<Vec<(u64, u64)>, LogError>;
371 async fn read_legacy_chunk(
377 &self,
378 name: &str,
379 version: u64,
380 chunk: u64,
381 ) -> Result<Option<Vec<u8>>, LogError>;
382 async fn list_legacy_chunks(&self, name: &str) -> Result<Vec<(u64, u64)>, LogError>;
390 async fn delete_all(&self, name: &str) -> Result<(), LogError>;
395}
396
397pub struct NativeVersionedLog<S: ?Sized> {
406 storage: Arc<S>,
407 name: String,
408 write_version: u64,
409 read_only: bool,
410 next_t: tokio::sync::Mutex<u64>,
412}
413
414impl<S: NativeLogStorage + ?Sized + 'static> NativeVersionedLog<S> {
415 pub async fn open(storage: Arc<S>, name: &str, write_version: u64) -> Result<Self, LogError> {
420 let records = read_native_merged(storage.as_ref(), name).await?;
424 let next_t = records.last().map_or(1, |r| r.t + 1);
425 Ok(Self {
426 storage,
427 name: name.to_owned(),
428 write_version,
429 read_only: false,
430 next_t: tokio::sync::Mutex::new(next_t),
431 })
432 }
433
434 #[must_use]
437 pub fn open_read_only(storage: Arc<S>, name: &str) -> Self {
438 Self {
439 storage,
440 name: name.to_owned(),
441 write_version: 0,
442 read_only: true,
443 next_t: tokio::sync::Mutex::new(0),
444 }
445 }
446}
447
448#[async_trait]
449impl<S: NativeLogStorage + ?Sized + 'static> TransactionLog for NativeVersionedLog<S> {
450 fn append(&self, record: &TxRecord) -> Result<(), LogError> {
451 let _ = record;
452 Err(LogError::AsyncOnly)
453 }
454
455 async fn append_async(&self, record: &TxRecord) -> Result<(), LogError> {
456 self.append_batch_async(std::slice::from_ref(record)).await
457 }
458
459 async fn append_batch_async(&self, records: &[TxRecord]) -> Result<(), LogError> {
460 if self.read_only {
461 return Err(LogError::Native("transaction log is read-only".into()));
462 }
463 if records.is_empty() {
464 return Ok(());
465 }
466 let mut next_t = self.next_t.lock().await;
467 for (offset, record) in records.iter().enumerate() {
469 if record.t != *next_t + offset as u64 {
470 return Err(LogError::Corrupt);
471 }
472 }
473 let framed = records
474 .iter()
475 .map(|record| {
476 let mut bytes = Vec::new();
477 append_framed_record(&mut bytes, record)?;
478 Ok((record.t, bytes))
479 })
480 .collect::<Result<Vec<_>, LogError>>()?;
481 if !self
486 .storage
487 .put_batch(&self.name, self.write_version, &framed)
488 .await?
489 {
490 return Err(LogError::Corrupt);
491 }
492 *next_t += records.len() as u64;
493 Ok(())
494 }
495
496 fn tx_range(&self, start: u64, end: Option<u64>) -> Result<Vec<TxRecord>, LogError> {
497 let _ = (start, end);
498 Err(LogError::AsyncOnly)
499 }
500
501 async fn tx_range_async(
502 &self,
503 start: u64,
504 end: Option<u64>,
505 ) -> Result<Vec<TxRecord>, LogError> {
506 let _guard = self.next_t.lock().await;
509 Ok(read_native_merged(self.storage.as_ref(), &self.name)
510 .await?
511 .into_iter()
512 .filter(|r| r.t >= start && end.is_none_or(|e| r.t < e))
513 .collect())
514 }
515}
516
517async fn read_native_merged<S: NativeLogStorage + ?Sized>(
518 storage: &S,
519 name: &str,
520) -> Result<Vec<TxRecord>, LogError> {
521 use std::collections::BTreeMap;
522
523 let mut per_version: BTreeMap<u64, Vec<TxRecord>> = BTreeMap::new();
526
527 let mut chunks = storage.list_legacy_chunks(name).await?;
531 chunks.sort_unstable();
532 for (version, chunk) in chunks {
533 let bytes = storage
534 .read_legacy_chunk(name, version, chunk)
535 .await?
536 .unwrap_or_default();
537 per_version
538 .entry(version)
539 .or_default()
540 .extend(decode_framed_records(&bytes)?);
541 }
542
543 let mut records = storage.list_records(name).await?;
545 records.sort_unstable();
546 for (version, t) in records {
547 let bytes = storage
548 .read_record(name, version, t)
549 .await?
550 .unwrap_or_default();
551 per_version
552 .entry(version)
553 .or_default()
554 .extend(decode_framed_records(&bytes)?);
555 }
556
557 let per_version: Vec<Vec<TxRecord>> = per_version
562 .into_values()
563 .map(|mut records| {
564 records.sort_by_key(|record| record.t);
565 records
566 })
567 .collect();
568 let merged = merge_versions(per_version);
569 for pair in merged.windows(2) {
570 if pair[1].t != pair[0].t + 1 {
571 return Err(LogError::Corrupt);
572 }
573 }
574 Ok(merged)
575}
576
577type VersionedRecords = Arc<Mutex<Vec<(u64, TxRecord)>>>;
580
581#[derive(Clone, Default)]
587pub struct MemLogRegistry {
588 logs: Arc<Mutex<HashMap<String, VersionedRecords>>>,
589}
590
591impl MemLogRegistry {
592 #[must_use]
594 pub fn new() -> Self {
595 Self::default()
596 }
597
598 fn entry(&self, name: &str) -> VersionedRecords {
599 Arc::clone(
600 self.logs
601 .lock()
602 .unwrap_or_else(std::sync::PoisonError::into_inner)
603 .entry(name.to_owned())
604 .or_default(),
605 )
606 }
607
608 #[must_use]
611 pub fn open(&self, name: &str, write_version: u64) -> MemVersionedLog {
612 let records = self.entry(name);
613 let next_t = {
614 let guard = records
615 .lock()
616 .unwrap_or_else(std::sync::PoisonError::into_inner);
617 MemVersionedLog::merged(&guard)
618 .last()
619 .map_or(1, |r| r.t + 1)
620 };
621 MemVersionedLog {
622 records,
623 write_version,
624 next_t: Mutex::new(next_t),
625 }
626 }
627
628 #[must_use]
630 pub fn exists(&self, name: &str) -> bool {
631 self.logs
632 .lock()
633 .unwrap_or_else(std::sync::PoisonError::into_inner)
634 .get(name)
635 .is_some_and(|entry| {
636 !entry
637 .lock()
638 .unwrap_or_else(std::sync::PoisonError::into_inner)
639 .is_empty()
640 })
641 }
642
643 pub fn delete_all(&self, name: &str) {
645 self.logs
646 .lock()
647 .unwrap_or_else(std::sync::PoisonError::into_inner)
648 .remove(name);
649 }
650}
651
652pub struct MemVersionedLog {
656 records: VersionedRecords,
657 write_version: u64,
658 next_t: Mutex<u64>,
662}
663
664impl MemVersionedLog {
665 fn merged(records: &[(u64, TxRecord)]) -> Vec<TxRecord> {
666 let mut versions: Vec<u64> = records.iter().map(|(version, _)| *version).collect();
667 versions.sort_unstable();
668 versions.dedup();
669 let per_version = versions
670 .into_iter()
671 .map(|version| {
672 records
673 .iter()
674 .filter(|(record_version, _)| *record_version == version)
675 .map(|(_, record)| record.clone())
676 .collect::<Vec<_>>()
677 })
678 .collect();
679 merge_versions(per_version)
680 }
681}
682
683impl TransactionLog for MemVersionedLog {
684 fn append(&self, record: &TxRecord) -> Result<(), LogError> {
685 let mut next_t = self
686 .next_t
687 .lock()
688 .unwrap_or_else(std::sync::PoisonError::into_inner);
689 if *next_t != record.t {
690 return Err(LogError::Corrupt);
691 }
692 self.records
693 .lock()
694 .unwrap_or_else(std::sync::PoisonError::into_inner)
695 .push((self.write_version, record.clone()));
696 *next_t += 1;
697 Ok(())
698 }
699
700 fn tx_range(&self, start: u64, end: Option<u64>) -> Result<Vec<TxRecord>, LogError> {
701 let records = self
702 .records
703 .lock()
704 .unwrap_or_else(std::sync::PoisonError::into_inner);
705 Ok(Self::merged(&records)
706 .into_iter()
707 .filter(|r| r.t >= start && end.is_none_or(|e| r.t < e))
708 .collect())
709 }
710}
711
712fn version_path(dir: &Path, name: &str, version: u64) -> PathBuf {
713 if version == 0 {
714 dir.join(format!("{name}.log"))
715 } else {
716 dir.join(format!("{name}.v{version}.log"))
717 }
718}
719
720fn version_files(dir: &Path, name: &str) -> Vec<(u64, PathBuf)> {
722 let mut files = Vec::new();
723 let legacy = version_path(dir, name, 0);
724 if legacy.is_file() {
725 files.push((0, legacy));
726 }
727 let prefix = format!("{name}.v");
728 if let Ok(entries) = fs::read_dir(dir) {
729 for entry in entries.flatten() {
730 let file_name = entry.file_name();
731 let Some(text) = file_name.to_str() else {
732 continue;
733 };
734 if let Some(version) = text
735 .strip_prefix(&prefix)
736 .and_then(|rest| rest.strip_suffix(".log"))
737 .and_then(|v| v.parse::<u64>().ok())
738 && version > 0
739 {
740 files.push((version, entry.path()));
741 }
742 }
743 }
744 files.sort_by_key(|(version, _)| *version);
745 files
746}
747
748fn read_merged(dir: &Path, name: &str) -> Result<Vec<TxRecord>, LogError> {
751 let files = version_files(dir, name);
752 let mut per_file: Vec<Vec<TxRecord>> = Vec::with_capacity(files.len());
753 for (_, path) in &files {
754 per_file.push(read_records(path)?.0);
755 }
756 let merged = merge_versions(per_file);
761 for pair in merged.windows(2) {
762 if pair[1].t != pair[0].t + 1 {
763 return Err(LogError::Corrupt);
764 }
765 }
766 Ok(merged)
767}
768
769fn encode_record(record: &TxRecord) -> Vec<u8> {
770 let mut out = Vec::new();
771 out.extend_from_slice(&record.t.to_be_bytes());
772 out.extend_from_slice(&record.tx_instant.to_be_bytes());
773 out.extend_from_slice(&(record.datoms.len() as u64).to_be_bytes());
774 for d in &record.datoms {
775 out.extend_from_slice(&d.e.raw().to_be_bytes());
776 out.extend_from_slice(&d.a.raw().to_be_bytes());
777 out.extend_from_slice(&d.tx.raw().to_be_bytes());
778 out.push(u8::from(d.added));
779 let v = encode_value(&d.v);
780 out.extend_from_slice(&(v.len() as u64).to_be_bytes());
781 out.extend_from_slice(&v);
782 }
783 out
784}
785fn decode_record(mut bytes: &[u8]) -> Result<TxRecord, LogError> {
786 fn take<'a>(bytes: &mut &'a [u8], n: usize) -> Result<&'a [u8], LogError> {
787 let value = bytes.get(..n).ok_or(LogError::Corrupt)?;
788 *bytes = &bytes[n..];
789 Ok(value)
790 }
791 fn u64_be(bytes: &mut &[u8]) -> Result<u64, LogError> {
792 Ok(u64::from_be_bytes(
793 take(bytes, 8)?.try_into().map_err(|_| LogError::Corrupt)?,
794 ))
795 }
796 let t = u64_be(&mut bytes)?;
797 let tx_instant = i64::from_be_bytes(
798 take(&mut bytes, 8)?
799 .try_into()
800 .map_err(|_| LogError::Corrupt)?,
801 );
802 let count = u64_be(&mut bytes)?;
803 let mut datoms = Vec::new();
804 for _ in 0..count {
805 let e = EntityId::from_raw(u64_be(&mut bytes)?);
806 let a = EntityId::from_raw(u64_be(&mut bytes)?);
807 let tx = EntityId::from_raw(u64_be(&mut bytes)?);
808 let added = take(&mut bytes, 1)?[0] != 0;
809 let len = usize::try_from(u64_be(&mut bytes)?).map_err(|_| LogError::Corrupt)?;
810 let raw = take(&mut bytes, len)?;
811 let (v, used) = decode_value(raw).map_err(|_| LogError::Corrupt)?;
812 if used != len {
813 return Err(LogError::Corrupt);
814 }
815 datoms.push(Datom { e, a, v, tx, added });
816 }
817 if !bytes.is_empty() {
818 return Err(LogError::Corrupt);
819 }
820 Ok(TxRecord {
821 t,
822 tx_instant,
823 datoms,
824 })
825}
826
827fn frame_header(payload_len: usize) -> Result<[u8; 8], LogError> {
828 let payload_len = u64::try_from(payload_len).map_err(|_| LogError::Corrupt)?;
829 if payload_len & CHECKSUMMED_FRAME != 0 {
830 return Err(LogError::Corrupt);
831 }
832 Ok((payload_len | CHECKSUMMED_FRAME).to_be_bytes())
833}
834
835fn frame_payload_len(header: [u8; 8]) -> Result<(usize, bool), LogError> {
836 let encoded = u64::from_be_bytes(header);
837 let checksummed = encoded & CHECKSUMMED_FRAME != 0;
838 let payload_len = encoded & !CHECKSUMMED_FRAME;
839 Ok((
840 usize::try_from(payload_len).map_err(|_| LogError::Corrupt)?,
841 checksummed,
842 ))
843}
844
845fn frame_checksum(header: [u8; 8], payload: &[u8]) -> u32 {
846 crc32c::crc32c_append(crc32c::crc32c(&header), payload)
847}
848
849fn read_records(path: &Path) -> Result<(Vec<TxRecord>, u64), LogError> {
857 let mut file = File::open(path)?;
858 let file_len = file.metadata()?.len();
859 let mut records = Vec::new();
860 let mut durable_len = 0_u64;
861 loop {
862 let mut len = [0; 8];
863 match file.read_exact(&mut len) {
864 Ok(()) => {}
865 Err(e) if e.kind() == io::ErrorKind::UnexpectedEof => break,
866 Err(e) => return Err(e.into()),
867 }
868 let (payload_len, checksummed) = frame_payload_len(len)?;
869 let frame_len = 8_u64
870 .checked_add(u64::try_from(payload_len).map_err(|_| LogError::Corrupt)?)
871 .and_then(|len| {
872 len.checked_add(if checksummed {
873 u64::try_from(FRAME_CHECKSUM_LEN).expect("checksum length fits u64")
874 } else {
875 0
876 })
877 })
878 .ok_or(LogError::Corrupt)?;
879 if file_len.saturating_sub(durable_len) < frame_len {
882 break;
883 }
884 let mut payload = vec![0; payload_len];
885 match file.read_exact(&mut payload) {
886 Ok(()) => {}
887 Err(e) if e.kind() == io::ErrorKind::UnexpectedEof => break,
888 Err(e) => return Err(e.into()),
889 }
890 if checksummed {
891 let mut stored_checksum = [0; FRAME_CHECKSUM_LEN];
892 match file.read_exact(&mut stored_checksum) {
893 Ok(()) => {}
894 Err(e) if e.kind() == io::ErrorKind::UnexpectedEof => break,
895 Err(e) => return Err(e.into()),
896 }
897 if u32::from_be_bytes(stored_checksum) != frame_checksum(len, &payload) {
898 return Err(LogError::Corrupt);
899 }
900 }
901 records.push(decode_record(&payload)?);
902 durable_len = durable_len
903 .checked_add(frame_len)
904 .ok_or(LogError::Corrupt)?;
905 }
906 Ok((records, durable_len))
907}
908
909pub fn append_framed_record(out: &mut Vec<u8>, record: &TxRecord) -> Result<(), LogError> {
918 let payload = encode_record(record);
919 let header = frame_header(payload.len())?;
920 out.extend_from_slice(&header);
921 out.extend_from_slice(&payload);
922 out.extend_from_slice(&frame_checksum(header, &payload).to_be_bytes());
923 Ok(())
924}
925
926pub fn decode_framed_records(mut bytes: &[u8]) -> Result<Vec<TxRecord>, LogError> {
936 let mut records = Vec::new();
937 while !bytes.is_empty() {
938 if bytes.len() < 8 {
939 return Err(LogError::Corrupt);
940 }
941 let header: [u8; 8] = bytes[..8].try_into().map_err(|_| LogError::Corrupt)?;
942 let (payload_len, checksummed) = frame_payload_len(header)?;
943 bytes = &bytes[8..];
944 let payload = bytes.get(..payload_len).ok_or(LogError::Corrupt)?;
945 bytes = &bytes[payload_len..];
946 if checksummed {
947 let stored_checksum = u32::from_be_bytes(
948 bytes
949 .get(..FRAME_CHECKSUM_LEN)
950 .ok_or(LogError::Corrupt)?
951 .try_into()
952 .map_err(|_| LogError::Corrupt)?,
953 );
954 if stored_checksum != frame_checksum(header, payload) {
955 return Err(LogError::Corrupt);
956 }
957 bytes = &bytes[FRAME_CHECKSUM_LEN..];
958 }
959 records.push(decode_record(payload)?);
960 }
961 Ok(records)
962}