use crate::algebra::ReindexPacker;
use crate::repr::Batch;
use crate::schema::{ColumnLocator, SchemaColumn, SchemaDescriptor, TypeCode};
use gnitz_wire::OrderKey;
use gnitz_wire::RowSource;
use crate::algebra::{image_slot_col, ImageCol, IMAGE_COL};
const RANK_COL: SchemaColumn = SchemaColumn::new(TypeCode::U8, false);
struct OrderSpec {
image: ImageCol,
nulls_first: bool,
nullable: bool,
}
impl OrderSpec {
#[inline]
fn rank(&self, src: &impl RowSource, row: usize) -> Option<(u8, bool)> {
self.nullable.then(|| {
let is_null = self.image.loc.is_null(src, row);
((is_null != self.nulls_first) as u8, is_null)
})
}
#[inline]
fn append_image(&self, src: &impl RowSource, row: usize, out: &mut Vec<u8>) {
let rank = self.rank(src, row);
out.extend(rank.map(|(rank, _)| rank));
if !matches!(rank, Some((_, true))) {
self.image.append(src, row, out);
}
}
#[inline]
fn write_lead(&self, src: &impl RowSource, row: usize, lead: &mut [u8], image: &mut Vec<u8>) {
let slot = match self.rank(src, row) {
None => lead,
Some((rank, is_null)) => {
let (first, slot) = lead.split_at_mut(1);
first[0] = rank;
if is_null {
slot.fill(0);
return;
}
slot
}
};
self.image.write_slot(src, row, slot, image);
}
fn lead_cols(&self) -> impl Iterator<Item = SchemaColumn> {
let slot = image_slot_col(self.image.is_wide());
self.nullable.then_some(RANK_COL).into_iter().chain([slot])
}
}
pub(super) struct TopNIndex {
key_packer: ReindexPacker,
pub(super) schema: SchemaDescriptor,
lead: OrderSpec,
rest: Vec<OrderSpec>,
pub(super) carried_in_index: Vec<ColumnLocator>,
lead_bytes: usize,
}
impl TopNIndex {
pub(super) fn new(
input: &SchemaDescriptor,
group_cols: &[u32],
order: &[OrderKey],
output: &SchemaDescriptor,
) -> Result<Self, String> {
let mut order = order.iter().map(|key| {
let (col, loc) = input.wire_col("top-n: order column", key.col as u32)?;
Ok(OrderSpec {
image: ImageCol::new(loc, key.desc),
nulls_first: key.nulls_first,
nullable: col.nullable,
})
});
let lead = order.next().ok_or_else(|| "top-n: no order keys".to_string())??;
let rest = order.collect::<Result<Vec<_>, String>>()?;
let suffix: Vec<SchemaColumn> = lead.lead_cols().collect();
let (key_packer, mut b) = ReindexPacker::new_group_key(input, group_cols, &suffix)?;
for _ in 0..rest.len() + usize::from(!lead.image.fits_slot()) {
b.push(IMAGE_COL);
}
b.push_payload_of(output);
let schema = b.finish().map_err(|e| format!("top-n: index {e}"))?;
let tail = schema.num_columns() - output.num_payload_cols();
Ok(TopNIndex {
key_packer,
lead,
rest,
carried_in_index: (tail..schema.num_columns()).map(|c| schema.locate(c)).collect(),
schema,
lead_bytes: suffix.iter().map(|c| c.size()).sum(),
})
}
#[inline]
pub(super) fn group_bytes(&self) -> usize {
self.key_packer.out_stride
}
pub(super) fn batch(&self, delta: &Batch, carried: &[ColumnLocator]) -> Batch {
let mb = delta.as_mem_batch();
let mut out = Batch::with_capacity_blob(&self.schema, delta.count.max(1), delta.blob().len());
let stride = self.key_packer.out_stride;
let mut image = Vec::new();
self.key_packer.for_each_key(&mb, stride + self.lead_bytes, |row, key| {
let weight = mb.get_weight(row);
image.clear();
self.lead
.write_lead(&mb, row, &mut key[stride..stride + self.lead_bytes], &mut image);
out.begin_row(key, weight);
let mut col = 0;
if !self.lead.image.fits_slot() {
out.extend_col_blob(col, &image);
col += 1;
}
for spec in &self.rest {
image.clear();
spec.append_image(&mb, row, &mut image);
out.extend_col_blob(col, &image);
col += 1;
}
out.append_cells_from(col, carried, &mb, row);
out.commit_row();
});
out
}
}