use crate::decoding::StreamingDecoder;
use crate::encoding::{
CompressionContext, CompressionLevel, EncoderDictionary, Matcher, Sequence, StreamingEncoder,
};
#[test]
fn a_reused_context_writes_the_frames_fresh_encoders_write() {
fn write_all(context: &mut CompressionContext, frame: &mut Vec<u8>, mut data: &[u8]) {
while !data.is_empty() {
let taken = context.write(frame, data).expect("write");
data = &data[taken..];
}
}
let mut state = 0x9E37_79B9u32;
let mut noise = |len: usize, alphabet: u32| -> Vec<u8> {
(0..len)
.map(|_| {
state = state.wrapping_mul(1_664_525).wrapping_add(1_013_904_223);
((state >> 24) % alphabet) as u8
})
.collect()
};
let text: Vec<u8> = (0..3_000u32)
.flat_map(|i| alloc::format!("row {} key {} val {}\n", i % 97, i % 13, i % 7).into_bytes())
.collect();
let lines: Vec<u8> = (0..4_000u32)
.flat_map(|i| alloc::format!("line {} of {}\n", i % 89, i % 7).into_bytes())
.collect();
let payloads: Vec<Vec<u8>> = vec![
text[..5_000].to_vec(),
noise(20_000, 256),
Vec::new(),
text.clone(),
noise(9_000, 16),
text[1_000..1_700].to_vec(),
text.clone(),
text[2_000..2_600].to_vec(),
text[..40_000].to_vec(),
lines[..3_000].to_vec(),
lines.clone(),
];
let dictionary = EncoderDictionary::from_serialized_or_raw_content(&text[..8_192])
.expect("raw content is a dictionary");
let mut diverged = Vec::new();
for level in [-3, 1, 2, 3, 5, 6, 9, 12, 14, 16, 17, 19, 22] {
for with_dictionary in [false, true] {
let compression_level = CompressionLevel::from_level(level);
let mut context = CompressionContext::new(compression_level);
context.set_content_checksum(true).unwrap();
if with_dictionary {
context.set_encoder_dictionary(dictionary.clone()).unwrap();
}
for (index, payload) in payloads.iter().enumerate() {
let pledged = index % 2 == 0;
let mut fresh = StreamingEncoder::new(Vec::new(), compression_level);
fresh.set_content_checksum(true).unwrap();
if with_dictionary {
fresh.set_encoder_dictionary(dictionary.clone()).unwrap();
}
if pledged {
fresh
.set_pledged_content_size(payload.len() as u64)
.unwrap();
}
fresh.write_all(payload).unwrap();
let expected = fresh.finish().unwrap();
let mut frame = Vec::new();
if pledged {
context
.set_pledged_content_size(payload.len() as u64)
.unwrap();
}
write_all(&mut context, &mut frame, payload);
context.finish_frame(&mut frame).unwrap();
if frame != expected {
diverged.push(alloc::format!(
"level {level}, dictionary {with_dictionary}, frame {index}: \
{} bytes reused against {} fresh",
frame.len(),
expected.len()
));
}
}
}
}
assert!(diverged.is_empty(), "{diverged:#?}");
}
#[test]
fn a_frame_that_did_not_finish_is_not_continued_by_the_next_encoder() {
let level = CompressionLevel::Default;
let payload = b"the next frame, whole and on its own".repeat(64);
let fresh = {
let mut encoder = StreamingEncoder::new(Vec::new(), level);
encoder.write_all(&payload).unwrap();
encoder.finish().unwrap()
};
let mut context = CompressionContext::new(level);
let next_frame = |context: &mut CompressionContext, case: &str| {
let mut encoder = StreamingEncoder::with_context(Vec::new(), context);
encoder.write_all(&payload).unwrap();
let frame = encoder.finish().unwrap();
assert!(frame == fresh, "{case}: the next frame is not a fresh one");
};
let mut encoder = StreamingEncoder::with_context(Vec::new(), &mut context);
encoder.set_pledged_content_size(1000).unwrap();
encoder.write_all(&[7u8; 600]).unwrap();
assert!(encoder.finish().is_err());
next_frame(&mut context, "short of the pledge");
let mut encoder = StreamingEncoder::with_context(Vec::new(), &mut context);
encoder.write_all(&[9u8; 300 * 1024]).unwrap();
drop(encoder);
next_frame(&mut context, "dropped mid-frame");
let mut encoder = StreamingEncoder::with_context(Vec::new(), &mut context);
encoder.set_pledged_content_size(5).unwrap();
assert!(encoder.finish().is_err());
next_frame(&mut context, "pledged and never written");
}
#[test]
fn encoder_settings_reach_its_context() {
let dict_raw = include_bytes!("../../../dict_tests/dictionary");
let payload: Vec<u8> = (0..400u32)
.flat_map(|i| alloc::format!("tenant=demo table=orders key={i} region=eu\n").into_bytes())
.collect();
type Setting = fn(&mut CompressionContext) -> Result<(), crate::io::Error>;
type EncoderSetting = fn(&mut StreamingEncoder<Vec<u8>>) -> Result<(), crate::io::Error>;
let settings: [(&str, Setting, EncoderSetting); 3] = [
(
"target block size",
|context| context.set_target_block_size(Some(2048)),
|encoder| encoder.set_target_block_size(Some(2048)),
),
(
"content size flag",
|context| context.set_content_size_flag(false),
|encoder| encoder.set_content_size_flag(false),
),
(
"dictionary ID flag",
|context| context.set_dictionary_id_flag(false),
|encoder| encoder.set_dictionary_id_flag(false),
),
];
let through_context = |setting: Setting| {
let mut context = CompressionContext::new(CompressionLevel::Default);
context.set_dictionary_from_bytes(dict_raw).unwrap();
assert!(context.dictionary().is_some());
setting(&mut context).unwrap();
let mut frame = Vec::new();
context
.set_pledged_content_size(payload.len() as u64)
.unwrap();
context.write(&mut frame, &payload).unwrap();
context.finish_frame(&mut frame).unwrap();
frame
};
for (name, setting, encoder_setting) in settings {
let mut encoder = StreamingEncoder::new(Vec::new(), CompressionLevel::Default);
encoder.set_dictionary_from_bytes(dict_raw).unwrap();
encoder_setting(&mut encoder).unwrap();
encoder
.set_pledged_content_size(payload.len() as u64)
.unwrap();
encoder.write_all(&payload).unwrap();
let frame = encoder.finish().unwrap();
assert!(
frame == through_context(setting),
"{name}: not the context's frame"
);
assert!(
frame != through_context(|_| Ok(())),
"{name}: the setting changed nothing"
);
}
}
#[test]
fn consecutive_unsized_frames_reuse_a_btultra2_context_like_fresh_ones() {
let lines: Vec<u8> = (0..4_000u32)
.flat_map(|i| alloc::format!("line {} of {}\n", i % 89, i % 7).into_bytes())
.collect();
let mut state = 0x2545_F491u32;
let mut noise = |len: usize| -> Vec<u8> {
(0..len)
.map(|_| {
state = state.wrapping_mul(1_664_525).wrapping_add(1_013_904_223);
b'a' + ((state >> 24) % 16) as u8
})
.collect()
};
let payloads = [noise(50_000), noise(60_000), lines, noise(30_000)];
let mut diverged = Vec::new();
for level in [19, 20, 22] {
let level = CompressionLevel::from_level(level);
let mut context = CompressionContext::new(level);
for (index, payload) in payloads.iter().enumerate() {
let mut fresh = StreamingEncoder::new(Vec::new(), level);
fresh.write_all(payload).unwrap();
let expected = fresh.finish().unwrap();
let mut frame = Vec::new();
let mut rest = payload.as_slice();
while !rest.is_empty() {
let taken = context.write(&mut frame, rest).unwrap();
rest = &rest[taken..];
}
context.finish_frame(&mut frame).unwrap();
if frame != expected {
diverged.push(alloc::format!(
"{level:?}, frame {index}: {} bytes reused against {} fresh",
frame.len(),
expected.len()
));
}
}
}
assert!(diverged.is_empty(), "{diverged:#?}");
}
#[test]
fn a_replaced_dictionary_leaves_nothing_of_the_old_one_behind() {
let text: Vec<u8> = (0..3_000u32)
.flat_map(|i| alloc::format!("row {} key {} val {}\n", i % 97, i % 13, i % 7).into_bytes())
.collect();
let other: Vec<u8> = text.iter().rev().copied().collect();
let first = EncoderDictionary::from_serialized_or_raw_content(&text[..8_192]).unwrap();
let second = EncoderDictionary::from_serialized_or_raw_content(&other[..8_192]).unwrap();
let payloads = [&text[..], &text[..900]];
let mut diverged = Vec::new();
for level in [1, 3, 5, 12, 16, 19] {
let level = CompressionLevel::from_level(level);
let mut context = CompressionContext::new(level);
context.set_encoder_dictionary(first.clone()).unwrap();
for payload in payloads {
context
.set_pledged_content_size(payload.len() as u64)
.unwrap();
context.write(&mut Vec::new(), payload).unwrap();
context.finish_frame(&mut Vec::new()).unwrap();
}
for replacement in [Some(&second), None] {
match replacement {
Some(dictionary) => context.set_encoder_dictionary(dictionary.clone()).unwrap(),
None => context.set_dictionary_from_bytes(&[]).unwrap(),
}
for (index, payload) in payloads.into_iter().enumerate() {
let mut fresh = StreamingEncoder::new(Vec::new(), level);
if let Some(dictionary) = replacement {
fresh.set_encoder_dictionary(dictionary.clone()).unwrap();
}
fresh
.set_pledged_content_size(payload.len() as u64)
.unwrap();
fresh.write_all(payload).unwrap();
let expected = fresh.finish().unwrap();
let mut frame = Vec::new();
context
.set_pledged_content_size(payload.len() as u64)
.unwrap();
context.write(&mut frame, payload).unwrap();
context.finish_frame(&mut frame).unwrap();
if frame != expected {
diverged.push(alloc::format!(
"{level:?}, dictionary {}, frame {index}: {} bytes reused against {} fresh",
replacement.is_some(),
frame.len(),
expected.len()
));
}
}
}
}
assert!(diverged.is_empty(), "{diverged:#?}");
}
#[test]
fn a_new_level_on_a_reused_context_drops_the_old_tuning() {
use crate::encoding::{CompressionParameters, Strategy};
let text: Vec<u8> = (0..3_000u32)
.flat_map(|i| alloc::format!("row {} key {} val {}\n", i % 97, i % 13, i % 7).into_bytes())
.collect();
let tuned = CompressionParameters::builder(CompressionLevel::from_level(3))
.strategy(Strategy::Btopt)
.window_log(17)
.build()
.unwrap();
let mut context = CompressionContext::new(CompressionLevel::from_level(3));
let mut diverged = Vec::new();
for level in [19, 1, 12, -2, 5, 22, 3] {
context.set_parameters(&tuned).unwrap();
context.write(&mut Vec::new(), &text).unwrap();
context.finish_frame(&mut Vec::new()).unwrap();
let level = CompressionLevel::from_level(level);
context.set_compression_level(level).unwrap();
let mut frame = Vec::new();
context.write(&mut frame, &text).unwrap();
context.finish_frame(&mut frame).unwrap();
let mut fresh = StreamingEncoder::new(Vec::new(), level);
fresh.write_all(&text).unwrap();
let expected = fresh.finish().unwrap();
if frame != expected {
diverged.push(alloc::format!(
"{level:?}: {} bytes reused against {} fresh",
frame.len(),
expected.len()
));
}
}
assert!(diverged.is_empty(), "{diverged:#?}");
}
#[test]
fn the_reported_footprint_covers_what_compressing_retained() {
let mut enc = StreamingEncoder::new(Vec::new(), CompressionLevel::Level(3));
let before = enc.heap_size();
let payload: Vec<u8> = (0..300_000u32)
.map(|i| ((i.wrapping_mul(2654435761) >> 24) % 32) as u8)
.collect();
enc.write_all(&payload).expect("write");
enc.flush().expect("flush");
let after = enc.heap_size();
assert!(
after > before,
"compressing retained buffers the footprint does not report: {before} -> {after}",
);
let scratch = core::mem::take(&mut enc.context.state.huff_weights);
let scratch_heap = scratch.heap_size();
assert!(
scratch_heap > 0,
"the weight builder's buffers should be populated after compressing",
);
let without_scratch = enc.heap_size();
assert_eq!(
after - without_scratch,
scratch_heap,
"the weight scratch is retained across blocks but is not counted in \
the reported footprint",
);
enc.context.state.huff_weights = scratch;
assert_eq!(enc.heap_size(), after, "restoring must undo the removal");
let entropy = enc.context.state.fse_tables.heap_size()
+ enc
.context
.state
.huff_rollback
.as_ref()
.map_or(0, |table| table.heap_size())
+ enc
.context
.state
.block_scratch
.huff_rollback
.as_ref()
.map_or(0, |table| table.heap_size());
assert!(
entropy > 0,
"compressing should have left entropy state retained",
);
let with_entropy = enc.heap_size();
let saved_tables = core::mem::replace(
&mut enc.context.state.fse_tables,
crate::encoding::frame_compressor::FseTables::new(),
);
let saved_rollback = enc.context.state.huff_rollback.take();
let saved_scratch_rollback = enc.context.state.block_scratch.huff_rollback.take();
assert_eq!(
with_entropy - enc.heap_size(),
entropy,
"the entropy state is retained but not counted in the reported footprint",
);
enc.context.state.fse_tables = saved_tables;
enc.context.state.huff_rollback = saved_rollback;
enc.context.state.block_scratch.huff_rollback = saved_scratch_rollback;
assert_eq!(enc.heap_size(), with_entropy, "restoring must undo removal");
}
use crate::io::{Error, ErrorKind, Read, Write};
use alloc::vec;
use alloc::vec::Vec;
struct TinyMatcher {
last_space: Vec<u8>,
window_size: u64,
}
impl TinyMatcher {
fn new(window_size: u64) -> Self {
Self {
last_space: Vec::new(),
window_size,
}
}
}
impl Matcher for TinyMatcher {
fn get_next_space(&mut self) -> Vec<u8> {
vec![0; self.window_size as usize]
}
fn get_last_space(&mut self) -> &[u8] {
self.last_space.as_slice()
}
fn commit_space(&mut self, space: Vec<u8>) {
self.last_space = space;
}
fn skip_matching(&mut self) {}
fn start_matching(&mut self, mut handle_sequence: impl for<'a> FnMut(Sequence<'a>)) {
handle_sequence(Sequence::Literals {
literals: self.last_space.as_slice(),
});
}
fn reset(&mut self, _level: CompressionLevel) {
self.last_space.clear();
}
fn window_size(&self) -> u64 {
self.window_size
}
}
struct FailingWriteOnce {
writes: usize,
fail_on_write_number: usize,
sink: Vec<u8>,
}
impl FailingWriteOnce {
fn new(fail_on_write_number: usize) -> Self {
Self {
writes: 0,
fail_on_write_number,
sink: Vec::new(),
}
}
}
impl Write for FailingWriteOnce {
fn write(&mut self, buf: &[u8]) -> Result<usize, Error> {
self.writes += 1;
if self.writes == self.fail_on_write_number {
return Err(super::other_error("injected write failure"));
}
self.sink.extend_from_slice(buf);
Ok(buf.len())
}
fn flush(&mut self) -> Result<(), Error> {
Ok(())
}
}
struct FailingWithKind {
writes: usize,
fail_on_write_number: usize,
kind: ErrorKind,
}
impl FailingWithKind {
fn new(fail_on_write_number: usize, kind: ErrorKind) -> Self {
Self {
writes: 0,
fail_on_write_number,
kind,
}
}
}
impl Write for FailingWithKind {
fn write(&mut self, buf: &[u8]) -> Result<usize, Error> {
self.writes += 1;
if self.writes == self.fail_on_write_number {
return Err(Error::from(self.kind));
}
Ok(buf.len())
}
fn flush(&mut self) -> Result<(), Error> {
Ok(())
}
}
struct PartialThenFailWriter {
writes: usize,
fail_on_write_number: usize,
partial_prefix_len: usize,
terminal_failure: bool,
sink: Vec<u8>,
}
impl PartialThenFailWriter {
fn new(fail_on_write_number: usize, partial_prefix_len: usize) -> Self {
Self {
writes: 0,
fail_on_write_number,
partial_prefix_len,
terminal_failure: false,
sink: Vec::new(),
}
}
}
impl Write for PartialThenFailWriter {
fn write(&mut self, buf: &[u8]) -> Result<usize, Error> {
if self.terminal_failure {
return Err(super::other_error("injected terminal write failure"));
}
self.writes += 1;
if self.writes == self.fail_on_write_number {
let written = core::cmp::min(self.partial_prefix_len, buf.len());
if written > 0 {
self.sink.extend_from_slice(&buf[..written]);
self.terminal_failure = true;
return Ok(written);
}
return Err(super::other_error("injected terminal write failure"));
}
self.sink.extend_from_slice(buf);
Ok(buf.len())
}
fn flush(&mut self) -> Result<(), Error> {
Ok(())
}
}
#[test]
fn streaming_encoder_pre_splits_full_blocks_like_the_frame_compressor() {
let path = concat!(env!("CARGO_MANIFEST_DIR"), "/decodecorpus_files/z000033");
let data = std::fs::read(path).unwrap();
for level in [6, 16] {
let mut streamed = Vec::new();
let mut enc = StreamingEncoder::new(&mut streamed, CompressionLevel::Level(level));
enc.set_pledged_content_size(data.len() as u64).unwrap();
for chunk in data.chunks(8192) {
enc.write_all(chunk).unwrap();
}
enc.finish().unwrap();
let mut read = Vec::new();
let mut fc: crate::encoding::FrameCompressor<&[u8], &mut Vec<u8>> =
crate::encoding::FrameCompressor::new(CompressionLevel::Level(level));
fc.set_source_size_hint(data.len() as u64);
fc.set_source(&data[..]);
fc.set_drain(&mut read);
fc.compress();
let (_, streamed_header) =
crate::decoding::frame::read_frame_header(&streamed[..]).unwrap();
let (_, read_header) = crate::decoding::frame::read_frame_header(&read[..]).unwrap();
assert_eq!(
streamed[usize::from(streamed_header)..],
read[usize::from(read_header)..],
"level {level}: streaming blocks must be pre-split like the reader path"
);
}
}
#[test]
fn streaming_encoder_matcher_and_gates_resolve_from_one_size() {
let mut enc = StreamingEncoder::new(Vec::new(), CompressionLevel::Level(13));
enc.set_pledged_content_size(4096).unwrap();
enc.set_source_size_hint(1 << 20).unwrap();
enc.write_all(&[0u8; 4096]).unwrap();
assert_eq!(
enc.context.state.matcher.active_backend(),
enc.context.state.strategy_tag.backend(),
"matcher backend must match the synchronized strategy ({:?})",
enc.context.state.strategy_tag,
);
enc.finish().unwrap();
}
#[test]
fn streaming_periodic_btlazy2_roundtrips() {
const LINES: &[&[u8]] = &[
b"ts=2026-03-26T21:39:28Z level=INFO msg=\"flush memtable\" tenant=demo table=orders region=eu-west\n",
b"ts=2026-03-26T21:39:29Z level=INFO msg=\"rotate segment\" tenant=demo table=orders region=eu-west\n",
b"ts=2026-03-26T21:39:30Z level=INFO msg=\"compact level\" tenant=demo table=orders region=eu-west\n",
b"ts=2026-03-26T21:39:31Z level=INFO msg=\"write block\" tenant=demo table=orders region=eu-west\n",
];
let target = 6 * 1024 * 1024usize;
let mut data = Vec::with_capacity(target);
'fill: loop {
for line in LINES {
if data.len() + line.len() > target {
break 'fill;
}
data.extend_from_slice(line);
}
}
for level in [13, 15] {
let mut out = Vec::new();
let mut enc = StreamingEncoder::new(&mut out, CompressionLevel::Level(level));
for chunk in data.chunks(64 * 1024) {
enc.write_all(chunk).unwrap();
}
enc.finish().unwrap();
let mut decoder = crate::decoding::FrameDecoder::new();
let mut round = Vec::with_capacity(data.len());
decoder
.decode_all_to_vec(&out, &mut round)
.unwrap_or_else(|e| panic!("L{level} decode failed: {e:?}"));
assert_eq!(round, data, "L{level} streamed periodic roundtrip");
}
}
#[test]
fn streaming_encoder_literal_gate_follows_the_effective_target_length() {
use crate::encoding::CompressionParameters;
let params = CompressionParameters::builder(CompressionLevel::Level(1))
.target_length(8)
.build()
.expect("valid override");
let mut plain = StreamingEncoder::new(Vec::new(), CompressionLevel::Level(1));
plain.set_parameters(¶ms).unwrap();
plain.write_all(b"plain frame payload").unwrap();
assert!(plain.context.state.literal_compression_disabled);
let dict: Vec<u8> = (0..4096u32)
.map(|i| (i.wrapping_mul(2_654_435_761) >> 13) as u8)
.collect();
let mut with_dict = StreamingEncoder::new(Vec::new(), CompressionLevel::Level(1));
with_dict.set_parameters(¶ms).unwrap();
with_dict
.set_encoder_dictionary(crate::encoding::EncoderDictionary::from_dictionary(
crate::decoding::Dictionary::from_raw_content(0xD1C7_0018, dict).unwrap(),
))
.unwrap();
with_dict.write_all(b"dictionary frame payload").unwrap();
assert!(with_dict.context.state.literal_compression_disabled);
let fast = CompressionParameters::builder(CompressionLevel::Level(22))
.strategy(crate::encoding::Strategy::Fast)
.build()
.expect("valid override");
let dict: Vec<u8> = (0..4096u32)
.map(|i| (i.wrapping_mul(2_654_435_761) >> 13) as u8)
.collect();
let mut moved = StreamingEncoder::new(Vec::new(), CompressionLevel::Level(22));
moved.set_parameters(&fast).unwrap();
moved
.set_encoder_dictionary(crate::encoding::EncoderDictionary::from_dictionary(
crate::decoding::Dictionary::from_raw_content(0xD1C7_001D, dict).unwrap(),
))
.unwrap();
moved.write_all(b"dictionary frame payload").unwrap();
assert!(moved.context.state.literal_compression_disabled);
}
#[test]
fn streaming_encoder_set_magicless_before_write_omits_magic_and_roundtrips() {
use crate::common::MAGIC_NUM;
let payload = b"streaming-magicless-roundtrip-".repeat(64);
let mut encoder = StreamingEncoder::new(Vec::new(), CompressionLevel::Fastest);
encoder
.set_magicless(true)
.expect("set_magicless pre-write");
encoder.write_all(&payload).unwrap();
let compressed = encoder.finish().unwrap();
assert!(
!compressed.starts_with(&MAGIC_NUM.to_le_bytes()),
"magicless frame must omit the 4-byte magic prefix",
);
let mut decoder = crate::decoding::FrameDecoder::new();
decoder.set_magicless(true);
let mut cursor: &[u8] = compressed.as_slice();
decoder.init(&mut cursor).expect("magicless init");
decoder
.decode_blocks(&mut cursor, crate::decoding::BlockDecodingStrategy::All)
.expect("decode_blocks");
let mut decoded: Vec<u8> = Vec::new();
decoder
.collect_to_writer(&mut decoded)
.expect("collect_to_writer");
assert_eq!(decoded, payload);
}
#[test]
fn streaming_encoder_set_magicless_after_first_write_errors() {
let mut encoder = StreamingEncoder::new(Vec::new(), CompressionLevel::Fastest);
encoder.write_all(b"first-block").unwrap();
let err = encoder
.set_magicless(true)
.expect_err("set_magicless after first write must error");
assert_eq!(
err.kind(),
crate::io::ErrorKind::InvalidInput,
"expected InvalidInput when setting magicless after frame_started, got {err:?}",
);
}
#[test]
fn streaming_encoder_roundtrip_multiple_writes() {
let payload = b"streaming-encoder-roundtrip-".repeat(1024);
let mut encoder = StreamingEncoder::new(Vec::new(), CompressionLevel::Fastest);
for chunk in payload.chunks(313) {
encoder.write_all(chunk).unwrap();
}
let compressed = encoder.finish().unwrap();
let mut decoder = StreamingDecoder::new(compressed.as_slice()).unwrap();
let mut decoded = Vec::new();
decoder.read_to_end(&mut decoded).unwrap();
assert_eq!(decoded, payload);
}
#[test]
fn flush_emits_nonempty_partial_output() {
let mut encoder = StreamingEncoder::new(Vec::new(), CompressionLevel::Fastest);
encoder.write_all(b"partial-block").unwrap();
encoder.flush().unwrap();
let flushed_len = encoder.get_ref().len();
assert!(
flushed_len > 0,
"flush should emit header+partial block bytes"
);
let compressed = encoder.finish().unwrap();
let mut decoder = StreamingDecoder::new(compressed.as_slice()).unwrap();
let mut decoded = Vec::new();
decoder.read_to_end(&mut decoded).unwrap();
assert_eq!(decoded, b"partial-block");
}
#[test]
fn flush_without_writes_does_not_emit_frame_header() {
let mut encoder = StreamingEncoder::new(Vec::new(), CompressionLevel::Fastest);
encoder.flush().unwrap();
assert!(encoder.get_ref().is_empty());
}
#[test]
fn block_boundary_write_emits_block_in_same_call() {
let mut boundary = StreamingEncoder::new_with_matcher(
TinyMatcher::new(4),
Vec::new(),
CompressionLevel::Uncompressed,
);
let mut below = StreamingEncoder::new_with_matcher(
TinyMatcher::new(4),
Vec::new(),
CompressionLevel::Uncompressed,
);
boundary.write_all(b"ABCD").unwrap();
below.write_all(b"ABC").unwrap();
let boundary_len = boundary.get_ref().len();
let below_len = below.get_ref().len();
assert!(
boundary_len > below_len,
"full block should be emitted immediately at block boundary"
);
}
#[test]
fn finish_consumes_encoder_and_emits_frame() {
let mut encoder = StreamingEncoder::new(Vec::new(), CompressionLevel::Fastest);
encoder.write_all(b"abc").unwrap();
let compressed = encoder.finish().unwrap();
let mut decoder = StreamingDecoder::new(compressed.as_slice()).unwrap();
let mut decoded = Vec::new();
decoder.read_to_end(&mut decoded).unwrap();
assert_eq!(decoded, b"abc");
}
#[test]
fn finish_without_writes_emits_empty_frame() {
let encoder = StreamingEncoder::new(Vec::new(), CompressionLevel::Fastest);
let compressed = encoder.finish().unwrap();
let mut decoder = StreamingDecoder::new(compressed.as_slice()).unwrap();
let mut decoded = Vec::new();
decoder.read_to_end(&mut decoded).unwrap();
assert!(decoded.is_empty());
}
#[test]
fn write_empty_buffer_returns_zero() {
let mut encoder = StreamingEncoder::new(Vec::new(), CompressionLevel::Fastest);
assert_eq!(encoder.write(&[]).unwrap(), 0);
let _ = encoder.finish().unwrap();
}
#[test]
fn uncompressed_level_roundtrip() {
let payload = b"uncompressed-streaming-roundtrip".repeat(64);
let mut encoder = StreamingEncoder::new(Vec::new(), CompressionLevel::Uncompressed);
for chunk in payload.chunks(41) {
encoder.write_all(chunk).unwrap();
}
let compressed = encoder.finish().unwrap();
let mut decoder = StreamingDecoder::new(compressed.as_slice()).unwrap();
let mut decoded = Vec::new();
decoder.read_to_end(&mut decoded).unwrap();
assert_eq!(decoded, payload);
}
#[test]
fn better_level_streaming_roundtrip() {
let payload = b"better-level-streaming-test".repeat(256);
let mut encoder = StreamingEncoder::new(Vec::new(), CompressionLevel::Better);
for chunk in payload.chunks(53) {
encoder.write_all(chunk).unwrap();
}
let compressed = encoder.finish().unwrap();
let mut decoder = StreamingDecoder::new(compressed.as_slice()).unwrap();
let mut decoded = Vec::new();
decoder.read_to_end(&mut decoded).unwrap();
assert_eq!(decoded, payload);
}
#[test]
fn zero_window_matcher_returns_invalid_input_error() {
let mut encoder = StreamingEncoder::new_with_matcher(
TinyMatcher::new(0),
Vec::new(),
CompressionLevel::Fastest,
);
let err = encoder.write_all(b"payload").unwrap_err();
assert_eq!(err.kind(), ErrorKind::InvalidInput);
}
#[test]
fn best_level_streaming_roundtrip() {
let payload = b"best-level-streaming-test".repeat(8 * 1024);
let mut encoder = StreamingEncoder::new(Vec::new(), CompressionLevel::Best);
for chunk in payload.chunks(53) {
encoder.write_all(chunk).unwrap();
}
let compressed = encoder.finish().unwrap();
let mut decoder = StreamingDecoder::new(compressed.as_slice()).unwrap();
let mut decoded = Vec::new();
decoder.read_to_end(&mut decoded).unwrap();
assert_eq!(decoded, payload);
}
#[test]
fn write_failure_poisoning_is_sticky() {
let mut encoder = StreamingEncoder::new_with_matcher(
TinyMatcher::new(4),
FailingWriteOnce::new(1),
CompressionLevel::Uncompressed,
);
assert!(encoder.write_all(b"ABCD").is_err());
assert!(encoder.flush().is_err());
assert!(encoder.write_all(b"EFGH").is_err());
assert_eq!(encoder.get_ref().sink.len(), 0);
assert!(encoder.finish().is_err());
}
#[test]
fn poisoned_encoder_returns_original_error_kind() {
let mut encoder = StreamingEncoder::new_with_matcher(
TinyMatcher::new(4),
FailingWithKind::new(1, ErrorKind::BrokenPipe),
CompressionLevel::Uncompressed,
);
let first_error = encoder.write_all(b"ABCD").unwrap_err();
assert_eq!(first_error.kind(), ErrorKind::BrokenPipe);
let second_error = encoder.write_all(b"EFGH").unwrap_err();
assert_eq!(second_error.kind(), ErrorKind::BrokenPipe);
}
#[test]
fn write_reports_progress_but_poisoning_is_sticky_after_later_block_failure() {
let payload = b"ABCDEFGHIJKL";
let mut encoder = StreamingEncoder::new_with_matcher(
TinyMatcher::new(4),
FailingWriteOnce::new(3),
CompressionLevel::Uncompressed,
);
let first_write = encoder.write(payload).unwrap();
assert_eq!(first_write, 8);
assert!(encoder.write(&payload[first_write..]).is_err());
assert!(encoder.flush().is_err());
assert!(encoder.write_all(b"EFGH").is_err());
}
#[test]
fn partial_write_failure_after_progress_poisons_encoder() {
let payload = b"ABCDEFGHIJKL";
let mut encoder = StreamingEncoder::new_with_matcher(
TinyMatcher::new(4),
PartialThenFailWriter::new(3, 1),
CompressionLevel::Uncompressed,
);
let first_write = encoder.write(payload).unwrap();
assert_eq!(first_write, 8);
let second_write = encoder.write(&payload[first_write..]);
assert!(second_write.is_err());
assert!(encoder.flush().is_err());
assert!(encoder.write_all(b"MNOP").is_err());
}
#[test]
fn new_with_matcher_and_get_mut_work() {
let matcher = TinyMatcher::new(128 * 1024);
let mut encoder =
StreamingEncoder::new_with_matcher(matcher, Vec::new(), CompressionLevel::Fastest);
encoder.get_mut().extend_from_slice(b"");
encoder.write_all(b"custom-matcher").unwrap();
let compressed = encoder.finish().unwrap();
let mut decoder = StreamingDecoder::new(compressed.as_slice()).unwrap();
let mut decoded = Vec::new();
decoder.read_to_end(&mut decoded).unwrap();
assert_eq!(decoded, b"custom-matcher");
}
#[test]
fn pledged_content_size_written_in_header() {
let payload = b"hello world, pledged size test";
let mut encoder = StreamingEncoder::new(Vec::new(), CompressionLevel::Fastest);
encoder
.set_pledged_content_size(payload.len() as u64)
.unwrap();
encoder.write_all(payload).unwrap();
let compressed = encoder.finish().unwrap();
let header = crate::decoding::frame::read_frame_header(compressed.as_slice())
.unwrap()
.0;
assert_eq!(header.frame_content_size(), payload.len() as u64);
let mut decoder = StreamingDecoder::new(compressed.as_slice()).unwrap();
let mut decoded = Vec::new();
decoder.read_to_end(&mut decoded).unwrap();
assert_eq!(decoded, payload);
}
#[test]
fn pledged_content_size_mismatch_returns_error() {
let mut encoder = StreamingEncoder::new(Vec::new(), CompressionLevel::Fastest);
encoder.set_pledged_content_size(100).unwrap();
encoder.write_all(b"short payload").unwrap(); let err = encoder.finish().unwrap_err();
assert_eq!(err.kind(), ErrorKind::InvalidInput);
}
#[test]
fn write_exceeding_pledge_returns_error() {
let mut encoder = StreamingEncoder::new(Vec::new(), CompressionLevel::Fastest);
encoder.set_pledged_content_size(5).unwrap();
let err = encoder.write_all(b"exceeds five bytes").unwrap_err();
assert_eq!(err.kind(), ErrorKind::InvalidInput);
}
#[test]
fn write_straddling_pledge_reports_partial_progress() {
let mut encoder = StreamingEncoder::new(Vec::new(), CompressionLevel::Fastest);
encoder.set_pledged_content_size(5).unwrap();
assert_eq!(encoder.write(b"abcdef").unwrap(), 5);
let err = encoder.write(b"g").unwrap_err();
assert_eq!(err.kind(), ErrorKind::InvalidInput);
}
#[test]
fn encoded_scratch_capacity_is_reused_across_blocks() {
let payload = vec![0xAB; 64 * 3];
let mut encoder = StreamingEncoder::new_with_matcher(
TinyMatcher::new(64),
Vec::new(),
CompressionLevel::Uncompressed,
);
encoder.write_all(&payload[..64]).unwrap();
let first_capacity = encoder.context.encoded_scratch.capacity();
assert!(
first_capacity >= 67,
"expected encoded scratch to keep block header + payload capacity",
);
encoder.write_all(&payload[64..128]).unwrap();
let second_capacity = encoder.context.encoded_scratch.capacity();
assert!(
second_capacity >= first_capacity,
"encoded scratch capacity should be reused across block emits",
);
encoder.write_all(&payload[128..]).unwrap();
let compressed = encoder.finish().unwrap();
let mut decoder = StreamingDecoder::new(compressed.as_slice()).unwrap();
let mut decoded = Vec::new();
decoder.read_to_end(&mut decoded).unwrap();
assert_eq!(decoded, payload);
}
#[test]
fn pledged_content_size_after_write_returns_error() {
let mut encoder = StreamingEncoder::new(Vec::new(), CompressionLevel::Fastest);
encoder.write_all(b"already writing").unwrap();
let err = encoder.set_pledged_content_size(15).unwrap_err();
assert_eq!(err.kind(), ErrorKind::InvalidInput);
}
#[test]
fn source_size_hint_directly_reduces_window_header() {
let payload = b"streaming-source-size-hint".repeat(64);
let mut no_hint = StreamingEncoder::new(Vec::new(), CompressionLevel::from_level(11));
no_hint.write_all(payload.as_slice()).unwrap();
let no_hint_frame = no_hint.finish().unwrap();
let no_hint_header = crate::decoding::frame::read_frame_header(no_hint_frame.as_slice())
.unwrap()
.0;
let no_hint_window = no_hint_header.window_size().unwrap();
let mut with_hint = StreamingEncoder::new(Vec::new(), CompressionLevel::from_level(11));
with_hint
.set_source_size_hint(payload.len() as u64)
.unwrap();
with_hint.write_all(payload.as_slice()).unwrap();
let late_hint_err = with_hint
.set_source_size_hint(payload.len() as u64)
.unwrap_err();
assert_eq!(late_hint_err.kind(), ErrorKind::InvalidInput);
let with_hint_frame = with_hint.finish().unwrap();
let with_hint_header = crate::decoding::frame::read_frame_header(with_hint_frame.as_slice())
.unwrap()
.0;
let with_hint_window = with_hint_header.window_size().unwrap();
assert!(
with_hint_window <= no_hint_window,
"source size hint should not increase advertised window"
);
let mut decoder = StreamingDecoder::new(with_hint_frame.as_slice()).unwrap();
let mut decoded = Vec::new();
decoder.read_to_end(&mut decoded).unwrap();
assert_eq!(decoded, payload);
}
#[test]
fn single_segment_requires_pledged_to_fit_matcher_window() {
let payload = b"streaming-window-gate-".repeat(60); let mut encoder = StreamingEncoder::new_with_matcher(
TinyMatcher::new(1024),
Vec::new(),
CompressionLevel::Fastest,
);
encoder
.set_pledged_content_size(payload.len() as u64)
.unwrap();
encoder.write_all(payload.as_slice()).unwrap();
let compressed = encoder.finish().unwrap();
let header = crate::decoding::frame::read_frame_header(compressed.as_slice())
.unwrap()
.0;
assert_eq!(header.frame_content_size(), payload.len() as u64);
assert!(
!header.descriptor.single_segment_flag(),
"single-segment must stay off when pledged content size exceeds matcher window"
);
assert!(
header.window_size().unwrap() >= 1024,
"window descriptor should be present when single-segment is disabled"
);
}
#[test]
fn ensure_frame_started_refreshes_stale_strategy_tag_at_reset() {
use crate::encoding::strategy::StrategyTag;
for level in [
CompressionLevel::Fastest,
CompressionLevel::Default,
CompressionLevel::Better,
CompressionLevel::Best,
] {
let expected = StrategyTag::for_compression_level(level);
let mut encoder = StreamingEncoder::new(Vec::new(), level);
let sentinel = StrategyTag::BtUltra2;
assert_ne!(
expected, sentinel,
"sentinel must differ from the legitimate tag at level {level:?}",
);
encoder.context.state.strategy_tag = sentinel;
encoder.write_all(b"x").unwrap();
assert_eq!(
encoder.context.state.strategy_tag, expected,
"reset-time strategy_tag sync missing at level {level:?}: \
sentinel survived `ensure_frame_started`",
);
let _ = encoder.finish().unwrap();
}
}
#[test]
fn level_22_streaming_window_roundtrips_in_our_decoder() {
let payload = b"level-22-streaming-window-cap-".repeat(512);
let mut encoder = StreamingEncoder::new(Vec::new(), CompressionLevel::from_level(22));
for chunk in payload.chunks(101) {
encoder.write_all(chunk).unwrap();
}
let compressed = encoder.finish().unwrap();
let header = crate::decoding::frame::read_frame_header(compressed.as_slice())
.unwrap()
.0;
let window = header.window_size().unwrap();
assert!(
window <= crate::common::MAXIMUM_ALLOWED_WINDOW_SIZE,
"L22 advertised window {window} exceeds decoder cap {}",
crate::common::MAXIMUM_ALLOWED_WINDOW_SIZE,
);
let mut decoder = StreamingDecoder::new(compressed.as_slice()).unwrap();
let mut decoded = Vec::new();
decoder.read_to_end(&mut decoded).unwrap();
assert_eq!(decoded, payload);
}
#[test]
fn streaming_encoder_set_content_checksum_false_clears_header_flag() {
let payload = b"streaming-checksum-toggle-".repeat(64);
let mut encoder = StreamingEncoder::new(Vec::new(), CompressionLevel::Fastest);
encoder
.set_content_checksum(false)
.expect("set_content_checksum pre-write");
encoder.write_all(&payload).unwrap();
let compressed = encoder.finish().unwrap();
let header = crate::decoding::frame::read_frame_header(compressed.as_slice())
.unwrap()
.0;
assert!(
!header.descriptor.content_checksum_flag(),
"content_checksum(false) must clear the frame header flag",
);
let mut decoder = StreamingDecoder::new(compressed.as_slice()).unwrap();
let mut decoded = Vec::new();
decoder.read_to_end(&mut decoded).unwrap();
assert_eq!(decoded, payload);
}
#[cfg(feature = "hash")]
#[test]
fn streaming_encoder_set_content_checksum_false_omits_trailer() {
let payload = b"streaming-checksum-trailer-".repeat(64);
let mut with = StreamingEncoder::new(Vec::new(), CompressionLevel::Fastest);
with.set_content_checksum(true)
.expect("set_content_checksum pre-write");
with.write_all(&payload).unwrap();
let with_checksum = with.finish().unwrap();
let mut without = StreamingEncoder::new(Vec::new(), CompressionLevel::Fastest);
without
.set_content_checksum(false)
.expect("set_content_checksum pre-write");
without.write_all(&payload).unwrap();
let without_checksum = without.finish().unwrap();
assert!(
crate::decoding::frame::read_frame_header(with_checksum.as_slice())
.unwrap()
.0
.descriptor
.content_checksum_flag(),
"default checksum-on frame must set the header flag",
);
assert_eq!(
with_checksum.len(),
without_checksum.len() + 4,
"checksum-on frame must carry exactly the 4-byte XXH64 trailer",
);
}
#[test]
fn streaming_encoder_set_content_checksum_after_first_write_errors() {
let mut encoder = StreamingEncoder::new(Vec::new(), CompressionLevel::Fastest);
encoder.write_all(b"first-block").unwrap();
let err = encoder
.set_content_checksum(false)
.expect_err("set_content_checksum after first write must error");
assert_eq!(
err.kind(),
ErrorKind::InvalidInput,
"expected InvalidInput when setting content checksum after frame_started, got {err:?}",
);
}
#[test]
fn no_pledged_size_omits_fcs_from_header() {
let mut encoder = StreamingEncoder::new(Vec::new(), CompressionLevel::Fastest);
encoder.write_all(b"no pledged size").unwrap();
let compressed = encoder.finish().unwrap();
let header = crate::decoding::frame::read_frame_header(compressed.as_slice())
.unwrap()
.0;
assert_eq!(header.frame_content_size(), 0);
assert_eq!(header.descriptor.frame_content_size_bytes().unwrap(), 0);
}
#[test]
fn streaming_encoder_with_dictionary_roundtrips_and_carries_dict_id() {
use alloc::format;
let dict_raw = include_bytes!("../../../dict_tests/dictionary");
let dict_id = crate::decoding::Dictionary::decode_dict(dict_raw)
.unwrap()
.id;
let mut payload = Vec::new();
for i in 0..400u32 {
payload.extend_from_slice(
format!("tenant=demo table=orders key={i} region=eu payload=aaaaabbbbbccccc\n")
.as_bytes(),
);
}
let mut encoder = StreamingEncoder::new(Vec::new(), CompressionLevel::Level(19));
encoder
.set_dictionary_from_bytes(dict_raw)
.expect("attach dictionary");
for chunk in payload.chunks(777) {
encoder.write_all(chunk).unwrap();
}
let compressed = encoder.finish().unwrap();
let header = crate::decoding::frame::read_frame_header(compressed.as_slice())
.unwrap()
.0;
assert_eq!(header.dictionary_id(), Some(dict_id));
let mut decoder =
StreamingDecoder::new_with_dictionary_bytes(compressed.as_slice(), dict_raw).unwrap();
let mut decoded = Vec::new();
decoder.read_to_end(&mut decoded).unwrap();
assert_eq!(decoded, payload);
let mut nodict = StreamingEncoder::new(Vec::new(), CompressionLevel::Level(19));
for chunk in payload.chunks(777) {
nodict.write_all(chunk).unwrap();
}
let nodict_frame = nodict.finish().unwrap();
assert!(
compressed.len() <= nodict_frame.len(),
"dict frame {} should not exceed no-dict frame {}",
compressed.len(),
nodict_frame.len()
);
}
#[test]
fn streaming_encoder_strategy_override_survives_frame_start() {
use crate::encoding::{CompressionParameters, Strategy};
let level = CompressionLevel::Fastest;
let level_tag = crate::encoding::strategy::StrategyTag::for_compression_level(level);
let override_tag = Strategy::Greedy.tag();
assert_ne!(
level_tag, override_tag,
"test needs an override that changes the derived tag"
);
let params = CompressionParameters::builder(level)
.strategy(Strategy::Greedy)
.build()
.unwrap();
let payload = b"override must outlive the frame header";
let mut encoder = StreamingEncoder::new(Vec::new(), level);
encoder.set_parameters(¶ms).unwrap();
encoder.write_all(payload).unwrap();
assert_eq!(
encoder.context.state.strategy_tag, override_tag,
"strategy override was discarded when the frame started"
);
let compressed = encoder.finish().unwrap();
let mut decoder = StreamingDecoder::new(compressed.as_slice()).unwrap();
let mut decoded = Vec::new();
decoder.read_to_end(&mut decoded).unwrap();
assert_eq!(decoded, payload);
}
#[test]
fn streaming_encoder_uncompressed_with_dictionary_omits_dict_id() {
let dict_raw = include_bytes!("../../../dict_tests/dictionary");
let payload = b"tenant=demo table=orders region=eu payload=aaaaabbbbbccccc";
let mut encoder = StreamingEncoder::new(Vec::new(), CompressionLevel::Uncompressed);
encoder
.set_dictionary_from_bytes(dict_raw)
.expect("attach dictionary");
encoder.write_all(payload).unwrap();
let compressed = encoder.finish().unwrap();
let header = crate::decoding::frame::read_frame_header(compressed.as_slice())
.unwrap()
.0;
assert_eq!(
header.dictionary_id(),
None,
"uncompressed frame must not require a dictionary at decode time"
);
let mut decoder = StreamingDecoder::new(compressed.as_slice()).unwrap();
let mut decoded = Vec::new();
decoder.read_to_end(&mut decoded).unwrap();
assert_eq!(decoded, payload);
}
#[test]
fn raw_dictionary_leaves_the_id_out_of_the_streaming_header() {
use crate::decoding::Dictionary;
use crate::encoding::EncoderDictionary;
let content: Vec<u8> = b"tenant=demo region=eu table=orders payload="
.iter()
.copied()
.cycle()
.take(2048)
.collect();
let mut payload = Vec::new();
while payload.len() < 8192 {
payload.extend_from_slice(&content);
}
let mut encoder = StreamingEncoder::new(Vec::new(), CompressionLevel::Default);
encoder
.set_encoder_dictionary(EncoderDictionary::from_dictionary(
Dictionary::from_raw_content(0, content.clone()).expect("a raw dictionary has no id"),
))
.expect("a raw dictionary must attach");
encoder.write_all(&payload).unwrap();
let compressed = encoder.finish().unwrap();
assert_eq!(
compressed[4] & 0b11,
0,
"a dictionary with no id must leave the Dictionary_ID field out, \
not store a zero in it"
);
let mut decoder = StreamingDecoder::new_with_dictionary_handle(
compressed.as_slice(),
&crate::decoding::DictionaryHandle::from_dictionary(
Dictionary::from_raw_content(0, content).expect("a raw dictionary has no id"),
),
)
.unwrap();
let mut decoded = Vec::new();
decoder.read_to_end(&mut decoded).unwrap();
assert_eq!(decoded, payload);
}
#[test]
fn set_dictionary_from_bytes_with_an_empty_buffer_clears_the_dictionary() {
let content: Vec<u8> = b"tenant=demo region=eu table=orders payload="
.iter()
.copied()
.cycle()
.take(2048)
.collect();
let mut payload = Vec::new();
while payload.len() < 8192 {
payload.extend_from_slice(&content);
}
let mut encoder = StreamingEncoder::new(Vec::new(), CompressionLevel::Default);
encoder
.set_dictionary_from_bytes(&content)
.expect("the dictionary attaches");
encoder
.set_dictionary_from_bytes(&[])
.expect("an empty buffer is how a caller says there is no dictionary");
encoder.write_all(&payload).unwrap();
let compressed = encoder.finish().unwrap();
let mut decoder = StreamingDecoder::new(compressed.as_slice()).unwrap();
let mut decoded = Vec::new();
decoder.read_to_end(&mut decoded).unwrap();
assert_eq!(decoded, payload);
}
#[test]
fn clearing_the_stream_dictionary_is_refused_where_attaching_is() {
let content: Vec<u8> = b"tenant=demo region=eu table=orders payload="
.iter()
.copied()
.cycle()
.take(1024)
.collect();
let mut enc = StreamingEncoder::new(Vec::new(), CompressionLevel::Default);
enc.set_dictionary_from_bytes(&content).expect("attach");
enc.write_all(b"the frame starts here").unwrap();
let err = enc
.set_dictionary_from_bytes(&[])
.expect_err("the frame's dictionary is already decided");
assert!(
alloc::format!("{err:?}").contains("before the first write"),
"unexpected error: {err:?}",
);
let mut enc = StreamingEncoder::new(FailingWriteOnce::new(1), CompressionLevel::Fastest);
let big = vec![b'x'; 512 * 1024];
let _ = enc.write_all(&big);
let _ = enc.flush();
assert!(
enc.set_dictionary_from_bytes(&[]).is_err(),
"a poisoned stream answers with its failure, not with success",
);
}
#[test]
fn clearing_the_stream_dictionary_gives_back_what_it_allocated() {
let content = include_bytes!("../../../dict_tests/dictionary").to_vec();
let mut enc = StreamingEncoder::new(Vec::new(), CompressionLevel::Default);
let empty = enc.heap_size();
enc.set_dictionary_from_bytes(&content)
.expect("the dictionary attaches");
let attached = enc.heap_size();
assert!(
attached > empty,
"attaching should have allocated something to give back: {empty} -> {attached}",
);
enc.set_dictionary_from_bytes(&[]).expect("clear");
assert_eq!(
enc.heap_size(),
empty,
"a cleared dictionary must leave the encoder holding no more than it did before",
);
}
#[test]
fn set_dictionary_from_bytes_takes_unmagicked_bytes_as_raw_content() {
use crate::decoding::Dictionary;
let content: Vec<u8> = b"tenant=demo region=eu table=orders payload="
.iter()
.copied()
.cycle()
.take(2048)
.collect();
assert_ne!(
content[..4],
crate::decoding::DICTIONARY_MAGIC,
"the fixture must not start with the dictionary magic",
);
let mut payload = Vec::new();
while payload.len() < 8192 {
payload.extend_from_slice(&content);
}
let mut encoder = StreamingEncoder::new(Vec::new(), CompressionLevel::Default);
encoder
.set_dictionary_from_bytes(&content)
.expect("raw content must load the way `zstd -D` loads it");
encoder.write_all(&payload).unwrap();
let compressed = encoder.finish().unwrap();
assert_eq!(compressed[4] & 0b11, 0);
let mut decoder = StreamingDecoder::new_with_dictionary_handle(
compressed.as_slice(),
&crate::decoding::DictionaryHandle::from_dictionary(
Dictionary::from_raw_content(0, content).expect("a raw dictionary has no id"),
),
)
.unwrap();
let mut decoded = Vec::new();
decoder.read_to_end(&mut decoded).unwrap();
assert_eq!(decoded, payload);
}