use std::fs::File;
use crate::error::{Error, ErrorContext, Result};
use super::budget::{ScanBudget, ScanMetrics};
use super::reader::Reader;
use super::scan::ScanPlan;
#[derive(Debug)]
pub struct Tail<'reader> {
reader: &'reader mut Reader,
plan: ScanPlan,
next_block: usize,
file: Option<File>,
budget: ScanBudget,
}
impl<'reader> Tail<'reader> {
pub(crate) fn new(reader: &'reader mut Reader) -> Self {
let limits = reader.limits();
Self {
next_block: reader.blocks().len(),
plan: ScanPlan::new(reader.schema()),
reader,
file: None,
budget: ScanBudget::new(limits.max_rows_per_scan(), limits.max_decoded_scan_bytes()),
}
}
pub fn project<I, S>(mut self, columns: I) -> Result<Self>
where
I: IntoIterator<Item = S>,
S: AsRef<str>,
{
self.plan.project(self.reader.schema(), columns)?;
Ok(self)
}
pub fn primary_range(mut self, range: crate::PrimaryRange) -> Result<Self> {
self.plan.primary_range(self.reader.schema(), range)?;
Ok(self)
}
pub fn file_order(self) -> Self {
self
}
pub fn poll_next(&mut self) -> Result<Option<crate::RecordBatch>> {
loop {
while self.next_block < self.reader.blocks().len() {
if let Some(batch) = self.decode_next()? {
return Ok(Some(batch));
}
}
let report = self.reader.refresh()?;
self.budget
.record_bytes_read(report.frame_bytes_scanned())?;
if report.blocks_added() == 0 {
return Ok(None);
}
}
}
pub fn metrics(&self) -> ScanMetrics {
self.budget.metrics()
}
fn decode_next(&mut self) -> Result<Option<crate::RecordBatch>> {
let index = self.next_block;
self.next_block += 1;
let block = self.reader.blocks()[index].clone();
let pruned = self.plan.should_prune(&block);
self.budget.record_block(pruned)?;
if pruned {
return Ok(None);
}
let selection = self.plan.selection(self.reader.schema_handle());
let file = match &mut self.file {
Some(file) => file,
None => match File::open(self.reader.path()) {
Ok(file) => self.file.insert(file),
Err(error) => {
return Err(Error::io(error, None).with_context(ErrorContext::File));
}
},
};
self.budget.charge_rows(block.row_count())?;
let decoded =
self.reader
.decode_selected_block_at(file, index, &selection, &mut self.budget)?;
let super::decode::DecodedBlock {
batch,
primary_values,
primary_sorted,
} = decoded;
let Some(range) = self.plan.range else {
self.budget.record_rows(batch.row_count())?;
return Ok(Some(batch));
};
let Some(primary_values) = primary_values else {
return Err(Error::internal(
"range tail did not decode its primary column",
));
};
match self
.plan
.filter_batch(primary_sorted, batch, &primary_values, range)?
{
Some(batch) => {
self.budget.record_rows(batch.row_count())?;
Ok(Some(batch))
}
None => Ok(None),
}
}
}