Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
23 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
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
7cf69fd
fix: clippy
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
9b35068
Merge branch 'main' into anais/encoder-v1-to-v04-and-refacto-2
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
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
57a490c
Merge branch 'main' into anais/encoder-v1-to-v04-and-refacto-2
anais-raison Jul 13, 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
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
1 change: 1 addition & 0 deletions Cargo.lock

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

4 changes: 2 additions & 2 deletions datadog-sidecar-ffi/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1715,7 +1715,7 @@ pub unsafe extern "C" fn ddog_send_traces_to_sidecar(
// Write traces to the shared memory
let mut shm_slice = mapped_shm.as_slice_mut();
let shm_slice_len = shm_slice.len();
let written = match msgpack_encoder::v04::write_to_slice(&mut shm_slice, traces) {
let written = match msgpack_encoder::v04::write_to_slice_from_v04(&mut shm_slice, traces) {
Ok(()) => shm_slice_len - shm_slice.len(),
Err(_) => {
tracing::error!("Failed serializing the traces");
Expand Down Expand Up @@ -1745,7 +1745,7 @@ pub unsafe extern "C" fn ddog_send_traces_to_sidecar(
match blocking::send_trace_v04_bytes(
&mut parameters.transport,
&parameters.instance_id,
msgpack_encoder::v04::to_vec_with_capacity(traces, written as u32),
msgpack_encoder::v04::to_vec_with_capacity_from_v04(traces, written as u32),
check!(
(&parameters.tracer_headers_tags).try_into(),
"Failed to convert tracer headers tags"
Expand Down
9 changes: 9 additions & 0 deletions libdd-data-pipeline/src/trace_exporter/builder.rs
Original file line number Diff line number Diff line change
Expand Up @@ -66,6 +66,7 @@ pub struct TraceExporterBuilder<R: SharedRuntime> {
instrumentation_scope_version: String,
git_commit_sha: String,
process_tags: String,
container_id: String,
input_format: TraceExporterInputFormat,
output_format: TraceExporterOutputFormat,
dogstatsd_url: Option<String>,
Expand Down Expand Up @@ -138,6 +139,7 @@ impl<R: SharedRuntime> TraceExporterBuilder<R> {
instrumentation_scope_version: String::new(),
git_commit_sha: String::new(),
process_tags: String::new(),
container_id: String::new(),
input_format: TraceExporterInputFormat::default(),
output_format: TraceExporterOutputFormat::default(),
dogstatsd_url: None,
Expand Down Expand Up @@ -238,6 +240,12 @@ impl<R: SharedRuntime> TraceExporterBuilder<R> {
self
}

/// Set the `Datadog-Container-Id` header
pub fn set_container_id(&mut self, container_id: &str) -> &mut Self {
container_id.clone_into(&mut self.container_id);
self
}

/// Set OTLP trace instrumentation scope metadata.
pub fn set_otlp_instrumentation_scope(&mut self, name: &str, version: &str) -> &mut Self {
name.clone_into(&mut self.instrumentation_scope_name);
Expand Down Expand Up @@ -821,6 +829,7 @@ impl<R: SharedRuntime> TraceExporterBuilder<R> {
app_version: self.app_version,
runtime_id,
service: self.service,
container_id: self.container_id,
},
input_format: self.input_format,
output_format: self.output_format,
Expand Down
34 changes: 17 additions & 17 deletions libdd-data-pipeline/src/trace_exporter/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1424,7 +1424,7 @@ mod tests {
span_id: 2,
..Default::default()
}]];
let data = msgpack_encoder::v04::to_vec(&traces);
let data = msgpack_encoder::v04::to_vec_from_v04(&traces);

let resp = exporter.send(data.as_ref()).unwrap();
assert!(matches!(resp, AgentResponse::Unchanged));
Expand Down Expand Up @@ -1465,7 +1465,7 @@ mod tests {
name: BytesString::from_slice(b"test").unwrap(),
..Default::default()
}]];
let data = msgpack_encoder::v04::to_vec(&traces);
let data = msgpack_encoder::v04::to_vec_from_v04(&traces);
// `send` is synchronous and, in log mode, returns after writing through the
// capability without initiating any HTTP; combined with the structural assert
// above this is deterministic (no background worker can race the mock).
Expand Down Expand Up @@ -1506,7 +1506,7 @@ mod tests {
..Default::default()
}],
];
let data = msgpack_encoder::v04::to_vec(&traces);
let data = msgpack_encoder::v04::to_vec_from_v04(&traces);

let _result = exporter.send(data.as_ref()).expect("failed to send trace");

Expand Down Expand Up @@ -1606,7 +1606,7 @@ mod tests {
name: BytesString::from_slice(b"test").unwrap(),
..Default::default()
}]];
let data = msgpack_encoder::v04::to_vec(&traces);
let data = msgpack_encoder::v04::to_vec_from_v04(&traces);
let result = exporter.send(data.as_ref());

assert!(result.is_err());
Expand Down Expand Up @@ -1714,7 +1714,7 @@ mod tests {
name: BytesString::from_slice(b"test").unwrap(),
..Default::default()
}]];
let data = msgpack_encoder::v04::to_vec(&traces);
let data = msgpack_encoder::v04::to_vec_from_v04(&traces);
let result = exporter.send(data.as_ref());

