use std::collections::BTreeMap;
use std::time::Instant;
use rudb_common::{Error, Result};
use crate::chooser::{Chooser, EXHAUSTIVE};
use crate::reader::Reader;
use crate::tally::{self, Family};
use crate::bitpack::{self, VALUES};
const MAX_DEPTH: u8 = 3;
const RUN: usize = 8;
const LONG_RUN: usize = 64;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum Kind {
Constant = 0,
Packed = 1,
Delta = 2,
Rle = 3,
Dict = 4,
Sparse = 5,
Strided = 6,
}
impl Kind {
pub const ALL: [Self; 7] = [
Self::Constant,
Self::Packed,
Self::Delta,
Self::Rle,
Self::Dict,
Self::Sparse,
Self::Strided,
];
fn tag(self) -> u8 {
self as u8
}
fn from_tag(tag: u8) -> Result<Self> {
match tag {
0 => Ok(Self::Constant),
1 => Ok(Self::Packed),
2 => Ok(Self::Delta),
3 => Ok(Self::Rle),
4 => Ok(Self::Dict),
5 => Ok(Self::Sparse),
6 => Ok(Self::Strided),
other => Err(Error::internal(format!("unknown encoding tag {other}"))),
}
}
#[must_use]
pub fn name(self) -> &'static str {
match self {
Self::Constant => "CONSTANT",
Self::Packed => "FOR+BITPACK",
Self::Delta => "DELTA",
Self::Rle => "RLE",
Self::Dict => "DICT",
Self::Sparse => "SPARSE",
Self::Strided => "STRIDE",
}
}
}
pub fn encode(values: &[i64]) -> Result<Vec<u8>> {
encode_with(values, &EXHAUSTIVE)
}
pub fn encode_with(values: &[i64], chooser: &dyn Chooser) -> Result<Vec<u8>> {
encode_at(values, 0, chooser)
}
pub fn decode(bytes: &[u8]) -> Result<Vec<i64>> {
let mut reader = Reader::new(bytes);
let values = with_decoding(|scratch| decode_chunk(&mut reader, scratch))?;
if reader.remaining() != 0 {
return Err(Error::internal(format!(
"{} bytes left over after decoding a chunk",
reader.remaining()
)));
}
Ok(values)
}
pub fn decode_as<T: Lane>(bytes: &[u8]) -> Result<Vec<T>> {
let mut reader = Reader::new(bytes);
let values = with_decoding(|scratch| decode_chunk_as(&mut reader, scratch))?;
if reader.remaining() != 0 {
return Err(Error::internal(format!(
"{} bytes left over after decoding a chunk",
reader.remaining()
)));
}
Ok(values)
}
pub trait Lane: Copy + Default {
fn fit(value: i64) -> Option<Self>;
fn wrap(value: i64) -> Self;
}
macro_rules! lanes {
($($ty:ty),* $(,)?) => {$(
impl Lane for $ty {
fn fit(value: i64) -> Option<Self> {
Self::try_from(value).ok()
}
#[allow(
clippy::cast_possible_truncation,
clippy::cast_sign_loss,
clippy::unnecessary_cast,
reason = "only called on a value the caller has checked fits"
)]
fn wrap(value: i64) -> Self {
value as Self
}
}
)*};
}
lanes!(i8, u8, i16, u16, i32, u32, i64, u64);
fn lane<T: Lane>(value: i64) -> Result<T> {
T::fit(value).ok_or_else(|| Error::internal(format!("{value} is outside the chunk's type")))
}
pub fn tally(bytes: &[u8]) -> Result<(usize, Vec<(i64, u64)>)> {
let mut counts = BTreeMap::<i64, u64>::new();
let rows = fold(bytes, |value, count| {
*counts.entry(value).or_default() += count;
Ok(())
})?;
Ok((rows, counts.into_iter().collect()))
}
pub fn fold(bytes: &[u8], mut emit: impl FnMut(i64, u64) -> Result<()>) -> Result<usize> {
let mut reader = Reader::new(bytes);
let kind = Kind::from_tag(reader.u8()?)?;
let count = reader.u32()? as usize;
match kind {
Kind::Constant => {
let value = reader.i64()?;
if count != 0 {
emit(value, count as u64)?;
}
}
Kind::Sparse => {
let dominant = reader.i64()?;
let exception_count = reader.u32()? as usize;
let (positions, values) = with_decoding(|scratch| -> Result<_> {
Ok((decode_chunk(&mut reader, scratch)?, decode_chunk(&mut reader, scratch)?))
})?;
if positions.len() != exception_count || values.len() != exception_count {
return Err(Error::internal("a sparse chunk disagrees about its exception count"));
}
let mut ordered = true;
let mut previous = None;
for &position in &positions {
let position = usize::try_from(position)
.ok()
.filter(|&position| position < count)
.ok_or_else(|| Error::internal("a sparse exception is outside the chunk"))?;
if previous.is_some_and(|last| position <= last) {
ordered = false;
}
previous = Some(position);
}
if ordered {
if count != exception_count {
emit(dominant, (count - exception_count) as u64)?;
}
for value in values {
emit(value, 1)?;
}
} else {
let mut exceptions = BTreeMap::<usize, i64>::new();
for (position, value) in positions.into_iter().zip(values) {
exceptions.insert(position as usize, value);
}
if count != exceptions.len() {
emit(dominant, (count - exceptions.len()) as u64)?;
}
for value in exceptions.into_values() {
emit(value, 1)?;
}
}
}
Kind::Rle => {
let (values, lengths) = with_decoding(|scratch| -> Result<_> {
Ok((decode_chunk(&mut reader, scratch)?, decode_chunk(&mut reader, scratch)?))
})?;
if values.len() != lengths.len() {
return Err(Error::internal("an RLE chunk has more runs than run lengths"));
}
let mut rows = 0_usize;
for (value, length) in values.into_iter().zip(lengths) {
let length = usize::try_from(length)
.map_err(|_| Error::internal("a negative RLE run length"))?;
rows = rows
.checked_add(length)
.filter(|&rows| rows <= count)
.ok_or_else(|| Error::internal("an RLE run ends past its chunk"))?;
if length != 0 {
emit(value, length as u64)?;
}
}
check_count(rows, count)?;
}
_ => {
reader = Reader::new(bytes);
let values = with_decoding(|scratch| decode_chunk(&mut reader, scratch))?;
check_count(values.len(), count)?;
for value in values {
emit(value, 1)?;
}
}
}
if reader.remaining() != 0 {
return Err(Error::internal(format!(
"{} bytes left over after counting a chunk",
reader.remaining()
)));
}
Ok(count)
}
const SPARSE: usize = 32;
pub fn decode_selected(bytes: &[u8], positions: &[usize]) -> Result<Vec<i64>> {
if positions.windows(2).any(|pair| pair[0] >= pair[1]) {
return Err(Error::internal("selected integer positions are not sorted and unique"));
}
let mut reader = Reader::new(bytes);
let values = with_decoding(|scratch| decode_selected_chunk(&mut reader, positions, scratch))?;
if reader.remaining() != 0 {
return Err(Error::internal(format!(
"{} bytes left over after decoding selected values",
reader.remaining()
)));
}
Ok(values)
}
#[must_use]
pub fn pointed(bytes: &[u8]) -> bool {
let simple = |bytes: &[u8]| {
bytes
.first()
.and_then(|&tag| Kind::from_tag(tag).ok())
.is_some_and(|kind| matches!(kind, Kind::Constant | Kind::Packed))
};
match bytes.first().and_then(|&tag| Kind::from_tag(tag).ok()) {
Some(Kind::Constant | Kind::Packed) => true,
Some(Kind::Strided) => bytes.get(1 + 4 + 8 + 8..).is_some_and(simple),
Some(Kind::Dict) => {
let mut reader = Reader::new(bytes.get(1 + 4..).unwrap_or_default());
skip_chunk(&mut reader).is_ok() && simple(reader.rest())
}
_ => false,
}
}
pub fn decode_prefix(bytes: &[u8]) -> Result<(Vec<i64>, usize)> {
let mut reader = Reader::new(bytes);
let values = with_decoding(|scratch| decode_chunk(&mut reader, scratch))?;
Ok((values, reader.used()))
}
pub fn describe_prefix(bytes: &[u8]) -> Result<(String, usize)> {
let mut reader = Reader::new(bytes);
let text = describe_chunk(&mut reader)?;
Ok((text, reader.used()))
}
pub fn candidate_sizes(values: &[i64]) -> Result<Vec<(Kind, usize)>> {
let mut sizes = Vec::new();
for kind in candidates(values, 0, &EXHAUSTIVE) {
if let Some(bytes) = encode_as(kind, values, 0, &EXHAUSTIVE)? {
sizes.push((kind, bytes.len()));
}
}
Ok(sizes)
}
#[must_use]
pub fn offered(values: &[i64]) -> Vec<Kind> {
candidates(values, 0, &EXHAUSTIVE)
}
pub fn encode_only(kind: Kind, values: &[i64]) -> Result<Option<Vec<u8>>> {
encode_as(kind, values, 0, &EXHAUSTIVE)
}
pub(crate) fn size_as(kind: Kind, values: &[i64], depth: u8) -> Result<Option<usize>> {
Ok(encode_as(kind, values, depth, &EXHAUSTIVE)?.map(|bytes| bytes.len()))
}
pub fn shape(bytes: &[u8]) -> Result<Vec<Kind>> {
let mut reader = Reader::new(bytes);
let mut kinds = Vec::new();
shape_chunk(&mut reader, &mut kinds)?;
Ok(kinds)
}
pub fn describe(bytes: &[u8]) -> Result<String> {
let mut reader = Reader::new(bytes);
describe_chunk(&mut reader)
}
fn encode_at(values: &[i64], depth: u8, chooser: &dyn Chooser) -> Result<Vec<u8>> {
let started = Instant::now();
let offered = candidates(values, depth, chooser);
let narrowed = chooser.narrow_integers(values, &offered, depth);
let counted = depth == 0;
if counted {
tally::chose(Family::Integer, started);
}
let mut best: Option<(Kind, Vec<u8>)> = None;
for kind in narrowed {
let encoded = if counted {
tally::offer(Family::Integer, kind.tag(), || encode_as(kind, values, depth, chooser))?
} else {
encode_as(kind, values, depth, chooser)?
};
let Some(bytes) = encoded else {
continue;
};
if best.as_ref().is_none_or(|(_, current)| bytes.len() < current.len()) {
best = Some((kind, bytes));
}
}
let (kind, bytes) = best.ok_or_else(|| Error::internal("no encoding applied to the chunk"))?;
if counted {
tally::kept(Family::Integer, kind.tag());
}
Ok(bytes)
}
fn candidates(values: &[i64], depth: u8, chooser: &dyn Chooser) -> Vec<Kind> {
let mut kinds = vec![Kind::Packed];
if depth >= MAX_DEPTH {
return kinds;
}
let Some(profile) = Profile::of(values) else {
return kinds;
};
if profile.runs == 1 {
return vec![Kind::Constant];
}
let considered = |kind| chooser.considers_integer(kind, depth);
let fits = profile.max.checked_sub(profile.min).is_some();
if considered(Kind::Delta)
&& (if fits { profile.deltas_pay(values) } else { deltas_fit(values) })
{
kinds.push(Kind::Delta);
}
if considered(Kind::Rle) && profile.runs * 4 <= values.len() * 3 {
kinds.push(Kind::Rle);
}
if considered(Kind::Dict) && profile.width() > 1 {
let distinct = spread_of(values).0;
if distinct * 2 <= values.len() && width_of(distinct as u64 - 1) < profile.width() {
kinds.push(Kind::Dict);
}
}
if considered(Kind::Sparse)
&& (profile.runs - 1) * 5 <= values.len() * 2
&& majority(values).is_some_and(|(_, count)| count * 10 >= values.len() * 8)
{
kinds.push(Kind::Sparse);
}
if considered(Kind::Strided) && stride_from(values, profile.min).is_some() {
kinds.push(Kind::Strided);
}
kinds
}
struct Profile {
min: i64,
max: i64,
runs: usize,
delta_low: u64,
delta_high: u64,
delta_runs: usize,
}
impl Profile {
fn of(values: &[i64]) -> Option<Self> {
let first = *values.first()?;
let (mut min, mut max, mut breaks) = (first, first, 0usize);
let mut last = values.get(1).map_or(0, |second| second.wrapping_sub(first));
let (mut delta_low, mut delta_high, mut turns) = (u64::MAX, 0u64, 0usize);
for (before, after) in values.iter().zip(&values[1..]) {
min = min.min(*after);
max = max.max(*after);
breaks += usize::from(before != after);
let delta = after.wrapping_sub(*before);
let zigzagged = zigzag(delta);
delta_low = delta_low.min(zigzagged);
delta_high = delta_high.max(zigzagged);
turns += usize::from(delta != last);
last = delta;
}
Some(Self { min, max, runs: breaks + 1, delta_low, delta_high, delta_runs: turns + 1 })
}
fn width(&self) -> u32 {
width_of(self.max.wrapping_sub(self.min) as u64)
}
fn deltas_pay(&self, values: &[i64]) -> bool {
let len = values.len();
let narrower = width_of(self.delta_high.wrapping_sub(self.delta_low)) < self.width();
let unruly = self.runs * 4 > len * 3;
let repeat = self.delta_runs * 4 <= len.saturating_sub(1) * 3;
narrower || (unruly && (repeat || few_deltas(values)))
}
}
fn few_deltas(values: &[i64]) -> bool {
let mut seen = [0i64; FEW_DELTAS];
let mut count = 0;
for pair in values.windows(2).take(FEW_DELTAS_SEEN) {
let delta = pair[1].wrapping_sub(pair[0]);
if seen[..count].contains(&delta) {
continue;
}
if count == FEW_DELTAS {
return false;
}
seen[count] = delta;
count += 1;
}
true
}
const FEW_DELTAS: usize = 4;
const FEW_DELTAS_SEEN: usize = 64;
fn width_of(range: u64) -> u32 {
u64::BITS - range.leading_zeros()
}
fn encode_as(
kind: Kind,
values: &[i64],
depth: u8,
chooser: &dyn Chooser,
) -> Result<Option<Vec<u8>>> {
let mut out = Vec::new();
put_u8(&mut out, kind.tag());
put_u32(&mut out, u32::try_from(values.len()).map_err(|_| too_long(values.len()))?);
match kind {
Kind::Constant => {
let Some(first) = values.first() else {
return Ok(None);
};
if values.iter().any(|value| value != first) {
return Ok(None);
}
put_i64(&mut out, *first);
}
Kind::Packed => encode_packed(values, &mut out)?,
Kind::Delta => {
let (Some(first), Some(deltas)) = (values.first(), deltas(values)) else {
return Ok(None);
};
put_i64(&mut out, *first);
out.extend_from_slice(&encode_at(&deltas, depth + 1, chooser)?);
}
Kind::Rle => {
let (run_values, run_lengths) = runs(values);
if run_values.is_empty() {
return Ok(None);
}
out.extend_from_slice(&encode_at(&run_values, depth + 1, chooser)?);
out.extend_from_slice(&encode_at(&run_lengths, depth + 1, chooser)?);
}
Kind::Dict => {
let dictionary = distinct_values(values);
if dictionary.is_empty() {
return Ok(None);
}
let codes = codes_over(values, &dictionary);
out.extend_from_slice(&encode_at(&dictionary, depth + 1, chooser)?);
out.extend_from_slice(&encode_at(&codes, depth + 1, chooser)?);
}
Kind::Sparse => {
let Some((value, _)) = majority(values).or_else(|| spread_of(values).1) else {
return Ok(None);
};
let mut positions = Vec::new();
let mut exceptions = Vec::new();
for (index, other) in values.iter().enumerate() {
if *other != value {
positions.push(index as i64);
exceptions.push(*other);
}
}
put_i64(&mut out, value);
put_u32(
&mut out,
u32::try_from(positions.len()).map_err(|_| too_long(positions.len()))?,
);
out.extend_from_slice(&encode_at(&positions, depth + 1, chooser)?);
out.extend_from_slice(&encode_at(&exceptions, depth + 1, chooser)?);
}
Kind::Strided => {
let (Some(base), Some(stride)) = (values.iter().min().copied(), stride_of(values))
else {
return Ok(None);
};
let mut steps = Vec::with_capacity(values.len());
for value in values {
let step = offset_from(*value, base) / stride;
let Ok(step) = i64::try_from(step) else {
return Ok(None);
};
steps.push(step);
}
put_i64(&mut out, base);
put_u64(&mut out, stride);
out.extend_from_slice(&encode_at(&steps, depth + 1, chooser)?);
}
}
Ok(Some(out))
}
fn encode_packed(values: &[i64], out: &mut Vec<u8>) -> Result<()> {
let mut offsets: Vec<u64> = Vec::with_capacity(VALUES);
let mut packed: Vec<u64> = vec![0; bitpack::packed_len::<u64>(64)];
let mut transposed = bitpack::Scratch::<u64>::new();
for unit in values.chunks(VALUES) {
let base = unit.iter().copied().min().unwrap_or(0);
offsets.clear();
offsets.extend(unit.iter().map(|value| offset_from(*value, base)));
let width = bitpack::required_width(&offsets);
put_i64(out, base);
put_u8(out, u8::try_from(width).map_err(|_| Error::internal("impossible width"))?);
if unit.len() == VALUES {
let words = bitpack::packed_len::<u64>(width);
bitpack::pack_with(&offsets, width, &mut packed[..words], &mut transposed)?;
for word in &packed[..words] {
put_u64(out, *word);
}
} else {
bitpack::pack_tail(&offsets, width, out)?;
}
}
Ok(())
}
struct Decoding {
packed: Vec<u64>,
}
thread_local! {
static DECODING: std::cell::RefCell<Decoding> =
const { std::cell::RefCell::new(Decoding::new()) };
}
fn with_decoding<T>(run: impl FnOnce(&mut Decoding) -> T) -> T {
DECODING.with(|cell| match cell.try_borrow_mut() {
Ok(mut scratch) => run(&mut scratch),
Err(_) => run(&mut Decoding::new()),
})
}
impl Decoding {
const fn new() -> Self {
Self { packed: Vec::new() }
}
fn ready(&mut self) {
if self.packed.len() != bitpack::packed_len::<u64>(64) {
self.packed.resize(bitpack::packed_len::<u64>(64), 0);
}
}
}
fn decode_chunk(reader: &mut Reader<'_>, scratch: &mut Decoding) -> Result<Vec<i64>> {
let kind = Kind::from_tag(reader.u8()?)?;
let count = reader.u32()? as usize;
match kind {
Kind::Constant => Ok(vec![reader.i64()?; count]),
Kind::Packed => {
let mut values = vec![0i64; count];
scratch.ready();
let mut done = 0;
while done < count {
let base = reader.i64()?;
let width = reader.u8()? as usize;
let wanted = (count - done).min(VALUES);
let into = &mut values[done..done + wanted];
if wanted == VALUES {
let words = bitpack::packed_len::<u64>(width);
for word in &mut scratch.packed[..words] {
*word = reader.u64()?;
}
bitpack::unpack_mapped(&scratch.packed[..words], width, into, |offset| {
value_from(offset, base)
})?;
} else {
let bytes = reader.bytes(bitpack::tail_len(wanted, width))?;
bitpack::unpack_tail_into(bytes, width, into, |offset| {
value_from(offset, base)
})?;
}
done += wanted;
}
Ok(values)
}
Kind::Delta => {
let first = reader.i64()?;
let deltas = decode_chunk(reader, scratch)?;
let mut values = Vec::with_capacity(count);
values.push(first);
let mut current = first;
for delta in deltas {
current = current.wrapping_add(unzigzag(delta as u64));
values.push(current);
}
check_count(values.len(), count)?;
Ok(values)
}
Kind::Rle => {
let run_values = decode_chunk(reader, scratch)?;
let run_lengths = decode_chunk(reader, scratch)?;
if run_values.len() != run_lengths.len() {
return Err(Error::internal("an RLE chunk has more runs than run lengths"));
}
let long = run_values.len().saturating_mul(LONG_RUN) <= count;
let mut values = if long { Vec::with_capacity(count) } else { vec![0; count + RUN] };
let mut at = 0usize;
for (value, length) in run_values.into_iter().zip(run_lengths) {
let length = usize::try_from(length)
.map_err(|_| Error::internal("a negative RLE run length"))?;
let end = at
.checked_add(length)
.filter(|end| *end <= count)
.ok_or_else(|| Error::internal("an RLE run ends past its chunk"))?;
if long {
values.resize(end, value);
} else {
let short =
if length <= RUN { values[at..].first_chunk_mut::<RUN>() } else { None };
match short {
Some(window) => window.fill(value),
None => values[at..end].fill(value),
}
}
at = end;
}
check_count(at, count)?;
values.truncate(count);
Ok(values)
}
Kind::Dict => {
let dictionary = decode_chunk(reader, scratch)?;
let codes = decode_chunk(reader, scratch)?;
let mut values = Vec::with_capacity(count);
for code in codes {
let index =
usize::try_from(code).ok().and_then(|index| dictionary.get(index)).ok_or_else(
|| Error::internal(format!("code {code} is not in the dictionary")),
)?;
values.push(*index);
}
check_count(values.len(), count)?;
Ok(values)
}
Kind::Sparse => {
let value = reader.i64()?;
let exception_count = reader.u32()? as usize;
let positions = decode_chunk(reader, scratch)?;
let exceptions = decode_chunk(reader, scratch)?;
if positions.len() != exception_count || exceptions.len() != exception_count {
return Err(Error::internal("a sparse chunk disagrees about its exception count"));
}
let mut values = vec![value; count];
for (position, exception) in positions.into_iter().zip(exceptions) {
let position = usize::try_from(position)
.ok()
.filter(|position| *position < count)
.ok_or_else(|| {
Error::internal(format!("exception at {position} is outside the chunk"))
})?;
values[position] = exception;
}
Ok(values)
}
Kind::Strided => {
let base = reader.i64()?;
let stride = reader.u64()?;
let steps = decode_chunk(reader, scratch)?;
check_count(steps.len(), count)?;
let mut values = Vec::with_capacity(count);
for step in steps {
let step = u64::try_from(step)
.map_err(|_| Error::internal("a negative number of strides"))?;
values.push(value_from(step.wrapping_mul(stride), base));
}
Ok(values)
}
}
}
fn decode_chunk_as<T: Lane>(reader: &mut Reader<'_>, scratch: &mut Decoding) -> Result<Vec<T>> {
let Some(&tag) = reader.rest().first() else {
return Err(Error::internal("a chunk ended before its encoding tag"));
};
if !matches!(Kind::from_tag(tag)?, Kind::Constant | Kind::Packed | Kind::Sparse | Kind::Rle) {
return decode_chunk(reader, scratch)?.into_iter().map(lane).collect();
}
let kind = Kind::from_tag(reader.u8()?)?;
let count = reader.u32()? as usize;
match kind {
Kind::Constant => Ok(vec![lane(reader.i64()?)?; count]),
Kind::Packed => {
let mut values = vec![T::default(); count];
scratch.ready();
let mut wide = [0i64; VALUES];
let mut done = 0;
while done < count {
let base = reader.i64()?;
let width = reader.u8()? as usize;
let wanted = (count - done).min(VALUES);
let into = &mut values[done..done + wanted];
let bytes = if wanted == VALUES {
let words = bitpack::packed_len::<u64>(width);
for word in &mut scratch.packed[..words] {
*word = reader.u64()?;
}
None
} else {
Some(reader.bytes(bitpack::tail_len(wanted, width))?)
};
let mask = if width >= 64 { u64::MAX } else { (1u64 << width) - 1 };
let top = i64::try_from(i128::from(base) + i128::from(mask)).ok();
if T::fit(base).is_some() && top.and_then(T::fit).is_some() {
let map = |offset| T::wrap(value_from(offset, base));
match bytes {
None => {
let words = bitpack::packed_len::<u64>(width);
bitpack::unpack_mapped(&scratch.packed[..words], width, into, map)?;
}
Some(bytes) => bitpack::unpack_tail_into(bytes, width, into, map)?,
}
} else {
let wide = &mut wide[..wanted];
let map = |offset| value_from(offset, base);
match bytes {
None => {
let words = bitpack::packed_len::<u64>(width);
bitpack::unpack_mapped(&scratch.packed[..words], width, wide, map)?;
}
Some(bytes) => bitpack::unpack_tail_into(bytes, width, wide, map)?,
}
for (value, &held) in into.iter_mut().zip(wide.iter()) {
*value = lane(held)?;
}
}
done += wanted;
}
Ok(values)
}
Kind::Rle => {
let run_values = decode_chunk(reader, scratch)?;
let run_lengths = decode_chunk(reader, scratch)?;
if run_values.len() != run_lengths.len() {
return Err(Error::internal("an RLE chunk has more runs than run lengths"));
}
let long = run_values.len().saturating_mul(LONG_RUN) <= count;
let mut values =
if long { Vec::with_capacity(count) } else { vec![T::default(); count + RUN] };
let mut at = 0usize;
for (value, length) in run_values.into_iter().zip(run_lengths) {
let value = lane::<T>(value)?;
let length = usize::try_from(length)
.map_err(|_| Error::internal("a negative RLE run length"))?;
let end = at
.checked_add(length)
.filter(|end| *end <= count)
.ok_or_else(|| Error::internal("an RLE run ends past its chunk"))?;
if long {
values.resize(end, value);
} else {
let short =
if length <= RUN { values[at..].first_chunk_mut::<RUN>() } else { None };
match short {
Some(window) => window.fill(value),
None => values[at..end].fill(value),
}
}
at = end;
}
check_count(at, count)?;
values.truncate(count);
Ok(values)
}
Kind::Sparse => {
let value = lane::<T>(reader.i64()?)?;
let exception_count = reader.u32()? as usize;
let positions = decode_chunk(reader, scratch)?;
let exceptions = decode_chunk(reader, scratch)?;
if positions.len() != exception_count || exceptions.len() != exception_count {
return Err(Error::internal("a sparse chunk disagrees about its exception count"));
}
let mut values = vec![value; count];
for (position, exception) in positions.into_iter().zip(exceptions) {
let position = usize::try_from(position)
.ok()
.filter(|position| *position < count)
.ok_or_else(|| {
Error::internal(format!("exception at {position} is outside the chunk"))
})?;
values[position] = lane(exception)?;
}
Ok(values)
}
Kind::Delta | Kind::Dict | Kind::Strided => {
Err(Error::internal("a chunk kind that decodes wide reached the narrow decoder"))
}
}
}
fn decode_selected_chunk(
reader: &mut Reader<'_>,
positions: &[usize],
scratch: &mut Decoding,
) -> Result<Vec<i64>> {
let Some(&tag) = reader.rest().first() else {
return Err(Error::internal("a chunk ended before its encoding tag"));
};
let kind = Kind::from_tag(tag)?;
if !matches!(kind, Kind::Constant | Kind::Packed | Kind::Rle | Kind::Strided | Kind::Dict) {
let values = decode_chunk(reader, scratch)?;
return positions
.iter()
.map(|&position| {
values.get(position).copied().ok_or_else(|| {
Error::internal(format!(
"selected integer position {position} is outside {} values",
values.len()
))
})
})
.collect();
}
let decoded = Kind::from_tag(reader.u8()?)?;
debug_assert_eq!(decoded, kind);
let count = reader.u32()? as usize;
if positions.last().is_some_and(|&position| position >= count) {
return Err(Error::internal(format!(
"selected integer position {} is outside {count} values",
positions.last().expect("a last position exists")
)));
}
match kind {
Kind::Constant => {
let value = reader.i64()?;
Ok(vec![value; positions.len()])
}
Kind::Packed => {
let mut out = Vec::with_capacity(positions.len());
let mut from = 0;
let mut done = 0;
while done < count {
let base = reader.i64()?;
let width = reader.u8()? as usize;
let wanted = (count - done).min(VALUES);
let upto = positions.partition_point(|&position| position < done + wanted);
if wanted == VALUES && (upto - from) * SPARSE > VALUES {
let words = bitpack::packed_len::<u64>(width);
scratch.ready();
for word in &mut scratch.packed[..words] {
*word = reader.u64()?;
}
let mut unit = [0_i64; VALUES];
bitpack::unpack_mapped(&scratch.packed[..words], width, &mut unit, |offset| {
value_from(offset, base)
})?;
out.extend(positions[from..upto].iter().map(|&position| unit[position - done]));
} else if wanted == VALUES {
let bytes = reader.bytes(bitpack::packed_len::<u64>(width) * 8)?;
for &position in &positions[from..upto] {
let offset = bitpack::unpack_u64_at(bytes, width, position - done)?;
out.push(value_from(offset, base));
}
} else {
let bytes = reader.bytes(bitpack::tail_len(wanted, width))?;
for &position in &positions[from..upto] {
let offset = bitpack::tail_at(bytes, width, position - done)?;
out.push(value_from(offset, base));
}
}
from = upto;
done += wanted;
}
Ok(out)
}
Kind::Rle => {
let run_value_bytes = reader.rest();
let mut run_value_reader = Reader::new(run_value_bytes);
let run_value_count = skip_chunk(&mut run_value_reader)?;
let run_value_len = run_value_reader.used();
reader.skip(run_value_len)?;
let run_lengths = decode_chunk(reader, scratch)?;
if run_value_count != run_lengths.len() {
return Err(Error::internal("an RLE chunk has more runs than run lengths"));
}
let mut wanted_runs = Vec::new();
let mut selected_per_run = Vec::new();
let mut selected = 0;
let mut at = 0usize;
for (run, length) in run_lengths.into_iter().enumerate() {
let length = usize::try_from(length)
.map_err(|_| Error::internal("a negative RLE run length"))?;
let end = at
.checked_add(length)
.filter(|end| *end <= count)
.ok_or_else(|| Error::internal("an RLE run ends past its chunk"))?;
let before = selected;
while selected < positions.len() && positions[selected] < end {
if positions[selected] < at {
return Err(Error::internal("selected integer positions went backwards"));
}
selected += 1;
}
if selected != before {
wanted_runs.push(run);
selected_per_run.push(selected - before);
}
at = end;
}
check_count(at, count)?;
if selected != positions.len() {
return Err(Error::internal("an RLE chunk ended before a selected position"));
}
let run_values = decode_selected(&run_value_bytes[..run_value_len], &wanted_runs)?;
let mut out = Vec::with_capacity(positions.len());
for (value, repeat) in run_values.into_iter().zip(selected_per_run) {
out.extend(std::iter::repeat_n(value, repeat));
}
Ok(out)
}
Kind::Strided => {
let base = reader.i64()?;
let stride = reader.u64()?;
let steps = decode_selected_chunk(reader, positions, scratch)?;
steps
.into_iter()
.map(|step| {
let step = u64::try_from(step)
.map_err(|_| Error::internal("a negative number of strides"))?;
Ok(value_from(step.wrapping_mul(stride), base))
})
.collect()
}
Kind::Dict => {
let dictionary = decode_chunk(reader, scratch)?;
let codes = decode_selected_chunk(reader, positions, scratch)?;
codes
.into_iter()
.map(|code| {
usize::try_from(code)
.ok()
.and_then(|index| dictionary.get(index))
.copied()
.ok_or_else(|| {
Error::internal(format!("code {code} is not in the dictionary"))
})
})
.collect()
}
_ => unreachable!("unsupported kinds used the full decoder"),
}
}
fn skip_chunk(reader: &mut Reader<'_>) -> Result<usize> {
let kind = Kind::from_tag(reader.u8()?)?;
let count = reader.u32()? as usize;
match kind {
Kind::Constant => reader.skip(8)?,
Kind::Packed => skip_packed(reader, count)?,
Kind::Delta => {
reader.skip(8)?;
skip_chunk(reader)?;
}
Kind::Rle | Kind::Dict => {
skip_chunk(reader)?;
skip_chunk(reader)?;
}
Kind::Sparse => {
reader.skip(12)?;
skip_chunk(reader)?;
skip_chunk(reader)?;
}
Kind::Strided => {
reader.skip(16)?;
skip_chunk(reader)?;
}
}
Ok(count)
}
fn skip_packed(reader: &mut Reader<'_>, count: usize) -> Result<()> {
let mut done = 0;
while done < count {
reader.skip(8)?;
let width = reader.u8()? as usize;
if width > 64 {
return Err(Error::internal(format!("a packed integer width of {width} is past 64")));
}
let wanted = (count - done).min(VALUES);
let bytes = if wanted == VALUES {
bitpack::packed_len::<u64>(width)
.checked_mul(8)
.ok_or_else(|| Error::internal("packed integer size overflow"))?
} else {
bitpack::tail_len(wanted, width)
};
reader.skip(bytes)?;
done += wanted;
}
Ok(())
}
fn shape_chunk(reader: &mut Reader<'_>, kinds: &mut Vec<Kind>) -> Result<()> {
let kind = Kind::from_tag(reader.u8()?)?;
let count = reader.u32()? as usize;
kinds.push(kind);
match kind {
Kind::Constant => reader.skip(8)?,
Kind::Packed => skip_packed(reader, count)?,
Kind::Delta => {
reader.skip(8)?;
shape_chunk(reader, kinds)?;
}
Kind::Rle | Kind::Dict => {
shape_chunk(reader, kinds)?;
shape_chunk(reader, kinds)?;
}
Kind::Sparse => {
reader.skip(12)?;
shape_chunk(reader, kinds)?;
shape_chunk(reader, kinds)?;
}
Kind::Strided => {
reader.skip(16)?;
shape_chunk(reader, kinds)?;
}
}
Ok(())
}
fn describe_chunk(reader: &mut Reader<'_>) -> Result<String> {
let kind = Kind::from_tag(reader.u8()?)?;
let count = reader.u32()? as usize;
Ok(match kind {
Kind::Constant => {
reader.i64()?;
"CONSTANT".to_string()
}
Kind::Packed => {
let mut widths = Vec::new();
let mut seen = 0;
while seen < count {
reader.i64()?;
let width = reader.u8()? as usize;
let wanted = (count - seen).min(VALUES);
if wanted == VALUES {
for _ in 0..bitpack::packed_len::<u64>(width) {
reader.u64()?;
}
} else {
reader.bytes(bitpack::tail_len(wanted, width))?;
}
widths.push(width);
seen += wanted;
}
let low = widths.iter().copied().min().unwrap_or(0);
let high = widths.iter().copied().max().unwrap_or(0);
if low == high {
format!("FOR+BITPACK[{low}]")
} else {
format!("FOR+BITPACK[{low}..{high}]")
}
}
Kind::Delta => {
reader.i64()?;
format!("DELTA({})", describe_chunk(reader)?)
}
Kind::Rle => {
let values = describe_chunk(reader)?;
let lengths = describe_chunk(reader)?;
format!("RLE({values}, {lengths})")
}
Kind::Dict => {
let dictionary = describe_chunk(reader)?;
let codes = describe_chunk(reader)?;
format!("DICT({dictionary}, {codes})")
}
Kind::Sparse => {
reader.i64()?;
reader.u32()?;
let positions = describe_chunk(reader)?;
let exceptions = describe_chunk(reader)?;
format!("SPARSE({positions}, {exceptions})")
}
Kind::Strided => {
reader.i64()?;
let stride = reader.u64()?;
format!("STRIDE[{stride}]({})", describe_chunk(reader)?)
}
})
}
fn stride_of(values: &[i64]) -> Option<u64> {
stride_from(values, values.iter().min().copied()?)
}
fn stride_from(values: &[i64], base: i64) -> Option<u64> {
let mut divisor = 0u64;
for value in values {
divisor = gcd(divisor, offset_from(*value, base));
if divisor == 1 {
return None;
}
}
(divisor > 1).then_some(divisor)
}
fn gcd(mut left: u64, mut right: u64) -> u64 {
if left == 0 {
return right;
}
if right == 0 {
return left;
}
let shift = (left | right).trailing_zeros();
left >>= left.trailing_zeros();
loop {
right >>= right.trailing_zeros();
if left > right {
std::mem::swap(&mut left, &mut right);
}
right -= left;
if right == 0 {
return left << shift;
}
}
}
fn offset_from(value: i64, base: i64) -> u64 {
(i128::from(value) - i128::from(base)) as u64
}
fn value_from(offset: u64, base: i64) -> i64 {
(i128::from(base) + i128::from(offset)) as i64
}
fn zigzag(value: i64) -> u64 {
((value << 1) ^ (value >> 63)) as u64
}
fn unzigzag(value: u64) -> i64 {
((value >> 1) as i64) ^ -((value & 1) as i64)
}
fn deltas_fit(values: &[i64]) -> bool {
values.windows(2).all(|pair| pair[1].checked_sub(pair[0]).is_some())
}
fn deltas(values: &[i64]) -> Option<Vec<i64>> {
let mut deltas = Vec::with_capacity(values.len().saturating_sub(1));
for pair in values.windows(2) {
let difference = pair[1].checked_sub(pair[0])?;
deltas.push(zigzag(difference) as i64);
}
Some(deltas)
}
fn runs(values: &[i64]) -> (Vec<i64>, Vec<i64>) {
let mut run_values: Vec<i64> = Vec::new();
let mut run_lengths: Vec<i64> = Vec::new();
let mut start = 0;
while let Some(&value) = values.get(start) {
let length = values[start..].iter().take_while(|other| **other == value).count();
run_values.push(value);
run_lengths.push(length as i64);
start += length;
}
(run_values, run_lengths)
}
fn spread_of(values: &[i64]) -> (usize, Option<(i64, usize)>) {
let mut sorted = values.to_vec();
sorted.sort_unstable();
let mut distinct = 0;
let mut best: Option<(i64, usize)> = None;
let mut index = 0;
while index < sorted.len() {
let value = sorted[index];
let mut end = index;
while end < sorted.len() && sorted[end] == value {
end += 1;
}
distinct += 1;
let count = end - index;
if best.is_none_or(|(_, seen)| count > seen) {
best = Some((value, count));
}
index = end;
}
(distinct, best)
}
fn majority(values: &[i64]) -> Option<(i64, usize)> {
let mut candidate = *values.first()?;
let mut lead = 0usize;
for value in values {
if lead == 0 {
candidate = *value;
lead = 1;
} else if *value == candidate {
lead += 1;
} else {
lead -= 1;
}
}
let count = values.iter().filter(|value| **value == candidate).count();
(count * 2 > values.len()).then_some((candidate, count))
}
fn distinct_values(values: &[i64]) -> Vec<i64> {
let mut distinct = values.to_vec();
distinct.sort_unstable();
distinct.dedup();
distinct
}
fn codes_over(values: &[i64], dictionary: &[i64]) -> Vec<i64> {
values
.iter()
.map(|value| {
dictionary
.binary_search(value)
.expect("the dictionary is the distinct values of this chunk") as i64
})
.collect()
}
fn check_count(actual: usize, expected: usize) -> Result<()> {
if actual == expected {
Ok(())
} else {
Err(Error::internal(format!(
"a chunk says it holds {expected} values and decoded to {actual}"
)))
}
}
fn too_long(len: usize) -> Error {
Error::internal(format!("a chunk of {len} values is longer than the format allows"))
}
fn put_u8(out: &mut Vec<u8>, value: u8) {
out.push(value);
}
fn put_u32(out: &mut Vec<u8>, value: u32) {
out.extend_from_slice(&value.to_le_bytes());
}
fn put_u64(out: &mut Vec<u8>, value: u64) {
out.extend_from_slice(&value.to_le_bytes());
}
fn put_i64(out: &mut Vec<u8>, value: i64) {
out.extend_from_slice(&value.to_le_bytes());
}
#[cfg(test)]
mod tests {
use super::*;
fn round_trip(values: &[i64]) -> Vec<u8> {
let bytes = encode(values).unwrap();
assert_eq!(decode(&bytes).unwrap(), values, "{}", describe(&bytes).unwrap());
bytes
}
fn kind_of(bytes: &[u8]) -> Kind {
Kind::from_tag(bytes[0]).unwrap()
}
struct Random(u64);
impl Random {
fn new() -> Self {
Self(0x9e37_79b9_7f4a_7c15)
}
fn next(&mut self) -> u64 {
self.0 ^= self.0 << 13;
self.0 ^= self.0 >> 7;
self.0 ^= self.0 << 17;
self.0
}
}
#[test]
fn the_dictionary_is_sorted_and_the_codes_point_back_at_the_values() {
let values = vec![30i64, 10, 30, 20, 10, -5];
let dictionary = distinct_values(&values);
let codes = codes_over(&values, &dictionary);
assert_eq!(dictionary, vec![-5, 10, 20, 30]);
assert_eq!(codes, vec![3, 1, 3, 2, 1, 0]);
for (code, value) in codes.iter().zip(&values) {
assert_eq!(dictionary[*code as usize], *value);
}
}
#[test]
fn one_sort_gives_the_distinct_count_and_the_most_frequent_value() {
let values = vec![7i64, 7, 7, 1, 2, 2];
assert_eq!(spread_of(&values), (3, Some((7, 3))));
assert_eq!(spread_of(&[]), (0, None));
assert_eq!(spread_of(&[9]), (1, Some((9, 1))));
assert_eq!(spread_of(&[4i64, 4, 8, 8]), (2, Some((4, 2))));
}
#[test]
fn the_majority_is_the_most_frequent_value_whenever_there_is_one() {
let chunks: Vec<Vec<i64>> = vec![
vec![],
vec![3],
vec![1, 2],
vec![1, 1, 2],
vec![2, 1, 1],
vec![4, 4, 8, 8],
vec![7, 1, 7, 2, 7, 3, 7],
vec![1, 2, 3, 9, 9, 9, 9],
(0..1000).map(|index| if index % 5 == 0 { index } else { -4 }).collect(),
(0..1000).map(|index| index % 3).collect(),
];
for chunk in chunks {
let (_, dominant) = spread_of(&chunk);
let expected = dominant.filter(|(_, count)| count * 2 > chunk.len());
assert_eq!(majority(&chunk), expected, "{chunk:?}");
}
}
#[test]
fn the_one_pass_offers_what_the_separate_tests_offered() {
let mut random = Random::new();
let mut chunks: Vec<Vec<i64>> = vec![
vec![],
vec![5],
vec![5, 5, 5],
vec![i64::MIN, i64::MAX],
vec![i64::MAX, i64::MIN, i64::MAX],
vec![i64::MIN, 0, i64::MAX],
vec![-1, i64::MAX],
(0..1000).map(|index| if index % 5 == 0 { index } else { -4 }).collect(),
(0..1000).map(|index| if index % 4 == 0 { index } else { -4 }).collect(),
(0..1000).map(|index| index / 7).collect(),
(0..1000).map(|index| index * 1_000_000).collect(),
];
for _ in 0..200 {
let len = (random.next() % 300) as usize;
let spread = 1 + random.next() % 8;
let common = (random.next() % 5) as i64;
chunks.push(
(0..len)
.map(|_| {
let draw = random.next();
if draw % 10 < spread { (draw >> 8) as i64 % 50 } else { common }
})
.collect(),
);
}
for chunk in chunks {
let mut expected = vec![Kind::Packed];
if !chunk.is_empty() {
if chunk.iter().all(|value| *value == chunk[0]) {
expected = vec![Kind::Constant];
} else {
let runs = 1 + chunk.windows(2).filter(|pair| pair[0] != pair[1]).count();
let low = *chunk.iter().min().unwrap();
let high = *chunk.iter().max().unwrap();
let bits = |range: u128| 128 - range.leading_zeros();
let width = bits((i128::from(high) - i128::from(low)) as u128);
let zigzags: Vec<u128> = chunk
.windows(2)
.map(|pair| u128::from(zigzag(pair[1].wrapping_sub(pair[0]))))
.collect();
let spread = zigzags.iter().max().unwrap() - zigzags.iter().min().unwrap();
let turns = 1 + zigzags.windows(2).filter(|pair| pair[0] != pair[1]).count();
let mut first: Vec<u128> = zigzags.iter().take(64).copied().collect();
first.sort_unstable();
first.dedup();
let pays = bits(spread) < width
|| (runs * 4 > chunk.len() * 3
&& (turns * 4 <= (chunk.len() - 1) * 3 || first.len() <= 4));
if deltas_fit(&chunk) && pays {
expected.push(Kind::Delta);
}
if runs * 4 <= chunk.len() * 3 {
expected.push(Kind::Rle);
}
let distinct = spread_of(&chunk).0;
if distinct * 2 <= chunk.len() && bits(distinct as u128 - 1) < width {
expected.push(Kind::Dict);
}
if majority(&chunk).is_some_and(|(_, count)| count * 10 >= chunk.len() * 8) {
expected.push(Kind::Sparse);
}
if stride_of(&chunk).is_some() {
expected.push(Kind::Strided);
}
}
}
assert_eq!(candidates(&chunk, 0, &EXHAUSTIVE), expected, "{chunk:?}");
}
}
#[test]
fn runs_are_every_stretch_of_equal_neighbours_in_order() {
assert_eq!(runs(&[]), (vec![], vec![]));
assert_eq!(runs(&[4]), (vec![4], vec![1]));
assert_eq!(runs(&[1, 1, 2, 1, 1, 1]), (vec![1, 2, 1], vec![2, 1, 3]));
}
#[test]
fn deltas_that_do_not_fit_are_refused_before_they_are_built() {
assert!(deltas_fit(&[1i64, 2, 3]));
assert!(deltas_fit(&[i64::MAX, i64::MAX]));
assert!(!deltas_fit(&[i64::MIN, i64::MAX]));
assert_eq!(deltas_fit(&[i64::MIN, i64::MAX]), deltas(&[i64::MIN, i64::MAX]).is_some());
assert_eq!(deltas_fit(&[1i64, 2, 3]), deltas(&[1i64, 2, 3]).is_some());
}
#[test]
fn what_the_chooser_returns_is_the_smallest_of_what_it_was_offered() {
let mut random = Random::new();
let noise: Vec<i64> = (0..2000).map(|_| (random.next() % 5000) as i64).collect();
let runs: Vec<i64> = (0..2000).map(|index: i64| index / 100).collect();
let climbing: Vec<i64> = (0..2000).map(|index| 1_700_000_000 + index).collect();
for values in [noise, runs, climbing, vec![7; 300], Vec::new()] {
let chosen = encode(&values).unwrap();
let mut smallest: Option<Vec<u8>> = None;
for kind in offered(&values) {
let Some(bytes) = encode_only(kind, &values).unwrap() else {
continue;
};
if smallest.as_ref().is_none_or(|best| bytes.len() < best.len()) {
smallest = Some(bytes);
}
}
assert_eq!(smallest.as_deref(), Some(chosen.as_slice()), "{}", values.len());
}
}
#[test]
fn a_column_of_whole_seconds_in_microseconds_pays_nothing_for_the_zeroes() {
let mut random = Random::new();
let day = 1_374_000_000_000_000i64;
let values: Vec<i64> =
(0..100_000).map(|_| day + (random.next() % 68_400) as i64 * 1_000_000).collect();
let bytes = round_trip(&values);
assert_eq!(kind_of(&bytes), Kind::Strided);
assert!(describe(&bytes).unwrap().starts_with("STRIDE[1000000]"), "{:?}", describe(&bytes));
let strided = 100_000 * 17 / 8;
assert!(bytes.len() < strided + 2000, "{} bytes for {strided} of payload", bytes.len());
let plain = encode_only(Kind::Packed, &values).unwrap().expect("packing always applies");
assert!(
bytes.len() * 2 < plain.len(),
"{} strided against {} packed",
bytes.len(),
plain.len()
);
}
#[test]
fn a_stride_is_the_common_factor_of_the_distances_from_the_smallest_value() {
assert_eq!(stride_of(&[10i64, 20, 40]), Some(10));
assert_eq!(stride_of(&[7i64, 17, 37]), Some(10));
assert_eq!(stride_of(&[10i64, 20, 23]), None);
assert_eq!(stride_of(&[5i64; 100]), None);
assert_eq!(stride_of(&[]), None);
assert_eq!(stride_of(&[i64::MIN, i64::MAX]), Some(u64::MAX));
}
#[test]
fn a_stride_across_the_whole_of_the_type_round_trips() {
for values in [vec![i64::MIN, i64::MAX], vec![i64::MIN, 0, i64::MAX]] {
let bytes = round_trip(&values);
assert_eq!(decode(&bytes).unwrap(), values);
}
}
#[test]
fn a_column_with_no_common_factor_is_not_offered_a_stride() {
let mut random = Random::new();
let values: Vec<i64> = (0..2000).map(|_| (random.next() % 1_000_000) as i64).collect();
assert!(!offered(&values).contains(&Kind::Strided));
assert!(encode_only(Kind::Strided, &values).unwrap().is_none());
}
#[test]
fn an_empty_chunk_round_trips() {
let bytes = round_trip(&[]);
assert_eq!(bytes.len(), 5);
}
#[test]
fn a_constant_column_costs_thirteen_bytes_however_long_it_is() {
let bytes = round_trip(&vec![42; 1_000_000]);
assert_eq!(kind_of(&bytes), Kind::Constant);
assert_eq!(bytes.len(), 13);
}
#[test]
fn a_narrow_range_is_packed_at_the_width_of_the_range_and_not_of_the_type() {
let mut random = Random::new();
let values: Vec<i64> = (0..100_000).map(|_| 1000 + (random.next() % 64) as i64).collect();
let bytes = round_trip(&values);
assert_eq!(kind_of(&bytes), Kind::Packed);
let packed = 100_000 * 6 / 8;
assert!(bytes.len() < packed + 2000, "{} bytes for {packed} of payload", bytes.len());
assert!(bytes.len() > packed, "{} bytes cannot hold {packed}", bytes.len());
}
#[test]
fn a_counter_becomes_deltas_and_then_a_constant() {
let values: Vec<i64> = (0..1_000_000).collect();
let bytes = round_trip(&values);
assert_eq!(kind_of(&bytes), Kind::Delta);
assert_eq!(describe(&bytes).unwrap(), "DELTA(CONSTANT)");
assert!(bytes.len() < 40, "{} bytes for a counter", bytes.len());
}
#[test]
fn a_column_that_counts_down_is_as_cheap_as_one_that_counts_up() {
let up: Vec<i64> = (0..100_000).collect();
let down: Vec<i64> = (0..100_000).rev().collect();
assert_eq!(round_trip(&up).len(), round_trip(&down).len());
}
#[test]
fn long_runs_become_rle() {
let mut values = Vec::new();
for run in 0..1000 {
values.extend(std::iter::repeat_n(run % 7, 200));
}
let bytes = round_trip(&values);
assert_eq!(kind_of(&bytes), Kind::Rle);
assert!(bytes.len() < 2000, "{} bytes for 1000 runs", bytes.len());
}
#[test]
fn a_low_cardinality_column_becomes_a_dictionary() {
let mut random = Random::new();
let dictionary: Vec<i64> =
(0..40).map(|_| 1_000_000_000 + (random.next() % (1 << 30)) as i64).collect();
let values: Vec<i64> =
(0..100_000).map(|_| dictionary[(random.next() % 40) as usize]).collect();
let bytes = round_trip(&values);
assert_eq!(kind_of(&bytes), Kind::Dict);
assert!(bytes.len() < 100_000, "{} bytes", bytes.len());
}
#[test]
fn a_nearly_constant_column_becomes_sparse() {
let mut values = vec![0i64; 100_000];
for index in 0..300 {
values[index * 331] = 1 << 40;
}
let bytes = round_trip(&values);
assert_eq!(kind_of(&bytes), Kind::Sparse);
assert!(bytes.len() < 3000, "{} bytes for 300 exceptions", bytes.len());
}
#[test]
fn encoded_counts_match_decoded_rows_across_integer_shapes() {
let mut sparse = vec![0_i64; 4096];
for (index, value) in [(7, -3), (91, 12), (1001, -3), (3000, 12)] {
sparse[index] = value;
}
let mut runs = Vec::new();
for value in [0, 7, 0, -5] {
runs.extend(std::iter::repeat_n(value, 500));
}
let mut random = Random::new();
let packed = (0..2000).map(|_| (random.next() % 251) as i64).collect::<Vec<_>>();
for values in [vec![0_i64; 1024], sparse, runs, packed] {
let bytes = encode(&values).unwrap();
let (rows, counts) = tally(&bytes).unwrap();
let mut expected = BTreeMap::<i64, u64>::new();
for value in decode(&bytes).unwrap() {
*expected.entry(value).or_default() += 1;
}
assert_eq!(rows, values.len());
assert_eq!(counts, expected.into_iter().collect::<Vec<_>>());
}
}
#[test]
fn folded_sparse_exceptions_keep_the_last_value_at_a_repeated_position() {
let mut bytes = vec![Kind::Sparse.tag()];
put_u32(&mut bytes, 10);
put_i64(&mut bytes, 0);
put_u32(&mut bytes, 2);
bytes.extend(encode(&[7, 7]).unwrap());
bytes.extend(encode(&[3, 5]).unwrap());
let mut counts = BTreeMap::<i64, u64>::new();
assert_eq!(
fold(&bytes, |value, count| {
*counts.entry(value).or_default() += count;
Ok(())
})
.unwrap(),
10
);
assert_eq!(counts, BTreeMap::from([(0, 9), (5, 1)]));
assert_eq!(decode(&bytes).unwrap()[7], 5);
}
#[test]
fn the_cascade_goes_more_than_one_level_deep() {
let mut values = Vec::new();
for index in 0..2000i64 {
values.extend(std::iter::repeat_n(1_000_000 + (index % 5) * 104_729, 100));
}
let bytes = round_trip(&values);
let shape = describe(&bytes).unwrap();
assert!(shape.contains('('), "{shape} is not a cascade");
assert!(bytes.len() < 4000, "{} bytes: {shape}", bytes.len());
}
#[test]
fn random_data_is_packed_at_full_width_and_costs_what_it_costs() {
let mut random = Random::new();
let values: Vec<i64> = (0..10_000).map(|_| random.next() as i64).collect();
let bytes = round_trip(&values);
assert_eq!(kind_of(&bytes), Kind::Packed);
assert!(bytes.len() < 10_000 * 8 + 1000, "{} bytes", bytes.len());
}
#[test]
fn the_extremes_of_the_type_survive() {
let values = vec![i64::MIN, i64::MAX, 0, -1, i64::MIN, i64::MAX];
round_trip(&values);
round_trip(&[i64::MIN; 3]);
round_trip(&[i64::MIN, i64::MIN + 1]);
}
#[test]
fn a_chunk_that_is_not_a_multiple_of_the_unit_round_trips() {
for len in [1, 2, 1023, 1024, 1025, 2047, 2049] {
let values: Vec<i64> = (0..len).map(|index| (index * 31 % 97) as i64).collect();
round_trip(&values);
}
}
#[test]
fn units_of_different_widths_in_one_chunk_do_not_read_each_others_leftovers() {
let mut random = Random::new();
let mut values = Vec::new();
for width in [40u32, 3, 61, 1, 17, 40] {
for _ in 0..1024 {
values.push((random.next() & ((1u64 << width) - 1)) as i64);
}
}
let bytes = encode(&values).unwrap();
let described = describe(&bytes).unwrap();
assert!(described.starts_with("FOR+BITPACK"), "expected one packed chunk, got {described}");
assert_eq!(decode(&bytes).unwrap(), values, "{described}");
}
#[test]
fn a_cascade_decodes_the_same_through_a_shared_scratch_as_through_its_own() {
let mut values = Vec::new();
for index in 0..8192i64 {
values.push(1_600_000_000 + index / 4 + (index % 7) * 1_000);
}
let bytes = encode(&values).unwrap();
let described = describe(&bytes).unwrap();
assert!(described.contains('('), "expected a cascade, got {described}");
assert_eq!(decode(&bytes).unwrap(), values, "{described}");
}
#[test]
fn selected_positions_agree_with_a_full_decode_for_every_kind_with_a_point_form() {
let positions = [0, 1, 17, 1023, 1024, 4097, 8191];
let packed: Vec<i64> = (0..8192).map(|index| index * 31 % 1_000_003).collect();
let mut runs = Vec::new();
for run in 0..160i64 {
runs.extend(std::iter::repeat_n(run * 13, (run as usize % 71) + 2));
}
runs.resize(8192, -7);
let strided: Vec<i64> = (0..8192).map(|index| 500 + index * 7 % 5003 * 100).collect();
let coded: Vec<i64> =
(0..8192).map(|index| [-9_000_000_000, 3, 77, 1 << 40][index % 4]).collect();
for (kind, values) in [
(Kind::Packed, packed),
(Kind::Rle, runs),
(Kind::Strided, strided),
(Kind::Dict, coded),
] {
let bytes = encode_only(kind, &values).unwrap().expect("encoding applies");
let selected = decode_selected(&bytes, &positions).unwrap();
let expected = positions.iter().map(|&position| values[position]).collect::<Vec<_>>();
assert_eq!(selected, expected, "{}", kind.name());
let shape = describe(&bytes).unwrap();
let simple = !shape.contains("RLE") && !shape.contains("DELTA");
assert_eq!(pointed(&bytes), simple, "{shape}");
}
}
#[test]
fn selected_positions_must_be_ordered_and_inside_the_chunk() {
let bytes = encode_only(Kind::Packed, &(0..2048).collect::<Vec<_>>())
.unwrap()
.expect("packed applies");
assert!(decode_selected(&bytes, &[7, 7]).is_err());
assert!(decode_selected(&bytes, &[8, 3]).is_err());
assert!(decode_selected(&bytes, &[2048]).is_err());
}
#[test]
fn a_partial_unit_costs_its_own_values_and_not_a_whole_unit() {
let values = vec![1i64 << 39, (1 << 39) + 7, 1 << 38];
let bytes = encode_only(Kind::Packed, &values).unwrap().unwrap();
assert_eq!(bytes.len(), 5 + 9 + 15);
assert_eq!(decode(&bytes).unwrap(), values);
}
#[test]
fn the_frame_of_reference_is_per_unit_and_not_per_chunk() {
let values: Vec<i64> =
(0..4096i64).map(|index| (index / 1024) * 1_000_000 + (index % 1024)).collect();
let bytes = encode_only(Kind::Packed, &values).unwrap().unwrap();
assert_eq!(describe(&bytes).unwrap(), "FOR+BITPACK[10]");
assert_eq!(decode(&bytes).unwrap(), values);
}
#[test]
fn every_candidate_that_applies_decodes_to_the_input() {
let mut values = vec![5i64; 3000];
for (index, value) in values.iter_mut().enumerate() {
if index % 500 == 0 {
*value = index as i64;
}
}
let applicable = candidates(&values, 0, &EXHAUSTIVE);
assert!(applicable.len() >= 4, "{applicable:?}");
for kind in applicable {
let bytes = encode_only(kind, &values).unwrap().unwrap();
assert_eq!(decode(&bytes).unwrap(), values, "{}", kind.name());
}
}
#[test]
fn every_kind_that_applies_decodes_to_what_it_was_given() {
let shapes: Vec<Vec<i64>> = vec![
Vec::new(),
vec![5; 1024],
vec![i64::MIN, i64::MAX, 0, -1],
(0..1024).map(|at| at * 7).collect(),
(0..1024).map(|at| at % 17).collect(),
(0..1024).map(|at| if at % 100 == 0 { at } else { 3 }).collect(),
(0..1024).map(|at| -at * 1_000_003).collect(),
(0..1024_i64)
.map(|at| {
at.wrapping_mul(6_364_136_223_846_793_005)
.wrapping_add(1_442_695_040_888_963_407)
})
.collect(),
];
let kinds =
[Kind::Constant, Kind::Packed, Kind::Delta, Kind::Rle, Kind::Dict, Kind::Sparse];
for values in &shapes {
for kind in kinds {
let Some(bytes) = encode_only(kind, values).unwrap() else {
continue;
};
assert_eq!(
&decode(&bytes).unwrap(),
values,
"{} over {} values",
kind.name(),
values.len()
);
}
}
}
#[test]
fn the_chooser_picks_the_smallest_candidate_rather_than_the_first_that_applies() {
let mut values = vec![5i64; 3000];
values[1500] = 9;
let chosen = encode(&values).unwrap();
for (_, size) in candidate_sizes(&values).unwrap() {
assert!(chosen.len() <= size);
}
}
#[test]
fn a_truncated_chunk_is_an_error_and_not_a_panic() {
let bytes = encode(&[1, 2, 3, 4, 5]).unwrap();
for len in 0..bytes.len() {
let error = decode(&bytes[..len]).unwrap_err();
assert!(error.message().contains("chunk"), "{error}");
}
}
#[test]
fn trailing_bytes_are_an_error() {
let mut bytes = encode(&[1, 2, 3]).unwrap();
bytes.push(0);
let error = decode(&bytes).unwrap_err();
assert!(error.message().contains("left over"), "{error}");
}
#[test]
fn an_unknown_tag_is_an_error() {
let error = decode(&[99, 0, 0, 0, 0]).unwrap_err();
assert!(error.message().contains("unknown encoding tag"), "{error}");
}
#[test]
fn a_dictionary_code_outside_the_dictionary_is_an_error() {
let mut bytes = vec![Kind::Dict.tag()];
put_u32(&mut bytes, 1);
bytes.extend_from_slice(&encode(&[10]).unwrap());
bytes.extend_from_slice(&encode(&[5]).unwrap());
let error = decode(&bytes).unwrap_err();
assert!(error.message().contains("not in the dictionary"), "{error}");
}
#[test]
fn a_negative_run_length_is_an_error() {
let mut bytes = vec![Kind::Rle.tag()];
put_u32(&mut bytes, 4);
bytes.extend_from_slice(&encode(&[7]).unwrap());
bytes.extend_from_slice(&encode(&[-4]).unwrap());
let error = decode(&bytes).unwrap_err();
assert!(error.message().contains("negative"), "{error}");
}
#[test]
fn a_run_that_runs_past_its_chunk_is_an_error() {
let mut bytes = vec![Kind::Rle.tag()];
put_u32(&mut bytes, 4);
bytes.extend_from_slice(&encode(&[7]).unwrap());
bytes.extend_from_slice(&encode(&[9]).unwrap());
let error = decode(&bytes).unwrap_err();
assert!(error.message().contains("past its chunk"), "{error}");
}
#[test]
fn runs_shorter_and_longer_than_the_width_they_are_written_in_all_come_back() {
let lengths = [1, 3, 7, 8, 9, 40, 2, 1];
let mut values = Vec::new();
for (at, length) in lengths.iter().enumerate() {
let value = i64::try_from(at).expect("eight runs") * 1000 - 3;
values.extend(std::iter::repeat_n(value, *length));
}
let bytes = encode(&values).expect("encodes");
assert_eq!(decode(&bytes).expect("decodes"), values, "runs around the write width");
let singles: Vec<i64> = (0..300).map(|index| index * 7 % 11).collect();
let bytes = encode(&singles).expect("encodes");
assert_eq!(decode(&bytes).expect("decodes"), singles, "no run longer than one");
}
#[test]
fn the_cascade_depth_is_bounded() {
let values: Vec<i64> = (0..50_000).map(|index| (index / 100) % 250).collect();
let bytes = round_trip(&values);
let shape = describe(&bytes).unwrap();
let depth = shape.matches('(').count();
assert!(depth <= MAX_DEPTH as usize, "{shape} is {depth} deep");
}
#[test]
fn candidate_sizes_reports_what_the_chooser_looked_at() {
let values: Vec<i64> = (0..5000).map(|index| index % 17 * 1000).collect();
let sizes = candidate_sizes(&values).unwrap();
assert!(sizes.iter().any(|(kind, _)| *kind == Kind::Dict));
assert!(sizes.iter().any(|(kind, _)| *kind == Kind::Packed));
assert!(sizes.iter().all(|(_, size)| *size > 0));
}
#[test]
fn a_chunk_can_be_read_from_the_front_of_a_longer_buffer() {
let first = encode(&[1, 2, 3]).unwrap();
let second: Vec<i64> = (0..3000).map(|index| index % 11).collect();
let second_bytes = encode(&second).unwrap();
let mut joined = first.clone();
joined.extend_from_slice(&second_bytes);
joined.extend_from_slice(b"and then something else");
let (values, used) = decode_prefix(&joined).unwrap();
assert_eq!(values, vec![1, 2, 3]);
assert_eq!(used, first.len());
let (more, used_again) = decode_prefix(&joined[used..]).unwrap();
assert_eq!(more, second);
assert_eq!(used_again, second_bytes.len());
let (text, described) = describe_prefix(&joined).unwrap();
assert_eq!(described, first.len());
assert_eq!(text, describe(&first).unwrap());
}
#[test]
fn a_truncated_chunk_is_still_an_error_when_read_as_a_prefix() {
let bytes = encode(&(0..2000).collect::<Vec<i64>>()).unwrap();
for len in 0..bytes.len() {
assert!(decode_prefix(&bytes[..len]).is_err(), "{len} bytes decoded");
}
}
#[test]
fn decoding_into_a_narrow_type_agrees_with_the_wide_decoder() {
let columns: Vec<Vec<i64>> = vec![
vec![7; 3000],
(0..3000).map(|i| (i * 37) % 200 - 100).collect(),
(0..3000).map(|i| if i % 97 == 0 { i % 50 } else { 0 }).collect(),
(0..3000).map(|i| i / 250).collect(),
(0..3000).map(|i| i * 3 + 11).collect(),
(0..3000).map(|i| [5, -9, 120][i as usize % 3]).collect(),
];
let kinds = [
Kind::Constant,
Kind::Packed,
Kind::Delta,
Kind::Rle,
Kind::Dict,
Kind::Sparse,
Kind::Strided,
];
for values in &columns {
for kind in kinds {
let Some(bytes) = encode_only(kind, values).unwrap() else { continue };
let wide = decode(&bytes).unwrap();
let as_i16: Vec<i64> =
decode_as::<i16>(&bytes).unwrap().into_iter().map(i64::from).collect();
let as_i32: Vec<i64> =
decode_as::<i32>(&bytes).unwrap().into_iter().map(i64::from).collect();
assert_eq!(as_i16, wide, "{kind:?} as i16");
assert_eq!(as_i32, wide, "{kind:?} as i32");
assert_eq!(decode_as::<i64>(&bytes).unwrap(), wide, "{kind:?} as i64");
}
let chosen = encode(values).unwrap();
let narrow: Vec<i64> =
decode_as::<i16>(&chosen).unwrap().into_iter().map(i64::from).collect();
assert_eq!(narrow, decode(&chosen).unwrap(), "the chosen cascade as i16");
}
}
#[test]
fn decoding_into_a_narrow_type_takes_what_fits_and_refuses_what_does_not() {
fn check<T: Lane + TryFrom<i64> + PartialEq + std::fmt::Debug>(edges: [i64; 2]) {
for edge in edges {
for value in [edge - 1, edge, edge + 1] {
for kind in [Kind::Constant, Kind::Packed, Kind::Rle, Kind::Sparse] {
let values = [value, value, edges[0].max(0).min(edges[1]), value];
let Some(bytes) = encode_only(kind, &values).unwrap() else { continue };
let fits = values.iter().all(|&value| T::try_from(value).is_ok());
assert_eq!(decode_as::<T>(&bytes).is_ok(), fits, "{value} {kind:?}");
}
}
}
}
check::<i8>([-128, 127]);
check::<u8>([0, 255]);
check::<i16>([-32_768, 32_767]);
check::<u16>([0, 65_535]);
check::<i32>([i64::from(i32::MIN), i64::from(i32::MAX)]);
check::<u32>([0, i64::from(u32::MAX)]);
}
#[test]
fn a_packed_block_wider_than_its_type_still_decodes_when_its_values_fit() {
let values: Vec<i64> = (0..1500).map(|i| if i % 2 == 0 { -5 } else { 32_767 }).collect();
let bytes = encode_only(Kind::Packed, &values).unwrap().expect("packing always applies");
let narrow: Vec<i64> =
decode_as::<i16>(&bytes).unwrap().into_iter().map(i64::from).collect();
assert_eq!(narrow, values);
let over: Vec<i64> = values.iter().map(|&value| value + 1).collect();
let bytes = encode_only(Kind::Packed, &over).unwrap().expect("packing always applies");
assert!(decode_as::<i16>(&bytes).is_err(), "32768 is not an i16");
}
}