Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

5 changes: 5 additions & 0 deletions libdd-trace-utils/src/msgpack_decoder/decode/buffer.rs
Original file line number Diff line number Diff line change
Expand Up @@ -45,6 +45,11 @@ impl<T: DeserializableTraceData> Buffer<T> {
self.0.borrow()
}

/// Returns the underlying owned bytes buffer.
pub fn bytes(&self) -> &T::Bytes {
&self.0
}

/// Tries to extract a slice of `bytes` from the buffer and advances the buffer.
pub fn try_slice_and_advance(&mut self, bytes: usize) -> Option<T::Bytes> {
T::try_slice_and_advance(&mut self.0, bytes)
Expand Down
162 changes: 158 additions & 4 deletions libdd-trace-utils/src/msgpack_decoder/v1/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -212,14 +212,64 @@ where
/// Consumes and discards the msgpack value at the current buffer position, regardless of its
/// type. Used to skip unknown keys for forward compatibility: if the V1 format gains new fields,
/// older decoders shouldn't reject the whole payload just because they don't recognize a key.
///
/// Any inline string encountered while skipping (at any nesting depth) is interned into `table`,
/// same as a recognized field would: skipping a value must not desync later back-references to
/// strings that happen to also appear inside it.
pub(super) fn skip_unknown_value<T: DeserializableTraceData>(
buf: &mut Buffer<T>,
) -> Result<(), DecodeError> {
rmpv::decode::read_value(buf.as_mut_slice())
table: &mut StringTable<T>,
) -> Result<(), DecodeError>
where
T::Text: Clone,
{
// Snapshot the buffer's owning handle *before* advancing past the skipped value: any string
// found inside it will be a substring of this exact allocation, so this is what
// `T::intern_skipped_str` must derive ownership from. Cloning is cheap (a refcount bump for
// `T::Bytes = Bytes`), unaffected by the lied `'static` lifetime `as_mut_slice` exposes.
let owner = buf.bytes().clone();
let value = rmpv::decode::read_value_ref(buf.as_mut_slice())
.map_err(|_| DecodeError::InvalidFormat("Failed to skip unknown V1 value".to_owned()))?;
record_strings_in_value_ref::<T>(&value, &owner, table);
Ok(())
}

/// Recursively walks a parsed [`rmpv::ValueRef`], interning every string it contains into
/// `table`. Strings with invalid UTF-8 are ignored: they can never have been produced by
/// [`read_interned_string`]'s own encoder-side counterpart, so they can't be the target of a
/// later back-reference either.
///
/// `owner` must be a snapshot of the buffer taken before it was advanced past `value`: the
/// strings inside `value` report a lied `'static` lifetime (see `Buffer::as_mut_slice`) but
/// really borrow from `owner`'s memory.
fn record_strings_in_value_ref<T: DeserializableTraceData>(
value: &rmpv::ValueRef<'static>,
Comment thread
anais-raison marked this conversation as resolved.
owner: &T::Bytes,
table: &mut StringTable<T>,
) where
T::Text: Clone,
{
match value {
rmpv::ValueRef::String(s) => {
if let Some(s) = (*s).into_str() {
table.record(&T::intern_skipped_str(owner, s));
}
}
rmpv::ValueRef::Array(items) => {
for item in items {
record_strings_in_value_ref::<T>(item, owner, table);
}
}
rmpv::ValueRef::Map(entries) => {
for (key, val) in entries {
record_strings_in_value_ref::<T>(key, owner, table);
record_strings_in_value_ref::<T>(val, owner, table);
}
}
_ => {}
}
}

/// Decodes the top-level V1 payload map: tracer metadata fields + chunks array.
fn decode_payload<T: DeserializableTraceData>(
buf: &mut Buffer<T>,
Expand Down Expand Up @@ -256,7 +306,7 @@ where
trace_key::ATTRIBUTES => {
payload.attributes = span::read_attributes_map(buf, table)?;
}
_unknown => skip_unknown_value(buf)?,
_unknown => skip_unknown_value(buf, table)?,
}
}

Expand Down Expand Up @@ -362,7 +412,7 @@ where
)
})?;
}
_unknown => skip_unknown_value(buf)?,
_unknown => skip_unknown_value(buf, table)?,
}
}

Expand Down Expand Up @@ -693,4 +743,108 @@ mod tests {
},
);
}

