use crate::{Error, Result};
pub const SYNC_MARKER_SIZE: usize = 8;
pub const RECORD_OVERHEAD: usize = 12;
#[derive(Debug)]
pub enum FrameStep<'a> {
Record(&'a [u8]),
End,
Truncated,
}
pub struct FrameWalker<'a> {
bytes: &'a [u8],
segment_id: i64,
marker_pos: usize,
cursor: usize,
section_end: usize,
in_section: bool,
truncated_end: bool,
done: bool,
}
impl<'a> FrameWalker<'a> {
pub fn new(bytes: &'a [u8], segment_id: i64, body_start: usize) -> Self {
Self {
bytes,
segment_id,
marker_pos: body_start,
cursor: 0,
section_end: 0,
in_section: false,
truncated_end: false,
done: false,
}
}
pub fn next_frame(&mut self) -> Result<FrameStep<'a>> {
if self.done {
return Ok(FrameStep::End);
}
loop {
if !self.in_section {
match self.open_section()? {
Some(()) => {}
None => {
self.done = true;
return Ok(if self.truncated_end {
FrameStep::Truncated
} else {
FrameStep::End
});
}
}
}
match self.read_record()? {
RecordOutcome::Body(b) => return Ok(FrameStep::Record(b)),
RecordOutcome::SectionDone => {
self.in_section = false;
self.marker_pos = self.section_end;
continue;
}
RecordOutcome::CleanEnd => {
self.done = true;
return Ok(if self.truncated_end {
FrameStep::Truncated
} else {
FrameStep::End
});
}
RecordOutcome::Truncated => {
self.done = true;
return Ok(FrameStep::Truncated);
}
}
}
}
fn open_section(&mut self) -> Result<Option<()>> {
let pos = self.marker_pos;
if pos + SYNC_MARKER_SIZE > self.bytes.len() {
self.truncated_end |= pos < self.bytes.len();
return Ok(None);
}
let next_marker = read_i32_be(self.bytes, pos);
if next_marker == 0 {
self.truncated_end = false;
return Ok(None);
}
let stored_crc = read_u32_be(self.bytes, pos + 4);
let computed = marker_crc(self.segment_id, pos);
if computed != stored_crc {
let all_zero = self.bytes[pos..pos + SYNC_MARKER_SIZE]
.iter()
.all(|&b| b == 0);
self.truncated_end |= !all_zero;
return Ok(None);
}
let next = next_marker as usize;
if next > self.bytes.len() {
self.section_end = self.bytes.len();
self.cursor = pos + SYNC_MARKER_SIZE;
self.in_section = true;
self.truncated_end = true;
return Ok(Some(()));
}
if next < pos + SYNC_MARKER_SIZE {
return Err(Error::CorruptCommitLogFrame(format!(
"sync marker at offset {pos} points backward to {next}"
)));
}
self.section_end = next;
self.cursor = pos + SYNC_MARKER_SIZE;
self.in_section = true;
Ok(Some(()))
}
fn read_record(&mut self) -> Result<RecordOutcome<'a>> {
let cur = self.cursor;
let len = self.bytes.len();
if cur >= self.section_end {
return Ok(RecordOutcome::SectionDone);
}
let is_final_section = self.section_end >= len;
if cur + 4 > self.section_end {
return Ok(if is_final_section {
RecordOutcome::Truncated
} else {
RecordOutcome::SectionDone
});
}
let size = read_i32_be(self.bytes, cur);
if size <= 0 {
return Ok(if is_final_section {
RecordOutcome::CleanEnd
} else {
RecordOutcome::SectionDone
});
}
let size = size as usize;
let end = cur + 4 + 4 + size + 4;
if end > len {
return Ok(RecordOutcome::Truncated);
}
if end > self.section_end {
return Err(Error::CorruptCommitLogFrame(format!(
"record at offset {cur} (len {size}) overruns section end {}",
self.section_end
)));
}
let size_crc = read_u32_be(self.bytes, cur + 4);
let size_bytes = &self.bytes[cur..cur + 4];
let mut hasher = crc32fast::Hasher::new();
hasher.update(size_bytes);
let computed_size_crc = hasher.finalize();
if computed_size_crc != size_crc {
return Err(Error::CorruptCommitLogFrame(format!(
"record size CRC mismatch at offset {cur}: stored={size_crc:#010x} \
computed={computed_size_crc:#010x}"
)));
}
let body_start = cur + 8;
let body = &self.bytes[body_start..body_start + size];
let body_crc = read_u32_be(self.bytes, body_start + size);
let mut hasher = crc32fast::Hasher::new();
hasher.update(size_bytes);
hasher.update(body);
let computed_body_crc = hasher.finalize();
if computed_body_crc != body_crc {
return Err(Error::CorruptCommitLogFrame(format!(
"record body CRC mismatch at offset {cur}: stored={body_crc:#010x} \
computed={computed_body_crc:#010x}"
)));
}
self.cursor = end;
Ok(RecordOutcome::Body(body))
}
}
enum RecordOutcome<'a> {
Body(&'a [u8]),
SectionDone,
CleanEnd,
Truncated,
}
pub fn marker_crc(segment_id: i64, marker_pos: usize) -> u32 {
let id = segment_id as u64;
let mut hasher = crc32fast::Hasher::new();
hasher.update(&((id & 0xFFFF_FFFF) as u32).to_be_bytes());
hasher.update(&((id >> 32) as u32).to_be_bytes());
hasher.update(&(marker_pos as u32).to_be_bytes());
hasher.finalize()
}
#[inline]
fn read_i32_be(b: &[u8], off: usize) -> i32 {
i32::from_be_bytes([b[off], b[off + 1], b[off + 2], b[off + 3]])
}
#[inline]
fn read_u32_be(b: &[u8], off: usize) -> u32 {
u32::from_be_bytes([b[off], b[off + 1], b[off + 2], b[off + 3]])
}
#[cfg(test)]
mod tests {
use super::*;
const SEGMENT_ID: i64 = 1234;
fn push_marker(buf: &mut Vec<u8>, pos: usize, next_marker: i32) {
buf.extend_from_slice(&next_marker.to_be_bytes());
buf.extend_from_slice(&marker_crc(SEGMENT_ID, pos).to_be_bytes());
}
fn push_record(buf: &mut Vec<u8>, body: &[u8]) {
let size = body.len() as i32;
let size_bytes = size.to_be_bytes();
let mut h = crc32fast::Hasher::new();
h.update(&size_bytes);
let size_crc = h.finalize();
buf.extend_from_slice(&size_bytes);
buf.extend_from_slice(&size_crc.to_be_bytes());
buf.extend_from_slice(body);
let mut h = crc32fast::Hasher::new();
h.update(&size_bytes);
h.update(body);
buf.extend_from_slice(&h.finalize().to_be_bytes());
}
#[test]
fn walks_records_across_multiple_sync_sections() {
let mut buf = Vec::new();
let record1 = b"section-one-record";
push_marker(&mut buf, 0, 0 );
push_record(&mut buf, record1);
buf.extend_from_slice(&[0u8; 8]); let section2_marker_pos = buf.len();
let next1 = section2_marker_pos as i32;
buf[0..4].copy_from_slice(&next1.to_be_bytes());
buf[4..8].copy_from_slice(&marker_crc(SEGMENT_ID, 0).to_be_bytes());
let record2 = b"section-two-record";
let record2_start = section2_marker_pos + SYNC_MARKER_SIZE;
let final_len = record2_start + 4 + 4 + record2.len() + 4;
push_marker(&mut buf, section2_marker_pos, final_len as i32);
push_record(&mut buf, record2);
assert_eq!(buf.len(), final_len, "test construction sanity check");
let mut walker = FrameWalker::new(&buf, SEGMENT_ID, 0);
let mut records = Vec::new();
loop {
match walker.next_frame().expect("no corruption in this fixture") {
FrameStep::Record(b) => records.push(b.to_vec()),
FrameStep::End => break,
FrameStep::Truncated => panic!("expected a clean end, not Truncated"),
}
}
assert_eq!(
records,
vec![record1.to_vec(), record2.to_vec()],
"both sections' records must be yielded — a regression that treats \
section 1's padding as segment-end would silently drop record2"
);
}
#[test]
fn torn_section_whose_record_ends_exactly_at_eof_still_reports_truncated() {
let mut buf = Vec::new();
push_marker(&mut buf, 0, 999_999); push_record(&mut buf, b"last-record-fills-the-buffer-exactly");
let mut walker = FrameWalker::new(&buf, SEGMENT_ID, 0);
match walker.next_frame().expect("valid record") {
FrameStep::Record(_) => {}
other => panic!("expected the one record, got {other:?}"),
}
match walker.next_frame().expect("no corruption") {
FrameStep::Truncated => {}
other => panic!(
"expected Truncated (torn section, record ends exactly at EOF), got {other:?}"
),
}
}
}