assert!(result.is_err());
Expand Down Expand Up @@ -1818,7 +1818,7 @@ mod tests {
name: BytesString::from_slice(b"test").unwrap(),
..Default::default()
}]];
let data = msgpack_encoder::v04::to_vec(&traces);
let data = msgpack_encoder::v04::to_vec_from_v04(&traces);

let _result = exporter.send(data.as_ref()).expect("failed to send trace");

Expand Down Expand Up @@ -1878,7 +1878,7 @@ mod tests {
name: BytesString::from_slice(b"test").unwrap(),
..Default::default()
}]];
let data = msgpack_encoder::v04::to_vec(&traces);
let data = msgpack_encoder::v04::to_vec_from_v04(&traces);
let result = exporter.send(data.as_ref()).unwrap();

assert_eq!(
Expand Down Expand Up @@ -1920,7 +1920,7 @@ mod tests {
name: BytesString::from_slice(b"test").unwrap(),
..Default::default()
}]];
let data = msgpack_encoder::v04::to_vec(&traces);
let data = msgpack_encoder::v04::to_vec_from_v04(&traces);
let code = match exporter.send(data.as_ref()).unwrap_err() {
TraceExporterError::Request(e) => Some(e.status()),
_ => None,
Expand Down Expand Up @@ -1955,7 +1955,7 @@ mod tests {
name: BytesString::from_slice(b"test").unwrap(),
..Default::default()
}]];
let data = msgpack_encoder::v04::to_vec(&traces);
let data = msgpack_encoder::v04::to_vec_from_v04(&traces);
let err = exporter.send(data.as_ref());

assert!(err.is_err());
Expand Down Expand Up @@ -2134,7 +2134,7 @@ mod tests {
..Default::default()
}];

let data = msgpack_encoder::v04::to_vec(&[trace_chunk]);
let data = msgpack_encoder::v04::to_vec_from_v04(&[trace_chunk]);

// Wait for the info fetcher to get the config
while mock_info.calls() == 0 {
Expand Down Expand Up @@ -2202,7 +2202,7 @@ mod tests {
error: 0,
..Default::default()
}]];
let data = msgpack_encoder::v04::to_vec(&traces);
let data = msgpack_encoder::v04::to_vec_from_v04(&traces);
let result = exporter.send(data.as_ref());

assert!(
Expand Down Expand Up @@ -2254,7 +2254,7 @@ mod tests {
error: 0,
..Default::default()
}]];
let data = msgpack_encoder::v04::to_vec(&traces);
let data = msgpack_encoder::v04::to_vec_from_v04(&traces);
let result = exporter.send(data.as_ref());

assert!(
Expand Down Expand Up @@ -2311,7 +2311,7 @@ mod tests {
duration: 1,
..Default::default()
}]];
let data = msgpack_encoder::v04::to_vec(&traces);
let data = msgpack_encoder::v04::to_vec_from_v04(&traces);
exporter.send(data.as_ref()).unwrap();
mock_intake.assert();
}
Expand Down Expand Up @@ -2564,7 +2564,7 @@ mod single_threaded_tests {
..Default::default()
}];

let data = msgpack_encoder::v04::to_vec(&[trace_chunk]);
let data = msgpack_encoder::v04::to_vec_from_v04(&[trace_chunk]);

// Wait for the info fetcher to get the config
while agent_info::get_agent_info().is_none() {
Expand Down Expand Up @@ -2665,7 +2665,7 @@ mod single_threaded_tests {
..Default::default()
}];

let data = msgpack_encoder::v04::to_vec(&[trace_chunk]);
let data = msgpack_encoder::v04::to_vec_from_v04(&[trace_chunk]);

// Wait for agent_info to be present so that sending a trace will trigger the stats worker
// to start
Expand Down Expand Up @@ -2764,7 +2764,7 @@ mod single_threaded_tests {
duration: 10,
..Default::default()
}];
let data = msgpack_encoder::v04::to_vec(&[trace_chunk]);
let data = msgpack_encoder::v04::to_vec_from_v04(&[trace_chunk]);
let _ = exporter.send(data.as_ref());