// ---------------------------------------------------------------------------------------------
// Forward-compatibility: unknown map keys must be skipped for every V1 map type. This test
// hand-builds wire bytes (the encoder never emits unknown keys) with `rmp::encode`, injecting
// a future/unknown key at every nesting level (payload, chunk, span, span_link, span_event),
// and asserts the surrounding known fields still decode correctly.
// ---------------------------------------------------------------------------------------------

use rmp::encode::{self, ByteBuf};

/// Writes a `u8` msgpack map key.
fn wkey(buf: &mut ByteBuf, k: u8) {
encode::write_uint(buf, k as u64).unwrap();
}

/// Exercises unknown-key skipping at every V1 nesting level in a single payload:
/// - payload: unknown field 99 carries the first occurrence of "ghost" (must be harvested as
/// table id 1, a scalar skip at the chunk level exercises the recursive skip, and a
/// subsequent span field back-references "prod" by id to prove the table wasn't desynced).
/// - chunk: unknown field 77 carries a nested `[uint, str, map]` value (recursive skip), whose
/// inline string "buried" must also be harvested (table id 3).
/// - span: unknown field 88 carries a scalar (f64) with no string to harvest.
/// - span_link / span_event: unknown fields 55 / 66 precede their known sibling field.
#[test]
fn unknown_keys_are_skipped_at_every_level() {
let mut span_link = ByteBuf::new();
encode::write_map_len(&mut span_link, 2).unwrap();
wkey(&mut span_link, 55); // unknown span_link key
encode::write_bool(&mut span_link, false).unwrap();
wkey(&mut span_link, span_link_key::SPAN_ID);
encode::write_uint(&mut span_link, 777).unwrap();

let mut span_event = ByteBuf::new();
encode::write_map_len(&mut span_event, 2).unwrap();
wkey(&mut span_event, 66); // unknown span_event key
encode::write_uint(&mut span_event, 999).unwrap();
wkey(&mut span_event, span_event_key::TIME);
encode::write_uint(&mut span_event, 123).unwrap();

let mut span = ByteBuf::new();
encode::write_map_len(&mut span, 6).unwrap();
wkey(&mut span, span_key::SPAN_ID);
encode::write_uint(&mut span, 42).unwrap();
wkey(&mut span, span_key::START);
encode::write_uint(&mut span, 100).unwrap();
wkey(&mut span, 88); // unknown span key: scalar, nothing to harvest
encode::write_f64(&mut span, 2.5).unwrap();
wkey(&mut span, span_key::SERVICE);
encode::write_uint(&mut span, 2).unwrap(); // back-reference to encoder id 2 ("prod")
wkey(&mut span, span_key::SPAN_LINKS);
encode::write_array_len(&mut span, 1).unwrap();
span.as_mut_vec().extend_from_slice(&span_link.into_vec());
wkey(&mut span, span_key::SPAN_EVENTS);
encode::write_array_len(&mut span, 1).unwrap();
span.as_mut_vec().extend_from_slice(&span_event.into_vec());

let mut chunk = ByteBuf::new();
encode::write_map_len(&mut chunk, 3).unwrap();
wkey(&mut chunk, 77); // unknown chunk key: nested value, recursive skip + string harvest
encode::write_array_len(&mut chunk, 3).unwrap();
encode::write_uint(&mut chunk, 1).unwrap();
encode::write_str(&mut chunk, "buried").unwrap(); // first occurrence -> table id 3
encode::write_map_len(&mut chunk, 1).unwrap();
encode::write_uint(&mut chunk, 5).unwrap();
encode::write_bool(&mut chunk, true).unwrap();
wkey(&mut chunk, chunk_key::TRACE_ID);
encode::write_bin(&mut chunk, &[9u8; 16]).unwrap();
wkey(&mut chunk, chunk_key::SPANS);
encode::write_array_len(&mut chunk, 1).unwrap();
chunk.as_mut_vec().extend_from_slice(&span.into_vec());

let mut buf = ByteBuf::new();
encode::write_map_len(&mut buf, 3).unwrap();
wkey(&mut buf, 99); // unknown payload key: first occurrence "ghost" -> table id 1
encode::write_str(&mut buf, "ghost").unwrap();
wkey(&mut buf, trace_key::ENV_REF);
encode::write_str(&mut buf, "prod").unwrap(); // first occurrence -> table id 2
wkey(&mut buf, trace_key::CHUNKS);
encode::write_array_len(&mut buf, 1).unwrap();
buf.as_mut_vec().extend_from_slice(&chunk.into_vec());
let buf = buf.into_vec();

let (decoded, consumed) =
from_bytes(Bytes::from(buf.clone())).expect("unknown keys must be skipped");
assert_eq!(
consumed,
buf.len(),
"decoder must consume every skipped value"
);
assert_eq!(decoded.env.as_str(), "prod");
let chunk = &decoded.chunks[0];
assert_eq!(chunk.trace_id, [9u8; 16]);
let span = &chunk.spans[0];
assert_eq!(span.span_id, 42);
assert_eq!(span.start, 100);
assert_eq!(
span.service.as_str(),
"prod",
"back-reference must still resolve correctly: harvesting \"ghost\" and \"buried\" \
while skipping unknown fields must not desync the string table"
);
assert_eq!(span.span_links[0].span_id, 777);
assert_eq!(span.span_events[0].time_unix_nano, 123);
}
}
6 changes: 3 additions & 3 deletions libdd-trace-utils/src/msgpack_decoder/v1/span.rs
Original file line number Diff line number Diff line change
Expand Up @@ -92,7 +92,7 @@ where
})?;
span.span_kind = SpanKind::from(kind);
}
_unknown => skip_unknown_value(buf)?,
_unknown => skip_unknown_value(buf, table)?,
}
}

