use std::fmt::Debug;
use std::io::Write;
use std::marker::PhantomData;
use crate::Flags;
use crate::bit_reader::BitReader;
use crate::bit_words::BitWords;
use crate::chunk_body_decompressor::ChunkBodyDecompressor;
use crate::chunk_metadata::{ChunkMetadata};
use crate::constants::{MAGIC_CHUNK_BYTE, MAGIC_HEADER, MAGIC_TERMINATION_BYTE, WORD_SIZE};
use crate::data_types::NumberLike;
use crate::errors::{ErrorKind, QCompressError, QCompressResult};
#[derive(Clone, Debug)]
pub struct DecompressorConfig {
pub numbers_limit_per_item: usize,
phantom: PhantomData<()>, }
impl Default for DecompressorConfig {
fn default() -> Self {
Self {
numbers_limit_per_item: 100000,
phantom: PhantomData,
}
}
}
impl DecompressorConfig {
pub fn with_numbers_limit_per_item(mut self, limit: usize) -> Self {
self.numbers_limit_per_item = limit;
self
}
}
#[derive(Clone, Debug)]
pub enum DecompressedItem<T: NumberLike> {
Flags(Flags),
ChunkMetadata(ChunkMetadata<T>),
Numbers(Vec<T>),
Footer,
}
#[derive(Clone, Debug, Default)]
struct State<T: NumberLike> {
bit_idx: usize,
flags: Option<Flags>,
chunk_body_decompressor: Option<ChunkBodyDecompressor<T>>,
terminated: bool,
}
pub(crate) fn read_header<T: NumberLike>(reader: &mut BitReader) -> QCompressResult<Flags> {
let bytes = reader.read_aligned_bytes(MAGIC_HEADER.len())?;
if bytes != MAGIC_HEADER {
return Err(QCompressError::corruption(format!(
"magic header does not match {:?}; instead found {:?}",
MAGIC_HEADER,
bytes,
)));
}
let bytes = reader.read_aligned_bytes(1)?;
let byte = bytes[0];
if byte != T::HEADER_BYTE {
return Err(QCompressError::corruption(format!(
"data type byte does not match {:?}; instead found {:?}",
T::HEADER_BYTE,
byte,
)));
}
Flags::parse_from(reader)
}
pub(crate) fn read_chunk_meta<T: NumberLike>(reader: &mut BitReader, flags: &Flags) -> QCompressResult<Option<ChunkMetadata<T>>> {
let magic_byte = reader.read_aligned_bytes(1)?[0];
if magic_byte == MAGIC_TERMINATION_BYTE {
return Ok(None);
} else if magic_byte != MAGIC_CHUNK_BYTE {
return Err(QCompressError::corruption(format!(
"invalid magic chunk byte: {}",
magic_byte
)));
}
let metadata = ChunkMetadata::parse_from(reader, flags)?;
reader.drain_empty_byte(|| QCompressError::corruption(
"nonzero bits in end of final byte of chunk metadata"
))?;
Ok(Some(metadata))
}
#[derive(Clone, Debug, Default)]
pub struct Decompressor<T> where T: NumberLike {
config: DecompressorConfig,
words: BitWords,
state: State<T>,
}
impl<T: NumberLike> Write for Decompressor<T> {
fn write(&mut self, buf: &[u8]) -> std::io::Result<usize> {
self.words.extend_bytes(buf);
Ok(buf.len())
}
fn flush(&mut self) -> std::io::Result<()> {
Ok(())
}
}
impl<T> Decompressor<T> where T: NumberLike {
pub fn from_config(config: DecompressorConfig) -> Self {
Self {
config,
..Default::default()
}
}
pub fn bit_idx(&self) -> usize {
self.state.bit_idx
}
fn with_reader<X, F>(&mut self, f: F) -> QCompressResult<X>
where F: FnOnce(&mut BitReader, &mut State<T>, &DecompressorConfig) -> QCompressResult<X> {
let mut reader = BitReader::from(&self.words);
reader.seek_to(self.state.bit_idx);
let res = f(&mut reader, &mut self.state, &self.config);
if res.is_ok() {
self.state.bit_idx = reader.bit_idx();
}
res
}
fn check_not_terminated(&self) -> QCompressResult<()> {
if self.state.terminated {
Err(QCompressError::invalid_argument("attempted to write to terminated decompressor"))
} else {
Ok(())
}
}
pub fn header(&mut self) -> QCompressResult<Flags> {
self.check_not_terminated()?;
if self.state.flags.is_some() {
return Err(QCompressError::invalid_argument(
"attempted to decompress header for the 2nd time"
))
}
self.with_reader(|reader, state, _| {
let flags = read_header::<T>(reader)?;
state.flags = Some(flags.clone());
Ok(flags)
})
}
pub fn chunk_metadata(&mut self) -> QCompressResult<Option<ChunkMetadata<T>>> {
self.check_not_terminated()?;
if self.state.flags.is_none() {
return Err(QCompressError::invalid_argument(
"attempted to decompress chunk metadata before header"
));
}
if self.state.chunk_body_decompressor.is_some() {
return Err(QCompressError::invalid_argument(
"attempted to decompress chunk metadata before chunk body was finished"
));
}
self.with_reader(|reader, state, _| {
let flags = state.flags.as_ref().unwrap();
let maybe_meta = read_chunk_meta(reader, flags)?;
if let Some(meta) = &maybe_meta {
state.chunk_body_decompressor = Some(ChunkBodyDecompressor::new(meta)?)
}
Ok(maybe_meta)
})
}
fn check_in_chunk_body(&self) -> QCompressResult<()> {
self.check_not_terminated()?;
if self.state.chunk_body_decompressor.is_none() {
return Err(QCompressError::invalid_argument(
"attempted to decompress chunk body before its chunk metadata"
));
}
Ok(())
}
pub fn skip_chunk_body(&mut self) -> QCompressResult<()> {
self.check_in_chunk_body()?;
let cbd = self.state.chunk_body_decompressor.as_ref().unwrap();
let skipped_bit_idx = self.state.bit_idx + cbd.bits_remaining();
if skipped_bit_idx <= self.words.total_bits {
self.state.bit_idx = skipped_bit_idx;
self.state.chunk_body_decompressor = None;
Ok(())
} else {
Err(QCompressError::insufficient_data(format!(
"unable to skip chunk body to bit index {} when only {} bits available",
skipped_bit_idx,
self.words.total_bits,
)))
}
}
pub fn chunk_body(&mut self) -> QCompressResult<Vec<T>> {
self.check_in_chunk_body()?;
self.with_reader(|reader, state, _| {
let chunk_body_decompressor = state.chunk_body_decompressor.as_mut().unwrap();
let numbers = chunk_body_decompressor.decompress_next_batch(
reader,
usize::MAX,
true,
)?;
state.chunk_body_decompressor = None;
Ok(numbers.nums)
})
}
pub fn simple_decompress(&mut self) -> QCompressResult<Vec<T>> {
let mut res: Option<Vec<T>> = None;
self.header()?;
while self.chunk_metadata()?.is_some() {
let nums = self.chunk_body()?;
res = match res {
Some(mut existing) => {
existing.extend(nums);
Some(existing)
}
None => {
Some(nums)
}
};
}
Ok(res.unwrap_or_default())
}
pub fn free_compressed_memory(&mut self) {
let words_to_free = self.state.bit_idx / WORD_SIZE;
if words_to_free > 0 {
self.words.truncate_left(words_to_free);
self.state.bit_idx -= words_to_free * WORD_SIZE;
}
}
}
impl<T: NumberLike> Iterator for &mut Decompressor<T> {
type Item = QCompressResult<DecompressedItem<T>>;
fn next(&mut self) -> Option<Self::Item> {
let res = self.with_reader(|reader, state, config| {
if state.terminated {
return Ok(None);
}
if state.flags.is_none() {
match read_header::<T>(reader) {
Ok(flags) => {
state.flags = Some(flags.clone());
Ok(Some(DecompressedItem::Flags(flags)))
},
Err(e) if matches!(e.kind, ErrorKind::InsufficientData) => Ok(None),
Err(e) => Err(e),
}
} else if state.chunk_body_decompressor.is_none() {
match read_chunk_meta::<T>(reader, state.flags.as_ref().unwrap()) {
Ok(Some(meta)) => {
match ChunkBodyDecompressor::new(&meta) {
Ok(cbd) => {
state.chunk_body_decompressor = Some(cbd);
Ok(Some(DecompressedItem::ChunkMetadata(meta)))
}
Err(e) => Err(e)
}
},
Ok(None) => {
state.terminated = true;
Ok(Some(DecompressedItem::Footer))
},
Err(e) if matches!(e.kind, ErrorKind::InsufficientData) => Ok(None),
Err(e) => Err(e),
}
} else {
let nums_result = state.chunk_body_decompressor.as_mut()
.unwrap()
.decompress_next_batch(reader, config.numbers_limit_per_item, false);
match nums_result {
Ok(numbers) => {
if numbers.nums.is_empty() {
Ok(None)
} else {
if numbers.finished_chunk_body {
state.chunk_body_decompressor = None;
}
Ok(Some(DecompressedItem::Numbers(numbers.nums)))
}
}
Err(e) => Err(e),
}
}
});
match res {
Ok(Some(x)) => Some(Ok(x)),
Ok(None) => None,
Err(e) => Some(Err(e))
}
}
}