Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
43 commits
Select commit Hold shift + click to select a range
d33370f
feat: encoder v1 to v04 + refacto
anais-raison Jun 22, 2026
02ab60b
refacto: renaming
anais-raison Jun 23, 2026
020b6cc
refacto name
anais-raison Jun 25, 2026
7fd7223
Merge branch 'main' into anais/encoder-v1-to-v04-and-refacto-2
anais-raison Jun 25, 2026
22c1044
fix: apply comments
anais-raison Jun 25, 2026
8936d5b
feat: add v1 decoder
anais-raison Jun 29, 2026
dd161cd
fix: comments
anais-raison Jun 29, 2026
5fbf201
Merge branch 'main' into anais/encoder-v1-to-v04-and-refacto-2
anais-raison Jun 29, 2026
a67ff41
Merge branch 'main' into anais/add-v1-decoder
anais-raison Jun 29, 2026
8c5d507
fix: clippy
anais-raison Jun 29, 2026
7cf69fd
fix: clippy
anais-raison Jun 29, 2026
5434161
Merge branch 'main' into anais/add-v1-decoder
anais-raison Jun 29, 2026
ffcd77f
Merge branch 'main' into anais/encoder-v1-to-v04-and-refacto-2
anais-raison Jun 29, 2026
8d3a26b
fix: comments
anais-raison Jul 1, 2026
51379b2
fix: comments
anais-raison Jul 2, 2026
4fbda9e
Merge branch 'main' into anais/encoder-v1-to-v04-and-refacto-2
anais-raison Jul 2, 2026
a077baa
fix: clippy
anais-raison Jul 2, 2026
67e02c2
fix: macro
anais-raison Jul 2, 2026
e315794
Merge branch 'main' into anais/add-v1-decoder
anais-raison Jul 2, 2026
9b35068
Merge branch 'main' into anais/encoder-v1-to-v04-and-refacto-2
anais-raison Jul 7, 2026
8b07b43
Merge branch 'main' into anais/add-v1-decoder
anais-raison Jul 7, 2026
e37e7fd
Merge branch 'main' into anais/encoder-v1-to-v04-and-refacto-2
anais-raison Jul 8, 2026
e33a348
Merge branch 'main' into anais/encoder-v1-to-v04-and-refacto-2
anais-raison Jul 8, 2026
05182fa
Merge branch 'main' into anais/add-v1-decoder
anais-raison Jul 8, 2026
bca3f24
fix: comments
anais-raison Jul 8, 2026
1081f93
Merge remote-tracking branch 'origin/anais/encoder-v1-to-v04-and-refa…
anais-raison Jul 8, 2026
dbbec42
fix: clippy
anais-raison Jul 8, 2026
1af6df1
fix: test
anais-raison Jul 8, 2026
dcf375e
fix: comments
anais-raison Jul 13, 2026
771e755
Merge branch 'main' into anais/encoder-v1-to-v04-and-refacto-2
anais-raison Jul 13, 2026
38da0ea
Merge remote-tracking branch 'origin/anais/encoder-v1-to-v04-and-refa…
anais-raison Jul 13, 2026
6c37c88
fix: merge
anais-raison Jul 13, 2026
57a490c
Merge branch 'main' into anais/encoder-v1-to-v04-and-refacto-2
anais-raison Jul 13, 2026
2391422
fix: comments
anais-raison Jul 15, 2026
41e75b1
fix: comments
anais-raison Jul 15, 2026
b44e5a6
Merge branch 'main' into anais/encoder-v1-to-v04-and-refacto-2
anais-raison Jul 15, 2026
cbca66a
fix: comments
anais-raison Jul 15, 2026
4ae04b8
Merge branch 'anais/encoder-v1-to-v04-and-refacto-2' into anais/add-v…
anais-raison Jul 15, 2026
7df05d6
Merge remote-tracking branch 'origin/main' into anais/add-v1-decoder
anais-raison Jul 15, 2026
196e0f3
fix: comments
anais-raison Jul 16, 2026
d418c61
fix: comments
anais-raison Jul 16, 2026
7c1161e
fix: comments
anais-raison Jul 17, 2026
8245c33
Merge branch 'main' into anais/add-v1-decoder
anais-raison Jul 17, 2026
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
59 changes: 46 additions & 13 deletions libdd-data-pipeline/src/trace_exporter/trace_serializer.rs
Original file line number Diff line number Diff line change
Expand Up @@ -17,7 +17,7 @@ use libdd_trace_utils::msgpack_encoder;
use libdd_trace_utils::span::{v04::Span, TraceData};
use libdd_trace_utils::trace_utils::{self, TracerHeaderTags};
use libdd_trace_utils::tracer_metadata::TracerMetadata;
use libdd_trace_utils::tracer_payload::{self, TraceEncoding};
use libdd_trace_utils::tracer_payload::{self};