let start = std::time::Instant::now();
Expand Down Expand Up @@ -2872,7 +2872,7 @@ mod single_threaded_tests {
duration: 10,
..Default::default()
}];
let data = msgpack_encoder::v04::to_vec(&[trace_chunk]);
let data = msgpack_encoder::v04::to_vec_from_v04(&[trace_chunk]);

// 1st send: /info has promoted v1_active=true, so this hits /v1.0/traces and 404s.
let result1 = exporter.send(&data);
Expand Down
4 changes: 2 additions & 2 deletions libdd-data-pipeline/src/trace_exporter/trace_serializer.rs
Original file line number Diff line number Diff line change
Expand Up @@ -120,7 +120,7 @@ impl TraceSerializer {
.max(MIN_BUFFER_CAPACITY);
let buff = match payload {
tracer_payload::TraceChunks::V04(p) => {
msgpack_encoder::v04::to_vec_with_capacity(p, capacity as u32)
msgpack_encoder::v04::to_vec_with_capacity_from_v04(p, capacity as u32)
}
tracer_payload::TraceChunks::V05(p) => {
let mut buff = Vec::with_capacity(capacity);
Expand All @@ -129,7 +129,7 @@ impl TraceSerializer {
buff
}
tracer_payload::TraceChunks::V1(p) => {
msgpack_encoder::v1::to_vec_with_capacity(p, capacity as u32, metadata)
msgpack_encoder::v1::to_vec_with_capacity_from_v04(p, capacity as u32, metadata)
}
};
self.previous_serialised_len
Expand Down
1 change: 1 addition & 0 deletions libdd-trace-utils/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,7 @@ required-features = ["bench-internals"]
anyhow = "1.0"
base64 = "0.22"
hex = "0.4"
itoa = "1.0"
hyper = { workspace = true, optional = true, default-features = false }
"http" = "1"
"http-body" = "1"
Expand Down
2 changes: 1 addition & 1 deletion libdd-trace-utils/benches/serialization.rs
Original file line number Diff line number Diff line change
Expand Up @@ -65,7 +65,7 @@ pub fn serialize_internal_to_msgpack(c: &mut Criterion) {
b.iter_batched(
|| vec![0u8; 12_000_000],
|mut vec| {
let _ = black_box(msgpack_encoder::v04::write_to_slice(
let _ = black_box(msgpack_encoder::v04::write_to_slice_from_v04(
&mut vec.as_mut_slice(),
black_box(&data),
));
Expand Down
26 changes: 26 additions & 0 deletions libdd-trace-utils/src/msgpack_encoder/mod.rs
Original file line number Diff line number Diff line change
@@ -1,6 +1,32 @@
// Copyright 2021-Present Datadog, Inc. https://www.datadoghq.com/
// SPDX-License-Identifier: Apache-2.0

//! # Encoder layout & naming convention
//!
//! ```text
//! msgpack_encoder/
//! ├── v04/
//! │ ├── mod.rs // public API + payload-level helpers
//! │ ├── span_v04.rs // v0.4 in-memory Span → v0.4 wire (native)
//! │ └── span_v1.rs // v1 in-memory Span → v0.4 wire (downgrade)
//! └── v1/
//! ├── mod.rs
//! ├── span_v04.rs // v0.4 in-memory Span → V1 wire (upgrade)
//! └── span_v1.rs // v1 in-memory Span → V1 wire (native)
//! ```
//!
//! - **Module (`v04`/`v1`) = output wire format.**
//! - **File suffix (`_v04`/`_v1`) = input span type.**
//! - **Public functions carry a `_from_<input>` suffix**, so a caller reads the *output* from the
//! module path and the *input* from the function name:
//!
//! | Module | Function | Input → Output |
//! |--------|----------|----------------|
//! | `v04::` | `to_vec_from_v04`, `write_to_slice_from_v04`, `to_encoded_byte_len_from_v04` | v04 → v0.4 (native) |
//! | `v04::` | `to_vec_from_v1`, `write_to_slice_from_v1`, `to_encoded_byte_len_from_v1` | v1 → v0.4 (downgrade) |
//! | `v1::` | `to_vec_from_v04`, `write_to_slice_from_v04`, `to_encoded_byte_len_from_v04` | v04 → V1 (upgrade) |
//! | `v1::` | `to_vec_from_v1`, `write_to_slice_from_v1`, `to_encoded_byte_len_from_v1` | v1 → V1 (native) |

pub mod v04;
pub mod v1;

Expand Down
Loading
Loading