From fef9b6fe64981c04a9a13d936d9f31e9d45deb9d Mon Sep 17 00:00:00 2001 From: David Roundy Date: Sun, 9 Jul 2023 04:41:46 -0700 Subject: [PATCH 1/7] start creating a new decompressor --- src/frame/mod.rs | 2 + src/frame/raw_decompress.rs | 389 ++++++++++++++++++++++++++++++++++++ 2 files changed, 391 insertions(+) create mode 100644 src/frame/raw_decompress.rs diff --git a/src/frame/mod.rs b/src/frame/mod.rs index 8acfad2c..9561e20e 100644 --- a/src/frame/mod.rs +++ b/src/frame/mod.rs @@ -23,10 +23,12 @@ use std::{fmt, io}; pub(crate) mod compress; #[cfg_attr(feature = "safe-decode", forbid(unsafe_code))] pub(crate) mod decompress; +pub(crate) mod raw_decompress; pub(crate) mod header; pub use compress::{AutoFinishEncoder, FrameEncoder}; pub use decompress::FrameDecoder; +pub use raw_decompress::Decoder; pub use header::{BlockMode, BlockSize, FrameInfo}; #[derive(Debug)] diff --git a/src/frame/raw_decompress.rs b/src/frame/raw_decompress.rs new file mode 100644 index 00000000..31015723 --- /dev/null +++ b/src/frame/raw_decompress.rs @@ -0,0 +1,389 @@ +use std::{ + convert::TryInto, + fmt, + hash::Hasher, + io::{self, ErrorKind}, + mem::size_of, +}; +use twox_hash::XxHash32; + +use super::header::{ + BlockInfo, BlockMode, FrameInfo, LZ4F_LEGACY_MAGIC_NUMBER, MAGIC_NUMBER_SIZE, + MAX_FRAME_INFO_SIZE, MIN_FRAME_INFO_SIZE, +}; +use super::Error; +use crate::{ + block::WINDOW_SIZE, + sink::{vec_sink_for_decompression, SliceSink}, +}; + +/// A reader for decompressing the LZ4 frame format +/// +/// This Decoder wraps any other reader that implements `io::Read`. +/// Bytes read will be decompressed according to the [LZ4 frame format]( +/// https://github.com/lz4/lz4/blob/dev/doc/lz4_Frame_format.md). +/// +/// # Example 1 +/// Deserializing json values out of a compressed file. +/// +/// ```no_run +/// let compressed_input = std::fs::File::open("datafile").unwrap(); +/// let mut decompressed_input = lz4_flex::frame::FrameDecoder::new(compressed_input); +/// let json: serde_json::Value = serde_json::from_reader(decompressed_input).unwrap(); +/// ``` +/// +/// # Example +/// Deserializing multiple json values out of a compressed file +/// +/// ```no_run +/// let compressed_input = std::fs::File::open("datafile").unwrap(); +/// let mut decompressed_input = lz4_flex::frame::FrameDecoder::new(compressed_input); +/// loop { +/// match serde_json::from_reader::<_, serde_json::Value>(&mut decompressed_input) { +/// Ok(json) => { println!("json {:?}", json); } +/// Err(e) if e.is_eof() => break, +/// Err(e) => panic!("{}", e), +/// } +/// } +/// ``` +pub struct Decoder { + /// The underlying reader. + r: R, + /// The FrameInfo of the frame currently being decoded. + /// It starts as `None` and is filled with the FrameInfo is read from the input. + /// It's reset to `None` once the frame EndMarker is read from the input. + current_frame_info: Option, + /// Xxhash32 used when content checksum is enabled. + content_hasher: XxHash32, + /// Total length of decompressed output for the current frame. + content_len: u64, + /// The compressed bytes buffer, taken from the underlying reader. + src: Vec, + /// The decompressed bytes buffer. Bytes are decompressed from src to dst + /// before being passed back to the caller. + dst: Vec, + /// Index into dst and length: starting point of bytes previously output + /// that are still part of the decompressor window. + ext_dict_offset: usize, + ext_dict_len: usize, + /// Index into dst: starting point of bytes not yet read by caller. + dst_start: usize, + /// Index into dst: ending point of bytes not yet read by caller. + dst_end: usize, +} + +impl Decoder { + /// Creates a new Decoder for the specified reader. + pub fn new(rdr: R) -> Decoder { + Decoder { + r: rdr, + src: Default::default(), + dst: Default::default(), + ext_dict_offset: 0, + ext_dict_len: 0, + dst_start: 0, + dst_end: 0, + current_frame_info: None, + content_hasher: XxHash32::with_seed(0), + content_len: 0, + } + } + + /// Gets a reference to the underlying reader in this decoder. + pub fn get_ref(&self) -> &R { + &self.r + } + + /// Gets a mutable reference to the underlying reader in this decoder. + /// + /// Note that mutation of the stream may result in surprising results if + /// this decoder is continued to be used. + pub fn get_mut(&mut self) -> &mut R { + &mut self.r + } + + /// Consumes the Decoder and returns the underlying reader. + pub fn into_inner(self) -> R { + self.r + } + + fn read_frame_info(&mut self) -> Result { + let mut buffer = [0u8; MAX_FRAME_INFO_SIZE]; + + match self.r.read(&mut buffer[..MAGIC_NUMBER_SIZE])? { + 0 => return Ok(0), + MAGIC_NUMBER_SIZE => (), + read => self.r.read_exact(&mut buffer[read..MAGIC_NUMBER_SIZE])?, + } + + if u32::from_le_bytes(buffer[0..MAGIC_NUMBER_SIZE].try_into().unwrap()) + != LZ4F_LEGACY_MAGIC_NUMBER + { + match self + .r + .read(&mut buffer[MAGIC_NUMBER_SIZE..MIN_FRAME_INFO_SIZE])? + { + 0 => return Ok(0), + MIN_FRAME_INFO_SIZE => (), + read => self + .r + .read_exact(&mut buffer[MAGIC_NUMBER_SIZE + read..MIN_FRAME_INFO_SIZE])?, + } + } + let required = FrameInfo::read_size(&buffer[..MIN_FRAME_INFO_SIZE])?; + if required != MIN_FRAME_INFO_SIZE && required != MAGIC_NUMBER_SIZE { + self.r + .read_exact(&mut buffer[MIN_FRAME_INFO_SIZE..required])?; + } + + let frame_info = FrameInfo::read(&buffer[..required])?; + if frame_info.dict_id.is_some() { + // Unsupported right now so it must be None + return Err(Error::DictionaryNotSupported.into()); + } + + let max_block_size = frame_info.block_size.get_size(); + let dst_size = if frame_info.block_mode == BlockMode::Linked { + // In linked mode we consume the output (bumping dst_start) but leave the + // beginning of dst to be used as a prefix in subsequent blocks. + // That is at least until we have at least `max_block_size + WINDOW_SIZE` + // bytes in dst, then we setup an ext_dict with the last WINDOW_SIZE bytes + // and the output goes to the beginning of dst again. + // Since we always want to be able to write a full block (up to max_block_size) + // we need a buffer with at least `max_block_size * 2 + WINDOW_SIZE` bytes. + max_block_size * 2 + WINDOW_SIZE + } else { + max_block_size + }; + self.src.clear(); + self.dst.clear(); + self.src.reserve_exact(max_block_size); + self.dst.reserve_exact(dst_size); + self.current_frame_info = Some(frame_info); + self.content_hasher = XxHash32::with_seed(0); + self.content_len = 0; + self.ext_dict_len = 0; + self.dst_start = 0; + self.dst_end = 0; + Ok(required) + } + + #[inline] + fn read_checksum(r: &mut R) -> Result { + let mut checksum_buffer = [0u8; size_of::()]; + r.read_exact(&mut checksum_buffer[..])?; + let checksum = u32::from_le_bytes(checksum_buffer); + Ok(checksum) + } + + #[inline] + fn check_block_checksum(data: &[u8], expected_checksum: u32) -> Result<(), io::Error> { + let mut block_hasher = XxHash32::with_seed(0); + block_hasher.write(data); + let calc_checksum = block_hasher.finish() as u32; + if calc_checksum != expected_checksum { + return Err(Error::BlockChecksumError.into()); + } + Ok(()) + } + + fn read_block(&mut self) -> io::Result { + debug_assert_eq!(self.dst_start, self.dst_end); + let frame_info = self.current_frame_info.as_ref().unwrap(); + + // Adjust dst buffer offsets to decompress the next block + let max_block_size = frame_info.block_size.get_size(); + if frame_info.block_mode == BlockMode::Linked { + // In linked mode we consume the output (bumping dst_start) but leave the + // beginning of dst to be used as a prefix in subsequent blocks. + // That is at least until we have at least `max_block_size + WINDOW_SIZE` + // bytes in dst, then we setup an ext_dict with the last WINDOW_SIZE bytes + // and the output goes to the beginning of dst again. + debug_assert_eq!(self.dst.capacity(), max_block_size * 2 + WINDOW_SIZE); + if self.dst_start + max_block_size > self.dst.capacity() { + // Output might not fit in the buffer. + // The ext_dict will become the last WINDOW_SIZE bytes + debug_assert!(self.dst_start >= max_block_size + WINDOW_SIZE); + self.ext_dict_offset = self.dst_start - WINDOW_SIZE; + self.ext_dict_len = WINDOW_SIZE; + // Output goes in the beginning of the buffer again. + self.dst_start = 0; + self.dst_end = 0; + } else if self.dst_start + self.ext_dict_len > WINDOW_SIZE { + // There's more than WINDOW_SIZE bytes of lookback adding the prefix and ext_dict. + // Since we have a limited buffer we must shrink ext_dict in favor of the prefix, + // so that we can fit up to max_block_size bytes between dst_start and ext_dict + // start. + let delta = self + .ext_dict_len + .min(self.dst_start + self.ext_dict_len - WINDOW_SIZE); + self.ext_dict_offset += delta; + self.ext_dict_len -= delta; + debug_assert!(self.dst_start + self.ext_dict_len >= WINDOW_SIZE) + } + } else { + debug_assert_eq!(self.ext_dict_len, 0); + debug_assert_eq!(self.dst.capacity(), max_block_size); + self.dst_start = 0; + self.dst_end = 0; + } + + // Read and decompress block + let block_info = { + let mut buffer = [0u8; 4]; + if let Err(err) = self.r.read_exact(&mut buffer) { + if err.kind() == ErrorKind::UnexpectedEof { + return Ok(0); + } else { + return Err(err); + } + } + BlockInfo::read(&buffer)? + }; + match block_info { + BlockInfo::Uncompressed(len) => { + let len = len as usize; + if len > max_block_size { + return Err(Error::BlockTooBig.into()); + } + // TODO: Attempt to avoid initialization of read buffer when + // https://github.com/rust-lang/rust/issues/42788 stabilizes + self.r.read_exact(vec_resize_and_get_mut( + &mut self.dst, + self.dst_start, + self.dst_start + len, + ))?; + if frame_info.block_checksums { + let expected_checksum = Self::read_checksum(&mut self.r)?; + Self::check_block_checksum( + &self.dst[self.dst_start..self.dst_start + len], + expected_checksum, + )?; + } + + self.dst_end += len; + self.content_len += len as u64; + } + BlockInfo::Compressed(len) => { + let len = len as usize; + if len > max_block_size { + return Err(Error::BlockTooBig.into()); + } + // TODO: Attempt to avoid initialization of read buffer when + // https://github.com/rust-lang/rust/issues/42788 stabilizes + self.r + .read_exact(vec_resize_and_get_mut(&mut self.src, 0, len))?; + if frame_info.block_checksums { + let expected_checksum = Self::read_checksum(&mut self.r)?; + Self::check_block_checksum(&self.src[..len], expected_checksum)?; + } + + let with_dict_mode = + frame_info.block_mode == BlockMode::Linked && self.ext_dict_len != 0; + let decomp_size = if with_dict_mode { + debug_assert!(self.dst_start + max_block_size <= self.ext_dict_offset); + let (head, tail) = self.dst.split_at_mut(self.ext_dict_offset); + let ext_dict = &tail[..self.ext_dict_len]; + + debug_assert!(head.len() - self.dst_start >= max_block_size); + crate::block::decompress::decompress_internal::( + &self.src[..len], + &mut SliceSink::new(head, self.dst_start), + ext_dict, + ) + } else { + // Independent blocks OR linked blocks with only prefix data + debug_assert!(self.dst.capacity() - self.dst_start >= max_block_size); + crate::block::decompress::decompress_internal::( + &self.src[..len], + &mut vec_sink_for_decompression( + &mut self.dst, + 0, + self.dst_start, + self.dst_start + max_block_size, + ), + b"", + ) + } + .map_err(Error::DecompressionError)?; + + self.dst_end += decomp_size; + self.content_len += decomp_size as u64; + } + + BlockInfo::EndMark => { + if let Some(expected) = frame_info.content_size { + if self.content_len != expected { + return Err(Error::ContentLengthError { + expected, + actual: self.content_len, + } + .into()); + } + } + if frame_info.content_checksum { + let expected_checksum = Self::read_checksum(&mut self.r)?; + let calc_checksum = self.content_hasher.finish() as u32; + if calc_checksum != expected_checksum { + return Err(Error::ContentChecksumError.into()); + } + } + self.current_frame_info = None; + return Ok(0); + } + } + + // Content checksum, if applicable + if frame_info.content_checksum { + self.content_hasher + .write(&self.dst[self.dst_start..self.dst_end]); + } + + Ok(self.dst_end - self.dst_start) + } + + fn read_more(&mut self) -> io::Result { + if self.current_frame_info.is_none() && self.read_frame_info()? == 0 { + return Ok(0); + } + self.read_block() + } +} + +impl Decoder { + /// Read the next block of decompressed data. + pub fn next_block(&mut self) -> io::Result<&[u8]> { + if self.dst_start == self.dst_end { + self.read_more()?; + } + let end = self.dst_end; + let start = self.dst_start; + self.dst_start = self.dst_end; + Ok(&self.dst[start..end]) + } +} + +impl fmt::Debug for Decoder { + fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result { + f.debug_struct("FrameDecoder") + .field("r", &self.r) + .field("content_hasher", &self.content_hasher) + .field("content_len", &self.content_len) + .field("src", &"[...]") + .field("dst", &"[...]") + .field("dst_end", &self.dst_end) + .field("ext_dict_offset", &self.ext_dict_offset) + .field("ext_dict_len", &self.ext_dict_len) + .field("current_frame_info", &self.current_frame_info) + .finish() + } +} + +/// Similar to `v.get_mut(start..end) but will adjust the len if needed. +#[inline] +fn vec_resize_and_get_mut(v: &mut Vec, start: usize, end: usize) -> &mut [u8] { + if end > v.len() { + v.resize(end, 0) + } + &mut v[start..end] +} From 1bc29ef9f2c32522866b5d275d32752dc213e788 Mon Sep 17 00:00:00 2001 From: David Roundy Date: Sun, 9 Jul 2023 04:43:07 -0700 Subject: [PATCH 2/7] remove some very reader-specific methods --- src/frame/raw_decompress.rs | 18 ------------------ 1 file changed, 18 deletions(-) diff --git a/src/frame/raw_decompress.rs b/src/frame/raw_decompress.rs index 31015723..d8767931 100644 --- a/src/frame/raw_decompress.rs +++ b/src/frame/raw_decompress.rs @@ -89,24 +89,6 @@ impl Decoder { } } - /// Gets a reference to the underlying reader in this decoder. - pub fn get_ref(&self) -> &R { - &self.r - } - - /// Gets a mutable reference to the underlying reader in this decoder. - /// - /// Note that mutation of the stream may result in surprising results if - /// this decoder is continued to be used. - pub fn get_mut(&mut self) -> &mut R { - &mut self.r - } - - /// Consumes the Decoder and returns the underlying reader. - pub fn into_inner(self) -> R { - self.r - } - fn read_frame_info(&mut self) -> Result { let mut buffer = [0u8; MAX_FRAME_INFO_SIZE]; From 3306091529ede551445b23dd3fb8a3b280e2737a Mon Sep 17 00:00:00 2001 From: David Roundy Date: Sun, 9 Jul 2023 04:49:44 -0700 Subject: [PATCH 3/7] switch to returning Error type --- src/frame/raw_decompress.rs | 17 ++++++++--------- 1 file changed, 8 insertions(+), 9 deletions(-) diff --git a/src/frame/raw_decompress.rs b/src/frame/raw_decompress.rs index d8767931..139dd395 100644 --- a/src/frame/raw_decompress.rs +++ b/src/frame/raw_decompress.rs @@ -169,7 +169,7 @@ impl Decoder { Ok(()) } - fn read_block(&mut self) -> io::Result { + fn read_block(&mut self) -> Result { debug_assert_eq!(self.dst_start, self.dst_end); let frame_info = self.current_frame_info.as_ref().unwrap(); @@ -217,7 +217,7 @@ impl Decoder { if err.kind() == ErrorKind::UnexpectedEof { return Ok(0); } else { - return Err(err); + return Err(err.into()); } } BlockInfo::read(&buffer)? @@ -226,7 +226,7 @@ impl Decoder { BlockInfo::Uncompressed(len) => { let len = len as usize; if len > max_block_size { - return Err(Error::BlockTooBig.into()); + return Err(Error::BlockTooBig); } // TODO: Attempt to avoid initialization of read buffer when // https://github.com/rust-lang/rust/issues/42788 stabilizes @@ -249,7 +249,7 @@ impl Decoder { BlockInfo::Compressed(len) => { let len = len as usize; if len > max_block_size { - return Err(Error::BlockTooBig.into()); + return Err(Error::BlockTooBig); } // TODO: Attempt to avoid initialization of read buffer when // https://github.com/rust-lang/rust/issues/42788 stabilizes @@ -299,15 +299,14 @@ impl Decoder { return Err(Error::ContentLengthError { expected, actual: self.content_len, - } - .into()); + }); } } if frame_info.content_checksum { let expected_checksum = Self::read_checksum(&mut self.r)?; let calc_checksum = self.content_hasher.finish() as u32; if calc_checksum != expected_checksum { - return Err(Error::ContentChecksumError.into()); + return Err(Error::ContentChecksumError); } } self.current_frame_info = None; @@ -324,7 +323,7 @@ impl Decoder { Ok(self.dst_end - self.dst_start) } - fn read_more(&mut self) -> io::Result { + fn read_more(&mut self) -> Result { if self.current_frame_info.is_none() && self.read_frame_info()? == 0 { return Ok(0); } @@ -334,7 +333,7 @@ impl Decoder { impl Decoder { /// Read the next block of decompressed data. - pub fn next_block(&mut self) -> io::Result<&[u8]> { + pub fn next_block(&mut self) -> Result<&[u8], Error> { if self.dst_start == self.dst_end { self.read_more()?; } From 970b09cbbd6974a4c0b4ac29fb58021806458dc5 Mon Sep 17 00:00:00 2001 From: David Roundy Date: Sun, 9 Jul 2023 04:55:53 -0700 Subject: [PATCH 4/7] simplify a bit more by removing read_more --- src/frame/raw_decompress.rs | 18 ++++++------------ 1 file changed, 6 insertions(+), 12 deletions(-) diff --git a/src/frame/raw_decompress.rs b/src/frame/raw_decompress.rs index 139dd395..31c49f33 100644 --- a/src/frame/raw_decompress.rs +++ b/src/frame/raw_decompress.rs @@ -170,6 +170,10 @@ impl Decoder { } fn read_block(&mut self) -> Result { + if self.current_frame_info.is_none() && self.read_frame_info()? == 0 { + return Ok(0); + } + debug_assert_eq!(self.dst_start, self.dst_end); let frame_info = self.current_frame_info.as_ref().unwrap(); @@ -322,25 +326,15 @@ impl Decoder { Ok(self.dst_end - self.dst_start) } - - fn read_more(&mut self) -> Result { - if self.current_frame_info.is_none() && self.read_frame_info()? == 0 { - return Ok(0); - } - self.read_block() - } } impl Decoder { /// Read the next block of decompressed data. pub fn next_block(&mut self) -> Result<&[u8], Error> { - if self.dst_start == self.dst_end { - self.read_more()?; - } - let end = self.dst_end; + self.read_block()?; let start = self.dst_start; self.dst_start = self.dst_end; - Ok(&self.dst[start..end]) + Ok(&self.dst[start..self.dst_end]) } } From de2e0c1cd829db5da67823bc6be97a248cb27a83 Mon Sep 17 00:00:00 2001 From: David Roundy Date: Sun, 9 Jul 2023 04:56:45 -0700 Subject: [PATCH 5/7] join impl blocks --- src/frame/raw_decompress.rs | 2 -- 1 file changed, 2 deletions(-) diff --git a/src/frame/raw_decompress.rs b/src/frame/raw_decompress.rs index 31c49f33..3b82bf25 100644 --- a/src/frame/raw_decompress.rs +++ b/src/frame/raw_decompress.rs @@ -326,9 +326,7 @@ impl Decoder { Ok(self.dst_end - self.dst_start) } -} -impl Decoder { /// Read the next block of decompressed data. pub fn next_block(&mut self) -> Result<&[u8], Error> { self.read_block()?; From 582c355450dc0a79538399c891b34340e92797bb Mon Sep 17 00:00:00 2001 From: David Roundy Date: Sun, 9 Jul 2023 05:38:54 -0700 Subject: [PATCH 6/7] make raw decoder no longer hold a reader --- src/frame/mod.rs | 7 ++- src/frame/raw_decompress.rs | 120 +++++++++++++++--------------------- 2 files changed, 53 insertions(+), 74 deletions(-) diff --git a/src/frame/mod.rs b/src/frame/mod.rs index 9561e20e..925f5f81 100644 --- a/src/frame/mod.rs +++ b/src/frame/mod.rs @@ -23,13 +23,13 @@ use std::{fmt, io}; pub(crate) mod compress; #[cfg_attr(feature = "safe-decode", forbid(unsafe_code))] pub(crate) mod decompress; -pub(crate) mod raw_decompress; pub(crate) mod header; +pub(crate) mod raw_decompress; pub use compress::{AutoFinishEncoder, FrameEncoder}; pub use decompress::FrameDecoder; -pub use raw_decompress::Decoder; pub use header::{BlockMode, BlockSize, FrameInfo}; +pub use raw_decompress::Decoder; #[derive(Debug)] #[non_exhaustive] @@ -71,6 +71,8 @@ pub enum Error { /// Actual content lenght. actual: u64, }, + /// The raw decompressor needs to be provided more compressed bytes. + NeedMoreInput(usize), } impl From for io::Error { @@ -80,6 +82,7 @@ impl From for io::Error { Error::CompressionError(_) | Error::DecompressionError(_) | Error::SkippableFrame(_) + | Error::NeedMoreInput(_) | Error::DictionaryNotSupported => io::Error::new(io::ErrorKind::Other, e), Error::WrongMagicNumber | Error::UnsupportedBlocksize(..) diff --git a/src/frame/raw_decompress.rs b/src/frame/raw_decompress.rs index 3b82bf25..b894b60a 100644 --- a/src/frame/raw_decompress.rs +++ b/src/frame/raw_decompress.rs @@ -1,10 +1,5 @@ -use std::{ - convert::TryInto, - fmt, - hash::Hasher, - io::{self, ErrorKind}, - mem::size_of, -}; +use alloc::collections::VecDeque; +use std::{convert::TryInto, hash::Hasher, io, mem::size_of}; use twox_hash::XxHash32; use super::header::{ @@ -46,9 +41,10 @@ use crate::{ /// } /// } /// ``` -pub struct Decoder { - /// The underlying reader. - r: R, +#[derive(Debug, Default)] +pub struct Decoder { + /// The bytes we have been given so far + raw: VecDeque, /// The FrameInfo of the frame currently being decoded. /// It starts as `None` and is filled with the FrameInfo is read from the input. /// It's reset to `None` once the frame EndMarker is read from the input. @@ -72,56 +68,54 @@ pub struct Decoder { dst_end: usize, } -impl Decoder { +impl Decoder { /// Creates a new Decoder for the specified reader. - pub fn new(rdr: R) -> Decoder { - Decoder { - r: rdr, - src: Default::default(), - dst: Default::default(), - ext_dict_offset: 0, - ext_dict_len: 0, - dst_start: 0, - dst_end: 0, - current_frame_info: None, - content_hasher: XxHash32::with_seed(0), - content_len: 0, + pub fn new() -> Decoder { + Decoder::default() + } + + /// Provide compressed data + pub fn push(&mut self, bytes: &[u8]) { + self.raw.extend(bytes); + } + fn read_exact(raw: &mut VecDeque, buffer: &mut [u8]) -> Result<(), Error> { + if raw.len() < buffer.len() { + return Err(Error::NeedMoreInput(buffer.len() - raw.len())); + } + for o in buffer.iter_mut() { + *o = raw.pop_front().unwrap(); } + Ok(()) } - fn read_frame_info(&mut self) -> Result { + fn read_frame_info(&mut self) -> Result { let mut buffer = [0u8; MAX_FRAME_INFO_SIZE]; - match self.r.read(&mut buffer[..MAGIC_NUMBER_SIZE])? { - 0 => return Ok(0), - MAGIC_NUMBER_SIZE => (), - read => self.r.read_exact(&mut buffer[read..MAGIC_NUMBER_SIZE])?, + if self.raw.is_empty() { + return Ok(0); } + Self::read_exact(&mut self.raw, &mut buffer[0..MAGIC_NUMBER_SIZE])?; if u32::from_le_bytes(buffer[0..MAGIC_NUMBER_SIZE].try_into().unwrap()) != LZ4F_LEGACY_MAGIC_NUMBER { - match self - .r - .read(&mut buffer[MAGIC_NUMBER_SIZE..MIN_FRAME_INFO_SIZE])? - { - 0 => return Ok(0), - MIN_FRAME_INFO_SIZE => (), - read => self - .r - .read_exact(&mut buffer[MAGIC_NUMBER_SIZE + read..MIN_FRAME_INFO_SIZE])?, + if self.raw.is_empty() { + return Ok(0); } + Self::read_exact( + &mut self.raw, + &mut buffer[MAGIC_NUMBER_SIZE..MIN_FRAME_INFO_SIZE], + )?; } let required = FrameInfo::read_size(&buffer[..MIN_FRAME_INFO_SIZE])?; if required != MIN_FRAME_INFO_SIZE && required != MAGIC_NUMBER_SIZE { - self.r - .read_exact(&mut buffer[MIN_FRAME_INFO_SIZE..required])?; + Self::read_exact(&mut self.raw, &mut buffer[MIN_FRAME_INFO_SIZE..required])?; } let frame_info = FrameInfo::read(&buffer[..required])?; if frame_info.dict_id.is_some() { // Unsupported right now so it must be None - return Err(Error::DictionaryNotSupported.into()); + return Err(Error::DictionaryNotSupported); } let max_block_size = frame_info.block_size.get_size(); @@ -151,9 +145,9 @@ impl Decoder { } #[inline] - fn read_checksum(r: &mut R) -> Result { + fn read_checksum(&mut self) -> Result { let mut checksum_buffer = [0u8; size_of::()]; - r.read_exact(&mut checksum_buffer[..])?; + Self::read_exact(&mut self.raw, &mut checksum_buffer[..])?; let checksum = u32::from_le_bytes(checksum_buffer); Ok(checksum) } @@ -175,7 +169,7 @@ impl Decoder { } debug_assert_eq!(self.dst_start, self.dst_end); - let frame_info = self.current_frame_info.as_ref().unwrap(); + let frame_info = self.current_frame_info.clone().unwrap(); // Adjust dst buffer offsets to decompress the next block let max_block_size = frame_info.block_size.get_size(); @@ -217,11 +211,11 @@ impl Decoder { // Read and decompress block let block_info = { let mut buffer = [0u8; 4]; - if let Err(err) = self.r.read_exact(&mut buffer) { - if err.kind() == ErrorKind::UnexpectedEof { + if let Err(err) = Self::read_exact(&mut self.raw, &mut buffer) { + if let Error::NeedMoreInput(_) = err { return Ok(0); } else { - return Err(err.into()); + return Err(err); } } BlockInfo::read(&buffer)? @@ -234,13 +228,12 @@ impl Decoder { } // TODO: Attempt to avoid initialization of read buffer when // https://github.com/rust-lang/rust/issues/42788 stabilizes - self.r.read_exact(vec_resize_and_get_mut( - &mut self.dst, - self.dst_start, - self.dst_start + len, - ))?; + Self::read_exact( + &mut self.raw, + vec_resize_and_get_mut(&mut self.dst, self.dst_start, self.dst_start + len), + )?; if frame_info.block_checksums { - let expected_checksum = Self::read_checksum(&mut self.r)?; + let expected_checksum = self.read_checksum()?; Self::check_block_checksum( &self.dst[self.dst_start..self.dst_start + len], expected_checksum, @@ -257,10 +250,9 @@ impl Decoder { } // TODO: Attempt to avoid initialization of read buffer when // https://github.com/rust-lang/rust/issues/42788 stabilizes - self.r - .read_exact(vec_resize_and_get_mut(&mut self.src, 0, len))?; + Self::read_exact(&mut self.raw, vec_resize_and_get_mut(&mut self.src, 0, len))?; if frame_info.block_checksums { - let expected_checksum = Self::read_checksum(&mut self.r)?; + let expected_checksum = self.read_checksum()?; Self::check_block_checksum(&self.src[..len], expected_checksum)?; } @@ -307,7 +299,7 @@ impl Decoder { } } if frame_info.content_checksum { - let expected_checksum = Self::read_checksum(&mut self.r)?; + let expected_checksum = self.read_checksum()?; let calc_checksum = self.content_hasher.finish() as u32; if calc_checksum != expected_checksum { return Err(Error::ContentChecksumError); @@ -336,22 +328,6 @@ impl Decoder { } } -impl fmt::Debug for Decoder { - fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result { - f.debug_struct("FrameDecoder") - .field("r", &self.r) - .field("content_hasher", &self.content_hasher) - .field("content_len", &self.content_len) - .field("src", &"[...]") - .field("dst", &"[...]") - .field("dst_end", &self.dst_end) - .field("ext_dict_offset", &self.ext_dict_offset) - .field("ext_dict_len", &self.ext_dict_len) - .field("current_frame_info", &self.current_frame_info) - .finish() - } -} - /// Similar to `v.get_mut(start..end) but will adjust the len if needed. #[inline] fn vec_resize_and_get_mut(v: &mut Vec, start: usize, end: usize) -> &mut [u8] { From bfa4fc8e57dff12249a47e5dceca9af275073e41 Mon Sep 17 00:00:00 2001 From: David Roundy Date: Sun, 16 Jul 2023 05:03:28 -0700 Subject: [PATCH 7/7] add docs --- src/frame/raw_decompress.rs | 37 ++++++++++++++----------------------- 1 file changed, 14 insertions(+), 23 deletions(-) diff --git a/src/frame/raw_decompress.rs b/src/frame/raw_decompress.rs index b894b60a..82a36ed9 100644 --- a/src/frame/raw_decompress.rs +++ b/src/frame/raw_decompress.rs @@ -12,34 +12,21 @@ use crate::{ sink::{vec_sink_for_decompression, SliceSink}, }; -/// A reader for decompressing the LZ4 frame format -/// -/// This Decoder wraps any other reader that implements `io::Read`. -/// Bytes read will be decompressed according to the [LZ4 frame format]( -/// https://github.com/lz4/lz4/blob/dev/doc/lz4_Frame_format.md). +/// A struct for decompressing the LZ4 frame format given pushed byte data /// /// # Example 1 /// Deserializing json values out of a compressed file. /// /// ```no_run -/// let compressed_input = std::fs::File::open("datafile").unwrap(); -/// let mut decompressed_input = lz4_flex::frame::FrameDecoder::new(compressed_input); -/// let json: serde_json::Value = serde_json::from_reader(decompressed_input).unwrap(); -/// ``` -/// -/// # Example -/// Deserializing multiple json values out of a compressed file -/// -/// ```no_run -/// let compressed_input = std::fs::File::open("datafile").unwrap(); -/// let mut decompressed_input = lz4_flex::frame::FrameDecoder::new(compressed_input); -/// loop { -/// match serde_json::from_reader::<_, serde_json::Value>(&mut decompressed_input) { -/// Ok(json) => { println!("json {:?}", json); } -/// Err(e) if e.is_eof() => break, -/// Err(e) => panic!("{}", e), -/// } +/// let mut decoder = lz4_flex::frame::Decoder::new(); +/// let compressed_input: Vec = std::fs::read("datafile").unwrap(); +/// decoder.push(&compressed_input); +/// let mut decompressed = Vec::new(); +/// while let Ok(v) = decoder.next_block() { +/// decompressed.extend(v); /// } +/// let json: serde_json::Value = serde_json::from_slice(&decompressed).unwrap(); +/// ```\ /// ``` #[derive(Debug, Default)] pub struct Decoder { @@ -74,7 +61,7 @@ impl Decoder { Decoder::default() } - /// Provide compressed data + /// Provide compressed data to decompress pub fn push(&mut self, bytes: &[u8]) { self.raw.extend(bytes); } @@ -320,6 +307,10 @@ impl Decoder { } /// Read the next block of decompressed data. + /// + /// When using this function, unless it is known that all data has already + /// been pushed, the user should check for an `Error::NeedMoreData` to see + /// if more data is needed. pub fn next_block(&mut self) -> Result<&[u8], Error> { self.read_block()?; let start = self.dst_start;