use crate::error::{Error, Result};
use std::cmp::Ordering;
use std::fmt::{self, Display};
#[derive(Debug, Clone, Copy, PartialEq, Default, bitcode::Encode, bitcode::Decode)]
#[repr(C)]
pub struct TSSample {
pub ts: u64,
pub v: f64,
}
impl TSSample {
pub const MAX_TIMESTAMP: u64 = u64::MAX;
pub const NAN_VALUE: f64 = f64::NAN;
#[inline]
pub const fn new(ts: u64, v: f64) -> Self {
Self { ts, v }
}
}
impl Display for TSSample {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
write!(f, "({}, {})", self.ts, self.v)
}
}
impl Eq for TSSample {}
impl PartialOrd for TSSample {
#[inline]
fn partial_cmp(&self, other: &Self) -> Option<Ordering> {
Some(self.cmp(other))
}
}
impl Ord for TSSample {
#[inline]
fn cmp(&self, other: &Self) -> Ordering {
self.ts.cmp(&other.ts)
}
}
#[derive(Debug, Default, Clone)]
pub struct BitWriter {
buf: Vec<u8>,
current_byte: u8,
bit_count: u8,
}
impl BitWriter {
#[inline]
pub fn new() -> Self {
Self::with_capacity(64)
}
#[inline]
pub fn with_capacity(capacity: usize) -> Self {
Self {
buf: Vec::with_capacity(capacity),
current_byte: 0,
bit_count: 0,
}
}
#[inline(always)]
pub fn write_bit(&mut self, bit: bool) {
if bit {
self.current_byte |= 1 << (7 - self.bit_count);
}
self.bit_count += 1;
if self.bit_count == 8 {
self.buf.push(self.current_byte);
self.current_byte = 0;
self.bit_count = 0;
}
}
#[inline(always)]
pub fn write_bits(&mut self, val: u64, mut num_bits: u8) {
if num_bits == 0 {
return;
}
while num_bits > 0 {
let space_in_byte = 8 - self.bit_count;
if num_bits <= space_in_byte {
let shift = space_in_byte - num_bits;
let mask = if num_bits == 64 {
u64::MAX
} else {
(1u64 << num_bits) - 1
};
self.current_byte |= ((val & mask) as u8) << shift;
self.bit_count += num_bits;
if self.bit_count == 8 {
self.buf.push(self.current_byte);
self.current_byte = 0;
self.bit_count = 0;
}
break;
} else {
let shift = num_bits - space_in_byte;
let chunk = (val >> shift) & ((1u64 << space_in_byte) - 1);
self.current_byte |= chunk as u8;
self.buf.push(self.current_byte);
self.current_byte = 0;
self.bit_count = 0;
num_bits -= space_in_byte;
}
}
}
#[inline]
pub fn finish(mut self) -> Vec<u8> {
if self.bit_count > 0 {
self.buf.push(self.current_byte);
self.current_byte = 0;
self.bit_count = 0;
}
self.buf
}
}
#[derive(Debug, Clone)]
pub struct BitReader<'a> {
data: &'a [u8],
byte_offset: usize,
bit_offset: u8,
}
impl<'a> BitReader<'a> {
#[inline]
pub const fn new(data: &'a [u8]) -> Self {
Self {
data,
byte_offset: 0,
bit_offset: 0,
}
}
#[inline(always)]
pub fn read_bit(&mut self) -> Option<bool> {
if self.byte_offset >= self.data.len() {
return None;
}
let b = self.data[self.byte_offset];
let bit = (b & (1 << (7 - self.bit_offset))) != 0;
self.bit_offset += 1;
if self.bit_offset == 8 {
self.byte_offset += 1;
self.bit_offset = 0;
}
Some(bit)
}
#[inline(always)]
pub fn read_bits(&mut self, mut num_bits: u8) -> Option<u64> {
if num_bits == 0 {
return Some(0);
}
let mut result = 0u64;
while num_bits > 0 {
if self.byte_offset >= self.data.len() {
return None;
}
let space_in_byte = 8 - self.bit_offset;
let b = self.data[self.byte_offset];
if num_bits <= space_in_byte {
let shift = space_in_byte - num_bits;
let chunk = (b >> shift) & (((1u16 << num_bits) - 1) as u8);
result = (result << num_bits) | (chunk as u64);
self.bit_offset += num_bits;
if self.bit_offset == 8 {
self.byte_offset += 1;
self.bit_offset = 0;
}
break;
} else {
let chunk = b & (((1u16 << space_in_byte) - 1) as u8);
result = (result << space_in_byte) | (chunk as u64);
self.byte_offset += 1;
self.bit_offset = 0;
num_bits -= space_in_byte;
}
}
Some(result)
}
}
pub fn gorilla_compress_samples(samples: &[TSSample]) -> Vec<u8> {
if samples.is_empty() {
return Vec::new();
}
let mut writer = BitWriter::with_capacity(samples.len() * 4);
let t0 = samples[0].ts;
writer.write_bits(t0, 64);
let v0_bits = samples[0].v.to_bits();
writer.write_bits(v0_bits, 64);
if samples.len() == 1 {
return writer.finish();
}
let mut prev_ts = t0;
let mut prev_delta = samples[1].ts.saturating_sub(prev_ts);
writer.write_bits(prev_delta, 32);
let mut prev_v_bits = v0_bits;
let mut prev_leading = 0xFF;
let mut prev_trailing = 0xFF;
compress_value(
&mut writer,
samples[1].v.to_bits(),
prev_v_bits,
&mut prev_leading,
&mut prev_trailing,
);
prev_ts = samples[1].ts;
prev_v_bits = samples[1].v.to_bits();
for s in &samples[2..] {
let cur_delta = s.ts.saturating_sub(prev_ts);
let dod = (cur_delta as i64) - (prev_delta as i64);
if dod == 0 {
writer.write_bit(false);
} else if (-63..=64).contains(&dod) {
writer.write_bits(0b10, 2);
writer.write_bits((dod + 63) as u64, 7);
} else if (-255..=256).contains(&dod) {
writer.write_bits(0b110, 3);
writer.write_bits((dod + 255) as u64, 9);
} else if (-2047..=2048).contains(&dod) {
writer.write_bits(0b1110, 4);
writer.write_bits((dod + 2047) as u64, 12);
} else {
writer.write_bits(0b1111, 4);
writer.write_bits(dod as u32 as u64, 32);
}
prev_delta = cur_delta;
prev_ts = s.ts;
compress_value(
&mut writer,
s.v.to_bits(),
prev_v_bits,
&mut prev_leading,
&mut prev_trailing,
);
prev_v_bits = s.v.to_bits();
}
writer.finish()
}
#[inline(always)]
fn compress_value(
writer: &mut BitWriter,
v_bits: u64,
prev_v_bits: u64,
prev_leading: &mut u8,
prev_trailing: &mut u8,
) {
let xor = v_bits ^ prev_v_bits;
if xor == 0 {
writer.write_bit(false);
} else {
writer.write_bit(true);
let leading = (xor.leading_zeros() as u8).min(31);
let trailing = xor.trailing_zeros() as u8;
if *prev_leading != 0xFF && leading >= *prev_leading && trailing >= *prev_trailing {
writer.write_bit(false);
let meaningful_len = 64 - *prev_leading - *prev_trailing;
let meaningful_bits = if meaningful_len == 64 {
xor
} else {
(xor >> *prev_trailing) & ((1u64 << meaningful_len) - 1)
};
writer.write_bits(meaningful_bits, meaningful_len);
} else {
writer.write_bit(true);
writer.write_bits(leading as u64, 5);
let meaningful_len = 64 - leading - trailing;
writer.write_bits((meaningful_len - 1) as u64, 6);
let meaningful_bits = if meaningful_len == 64 {
xor
} else {
(xor >> trailing) & ((1u64 << meaningful_len) - 1)
};
writer.write_bits(meaningful_bits, meaningful_len);
*prev_leading = leading;
*prev_trailing = trailing;
}
}
}
pub fn gorilla_decompress_into(
data: &[u8],
count: usize,
samples: &mut Vec<TSSample>,
) -> Result<()> {
if count == 0 || data.is_empty() {
return Ok(());
}
let mut reader = BitReader::new(data);
samples.reserve(count);
let t0 = reader
.read_bits(64)
.ok_or_else(|| Error::invalid_data("ERR TSDB: corrupted gorilla stream (missing t0)"))?;
let v0_bits = reader
.read_bits(64)
.ok_or_else(|| Error::invalid_data("ERR TSDB: corrupted gorilla stream (missing v0)"))?;
samples.push(TSSample::new(t0, f64::from_bits(v0_bits)));
if count == 1 {
return Ok(());
}
let delta0 = reader.read_bits(32).ok_or_else(|| {
Error::invalid_data("ERR TSDB: corrupted gorilla stream (missing delta0)")
})?;
let mut prev_ts = t0.saturating_add(delta0);
let mut prev_delta = delta0;
let mut prev_v_bits = v0_bits;
let mut prev_leading = 0xFF;
let mut prev_trailing = 0xFF;
let v1_bits = decompress_value(
&mut reader,
prev_v_bits,
&mut prev_leading,
&mut prev_trailing,
)?;
samples.push(TSSample::new(prev_ts, f64::from_bits(v1_bits)));
prev_v_bits = v1_bits;
for _ in 2..count {
let bit = reader
.read_bit()
.ok_or_else(|| Error::invalid_data("ERR TSDB: corrupted gorilla dod bit"))?;
let dod: i64 = if !bit {
0
} else {
let b2 = reader
.read_bit()
.ok_or_else(|| Error::invalid_data("ERR TSDB: corrupted gorilla bits"))?;
if !b2 {
let val = reader
.read_bits(7)
.ok_or_else(|| Error::invalid_data("ERR TSDB: corrupted gorilla bits"))?;
(val as i64) - 63
} else {
let b3 = reader
.read_bit()
.ok_or_else(|| Error::invalid_data("ERR TSDB: corrupted gorilla bits"))?;
if !b3 {
let val = reader
.read_bits(9)
.ok_or_else(|| Error::invalid_data("ERR TSDB: corrupted gorilla bits"))?;
(val as i64) - 255
} else {
let b4 = reader
.read_bit()
.ok_or_else(|| Error::invalid_data("ERR TSDB: corrupted gorilla bits"))?;
if !b4 {
let val = reader.read_bits(12).ok_or_else(|| {
Error::invalid_data("ERR TSDB: corrupted gorilla bits")
})?;
(val as i64) - 2047
} else {
let val = reader.read_bits(32).ok_or_else(|| {
Error::invalid_data("ERR TSDB: corrupted gorilla bits")
})?;
val as u32 as i32 as i64
}
}
}
};
let cur_delta = ((prev_delta as i64) + dod).max(0) as u64;
prev_ts = prev_ts.saturating_add(cur_delta);
prev_delta = cur_delta;
let cur_v_bits = decompress_value(
&mut reader,
prev_v_bits,
&mut prev_leading,
&mut prev_trailing,
)?;
samples.push(TSSample::new(prev_ts, f64::from_bits(cur_v_bits)));
prev_v_bits = cur_v_bits;
}
Ok(())
}
#[inline]
pub fn gorilla_decompress_samples(data: &[u8], count: usize) -> Result<Vec<TSSample>> {
let mut samples = Vec::with_capacity(count);
gorilla_decompress_into(data, count, &mut samples)?;
Ok(samples)
}
#[inline]
pub fn gorilla_decompress_last_timestamp(data: &[u8], count: usize) -> Option<u64> {
if count == 0 || data.is_empty() {
return None;
}
let mut reader = BitReader::new(data);
let t0 = reader.read_bits(64)?;
reader.read_bits(64)?;
if count == 1 {
return Some(t0);
}
let delta0 = reader.read_bits(32)?;
let mut prev_ts = t0.saturating_add(delta0);
let mut prev_delta = delta0;
let mut prev_leading = 0xFF;
let mut prev_trailing = 0xFF;
skip_value(&mut reader, &mut prev_leading, &mut prev_trailing)?;
if count == 2 {
return Some(prev_ts);
}
for _ in 2..count {
let bit = reader.read_bit()?;
let dod: i64 = if !bit {
0
} else {
let b2 = reader.read_bit()?;
if !b2 {
let val = reader.read_bits(7)?;
(val as i64) - 63
} else {
let b3 = reader.read_bit()?;
if !b3 {
let val = reader.read_bits(9)?;
(val as i64) - 255
} else {
let b4 = reader.read_bit()?;
if !b4 {
let val = reader.read_bits(12)?;
(val as i64) - 2047
} else {
let val = reader.read_bits(32)?;
val as u32 as i32 as i64
}
}
}
};
let cur_delta = ((prev_delta as i64) + dod).max(0) as u64;
prev_ts = prev_ts.saturating_add(cur_delta);
prev_delta = cur_delta;
skip_value(&mut reader, &mut prev_leading, &mut prev_trailing)?;
}
Some(prev_ts)
}
#[inline(always)]
fn skip_value(
reader: &mut BitReader<'_>,
prev_leading: &mut u8,
prev_trailing: &mut u8,
) -> Option<()> {
let bit = reader.read_bit()?;
if bit {
let reuse_header = !reader.read_bit()?;
if reuse_header {
if *prev_leading == 0xFF || *prev_trailing == 0xFF {
return None;
}
let meaningful_len = 64 - *prev_leading - *prev_trailing;
reader.read_bits(meaningful_len)?;
} else {
let lz = reader.read_bits(5)? as u8;
let meaningful_len = (reader.read_bits(6)? as u8) + 1;
reader.read_bits(meaningful_len)?;
let shift = 64 - lz - meaningful_len;
*prev_leading = lz;
*prev_trailing = shift;
}
}
Some(())
}
#[inline(always)]
fn decompress_value(
reader: &mut BitReader<'_>,
prev_v_bits: u64,
prev_leading: &mut u8,
prev_trailing: &mut u8,
) -> Result<u64> {
let bit = reader
.read_bit()
.ok_or_else(|| Error::invalid_data("ERR TSDB: corrupted gorilla value bit"))?;
if !bit {
Ok(prev_v_bits)
} else {
let reuse_header = !reader
.read_bit()
.ok_or_else(|| Error::invalid_data("ERR TSDB: corrupted gorilla reuse bit"))?;
if reuse_header {
if *prev_leading == 0xFF || *prev_trailing == 0xFF {
return Err(Error::invalid_data(
"ERR TSDB: corrupted gorilla stream (invalid reuse header)",
));
}
let meaningful_len = 64 - *prev_leading - *prev_trailing;
let bits = reader.read_bits(meaningful_len).ok_or_else(|| {
Error::invalid_data("ERR TSDB: corrupted gorilla meaningful bits")
})?;
let shift = *prev_trailing;
let xor = if meaningful_len == 64 {
bits
} else {
bits << shift
};
Ok(prev_v_bits ^ xor)
} else {
let lz = reader
.read_bits(5)
.ok_or_else(|| Error::invalid_data("ERR TSDB: corrupted gorilla lz"))?
as u8;
let meaningful_len = (reader
.read_bits(6)
.ok_or_else(|| Error::invalid_data("ERR TSDB: corrupted gorilla meaningful len"))?
as u8)
+ 1;
let bits = reader.read_bits(meaningful_len).ok_or_else(|| {
Error::invalid_data("ERR TSDB: corrupted gorilla meaningful bits")
})?;
let shift = 64 - lz - meaningful_len;
let xor = if meaningful_len == 64 {
bits
} else {
bits << shift
};
*prev_leading = lz;
*prev_trailing = shift;
Ok(prev_v_bits ^ xor)
}
}
}