use gnitz_wire::ClampKind;
use crate::repr::{materialize_carrying, Batch};
use crate::stream::OpenAt;
fn cap(kind: ClampKind) -> i64 {
match kind {
ClampKind::Distinct => 1,
ClampKind::PositivePart => i64::MAX,
}
}
pub fn op_weight_clamp(delta: &Batch, hist: OpenAt<'_>, kind: ClampKind) -> Batch {
debug_assert!(delta.is_consolidated());
if delta.is_empty() {
return Batch::empty_with_schema(delta.schema());
}
let cursor = &mut hist(delta.get_pk_bytes(0), delta.get_pk_bytes(delta.len() - 1));
let cap = cap(kind);
let mb = delta.as_mem_batch();
let mut rows: Vec<(u32, u32, i64)> = Vec::new();
let mut identity = true;
cursor.for_each_mem_row_weight(&mb, |i, w_old| {
let dw = mb.get_weight(i);
let out_w = w_old.wrapping_add(dw).clamp(0, cap) - w_old.clamp(0, cap);
identity &= out_w == dw;
if out_w != 0 {
rows.push((0, i as u32, out_w));
}
});
if identity {
return delta.clone();
}
let mut out = materialize_carrying(std::slice::from_ref(&mb), delta.schema(), &rows);
out.certify_consolidated();
out
}
#[cfg(test)]
#[path = "tests/clamp.rs"]
mod tests;
#[cfg(test)]
#[path = "benches/clamp.rs"]
mod bench;