use std::cmp::{max, min};
use std::fmt::Debug;
use std::marker::PhantomData;
use crate::{Flags, gcd_utils, huffman_encoding};
use crate::bit_writer::BitWriter;
use crate::chunk_metadata::{ChunkMetadata, PrefixMetadata};
use crate::compression_table::CompressionTable;
use crate::constants::*;
use crate::data_types::{NumberLike, UnsignedLike};
use crate::delta_encoding;
use crate::delta_encoding::DeltaMoments;
use crate::errors::{QCompressError, QCompressResult};
use crate::gcd_utils::{GcdOperator, GeneralGcdOp, TrivialGcdOp};
use crate::prefix::{Prefix, PrefixCompressionInfo, WeightedPrefix};
use crate::prefix_optimization;
const MIN_N_TO_USE_RUN_LEN: usize = 1001;
const MIN_FREQUENCY_TO_USE_RUN_LEN: f64 = 0.8;
const DEFAULT_CHUNK_SIZE: usize = 1000000;
struct JumpstartConfiguration {
weight: usize,
jumpstart: usize,
}
#[derive(Clone, Debug)]
pub struct CompressorConfig {
pub compression_level: usize,
pub delta_encoding_order: usize,
pub use_gcds: bool,
phantom: PhantomData<()>,
}
impl Default for CompressorConfig {
fn default() -> Self {
Self {
compression_level: DEFAULT_COMPRESSION_LEVEL,
delta_encoding_order: 0,
use_gcds: true,
phantom: PhantomData,
}
}
}
impl CompressorConfig {
pub fn with_compression_level(mut self, level: usize) -> Self {
self.compression_level = level;
self
}
pub fn with_delta_encoding_order(mut self, order: usize) -> Self {
self.delta_encoding_order = order;
self
}
pub fn with_use_gcds(mut self, use_gcds: bool) -> Self {
self.use_gcds = use_gcds;
self
}
}
#[derive(Clone, Debug)]
struct InternalCompressorConfig {
pub compression_level: usize,
}
impl From<&CompressorConfig> for InternalCompressorConfig {
fn from(config: &CompressorConfig) -> Self {
InternalCompressorConfig {
compression_level: config.compression_level,
}
}
}
impl Default for InternalCompressorConfig {
fn default() -> Self {
Self::from(&CompressorConfig::default())
}
}
fn choose_run_len_jumpstart(
count: usize,
n: usize,
) -> JumpstartConfiguration {
let freq = (count as f64) / (n as f64);
let non_freq = 1.0 - freq;
let jumpstart = min((-non_freq.log2()).ceil() as usize, MAX_JUMPSTART);
let expected_n_runs = (freq * non_freq * n as f64).ceil() as usize;
JumpstartConfiguration {
weight: expected_n_runs,
jumpstart,
}
}
struct PrefixBuffer<'a, T: NumberLike> {
pub seq: &'a mut Vec<WeightedPrefix<T>>,
pub prefix_idx: &'a mut usize,
pub max_n_pref: usize,
pub n_unsigneds: usize,
pub sorted: &'a [T::Unsigned],
pub use_gcd: bool,
}
fn push_pref<T: NumberLike>(
buffer: &mut PrefixBuffer<'_, T>,
i: usize,
j: usize,
) {
let sorted = buffer.sorted;
let n_unsigneds = buffer.n_unsigneds;
let count = j - i;
let frequency = count as f64 / buffer.n_unsigneds as f64;
let new_prefix_idx = max(*buffer.prefix_idx + 1, (j * buffer.max_n_pref) / n_unsigneds);
let lower = T::from_unsigned(sorted[i]);
let upper = T::from_unsigned(sorted[j - 1]);
let gcd = if buffer.use_gcd {
gcd_utils::gcd(&sorted[i..j])
} else {
T::Unsigned::ONE
};
if n_unsigneds < MIN_N_TO_USE_RUN_LEN || frequency < MIN_FREQUENCY_TO_USE_RUN_LEN || count == n_unsigneds {
buffer.seq.push(WeightedPrefix::new(
count,
count,
lower,
upper,
None,
gcd,
));
} else {
let config = choose_run_len_jumpstart(count, n_unsigneds);
buffer.seq.push(WeightedPrefix::new(
count,
config.weight,
lower,
upper,
Some(config.jumpstart),
gcd,
));
}
*buffer.prefix_idx = new_prefix_idx;
}
fn train_prefixes<T: NumberLike>(
unsigneds: Vec<T::Unsigned>,
internal_config: &InternalCompressorConfig,
flags: &Flags,
n: usize, ) -> QCompressResult<Vec<Prefix<T>>> {
if unsigneds.is_empty() {
return Ok(Vec::new());
}
let comp_level = internal_config.compression_level;
if comp_level > MAX_COMPRESSION_LEVEL {
return Err(QCompressError::invalid_argument(format!(
"compresion level may not exceed {} (was {})",
MAX_COMPRESSION_LEVEL,
comp_level,
)));
}
if n > MAX_ENTRIES {
return Err(QCompressError::invalid_argument(format!(
"count may not exceed {} per chunk (was {})",
MAX_ENTRIES,
n,
)));
}
let n_unsigneds = unsigneds.len();
let mut sorted = unsigneds;
sorted.sort_unstable();
let safe_comp_level = min(comp_level, (n_unsigneds as f64).log2() as usize);
let max_n_pref = 1_usize << safe_comp_level;
let mut raw_prefs: Vec<WeightedPrefix<T>> = Vec::new();
let mut pref_idx = 0_usize;
let use_gcd = flags.use_gcds;
let mut i = 0;
let mut backup_j = 0_usize;
let mut prefix_buffer = PrefixBuffer::<T> {
seq: &mut raw_prefs,
prefix_idx: &mut pref_idx,
max_n_pref,
n_unsigneds,
sorted: &sorted,
use_gcd,
};
for j in 0..n_unsigneds {
let target_j = ((*prefix_buffer.prefix_idx + 1) * n_unsigneds) / max_n_pref;
if j > 0 && sorted[j] == sorted[j - 1] {
if j >= target_j && j - target_j >= target_j - backup_j && backup_j > i {
push_pref(&mut prefix_buffer, i, backup_j);
i = backup_j;
}
} else {
backup_j = j;
if j >= target_j {
push_pref(&mut prefix_buffer, i, j);
i = j;
}
}
}
push_pref(&mut prefix_buffer, i, n_unsigneds);
let mut optimized_prefs = prefix_optimization::optimize_prefixes(
raw_prefs,
flags,
n,
);
huffman_encoding::make_huffman_code(&mut optimized_prefs);
let prefixes = optimized_prefs.iter()
.map(|wp| wp.prefix.clone())
.collect();
Ok(prefixes)
}
#[derive(Clone)]
struct TrainedChunkCompressor<U: UnsignedLike, GcdOp: GcdOperator<U>> {
pub table: CompressionTable<U>,
op: PhantomData<GcdOp>,
}
fn trained_compress_chunk_nums<T: NumberLike>(
prefixes: &[Prefix<T>],
unsigneds: &[T::Unsigned],
writer: &mut BitWriter,
) -> QCompressResult<()> {
let table = CompressionTable::from(prefixes);
if gcd_utils::use_gcd_arithmetic(prefixes) {
TrainedChunkCompressor::<T::Unsigned, GeneralGcdOp> { table, op: PhantomData }
.compress_nums(unsigneds, writer)
} else {
TrainedChunkCompressor::<T::Unsigned, TrivialGcdOp> { table, op: PhantomData }
.compress_nums(unsigneds, writer)
}
}
impl<U, GcdOp> TrainedChunkCompressor<U, GcdOp> where U: UnsignedLike, GcdOp: GcdOperator<U> {
fn compress_nums(&self, unsigneds: &[U], writer: &mut BitWriter) -> QCompressResult<()> {
let mut i = 0;
while i < unsigneds.len() {
let unsigned = unsigneds[i];
let p = self.table.search(unsigned)?;
writer.write_usize(p.code, p.code_len);
match p.run_len_jumpstart {
None => {
Self::compress_offset_bits_w_prefix(unsigned, p, writer);
i += 1;
}
Some(jumpstart) => {
let mut reps = 1;
for &other in unsigneds.iter().skip(i + 1) {
if p.contains(other) {
reps += 1;
} else {
break;
}
}
writer.write_varint(reps - 1, jumpstart);
for &unsigned in unsigneds.iter().skip(i).take(reps) {
Self::compress_offset_bits_w_prefix(unsigned, p, writer);
}
i += reps;
}
}
}
writer.finish_byte();
Ok(())
}
fn compress_offset_bits_w_prefix(
unsigned: U,
p: &PrefixCompressionInfo<U>,
writer: &mut BitWriter,
) {
let off = GcdOp::get_offset(unsigned - p.lower, p.gcd);
writer.write_diff(off, p.k);
if off < p.only_k_bits_lower || off > p.only_k_bits_upper {
writer.write_one((off & (U::ONE << p.k)) > U::ZERO);
}
}
}
#[derive(Clone, Debug, Default)]
struct State {
has_written_header: bool,
has_written_footer: bool,
}
#[derive(Clone, Debug)]
pub struct Compressor<T> where T: NumberLike {
internal_config: InternalCompressorConfig,
flags: Flags,
writer: BitWriter,
state: State,
phantom: PhantomData<T>,
}
impl<T: NumberLike> Default for Compressor<T> {
fn default() -> Self {
Self::from_config(CompressorConfig::default())
}
}
impl<T> Compressor<T> where T: NumberLike {
pub fn from_config(config: CompressorConfig) -> Self {
Self {
internal_config: InternalCompressorConfig::from(&config),
flags: Flags::from(&config),
writer: BitWriter::default(),
state: State::default(),
phantom: PhantomData,
}
}
pub fn flags(&self) -> &Flags {
&self.flags
}
pub fn header(&mut self) -> QCompressResult<()> {
if self.state.has_written_header {
return Err(QCompressError::invalid_argument(
"attempted to write second header with compressor"
));
}
if self.state.has_written_footer {
return Err(QCompressError::invalid_argument(
"attempted to write header after footer"
));
}
self.writer.write_aligned_bytes(&MAGIC_HEADER)?;
self.writer.write_aligned_byte(T::HEADER_BYTE)?;
self.flags.write(&mut self.writer)?;
self.state.has_written_header = true;
Ok(())
}
pub fn chunk(&mut self, nums: &[T]) -> QCompressResult<ChunkMetadata<T>> {
if !self.state.has_written_header {
return Err(QCompressError::invalid_argument(
"attempted to write chunk before header"
));
}
if self.state.has_written_footer {
return Err(QCompressError::invalid_argument(
"attempted to write chunk to terminated compressor"
));
}
if nums.is_empty() {
return Err(QCompressError::invalid_argument(
"cannot compress empty chunk"
));
}
self.writer.write_aligned_byte(MAGIC_CHUNK_BYTE)?;
let n = nums.len();
let pre_meta_bit_idx = self.writer.bit_size();
let order = self.flags.delta_encoding_order;
let (mut metadata, post_meta_byte_idx) = if order == 0 {
let unsigneds = nums.iter()
.map(|x| x.to_unsigned())
.collect::<Vec<_>>();
let prefixes = train_prefixes(
unsigneds.clone(),
&self.internal_config,
&self.flags,
n,
)?;
let prefix_metadata = PrefixMetadata::Simple {
prefixes: prefixes.clone(),
};
let metadata = ChunkMetadata {
n,
compressed_body_size: 0,
prefix_metadata,
phantom: PhantomData,
};
metadata.write_to(&mut self.writer, &self.flags);
let post_meta_idx = self.writer.byte_size();
trained_compress_chunk_nums(
&prefixes,
&unsigneds,
&mut self.writer
)?;
(metadata, post_meta_idx)
} else {
let delta_moments = DeltaMoments::from(nums, order);
let deltas = delta_encoding::nth_order_deltas(nums, order);
let unsigneds = deltas.iter()
.map(|x| x.to_unsigned())
.collect::<Vec<_>>();
let prefixes = train_prefixes(
unsigneds.clone(),
&self.internal_config,
&self.flags,
n,
)?;
let prefix_metadata = PrefixMetadata::Delta {
delta_moments,
prefixes: prefixes.clone(),
};
let metadata = ChunkMetadata {
n,
compressed_body_size: 0,
prefix_metadata,
phantom: PhantomData,
};
metadata.write_to(&mut self.writer, &self.flags);
let post_meta_idx = self.writer.byte_size();
trained_compress_chunk_nums(
&prefixes,
&unsigneds,
&mut self.writer
)?;
(metadata, post_meta_idx)
};
metadata.compressed_body_size = self.writer.byte_size() - post_meta_byte_idx;
metadata.update_write_compressed_body_size(&mut self.writer, pre_meta_bit_idx);
Ok(metadata)
}
pub fn footer(&mut self) -> QCompressResult<()> {
if !self.state.has_written_header {
return Err(QCompressError::invalid_argument(
"attempted to write footer before header"
));
}
if self.state.has_written_footer {
return Err(QCompressError::invalid_argument(
"attempted to write second footer"
));
}
self.writer.write_aligned_byte(MAGIC_TERMINATION_BYTE)?;
self.state.has_written_footer = true;
Ok(())
}
pub fn simple_compress(&mut self, nums: &[T]) -> Vec<u8> {
self.header().unwrap();
nums.chunks(DEFAULT_CHUNK_SIZE)
.for_each(|chunk| {
self.chunk(chunk).unwrap();
});
self.footer().unwrap();
self.drain_bytes()
}
pub fn drain_bytes(&mut self) -> Vec<u8> {
self.writer.drain_bytes()
}
pub fn byte_size(&mut self) -> usize {
self.writer.byte_size()
}
}