pub(super) struct RowAtATime<'a> {
input: Box<dyn uqa_execution::PhysicalOperator + 'a>,
schema: uqa_execution::RowSchema,
ordering: Vec<uqa_execution::PhysicalOrder>,
pending: std::vec::IntoIter<uqa_execution::PhysicalRow>,
}
impl<'a> RowAtATime<'a> {
pub(super) fn new(input: Box<dyn uqa_execution::PhysicalOperator + 'a>) -> Self {
let schema = input.row_schema().clone();
let ordering = input.output_ordering().to_vec();
Self {
input,
schema,
ordering,
pending: Vec::new().into_iter(),
}
}
}
impl uqa_execution::PhysicalOperator for RowAtATime<'_> {
fn row_schema(&self) -> &uqa_execution::RowSchema {
&self.schema
}
fn estimated_cardinality(&self) -> Option<u64> {
self.input.estimated_cardinality()
}
fn output_ordering(&self) -> &[uqa_execution::PhysicalOrder] {
&self.ordering
}
fn open(&mut self) -> uqa_execution::ExecResult<()> {
self.pending = Vec::new().into_iter();
self.input.open()
}
fn next(&mut self) -> uqa_execution::ExecResult<Option<uqa_execution::Batch>> {
loop {
if let Some(row) = self.pending.next() {
return Ok(Some(uqa_execution::Batch::from_physical_rows(
self.schema.clone(),
vec![row],
)));
}
let Some(batch) = self.input.next()? else {
return Ok(None);
};
if batch.schema != self.schema {
return Err(uqa_execution::ExecError::Other(format!(
"row-at-a-time input schema mismatch: expected {:?}, got {:?}",
self.schema, batch.schema
)));
}
self.pending = batch.rows.into_iter();
}
}
fn close(&mut self) -> uqa_execution::ExecResult<()> {
self.pending = Vec::new().into_iter();
self.input.close()
}
}