From 0b43c246ce6db3fe9da4736a0237dfcd2c182838 Mon Sep 17 00:00:00 2001 From: Roch Devost Date: Mon, 24 Aug 2026 11:02:24 -0400 Subject: [PATCH 1/4] fix(compression): align zstd behavior across targets --- libdd-profiling/src/pprof/test_utils.rs | 4 +- libdd-profiling/src/profiles/compressor.rs | 42 +++++++++++- .../src/send_with_retry/compression.rs | 60 ++++++++++++++--- tests/wasm/tests/compression.rs | 67 +++++++++++++++++-- 4 files changed, 151 insertions(+), 22 deletions(-) diff --git a/libdd-profiling/src/pprof/test_utils.rs b/libdd-profiling/src/pprof/test_utils.rs index d315352b65..ec6fafbc28 100644 --- a/libdd-profiling/src/pprof/test_utils.rs +++ b/libdd-profiling/src/pprof/test_utils.rs @@ -6,8 +6,8 @@ use libdd_profiling_protobuf::prost_impls::{Profile, Sample}; fn deserialize_compressed_pprof(encoded: &[u8]) -> anyhow::Result { use prost::Message; - // Native zstd uses FFI and is unavailable under Miri, where the default - // profile codec remains uncompressed. + // Miri cannot call native zstd's FFI, so DefaultProfileCodec resolves to + // NoopProfileCodec under cfg(miri) and these bytes are intentionally uncompressed. #[cfg(miri)] let buf = encoded.to_vec(); #[cfg(all(not(miri), not(target_arch = "wasm32")))] diff --git a/libdd-profiling/src/profiles/compressor.rs b/libdd-profiling/src/profiles/compressor.rs index cc90786838..24ae481634 100644 --- a/libdd-profiling/src/profiles/compressor.rs +++ b/libdd-profiling/src/profiles/compressor.rs @@ -3,6 +3,16 @@ use std::io::{self, BufWriter, Read, Write}; +const DEFAULT_COMPRESSION_LEVEL: i32 = 3; + +fn zstd_compression_level(level: i32) -> i32 { + if level == 0 { + DEFAULT_COMPRESSION_LEVEL + } else { + level + } +} + /// This type wraps a [`Vec`] to provide a [`Write`] interface that has a max /// capacity that won't be exceeded. Additionally, it gracefully handles /// out-of-memory conditions instead of panicking (unfortunately not compatible @@ -119,6 +129,11 @@ impl ProfileCodec for NoopProfileCodec { } } +/// A zstd-compatible profile codec. +/// +/// WASM accepts levels `-7..=4`. Native targets accept the range reported by +/// `zstd::compression_level_range()`, making `-7..=4` the portable range. Level +/// `0` selects level `3` on every target. #[allow(unused)] pub struct ZstdProfileCodec; @@ -132,7 +147,10 @@ impl ProfileCodec for ZstdProfileCodec { compression_level: i32, ) -> io::Result { let buffer = SizeRestrictedBuffer::try_new(size_hint, max_capacity)?; - zstd::Encoder::<'static, SizeRestrictedBuffer>::new(buffer, compression_level) + zstd::Encoder::<'static, SizeRestrictedBuffer>::new( + buffer, + zstd_compression_level(compression_level), + ) } fn finish(encoder: Self::Encoder) -> io::Result> { @@ -157,7 +175,8 @@ impl ProfileCodec for ZstdProfileCodec { compression_level: i32, ) -> io::Result { let buffer = SizeRestrictedBuffer::try_new(size_hint, max_capacity)?; - zrip::FrameEncoder::new(buffer, compression_level).map_err(io::Error::other) + zrip::FrameEncoder::new(buffer, zstd_compression_level(compression_level)) + .map_err(io::Error::other) } fn finish(encoder: Self::Encoder) -> io::Result> { @@ -267,7 +286,8 @@ impl Compressor { /// - `size_hint`: beginning capacity for the output buffer. This is a hint for the starting /// size, and the implementation may use something different. /// - `max_capacity`: the maximum size for the output buffer (hard limit). - /// - `compression_level`: must be supported by the target's zstd encoder. + /// - `compression_level`: passed to `C::new_encoder`. For the default [`ZstdProfileCodec`], see + /// its target-specific level documentation. pub fn try_new( size_hint: usize, max_capacity: usize, @@ -299,3 +319,19 @@ impl Write for Compressor { self.encoder.flush() } } + +#[cfg(all(test, not(target_arch = "wasm32")))] +mod tests { + use super::*; + + fn compress(level: i32) -> Vec { + let mut compressor = Compressor::::try_new(256, 4096, level).unwrap(); + compressor.write_all(b"hello profile").unwrap(); + compressor.finish().unwrap() + } + + #[test] + fn zero_uses_default_compression_level() { + assert_eq!(compress(0), compress(DEFAULT_COMPRESSION_LEVEL)); + } +} diff --git a/libdd-trace-utils/src/send_with_retry/compression.rs b/libdd-trace-utils/src/send_with_retry/compression.rs index 53f7bcf918..f3ec1fc28b 100644 --- a/libdd-trace-utils/src/send_with_retry/compression.rs +++ b/libdd-trace-utils/src/send_with_retry/compression.rs @@ -6,11 +6,42 @@ use std::io::Write as _; #[cfg(feature = "compression")] const CONTENT_ENCODING_ZSTD: http::HeaderValue = http::HeaderValue::from_static("zstd"); +#[cfg(feature = "compression")] +const DEFAULT_COMPRESSION_LEVEL: i32 = 3; + +#[cfg(all(feature = "compression", not(target_arch = "wasm32")))] +type ZstdEncoder = zstd::Encoder<'static, Vec>; +#[cfg(all(feature = "compression", target_arch = "wasm32"))] +type ZstdEncoder = zrip::FrameEncoder>; + +#[cfg(feature = "compression")] +fn zstd_compression_level(level: i32) -> i32 { + if level == 0 { + DEFAULT_COMPRESSION_LEVEL + } else { + level + } +} + +#[cfg(all(feature = "compression", not(target_arch = "wasm32")))] +fn new_zstd_encoder(writer: Vec, level: i32) -> std::io::Result { + zstd::Encoder::new(writer, zstd_compression_level(level)) +} + +#[cfg(all(feature = "compression", target_arch = "wasm32"))] +fn new_zstd_encoder(writer: Vec, level: i32) -> std::io::Result { + zrip::FrameEncoder::new(writer, zstd_compression_level(level)).map_err(std::io::Error::other) +} #[derive(Clone, Copy, Debug)] pub enum CompressionStrategy { None, #[cfg(feature = "compression")] + /// Zstd-compatible compression. + /// + /// WASM accepts levels `-7..=4`. Native targets accept the range reported by + /// `zstd::compression_level_range()`, making `-7..=4` the portable range. + /// Level `0` selects level `3` on every target. Zstd { level: i32, }, @@ -27,18 +58,10 @@ pub fn compress(data: Vec, strategy: CompressionStrategy) -> (Vec, Compr // Allocate 1/10th of the original buffer, so we shouldn't add too // much memory usage, and no less than 256 bytes let writer = Vec::with_capacity((data.len() / 10).max(256)); - #[cfg(not(target_arch = "wasm32"))] - let result = zstd::Encoder::new(writer, level).and_then(|mut e| { - e.write_all(&data)?; - Ok((e.finish()?, strategy)) + let result = new_zstd_encoder(writer, level).and_then(|mut encoder| { + encoder.write_all(&data)?; + Ok((encoder.finish()?, strategy)) }); - #[cfg(target_arch = "wasm32")] - let result = zrip::FrameEncoder::new(writer, level) - .map_err(std::io::Error::other) - .and_then(|mut e| { - e.write_all(&data)?; - Ok((e.finish()?, strategy)) - }); result.unwrap_or((data, CompressionStrategy::None)) } } @@ -72,4 +95,19 @@ mod tests { assert!(matches!(strategy, CompressionStrategy::Zstd { level: 1 })); assert_eq!(decompress(&compressed).unwrap(), data); } + + #[test] + fn zero_uses_default_compression_level() { + let data = b"hello zstd".repeat(100); + let (default_compressed, _) = + compress(data.clone(), CompressionStrategy::Zstd { level: 0 }); + let (level_three_compressed, _) = compress( + data, + CompressionStrategy::Zstd { + level: DEFAULT_COMPRESSION_LEVEL, + }, + ); + + assert_eq!(default_compressed, level_three_compressed); + } } diff --git a/tests/wasm/tests/compression.rs b/tests/wasm/tests/compression.rs index ca04a4b92a..e1c93dfdb1 100644 --- a/tests/wasm/tests/compression.rs +++ b/tests/wasm/tests/compression.rs @@ -15,6 +15,14 @@ mod profiling_compressor; use std::io::{Read, Write}; +fn compress_profile(payload: &[u8], level: i32) -> std::io::Result> { + use profiling_compressor::{Compressor, ZstdProfileCodec}; + + let mut compressor = Compressor::::try_new(256, 4096, level)?; + compressor.write_all(payload)?; + compressor.finish() +} + #[wasm_bindgen_test::wasm_bindgen_test] fn trace_compression_uses_zrip() { use trace_compression::{add_headers, compress, CompressionStrategy}; @@ -28,16 +36,41 @@ fn trace_compression_uses_zrip() { assert_eq!(headers["content-encoding"], "zstd"); } +#[wasm_bindgen_test::wasm_bindgen_test] +fn trace_compression_levels_are_portable() { + use trace_compression::{compress, CompressionStrategy}; + + let payload = b"hello zstd".repeat(100); + for level in [-7, 4] { + let (compressed, strategy) = compress(payload.clone(), CompressionStrategy::Zstd { level }); + assert!(matches!(strategy, CompressionStrategy::Zstd { level: actual } if actual == level)); + assert_eq!(zrip::decompress(&compressed).unwrap(), payload); + } + for level in [-8, 5] { + let (uncompressed, strategy) = + compress(payload.clone(), CompressionStrategy::Zstd { level }); + assert!(matches!(strategy, CompressionStrategy::None)); + assert_eq!(uncompressed, payload); + } +} + +#[wasm_bindgen_test::wasm_bindgen_test] +fn trace_compression_zero_uses_level_three() { + use trace_compression::{compress, CompressionStrategy}; + + let payload = b"hello zstd".repeat(100); + let (default_compressed, _) = compress(payload.clone(), CompressionStrategy::Zstd { level: 0 }); + let (level_three_compressed, _) = compress(payload, CompressionStrategy::Zstd { level: 3 }); + + assert_eq!(default_compressed, level_three_compressed); +} + #[wasm_bindgen_test::wasm_bindgen_test] fn profiling_codecs_use_zrip() { - use profiling_compressor::{ - Compressor, ObservationCodec, ZstdObservationCodec, ZstdProfileCodec, - }; + use profiling_compressor::{ObservationCodec, ZstdObservationCodec}; let payload = b"hello profile".repeat(100); - let mut compressor = Compressor::::try_new(256, 4096, 1).unwrap(); - compressor.write_all(&payload).unwrap(); - let compressed = compressor.finish().unwrap(); + let compressed = compress_profile(&payload, 1).unwrap(); assert_eq!(zrip::decompress(&compressed).unwrap(), payload); let mut encoder = ZstdObservationCodec::new_encoder(256, 4096).unwrap(); @@ -47,3 +80,25 @@ fn profiling_codecs_use_zrip() { decoder.read_to_end(&mut decoded).unwrap(); assert_eq!(decoded, payload); } + +#[wasm_bindgen_test::wasm_bindgen_test] +fn profiling_compression_levels_are_portable() { + let payload = b"hello profile".repeat(100); + for level in [-7, 4] { + let compressed = compress_profile(&payload, level).unwrap(); + assert_eq!(zrip::decompress(&compressed).unwrap(), payload); + } + for level in [-8, 5] { + assert!(compress_profile(&payload, level).is_err()); + } +} + +#[wasm_bindgen_test::wasm_bindgen_test] +fn profiling_compression_zero_uses_level_three() { + let payload = b"hello profile".repeat(100); + + assert_eq!( + compress_profile(&payload, 0).unwrap(), + compress_profile(&payload, 3).unwrap() + ); +} From 0ccc548be0e731772fc7aeee5175b2d968debddc Mon Sep 17 00:00:00 2001 From: Roch Devost Date: Mon, 24 Aug 2026 11:42:53 -0400 Subject: [PATCH 2/4] test(profiling): skip native zstd under miri --- libdd-profiling/src/profiles/compressor.rs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/libdd-profiling/src/profiles/compressor.rs b/libdd-profiling/src/profiles/compressor.rs index 24ae481634..b58bf91867 100644 --- a/libdd-profiling/src/profiles/compressor.rs +++ b/libdd-profiling/src/profiles/compressor.rs @@ -320,7 +320,7 @@ impl Write for Compressor { } } -#[cfg(all(test, not(target_arch = "wasm32")))] +#[cfg(all(test, not(miri), not(target_arch = "wasm32")))] mod tests { use super::*; From 2fa7bbb49882e1549fe9f13904b8d3ac5e99095d Mon Sep 17 00:00:00 2001 From: Roch Devost Date: Mon, 24 Aug 2026 12:27:15 -0400 Subject: [PATCH 3/4] fix(compression): clamp zrip compression levels --- libdd-profiling/src/profiles/compressor.rs | 25 +++++++++--- .../src/send_with_retry/compression.rs | 38 +++++++++++++++---- tests/wasm/tests/compression.rs | 26 ++++++++----- 3 files changed, 66 insertions(+), 23 deletions(-) diff --git a/libdd-profiling/src/profiles/compressor.rs b/libdd-profiling/src/profiles/compressor.rs index b58bf91867..3a340771ed 100644 --- a/libdd-profiling/src/profiles/compressor.rs +++ b/libdd-profiling/src/profiles/compressor.rs @@ -4,13 +4,22 @@ use std::io::{self, BufWriter, Read, Write}; const DEFAULT_COMPRESSION_LEVEL: i32 = 3; +#[cfg(target_arch = "wasm32")] +const MIN_ZRIP_COMPRESSION_LEVEL: i32 = -7; +#[cfg(target_arch = "wasm32")] +const MAX_ZRIP_COMPRESSION_LEVEL: i32 = 4; fn zstd_compression_level(level: i32) -> i32 { - if level == 0 { + let level = if level == 0 { DEFAULT_COMPRESSION_LEVEL } else { level - } + }; + + #[cfg(target_arch = "wasm32")] + let level = level.clamp(MIN_ZRIP_COMPRESSION_LEVEL, MAX_ZRIP_COMPRESSION_LEVEL); + + level } /// This type wraps a [`Vec`] to provide a [`Write`] interface that has a max @@ -131,9 +140,9 @@ impl ProfileCodec for NoopProfileCodec { /// A zstd-compatible profile codec. /// -/// WASM accepts levels `-7..=4`. Native targets accept the range reported by -/// `zstd::compression_level_range()`, making `-7..=4` the portable range. Level -/// `0` selects level `3` on every target. +/// Native targets accept the range reported by `zstd::compression_level_range()`. +/// WASM clamps levels to zrip's supported range of `-7..=4`. Level `0` selects +/// level `3` on every target. #[allow(unused)] pub struct ZstdProfileCodec; @@ -334,4 +343,10 @@ mod tests { fn zero_uses_default_compression_level() { assert_eq!(compress(0), compress(DEFAULT_COMPRESSION_LEVEL)); } + + #[test] + fn native_compression_level_is_not_clamped() { + assert_eq!(zstd_compression_level(22), 22); + assert!(!compress(22).is_empty()); + } } diff --git a/libdd-trace-utils/src/send_with_retry/compression.rs b/libdd-trace-utils/src/send_with_retry/compression.rs index f3ec1fc28b..466da0a7b5 100644 --- a/libdd-trace-utils/src/send_with_retry/compression.rs +++ b/libdd-trace-utils/src/send_with_retry/compression.rs @@ -8,6 +8,10 @@ use std::io::Write as _; const CONTENT_ENCODING_ZSTD: http::HeaderValue = http::HeaderValue::from_static("zstd"); #[cfg(feature = "compression")] const DEFAULT_COMPRESSION_LEVEL: i32 = 3; +#[cfg(all(feature = "compression", target_arch = "wasm32"))] +const MIN_ZRIP_COMPRESSION_LEVEL: i32 = -7; +#[cfg(all(feature = "compression", target_arch = "wasm32"))] +const MAX_ZRIP_COMPRESSION_LEVEL: i32 = 4; #[cfg(all(feature = "compression", not(target_arch = "wasm32")))] type ZstdEncoder = zstd::Encoder<'static, Vec>; @@ -16,21 +20,26 @@ type ZstdEncoder = zrip::FrameEncoder>; #[cfg(feature = "compression")] fn zstd_compression_level(level: i32) -> i32 { - if level == 0 { + let level = if level == 0 { DEFAULT_COMPRESSION_LEVEL } else { level - } + }; + + #[cfg(all(feature = "compression", target_arch = "wasm32"))] + let level = level.clamp(MIN_ZRIP_COMPRESSION_LEVEL, MAX_ZRIP_COMPRESSION_LEVEL); + + level } #[cfg(all(feature = "compression", not(target_arch = "wasm32")))] fn new_zstd_encoder(writer: Vec, level: i32) -> std::io::Result { - zstd::Encoder::new(writer, zstd_compression_level(level)) + zstd::Encoder::new(writer, level) } #[cfg(all(feature = "compression", target_arch = "wasm32"))] fn new_zstd_encoder(writer: Vec, level: i32) -> std::io::Result { - zrip::FrameEncoder::new(writer, zstd_compression_level(level)).map_err(std::io::Error::other) + zrip::FrameEncoder::new(writer, level).map_err(std::io::Error::other) } #[derive(Clone, Copy, Debug)] @@ -39,9 +48,9 @@ pub enum CompressionStrategy { #[cfg(feature = "compression")] /// Zstd-compatible compression. /// - /// WASM accepts levels `-7..=4`. Native targets accept the range reported by - /// `zstd::compression_level_range()`, making `-7..=4` the portable range. - /// Level `0` selects level `3` on every target. + /// Native targets accept the range reported by `zstd::compression_level_range()`. + /// WASM clamps levels to zrip's supported range of `-7..=4`. Level `0` selects + /// level `3` on every target. Zstd { level: i32, }, @@ -54,6 +63,8 @@ pub fn compress(data: Vec, strategy: CompressionStrategy) -> (Vec, Compr CompressionStrategy::None => (data, CompressionStrategy::None), #[cfg(feature = "compression")] CompressionStrategy::Zstd { level } => { + let level = zstd_compression_level(level); + let strategy = CompressionStrategy::Zstd { level }; // Start with an initial buffer // Allocate 1/10th of the original buffer, so we shouldn't add too // much memory usage, and no less than 256 bytes @@ -99,7 +110,7 @@ mod tests { #[test] fn zero_uses_default_compression_level() { let data = b"hello zstd".repeat(100); - let (default_compressed, _) = + let (default_compressed, strategy) = compress(data.clone(), CompressionStrategy::Zstd { level: 0 }); let (level_three_compressed, _) = compress( data, @@ -109,5 +120,16 @@ mod tests { ); assert_eq!(default_compressed, level_three_compressed); + assert!(matches!(strategy, CompressionStrategy::Zstd { level: 3 })); + } + + #[test] + fn native_compression_level_is_not_clamped() { + let data = b"hello zstd".repeat(100); + let (compressed, strategy) = + compress(data.clone(), CompressionStrategy::Zstd { level: 22 }); + + assert!(matches!(strategy, CompressionStrategy::Zstd { level: 22 })); + assert_eq!(decompress(&compressed).unwrap(), data); } } diff --git a/tests/wasm/tests/compression.rs b/tests/wasm/tests/compression.rs index e1c93dfdb1..abca3cc384 100644 --- a/tests/wasm/tests/compression.rs +++ b/tests/wasm/tests/compression.rs @@ -37,7 +37,7 @@ fn trace_compression_uses_zrip() { } #[wasm_bindgen_test::wasm_bindgen_test] -fn trace_compression_levels_are_portable() { +fn trace_compression_levels_are_clamped() { use trace_compression::{compress, CompressionStrategy}; let payload = b"hello zstd".repeat(100); @@ -46,11 +46,12 @@ fn trace_compression_levels_are_portable() { assert!(matches!(strategy, CompressionStrategy::Zstd { level: actual } if actual == level)); assert_eq!(zrip::decompress(&compressed).unwrap(), payload); } - for level in [-8, 5] { - let (uncompressed, strategy) = - compress(payload.clone(), CompressionStrategy::Zstd { level }); - assert!(matches!(strategy, CompressionStrategy::None)); - assert_eq!(uncompressed, payload); + for (level, expected) in [(-8, -7), (5, 4), (22, 4)] { + let (compressed, strategy) = compress(payload.clone(), CompressionStrategy::Zstd { level }); + assert!( + matches!(strategy, CompressionStrategy::Zstd { level: actual } if actual == expected) + ); + assert_eq!(zrip::decompress(&compressed).unwrap(), payload); } } @@ -59,10 +60,12 @@ fn trace_compression_zero_uses_level_three() { use trace_compression::{compress, CompressionStrategy}; let payload = b"hello zstd".repeat(100); - let (default_compressed, _) = compress(payload.clone(), CompressionStrategy::Zstd { level: 0 }); + let (default_compressed, strategy) = + compress(payload.clone(), CompressionStrategy::Zstd { level: 0 }); let (level_three_compressed, _) = compress(payload, CompressionStrategy::Zstd { level: 3 }); assert_eq!(default_compressed, level_three_compressed); + assert!(matches!(strategy, CompressionStrategy::Zstd { level: 3 })); } #[wasm_bindgen_test::wasm_bindgen_test] @@ -82,14 +85,17 @@ fn profiling_codecs_use_zrip() { } #[wasm_bindgen_test::wasm_bindgen_test] -fn profiling_compression_levels_are_portable() { +fn profiling_compression_levels_are_clamped() { let payload = b"hello profile".repeat(100); for level in [-7, 4] { let compressed = compress_profile(&payload, level).unwrap(); assert_eq!(zrip::decompress(&compressed).unwrap(), payload); } - for level in [-8, 5] { - assert!(compress_profile(&payload, level).is_err()); + for (level, expected) in [(-8, -7), (5, 4), (22, 4)] { + assert_eq!( + compress_profile(&payload, level).unwrap(), + compress_profile(&payload, expected).unwrap() + ); } } From e5dc6ef0dd4b2c3bdec14111320c702e070be8ff Mon Sep 17 00:00:00 2001 From: Roch Devost Date: Mon, 24 Aug 2026 14:12:32 -0400 Subject: [PATCH 4/4] docs(profiling): restore miri explanation --- libdd-profiling/src/pprof/test_utils.rs | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/libdd-profiling/src/pprof/test_utils.rs b/libdd-profiling/src/pprof/test_utils.rs index ec6fafbc28..42a15b30e3 100644 --- a/libdd-profiling/src/pprof/test_utils.rs +++ b/libdd-profiling/src/pprof/test_utils.rs @@ -6,8 +6,8 @@ use libdd_profiling_protobuf::prost_impls::{Profile, Sample}; fn deserialize_compressed_pprof(encoded: &[u8]) -> anyhow::Result { use prost::Message; - // Miri cannot call native zstd's FFI, so DefaultProfileCodec resolves to - // NoopProfileCodec under cfg(miri) and these bytes are intentionally uncompressed. + // The zstd bindings use FFI so they don't work under miri. This means the + // buffer isn't compressed, so simply convert to a vec. #[cfg(miri)] let buf = encoded.to_vec(); #[cfg(all(not(miri), not(target_arch = "wasm32")))]