From d6134edc69322a8df11c69d1de26bd1807cbe44f Mon Sep 17 00:00:00 2001 From: Pure White Date: Wed, 14 Sep 2022 14:13:11 +0800 Subject: [PATCH] feat(frame): add reset API for frame codec --- bench/Cargo.toml | 1 + bench/src/bench.rs | 332 ++++++++++++++++++++++++++++++++++++++++++++- src/compress.rs | 9 ++ src/read.rs | 48 ++++++- test/tests.rs | 24 ++++ 5 files changed, 411 insertions(+), 3 deletions(-) diff --git a/bench/Cargo.toml b/bench/Cargo.toml index fa429c6..e6c14b1 100644 --- a/bench/Cargo.toml +++ b/bench/Cargo.toml @@ -19,6 +19,7 @@ path = "src/bench.rs" [features] cpp = ["snappy-cpp"] +reuse = [] [dependencies] criterion = "0.3.1" diff --git a/bench/src/bench.rs b/bench/src/bench.rs index 7a9b08e..8e10b40 100644 --- a/bench/src/bench.rs +++ b/bench/src/bench.rs @@ -1,8 +1,9 @@ -use std::time::Duration; +use std::{io::Read, time::Duration}; use criterion::{ criterion_group, criterion_main, Bencher, Benchmark, Criterion, Throughput, }; +use snap::read::FrameEncoder; const CORPUS_HTML: &'static [u8] = include_bytes!("../../data/html"); const CORPUS_URLS_10K: &'static [u8] = include_bytes!("../../data/urls.10K"); @@ -21,6 +22,8 @@ const CORPUS_GEOPROTO: &'static [u8] = include_bytes!("../../data/geo.protodata"); const CORPUS_KPPKN: &'static [u8] = include_bytes!("../../data/kppkn.gtb"); +const FRAME_HEADER_SIZE: usize = 16; + macro_rules! compress { ($c:expr, $comp:expr, $group:expr, $name:expr, $corpus:expr) => { compress!($c, $comp, $group, $name, $corpus, 0); @@ -39,6 +42,66 @@ macro_rules! compress { }; } +macro_rules! frame_compress { + ($c:expr, $comp:expr, $group:expr, $name:expr, $corpus:expr) => { + frame_compress!($c, $comp, $group, $name, $corpus, 0); + }; + ($c:expr, $comp:expr, $group:expr, $name:expr, $corpus:expr, $size:expr) => { + let mut corpus = $corpus; + if $size > 0 { + corpus = &corpus[..$size]; + } + let mut dst = vec![ + 0; + snap::raw::max_compress_len(corpus.len()) + + FRAME_HEADER_SIZE + ]; + define( + $c, + $group, + &format!("frame_compress/{}", $name), + corpus, + move |b| { + b.iter(|| { + dst.clear(); + $comp(corpus, &mut dst); + }); + }, + ); + }; +} + +macro_rules! frame_compress_reuse { + ($c:expr, $group:expr, $name:expr, $corpus:expr) => { + frame_compress_reuse!($c, $group, $name, $corpus, 0); + }; + ($c:expr, $group:expr, $name:expr, $corpus:expr, $size:expr) => { + let mut corpus = $corpus; + if $size > 0 { + corpus = &corpus[..$size]; + } + let mut dst = vec![ + 0; + snap::raw::max_compress_len(corpus.len()) + + FRAME_HEADER_SIZE + ]; + let mut encoder = snap::read::FrameEncoder::new(corpus); + define( + $c, + $group, + &format!("frame_compress_reuse/{}", $name), + corpus, + move |b| { + b.iter(|| { + dst.clear(); + encoder.reset(corpus); + encoder.read_to_end(&mut dst).unwrap(); + }); + }, + ); + }; +} + macro_rules! decompress { ($c:expr, $decomp:expr, $group:expr, $name:expr, $corpus:expr) => { decompress!($c, $decomp, $group, $name, $corpus, 0); @@ -65,6 +128,71 @@ macro_rules! decompress { }; } +macro_rules! frame_decompress { + ($c:expr, $decomp:expr, $group:expr, $name:expr, $corpus:expr) => { + frame_decompress!($c, $decomp, $group, $name, $corpus, 0); + }; + ($c:expr, $decomp:expr, $group:expr, $name:expr, $corpus:expr, $size:expr) => { + let mut corpus = $corpus; + if $size > 0 { + corpus = &corpus[..$size]; + } + let mut compressed = Vec::with_capacity( + snap::raw::max_compress_len(corpus.len()) + FRAME_HEADER_SIZE, + ); + snap::read::FrameEncoder::new(corpus) + .read_to_end(&mut compressed) + .unwrap(); + let mut dst = vec![0; corpus.len() + FRAME_HEADER_SIZE]; + define( + $c, + $group, + &format!("frame_decompress/{}", $name), + corpus, + move |b| { + b.iter(|| { + dst.clear(); + $decomp(&compressed, &mut dst); + }); + }, + ); + }; +} + +macro_rules! frame_decompress_reuse { + ($c:expr, $group:expr, $name:expr, $corpus:expr) => { + frame_decompress_reuse!($c, $group, $name, $corpus, 0); + }; + ($c:expr, $group:expr, $name:expr, $corpus:expr, $size:expr) => { + let mut corpus = $corpus; + if $size > 0 { + corpus = &corpus[..$size]; + } + let mut compressed = Vec::with_capacity( + snap::raw::max_compress_len(corpus.len()) + FRAME_HEADER_SIZE, + ); + snap::read::FrameEncoder::new(corpus) + .read_to_end(&mut compressed) + .unwrap(); + let mut dst = vec![0; corpus.len() + FRAME_HEADER_SIZE]; + define( + $c, + $group, + &format!("frame_decompress_reuse/{}", $name), + corpus, + move |b| { + let src: &[u8] = &compressed; + let mut decoder = snap::read::FrameDecoder::new(src); + b.iter(|| { + dst.clear(); + decoder.reset(src); + decoder.read_to_end(&mut dst).unwrap(); + }); + }, + ); + }; +} + fn all(c: &mut Criterion) { rust(c); #[cfg(feature = "cpp")] @@ -76,10 +204,20 @@ fn rust(c: &mut Criterion) { snap::raw::Encoder::new().compress(input, output) } + fn frame_compress(input: &[u8], output: &mut Vec) { + let mut encoder = snap::read::FrameEncoder::new(input); + encoder.read_to_end(output).unwrap(); + } + fn decompress(input: &[u8], output: &mut [u8]) -> snap::Result { snap::raw::Decoder::new().decompress(input, output) } + fn frame_decompress(input: &[u8], output: &mut Vec) { + let mut decoder = snap::read::FrameDecoder::new(input); + decoder.read_to_end(output).unwrap(); + } + compress!(c, compress, "snap", "zflat00_html", CORPUS_HTML); compress!(c, compress, "snap", "zflat01_urls", CORPUS_URLS_10K); compress!(c, compress, "snap", "zflat02_jpg", CORPUS_FIREWORKS); @@ -93,6 +231,90 @@ fn rust(c: &mut Criterion) { compress!(c, compress, "snap", "zflat10_pb", CORPUS_GEOPROTO); compress!(c, compress, "snap", "zflat11_gaviota", CORPUS_KPPKN); + frame_compress!(c, frame_compress, "snap", "zflat00_html", CORPUS_HTML); + frame_compress!( + c, + frame_compress, + "snap", + "zflat01_urls", + CORPUS_URLS_10K + ); + frame_compress!( + c, + frame_compress, + "snap", + "zflat02_jpg", + CORPUS_FIREWORKS + ); + frame_compress!( + c, + frame_compress, + "snap", + "zflat03_jpg_200", + CORPUS_FIREWORKS, + 200 + ); + frame_compress!( + c, + frame_compress, + "snap", + "zflat04_pdf", + CORPUS_PAPER_100K + ); + frame_compress!( + c, + frame_compress, + "snap", + "zflat05_html4", + CORPUS_HTML_X_4 + ); + frame_compress!(c, frame_compress, "snap", "zflat06_txt1", CORPUS_ALICE29); + frame_compress!( + c, + frame_compress, + "snap", + "zflat07_txt2", + CORPUS_ASYOULIK + ); + frame_compress!(c, frame_compress, "snap", "zflat08_txt3", CORPUS_LCET10); + frame_compress!( + c, + frame_compress, + "snap", + "zflat09_txt4", + CORPUS_PLRABN12 + ); + frame_compress!(c, frame_compress, "snap", "zflat10_pb", CORPUS_GEOPROTO); + frame_compress!( + c, + frame_compress, + "snap", + "zflat11_gaviota", + CORPUS_KPPKN + ); + + #[cfg(feature = "reuse")] + { + frame_compress_reuse!(c, "snap", "zflat00_html", CORPUS_HTML); + frame_compress_reuse!(c, "snap", "zflat01_urls", CORPUS_URLS_10K); + frame_compress_reuse!(c, "snap", "zflat02_jpg", CORPUS_FIREWORKS); + frame_compress_reuse!( + c, + "snap", + "zflat03_jpg_200", + CORPUS_FIREWORKS, + 200 + ); + frame_compress_reuse!(c, "snap", "zflat04_pdf", CORPUS_PAPER_100K); + frame_compress_reuse!(c, "snap", "zflat05_html4", CORPUS_HTML_X_4); + frame_compress_reuse!(c, "snap", "zflat06_txt1", CORPUS_ALICE29); + frame_compress_reuse!(c, "snap", "zflat07_txt2", CORPUS_ASYOULIK); + frame_compress_reuse!(c, "snap", "zflat08_txt3", CORPUS_LCET10); + frame_compress_reuse!(c, "snap", "zflat09_txt4", CORPUS_PLRABN12); + frame_compress_reuse!(c, "snap", "zflat10_pb", CORPUS_GEOPROTO); + frame_compress_reuse!(c, "snap", "zflat11_gaviota", CORPUS_KPPKN); + } + decompress!(c, decompress, "snap", "uflat00_html", CORPUS_HTML); decompress!(c, decompress, "snap", "uflat01_urls", CORPUS_URLS_10K); decompress!(c, decompress, "snap", "uflat02_jpg", CORPUS_FIREWORKS); @@ -112,6 +334,114 @@ fn rust(c: &mut Criterion) { decompress!(c, decompress, "snap", "uflat09_txt4", CORPUS_PLRABN12); decompress!(c, decompress, "snap", "uflat10_pb", CORPUS_GEOPROTO); decompress!(c, decompress, "snap", "uflat11_gaviota", CORPUS_KPPKN); + + frame_decompress!( + c, + frame_decompress, + "snap", + "uflat00_html", + CORPUS_HTML + ); + frame_decompress!( + c, + frame_decompress, + "snap", + "uflat01_urls", + CORPUS_URLS_10K + ); + frame_decompress!( + c, + frame_decompress, + "snap", + "uflat02_jpg", + CORPUS_FIREWORKS + ); + frame_decompress!( + c, + frame_decompress, + "snap", + "uflat03_jpg_200", + CORPUS_FIREWORKS, + 200 + ); + frame_decompress!( + c, + frame_decompress, + "snap", + "uflat04_pdf", + CORPUS_PAPER_100K + ); + frame_decompress!( + c, + frame_decompress, + "snap", + "uflat05_html4", + CORPUS_HTML_X_4 + ); + frame_decompress!( + c, + frame_decompress, + "snap", + "uflat06_txt1", + CORPUS_ALICE29 + ); + frame_decompress!( + c, + frame_decompress, + "snap", + "uflat07_txt2", + CORPUS_ASYOULIK + ); + frame_decompress!( + c, + frame_decompress, + "snap", + "uflat08_txt3", + CORPUS_LCET10 + ); + frame_decompress!( + c, + frame_decompress, + "snap", + "uflat09_txt4", + CORPUS_PLRABN12 + ); + frame_decompress!( + c, + frame_decompress, + "snap", + "uflat10_pb", + CORPUS_GEOPROTO + ); + frame_decompress!( + c, + frame_decompress, + "snap", + "uflat11_gaviota", + CORPUS_KPPKN + ); + + #[cfg(feature = "reuse")] + { + frame_decompress_reuse!(c, "snap", "uflat00_html", CORPUS_HTML); + frame_decompress_reuse!(c, "snap", "uflat01_urls", CORPUS_URLS_10K); + frame_decompress_reuse!(c, "snap", "uflat02_jpg", CORPUS_FIREWORKS); + frame_decompress_reuse!( + c, + "snap", + "uflat03_jpg_200", + CORPUS_FIREWORKS, + 200 + ); + frame_decompress_reuse!(c, "snap", "uflat04_pdf", CORPUS_PAPER_100K); + frame_decompress_reuse!(c, "snap", "uflat05_html4", CORPUS_HTML_X_4); + frame_decompress_reuse!(c, "snap", "uflat06_txt1", CORPUS_ALICE29); + frame_decompress_reuse!(c, "snap", "uflat07_txt2", CORPUS_ASYOULIK); + frame_decompress_reuse!(c, "snap", "uflat08_txt3", CORPUS_LCET10); + frame_decompress_reuse!(c, "snap", "uflat09_txt4", CORPUS_PLRABN12); + frame_decompress_reuse!(c, "snap", "uflat10_pb", CORPUS_GEOPROTO); + frame_decompress_reuse!(c, "snap", "uflat11_gaviota", CORPUS_KPPKN); + } } #[cfg(feature = "cpp")] diff --git a/src/compress.rs b/src/compress.rs index 1a6638d..2a19254 100644 --- a/src/compress.rs +++ b/src/compress.rs @@ -167,6 +167,15 @@ impl Encoder { buf.truncate(n); Ok(buf) } + + /// Resets the internal state of this encoder. + /// + /// This can make the encoder reusable and thus reduce the overhead of + /// memory allocation. + pub fn reset(&mut self) { + // small doesn't need to be reset, since it's reset when it's used. + self.big.clear(); + } } struct Block<'s, 'd> { diff --git a/src/read.rs b/src/read.rs index a924bf9..39a3f9e 100644 --- a/src/read.rs +++ b/src/read.rs @@ -94,6 +94,24 @@ impl FrameDecoder { pub fn get_mut(&mut self) -> &mut R { &mut self.r } + + /// Resets the internal state of this decoder and sets a new underlying + /// reader. + /// + /// This can make the decoder reusable and thus reduce the overhead of + /// memory allocation. + pub fn reset(&mut self, rdr: R) { + self.r = rdr; + unsafe { + self.src.set_len(MAX_COMPRESS_BLOCK_SIZE); + } + unsafe { + self.dst.set_len(MAX_BLOCK_SIZE); + } + self.dsts = 0; + self.dste = 0; + self.read_stream_ident = false; + } } impl io::Read for FrameDecoder { @@ -130,12 +148,12 @@ impl io::Read for FrameDecoder { } let len = len64 as usize; match ty { - Err(b) if 0x02 <= b && b <= 0x7F => { + Err(b) if (0x02..=0x7F).contains(&b) => { // Spec says that chunk types 0x02-0x7F are reserved and // conformant decoders must return an error. fail!(Error::UnsupportedChunkType { byte: b }); } - Err(b) if 0x80 <= b && b <= 0xFD => { + Err(b) if (0x80..=0xFD).contains(&b) => { // Spec says that chunk types 0x80-0xFD are reserved but // skippable. self.r.read_exact(&mut self.src[0..len])?; @@ -331,6 +349,18 @@ impl FrameEncoder { self.dsts += count; count } + + /// Resets the internal state of this encoder and sets a new underlying + /// reader. + /// + /// This can make the encoder reusable and thus reduce the overhead of + /// memory allocation. + pub fn reset(&mut self, rdr: R) { + self.inner.reset(rdr); + unsafe { self.dst.set_len(MAX_READ_FRAME_ENCODER_BLOCK_SIZE) }; + self.dsts = 0; + self.dste = 0; + } } impl io::Read for FrameEncoder { @@ -402,6 +432,20 @@ impl Inner { )?; Ok(dst_write_start + frame_data.len()) } + + /// Resets the internal state of this inner encoder and sets a new underlying + /// reader. + /// + /// This can make the inner encoder reusable and thus reduce the overhead of + /// memory allocation. + pub fn reset(&mut self, rdr: R) { + self.r = rdr; + self.enc.reset(); + unsafe { + self.src.set_len(MAX_BLOCK_SIZE); + } + self.wrote_stream_ident = false; + } } impl fmt::Debug for FrameEncoder { diff --git a/test/tests.rs b/test/tests.rs index e23dafa..e930692 100644 --- a/test/tests.rs +++ b/test/tests.rs @@ -495,6 +495,30 @@ fn qc_roundtrip_stream() { .quickcheck(p as fn(_) -> _); } +#[test] +fn qc_roundtrip_stream_reuse() { + fn p(bytes: Vec) -> TestResult { + if bytes.is_empty() { + return TestResult::discard(); + } + use snap::read; + use std::io::Read; + + let compressed = write_frame_press(&bytes); + let mut decoder = read::FrameDecoder::new(&compressed[..]); + let mut buf = vec![]; + decoder.read_to_end(&mut buf).unwrap(); + buf.clear(); + decoder.reset(&compressed); + decoder.read_to_end(&mut buf).unwrap(); + TestResult::from_bool(buf == bytes) + } + QuickCheck::new() + .gen(StdGen::new(rand::thread_rng(), 10_000)) + .tests(1_000) + .quickcheck(p as fn(_) -> _); +} + #[test] fn test_short_input() { // Regression test for https://github.com/BurntSushi/rust-snappy/issues/42