Expand Down Expand Up @@ -276,7 +276,7 @@ where
DecodeError::InvalidFormat(format!("V1 span_link flags {v} exceeds u32::MAX"))
})?;
}
_unknown => skip_unknown_value(buf)?,
_unknown => skip_unknown_value(buf, table)?,
}
}
Ok(link)
Expand Down Expand Up @@ -325,7 +325,7 @@ where
span_event_key::ATTRIBUTES => {
event.attributes = read_attributes_map(buf, table)?;
}
_unknown => skip_unknown_value(buf)?,
_unknown => skip_unknown_value(buf, table)?,
}
}
Ok(event)
Expand Down
21 changes: 20 additions & 1 deletion libdd-trace-utils/src/span/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -67,7 +67,7 @@ impl SpanText for BytesString {
}
}

pub trait SpanBytes: Debug + Eq + Hash + Borrow<[u8]> + Serialize + Default {
pub trait SpanBytes: Debug + Eq + Hash + Borrow<[u8]> + Serialize + Default + Clone {
fn from_static_bytes(value: &'static [u8]) -> Self;
}

Expand Down Expand Up @@ -98,6 +98,13 @@ pub trait DeserializableTraceData: TraceData {
fn try_slice_and_advance(buf: &mut Self::Bytes, bytes: usize) -> Option<Self::Bytes>;

fn read_string(buf: &mut Self::Bytes) -> Result<Self::Text, DecodeError>;

/// Interns a string found while walking a value through `get_mut_slice`'s lied `'static`
/// view (e.g. skipping an unrecognized V1 field for forward compatibility). `s` really
/// borrows from `owner`'s memory, not `'static`: implementations must derive `Self::Text`
/// from `owner` itself rather than trusting that lifetime, so a refcounted backing
/// allocation isn't freed out from under the interned string.
fn intern_skipped_str(owner: &Self::Bytes, s: &'static str) -> Self::Text;
}

/// TraceData implementation using `Bytes` and `BytesString`.
Expand Down Expand Up @@ -149,6 +156,11 @@ impl DeserializableTraceData for BytesData {
}
Ok(string)
}

#[inline]
fn intern_skipped_str(owner: &Bytes, s: &'static str) -> BytesString {
BytesString::from_bytes_slice(owner, s)
}
}

/// TraceData implementation using `&str` and `&[u8]`.
Expand Down Expand Up @@ -179,6 +191,13 @@ impl<'a> DeserializableTraceData for SliceData<'a> {
str
})
}

#[inline]
fn intern_skipped_str(_owner: &&'a [u8], s: &'static str) -> &'a str {
// No refcounted allocation to preserve here: `s` borrows from a plain slice the
// caller owns for `'a`, and a `'static` reference is always a valid `'a` reference.
s
}
}

#[derive(Debug)]
Expand Down
Loading