use std::fs::File;
use std::path::Path;
use crate::error::{Error, ErrorContext, Result};
use crate::limits::Limits;
use crate::schema::Schema;
use super::constants::{
FIRST_DATA_FRAME_SEQUENCE, PROLOGUE_SIZE, ROW_IDS_FEATURE, SCHEMA_FRAME_SEQUENCE,
};
use super::data_frame::{self, DataFrameMetadata};
use super::frame::{self, FrameMetadata, FrameRead};
use super::prologue::{self, Prologue};
use super::schema_frame;
pub(crate) struct FileScan {
file: File,
file_size: u64,
limits: Limits,
prologue: Prologue,
schema: Schema,
first_data_frame_offset: u64,
}
pub(crate) struct Walk {
pub(crate) frame_count: u64,
pub(crate) next_sequence: u64,
pub(crate) last_good_offset: u64,
pub(crate) incomplete_tail: bool,
pub(crate) next_row_id: Option<u64>,
pub(crate) frame_bytes_scanned: u64,
}
pub(crate) struct ResumePoint {
pub(crate) offset: u64,
pub(crate) sequence: u64,
pub(crate) expected_base_row_id: Option<u64>,
pub(crate) frame_count: u64,
}
impl FileScan {
pub(crate) fn open(path: &Path, limits: Limits) -> Result<Self> {
let file = File::open(path)
.map_err(|error| Error::io(error, None).with_context(ErrorContext::File))?;
Self::from_file(file, limits)
}
pub(crate) fn from_file(mut file: File, limits: Limits) -> Result<Self> {
let file_size = file
.metadata()
.map(|metadata| metadata.len())
.map_err(|error| Error::io(error, None).with_context(ErrorContext::File))?;
let prologue = prologue::read_from_file(&mut file, file_size)?;
let schema_offset = PROLOGUE_SIZE as u64;
let frame = match frame::read_frame(
&mut file,
file_size,
schema_offset,
SCHEMA_FRAME_SEQUENCE,
limits,
)? {
FrameRead::Complete(frame) => frame,
FrameRead::IncompleteTail => return Err(incomplete_schema_frame(schema_offset)),
};
let schema = schema_frame::parse(&mut file, &frame, limits)?;
let first_data_frame_offset = next_offset(frame.frame_offset, frame.total_length)?;
Ok(Self {
file,
file_size,
limits,
prologue,
schema,
first_data_frame_offset,
})
}
pub(crate) fn prologue(&self) -> Prologue {
self.prologue
}
pub(crate) fn schema(&self) -> &Schema {
&self.schema
}
pub(crate) fn file_size(&self) -> u64 {
self.file_size
}
pub(crate) fn into_file(self) -> File {
self.file
}
pub(crate) fn walk_data_frames(
&mut self,
visit: impl FnMut(&FrameMetadata, &DataFrameMetadata) -> Result<()>,
) -> Result<Walk> {
self.walk_data_frames_from(self.initial_resume_point(), visit)
}
pub(crate) fn walk_data_frames_from(
&mut self,
resume: ResumePoint,
mut visit: impl FnMut(&FrameMetadata, &DataFrameMetadata) -> Result<()>,
) -> Result<Walk> {
let mut offset = resume.offset;
let mut sequence = resume.sequence;
let mut frame_count = resume.frame_count;
let mut expected_base_row_id = resume.expected_base_row_id;
let mut incomplete_tail = false;
let mut frame_bytes_scanned = 0_u64;
while offset < self.file_size {
let frame = match frame::read_frame(
&mut self.file,
self.file_size,
offset,
sequence,
self.limits,
)? {
FrameRead::Complete(frame) => frame,
FrameRead::IncompleteTail => {
incomplete_tail = true;
break;
}
};
let block = data_frame::parse(
&mut self.file,
&frame,
&self.schema,
self.prologue.feature_flags,
expected_base_row_id,
)?;
if let Some(base) = block.base_row_id {
expected_base_row_id =
Some(base.checked_add(block.row_count).ok_or_else(|| {
Error::corruption("base row ID overflow", Some(frame.frame_offset))
.with_context(ErrorContext::Frame {
sequence: frame.sequence,
})
})?);
}
visit(&frame, &block)?;
offset = next_offset(offset, frame.total_length)?;
let (Some(advanced_sequence), Some(advanced_count), Some(advanced_bytes)) = (
sequence.checked_add(1),
frame_count.checked_add(1),
frame_bytes_scanned.checked_add(frame.total_length),
) else {
return Err(
Error::corruption("frame sequence overflow", Some(frame.frame_offset))
.with_context(ErrorContext::File),
);
};
sequence = advanced_sequence;
frame_count = advanced_count;
frame_bytes_scanned = advanced_bytes;
}
Ok(Walk {
frame_count,
next_sequence: sequence,
last_good_offset: offset,
incomplete_tail,
next_row_id: expected_base_row_id,
frame_bytes_scanned,
})
}
fn initial_resume_point(&self) -> ResumePoint {
ResumePoint {
offset: self.first_data_frame_offset,
sequence: FIRST_DATA_FRAME_SEQUENCE,
expected_base_row_id: self.row_ids_enabled().then_some(0_u64),
frame_count: 1,
}
}
fn row_ids_enabled(&self) -> bool {
self.prologue.feature_flags & ROW_IDS_FEATURE != 0
}
}
fn next_offset(offset: u64, length: u64) -> Result<u64> {
offset.checked_add(length).ok_or_else(|| {
Error::corruption("next frame offset overflow", Some(offset))
.with_context(ErrorContext::File)
})
}
fn incomplete_schema_frame(offset: u64) -> Error {
Error::incomplete_tail(
"the file ends before its schema frame is complete",
Some(offset),
)
.with_context(ErrorContext::Frame {
sequence: SCHEMA_FRAME_SEQUENCE,
})
}