use crate::repr::Batch;
use crate::schema::ColumnTable;
use crate::schema::SchemaDescriptor;
use gnitz_wire::OrderKey;
use super::index::TopNIndex;
use crate::algebra::GroupOutKey;
pub struct TopNPlan {
pub output_schema: SchemaDescriptor,
pub(super) key: GroupOutKey,
pub(super) offset: u64,
pub(super) limit: u64,
pub index: TopNIndex,
}
impl TopNPlan {
pub fn from_wire(
input_schema: &SchemaDescriptor,
group_cols: &[u32],
order: &[OrderKey],
limit: u64,
offset: u64,
) -> Result<Self, String> {
if limit == 0 {
return Err("top-n: a zero limit selects nothing".to_string());
}
let (key, prefix) = GroupOutKey::new(input_schema, group_cols, 0..input_schema.num_columns() as u32)?;
let output_schema = prefix.finish().map_err(|e| format!("top-n: output {e}"))?;
let index = TopNIndex::new(input_schema, group_cols, order, &output_schema)?;
Ok(TopNPlan { output_schema, key, offset, limit, index })
}
pub fn index_batch(&self, delta: &Batch) -> Batch {
self.index.batch(delta, self.key.carried())
}
pub fn partial(
input_schema: &SchemaDescriptor,
order: &[OrderKey],
limit: u64,
offset: u64,
) -> Result<Self, String> {
Self::from_wire(input_schema, &[], order, limit.saturating_add(offset), 0)
}
pub fn combine(partials: &SchemaDescriptor, order: &[OrderKey], limit: u64, offset: u64) -> Result<Self, String> {
let key = partials.pk_cols();
let shifted: Vec<OrderKey> = order
.iter()
.map(|k| OrderKey {
col: k.col.saturating_add(key.len() as u16),
..*k
})
.collect();
let plan = Self::from_wire(partials, key, &shifted, limit, offset)?;
debug_assert_eq!(plan.output_schema, *partials);
Ok(plan)
}
}