/// Minimal capacity of fresh buffers allocated to encode traces, in bytes.
const MIN_BUFFER_CAPACITY: usize = 1024;
Expand Down Expand Up @@ -59,7 +59,7 @@ impl TraceSerializer {
let chunks = payload.size();
let headers =
self.build_traces_headers(header_tags, chunks, agent_payload_response_version);
let mp_payload = self.serialize_payload(&payload, metadata)?;
let mp_payload = self.serialize_payload(&payload, metadata, output_format)?;

Ok(PreparedTracesPayload {
data: mp_payload,
Expand All @@ -78,12 +78,19 @@ impl TraceSerializer {
TraceExporterError::Deserialization(DecodeError::InvalidFormat(e.to_string()))
};
match output_format {
TraceExporterOutputFormat::V1 => Ok(tracer_payload::TraceChunks::V1(traces)),
TraceExporterOutputFormat::V04 => {
trace_utils::collect_trace_chunks(traces, TraceEncoding::V04).map_err(map_err)
// v0.4 input spans are kept as-is in `TraceChunks::V04`. Whether they go out as v0.4
// or are cross-encoded into V1 on the wire is decided in `serialize_payload`.
//
// APMSP-2812 - TODO: when the data-pipeline gains a V1-native input model (its own
// `v1::Span`-shaped builder), route `OutputFormat::V1` to
// `TraceChunks::V1(v1::TracerPayload)` instead and serialize via
// `to_vec_from_payload_v1`. A `StatSpan` impl on `v1::Span<T>` will also be needed
// if client-side stats are enabled on the V1-native path.
TraceExporterOutputFormat::V04 | TraceExporterOutputFormat::V1 => {
Ok(tracer_payload::TraceChunks::V04(traces))
}
TraceExporterOutputFormat::V05 => {
trace_utils::collect_trace_chunks(traces, TraceEncoding::V05).map_err(map_err)
trace_utils::convert_trace_chunks_v04_to_v05(traces).map_err(map_err)
}
}
}
Expand Down Expand Up @@ -113,23 +120,41 @@ impl TraceSerializer {
&self,
payload: &tracer_payload::TraceChunks<T>,
metadata: &TracerMetadata,
output_format: TraceExporterOutputFormat,
) -> Result<Vec<u8>, TraceExporterError> {
let capacity = self
.previous_serialised_len
.load(Ordering::Relaxed)
.max(MIN_BUFFER_CAPACITY);
let buff = match payload {
tracer_payload::TraceChunks::V04(p) => {
let buff = match (payload, output_format) {
(tracer_payload::TraceChunks::V04(p), TraceExporterOutputFormat::V04) => {
msgpack_encoder::v04::to_vec_with_capacity_from_v04(p, capacity as u32)
}
tracer_payload::TraceChunks::V05(p) => {
// v0.4 spans cross-encoded as V1 on the wire (used when the agent advertises
// /v1.0/traces).
(tracer_payload::TraceChunks::V04(p), TraceExporterOutputFormat::V1) => {
msgpack_encoder::v1::to_vec_with_capacity_from_v04(p, capacity as u32, metadata)
}
(tracer_payload::TraceChunks::V05(p), TraceExporterOutputFormat::V05) => {
let mut buff = Vec::with_capacity(capacity);
rmp_serde::encode::write(&mut buff, p)
.map_err(TraceExporterError::Serialization)?;
buff
}
tracer_payload::TraceChunks::V1(p) => {
msgpack_encoder::v1::to_vec_with_capacity_from_v04(p, capacity as u32, metadata)
// Native V1 input model, serialized directly — the payload carries its own
// tracer-level metadata, so no `TracerMetadata` is needed here.
(tracer_payload::TraceChunks::V1(p), TraceExporterOutputFormat::V1) => {
msgpack_encoder::v1::to_vec_with_capacity_from_v1(p, capacity as u32)
}
// `collect_and_process_traces` only produces (V04, V04|V1), (V05, V05),
// or (V1, V1) — any other combination here is a programming error.
_ => {
return Err(TraceExporterError::Deserialization(
DecodeError::InvalidFormat(
"Unsupported (TraceChunks, OutputFormat) combination for serialization"
.to_owned(),
),
));
}
};
self.previous_serialised_len
Expand Down Expand Up @@ -277,7 +302,11 @@ mod tests {
.collect_and_process_traces(original_traces.clone(), TraceExporterOutputFormat::V04)
.unwrap();

let result = serializer.serialize_payload(&payload, &TracerMetadata::default());
let result = serializer.serialize_payload(
&payload,
&TracerMetadata::default(),
TraceExporterOutputFormat::V04,
);
assert!(result.is_ok());

let serialized = result.unwrap();
Expand Down Expand Up @@ -312,7 +341,11 @@ mod tests {
.collect_and_process_traces(original_traces.clone(), TraceExporterOutputFormat::V05)
.unwrap();

let result = serializer.serialize_payload(&payload, &TracerMetadata::default());
let result = serializer.serialize_payload(
&payload,
&TracerMetadata::default(),
TraceExporterOutputFormat::V05,
);
assert!(result.is_ok());

let serialized = result.unwrap();
Expand Down
13 changes: 13 additions & 0 deletions libdd-trace-utils/src/msgpack_decoder/decode/buffer.rs
Original file line number Diff line number Diff line change
Expand Up @@ -40,6 +40,11 @@ impl<T: DeserializableTraceData> Buffer<T> {
T::get_mut_slice(&mut self.0)
}

/// Returns an immutable reference to the underlying slice, without advancing the buffer.
pub fn as_slice(&self) -> &[u8] {
self.0.borrow()
}

/// 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 All @@ -52,6 +57,14 @@ impl<T: DeserializableTraceData> Buffer<T> {
pub fn read_string(&mut self) -> Result<T::Text, DecodeError> {
T::read_string(&mut self.0)
}

/// Caps a decoded element count at the bytes remaining in the buffer. Each msgpack
/// element needs >=1 byte on the wire, so a length prefix can't legitimately exceed
/// the remaining bytes — this prevents a malicious count (e.g. 0xFFFFFFFF) from
/// forcing a huge pre-allocation before any element is read.
pub fn capped_capacity(&self, count: usize) -> usize {
count.min(self.len())
}
}

impl<T: DeserializableTraceData> Deref for Buffer<T> {
Expand Down
1 change: 1 addition & 0 deletions libdd-trace-utils/src/msgpack_decoder/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -4,3 +4,4 @@
pub mod decode;
pub mod v04;
pub mod v05;
pub mod v1;
Loading
Loading