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
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.

1 change: 1 addition & 0 deletions datadog-opentelemetry/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -64,6 +64,7 @@ libdd-trace-utils = { workspace = true }
libdd-telemetry = { workspace = true }
libdd-common = { workspace = true }
libdd-tinybytes = { workspace = true }
percent-encoding = "2.3.1"

[dev-dependencies]
assert_unordered = "0.3"
Expand Down
5 changes: 4 additions & 1 deletion datadog-opentelemetry/examples/propagator/src/server.rs
Original file line number Diff line number Diff line change
Expand Up @@ -119,7 +119,10 @@ async fn send_request(
) -> std::result::Result<(), Box<dyn std::error::Error + Send + Sync + 'static>> {
let client = Client::builder(TokioExecutor::new()).build_http();

let cx = Context::current();
let cx = Context::current().with_baggage(vec![
KeyValue::new("request-id", "xyz-123"),
KeyValue::new("caller", "rust-propagator-example"),
]);

let mut req = hyper::Request::builder().uri(url);
global::get_text_map_propagator(|propagator| {
Expand Down
47 changes: 47 additions & 0 deletions datadog-opentelemetry/src/core/configuration/configuration.rs
Original file line number Diff line number Diff line change
Expand Up @@ -773,6 +773,8 @@ pub enum TracePropagationStyle {
Datadog,
/// W3C Trace Context propagation format using `traceparent` and `tracestate` headers.
TraceContext,
/// W3C Baggage propagation format using the `baggage` header.
Baggage,
/// No propagation - trace context is not propagated.
None,
}
Expand Down Expand Up @@ -804,6 +806,7 @@ impl FromStr for TracePropagationStyle {
match s.trim().to_lowercase().as_str() {
"datadog" => Ok(TracePropagationStyle::Datadog),
"tracecontext" => Ok(TracePropagationStyle::TraceContext),
"baggage" => Ok(TracePropagationStyle::Baggage),
"none" => Ok(TracePropagationStyle::None),
_ => Err(format!("Unknown trace propagation style: '{s}'")),
}
Expand All @@ -815,6 +818,7 @@ impl Display for TracePropagationStyle {
let style = match self {
TracePropagationStyle::Datadog => "datadog",
TracePropagationStyle::TraceContext => "tracecontext",
TracePropagationStyle::Baggage => "baggage",
TracePropagationStyle::None => "none",
};
write!(f, "{style}")
Expand Down Expand Up @@ -1822,6 +1826,7 @@ fn default_config() -> Config {
Some(vec![
TracePropagationStyle::Datadog,
TracePropagationStyle::TraceContext,
TracePropagationStyle::Baggage,

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Since this is a change in behavior, it should be flagged as a breaking change feat(baggage)!:...

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

We should have made this enum #[non_exhaustive] 🫤

]),
),
trace_propagation_style_extract: ConfigItem::new(
Expand Down Expand Up @@ -2726,6 +2731,7 @@ mod tests {
Some(vec![
TracePropagationStyle::Datadog,
TracePropagationStyle::TraceContext,
TracePropagationStyle::Baggage,
])
.as_deref()
);
Expand All @@ -2737,6 +2743,47 @@ mod tests {
assert!(config.trace_propagation_extract_first());
}

#[test]
fn test_propagation_style_baggage_parsed_from_env() {
// "baggage" is recognised as a valid style value (case-insensitive) in all three env vars.
let mut sources = CompositeSource::new();
sources.add_source(HashMapSource::from_iter(
[
("DD_TRACE_PROPAGATION_STYLE", "datadog,tracecontext,baggage"),
("DD_TRACE_PROPAGATION_STYLE_EXTRACT", "Baggage,datadog"),
("DD_TRACE_PROPAGATION_STYLE_INJECT", "BAGGAGE,tracecontext"),
],
ConfigSourceOrigin::EnvVar,
));
let config = Config::builder_with_sources(&sources).build();

assert_eq!(
config.trace_propagation_style(),
Some(vec![
TracePropagationStyle::Datadog,
TracePropagationStyle::TraceContext,
TracePropagationStyle::Baggage,
])
.as_deref()
);
assert_eq!(
config.trace_propagation_style_extract(),
Some(vec![
TracePropagationStyle::Baggage,
TracePropagationStyle::Datadog,
])
.as_deref()
);
assert_eq!(
config.trace_propagation_style_inject(),
Some(vec![
TracePropagationStyle::Baggage,
TracePropagationStyle::TraceContext,
])
.as_deref()
);
}

#[test]
fn test_stats_computation_enabled_config() {
let mut sources = CompositeSource::new();
Expand Down
270 changes: 270 additions & 0 deletions datadog-opentelemetry/src/propagation/baggage.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,270 @@
// Copyright 2025-Present Datadog, Inc. https://www.datadoghq.com/
// SPDX-License-Identifier: Apache-2.0

//! W3C Baggage propagation (`baggage` header).
//!
//! This module exposes the header key so the composite propagator can include it in its `fields()`
//! list. It contains only extraction code.
//!
//! Injection is performed by [`opentelemetry_sdk::propagation::BaggagePropagator`]
//! at the [`DatadogPropagator`](crate::text_map_propagator::DatadogPropagator) layer, which has
//! access to the OTel [`Context`](opentelemetry::Context) that carries baggage.
use std::sync::LazyLock;

use opentelemetry::baggage::{Baggage, KeyValueMetadata};
use percent_encoding::percent_decode_str;

use crate::{dd_warn, propagation::carrier::Extractor};

/// The W3C `baggage` header name.
pub const BAGGAGE_KEY: &str = "baggage";

/// Cap on the number of baggage members extracted. Entries beyond this are dropped.
const MAX_BAGGAGE_MEMBERS: usize = 32;

/// Cap on the cumulative byte length of baggage members parsed. Once exceeded, remaining entries
/// are dropped.
const MAX_BAGGAGE_LENGTH: usize = 1024;

static BAGGAGE_HEADER_KEYS: LazyLock<[String; 1]> = LazyLock::new(|| [BAGGAGE_KEY.to_owned()]);

/// Returns the header keys used by the W3C baggage propagator.
pub fn keys() -> &'static [String] {
BAGGAGE_HEADER_KEYS.as_slice()
}

fn parse_baggage_member(baggage_member: &str) -> Option<KeyValueMetadata> {
let mut member = baggage_member.split(';');
let Some(name_and_value) = member.next() else {
dd_warn!("Propagator (baggage): invalid format");
return None;
};
let mut iter = name_and_value.split('=');
let (Some(name), Some(value)) = (iter.next(), iter.next()) else {
dd_warn!("Propagator (baggage): invalid key-value format");
return None;
};
let decode_name = percent_decode_str(name).decode_utf8();
let decode_value = percent_decode_str(value).decode_utf8();

let (Ok(name), Ok(value)) = (decode_name, decode_value) else {
dd_warn!("Propagator (baggage): invalid percent encoded UTF8 string in key values");
return None;
};

// decode and trim metadata entries associated with the key-value
let decoded_props = member
.flat_map(|prop| percent_decode_str(prop).decode_utf8())
.enumerate()
.fold(String::new(), |mut acc, (i, prop)| {
if i != 0 {
acc.push(';');
}
acc.push_str(prop.trim());
acc
});

Some(KeyValueMetadata::new(
name.trim().to_owned(),
value.trim().to_string(),
decoded_props,
))
}

pub(crate) fn extract_baggage(extractor: &dyn Extractor) -> Option<Baggage> {
let header_value = extractor.get(BAGGAGE_KEY)?;
let mut members = 0;
let mut allocated_size = 0;
let baggage = header_value.split(',')
.take_while(|member| {
allocated_size += member.len();
let drop_entry = allocated_size > MAX_BAGGAGE_LENGTH;
if drop_entry {
dd_warn!("Propagator (baggage): ignored baggage key-values, only first {} bytes propagated", MAX_BAGGAGE_LENGTH)
}
!drop_entry
})
.filter_map(parse_baggage_member)
.take_while(|_| {
members +=1;
let drop_entry = members > MAX_BAGGAGE_MEMBERS;
if drop_entry {
dd_warn!("Propagator (baggage): ignored baggage key-values, only first {} propagated", MAX_BAGGAGE_MEMBERS)
}
!drop_entry
});
Some(Baggage::from(baggage))
}

#[cfg(test)]
mod tests {
use super::*;
use opentelemetry::{baggage::BaggageMetadata, Key, StringValue};
use std::collections::HashMap;

fn valid_extract_data() -> Vec<(&'static str, HashMap<Key, StringValue>)> {
vec![
// "valid w3cHeader"
(
"key1=val1,key2=val2",
vec![
(Key::new("key1"), StringValue::from("val1")),
(Key::new("key2"), StringValue::from("val2")),
]
.into_iter()
.collect(),
),
// "valid w3cHeader with spaces"
(
"key1 = val1, key2 =val2 ",
vec![
(Key::new("key1"), StringValue::from("val1")),
(Key::new("key2"), StringValue::from("val2")),
]
.into_iter()
.collect(),
),
// "valid header with url-escaped comma"
(
"key1=val1,key2=val2%2Cval3",
vec![
(Key::new("key1"), StringValue::from("val1")),
(Key::new("key2"), StringValue::from("val2,val3")),
]
.into_iter()
.collect(),
),
// "valid header with an invalid header"
(
"key1=val1,key2=val2,a,val3",
vec![
(Key::new("key1"), StringValue::from("val1")),
(Key::new("key2"), StringValue::from("val2")),
]
.into_iter()
.collect(),
),
// "valid header with no value"
(
"key1=,key2=val2",
vec![
(Key::new("key1"), StringValue::from("")),
(Key::new("key2"), StringValue::from("val2")),
]
.into_iter()
.collect(),
),
]
}

#[allow(clippy::type_complexity)]
fn valid_extract_data_with_metadata(
) -> Vec<(&'static str, HashMap<Key, (StringValue, BaggageMetadata)>)> {
vec![
// "valid w3cHeader with properties"
("key1=val1,key2=val2;prop=1", vec![(Key::new("key1"), (StringValue::from("val1"), BaggageMetadata::default())), (Key::new("key2"), (StringValue::from("val2"), BaggageMetadata::from("prop=1")))].into_iter().collect()),
// prop can don't need to be key value pair
("key1=val1,key2=val2;prop1", vec![(Key::new("key1"), (StringValue::from("val1"), BaggageMetadata::default())), (Key::new("key2"), (StringValue::from("val2"), BaggageMetadata::from("prop1")))].into_iter().collect()),
("key1=value1;property1;property2, key2 = value2, key3=value3; propertyKey=propertyValue",
vec![
(Key::new("key1"), (StringValue::from("value1"), BaggageMetadata::from("property1;property2"))),
(Key::new("key2"), (StringValue::from("value2"), BaggageMetadata::default())),
(Key::new("key3"), (StringValue::from("value3"), BaggageMetadata::from("propertyKey=propertyValue"))),
].into_iter().collect()),
]
}

#[test]
fn test_extract_baggage() {
for (header_value, kvs) in valid_extract_data() {
let mut extractor: HashMap<String, String> = HashMap::new();
extractor.insert(BAGGAGE_KEY.to_string(), header_value.to_string());
let baggage = extract_baggage(&extractor).expect("baggage extracted");

assert_eq!(kvs.len(), baggage.len());
for (key, (value, _metadata)) in &baggage {
assert_eq!(Some(value), kvs.get(key))
}
}
}

#[test]
fn test_extract_baggage_with_metadata() {
for (header_value, kvm) in valid_extract_data_with_metadata() {
let mut extractor: HashMap<String, String> = HashMap::new();
extractor.insert(BAGGAGE_KEY.to_string(), header_value.to_string());
let baggage = extract_baggage(&extractor).expect("baggage extracted");

assert_eq!(kvm.len(), baggage.len());
for (key, value_and_prop) in &baggage {
assert_eq!(Some(value_and_prop), kvm.get(key))
}
}
}

#[test]
fn test_extract_baggage_respects_max_members() {
let total = MAX_BAGGAGE_MEMBERS + 8;
let header_value = (0..total)
.map(|i| format!("k{i}=v{i}"))
.collect::<Vec<_>>()
.join(",");

let mut extractor: HashMap<String, String> = HashMap::new();
extractor.insert(BAGGAGE_KEY.to_string(), header_value);
let baggage = extract_baggage(&extractor).expect("baggage extracted");

assert_eq!(baggage.len(), MAX_BAGGAGE_MEMBERS);
for i in 0..MAX_BAGGAGE_MEMBERS {
let key = Key::new(format!("k{i}"));
assert_eq!(
baggage.get(&key).map(|v| v.as_str()),
Some(format!("v{i}").as_str())
);
}
for i in MAX_BAGGAGE_MEMBERS..total {
let key = Key::new(format!("k{i}"));
assert!(baggage.get(&key).is_none());
}
}

#[test]
fn test_extract_baggage_respects_max_length() {
// Each member is exactly 50 bytes: "k{NN}={padding...}".
// 21 members * 50 = 1050 bytes (separators excluded by the
// implementation), so the 21st entry pushes allocated_size past
// MAX_BAGGAGE_LENGTH (1024) and must be dropped.
let member_size = 50;
let total = 25;
let members: Vec<String> = (0..total)
.map(|i| {
let prefix = format!("k{i:02}=");
let pad = "x".repeat(member_size - prefix.len());
format!("{prefix}{pad}")
})
.collect();
assert!(members.iter().map(|m| m.len()).sum::<usize>() > MAX_BAGGAGE_LENGTH);
let header_value = members.join(",");

let mut extractor: HashMap<String, String> = HashMap::new();
extractor.insert(BAGGAGE_KEY.to_string(), header_value);
let baggage = extract_baggage(&extractor).expect("baggage extracted");

let expected_kept = MAX_BAGGAGE_LENGTH / member_size;
assert_eq!(baggage.len(), expected_kept);
for i in 0..expected_kept {
let key = Key::new(format!("k{i:02}"));
assert!(
baggage.get(&key).is_some(),
"expected key k{i:02} to be present"
);
}
for i in expected_kept..total {
let key = Key::new(format!("k{i:02}"));
assert!(
baggage.get(&key).is_none(),
"expected key k{i:02} to be dropped"
);
}
}
}
Loading
Loading