From aa3cf45fbcd5079b416ce9968f5ca452e8adfefb Mon Sep 17 00:00:00 2001 From: jsurany Date: Thu, 21 May 2026 15:26:19 -0400 Subject: [PATCH 01/13] add prototype for EnvVarsCarrier --- .../src/propagation/env_vars_carrier.rs | 135 ++++++++++++++++++ opentelemetry-sdk/src/propagation/mod.rs | 2 + 2 files changed, 137 insertions(+) create mode 100644 opentelemetry-sdk/src/propagation/env_vars_carrier.rs diff --git a/opentelemetry-sdk/src/propagation/env_vars_carrier.rs b/opentelemetry-sdk/src/propagation/env_vars_carrier.rs new file mode 100644 index 0000000000..5fedb4837a --- /dev/null +++ b/opentelemetry-sdk/src/propagation/env_vars_carrier.rs @@ -0,0 +1,135 @@ +use opentelemetry::propagation::{Extractor, Injector}; +use std::collections::HashMap; +/// Propagates name-value pairs via environment variables. +/// +/// This propagator provides a mechanism for propagating context information +/// across process boundaries using environment variables, usually for when +/// network protocols are not applicable. +/// +/// Note that to comply with environment variable naming conventions, all keys +/// are normalized to be compatible with the [POSIX.1-2024] standard. The +/// normalization process follows these rules: +/// +/// - uppercase ASCII letters +/// - replace non-alphanumeric/non-underscore characters with an underscore +/// - prefix name with an underscore if it otherwise starts with a digit +/// +/// # Examples +/// ``` +/// use opentelemetry::propagation::{Extractor, Injector}; +/// use opentelemetry_sdk::propagation::EnvVarsCarrier; +/// +/// // Builds the carrier, fetching the environment into the carrier mapping. +/// let mut carrier = EnvVarsCarrier::new(); +/// +/// // Looks for the normalized "FOO" value in the carrier mapping +/// let val = carrier.get("foo") +/// +/// // Sets the value for normalized "FOO" to "bar", does NOT set env vars +/// carrier.set("foo", String::from("bar")); +/// ``` +/// +/// [POSIX.1-2024]: https://pubs.opengroup.org/onlinepubs/9799919799/basedefs/V1_chap08.html +#[derive(Debug, Clone)] +pub struct EnvVarsCarrier { + map: HashMap, +} + +impl Default for EnvVarsCarrier { + fn default() -> Self { + Self::new() + } +} + +impl EnvVarsCarrier { + /// Create a new `EnvVarsCarrier` object, built from environment variables. + /// Environment variables are fetched and normalized at construction time. + pub fn new() -> Self { + let map = std::env::vars().map(|(k, v)| (normalize(&k), v)).collect(); + + Self { map } + } + + /// Create a new `EnvVarsCarrier` object, internally empty. Useful for + /// testing and for setting up environment mapping for subprocesses from + /// scratch. + pub fn new_empty() -> Self { + Self { + map: HashMap::new(), + } + } +} + +impl Injector for EnvVarsCarrier { + /// Set the value for the normalized key in the carrier mapping. + /// Does NOT set environment variables, but may + fn set(&mut self, key: &str, value: String) { + self.map.insert(normalize(key), value); + } +} + +impl Extractor for EnvVarsCarrier { + /// Get the value for the normalized key from the carrier mapping. + fn get(&self, key: &str) -> Option<&str> { + self.map.get(&normalize(key)).map(|s| s.as_str()) + } + + /// List all of the internal mapping keys in their normalized form. + fn keys(&self) -> Vec<&str> { + self.map.keys().map(|k| k.as_str()).collect() + } +} + +#[inline(always)] +fn normalize_char(c: char) -> char { + if c.is_ascii_alphanumeric() { + c.to_ascii_uppercase() + } else { + '_' + } +} + +fn normalize(name: &str) -> String { + let mut bytes = name.chars().peekable(); + let needs_prefix = bytes.peek().is_some_and(|b| b.is_ascii_digit()); + + needs_prefix + .then_some('_') + .into_iter() + .chain(bytes.map(normalize_char)) + .collect() +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn test_normalize() { + assert_eq!(normalize("foo.bar"), "FOO_BAR"); + assert_eq!(normalize("3abc"), "_3ABC"); + assert_eq!(normalize("HELLO_WORLD"), "HELLO_WORLD"); + assert_eq!(normalize("a.b.c"), "A_B_C"); + assert_eq!(normalize("key with spaces"), "KEY_WITH_SPACES"); + assert_eq!(normalize("ⵕu⾫tⅭf⼤8"), "_U_T_F_8"); + } + + #[test] + fn test_env_vars_carrier_injector() { + let mut carrier = EnvVarsCarrier::new_empty(); + carrier.set("1foo.barᎁbaz", "bar".to_string()); + + let entry = carrier.map.get("_1FOO_BAR_BAZ").unwrap(); + assert_eq!(entry, "bar"); + } + + #[test] + fn test_env_vars_carrier_extractor() { + let mut carrier = EnvVarsCarrier::new_empty(); + carrier + .map + .insert("FOO_BAR".to_string(), "value".to_string()); + + assert_eq!(Extractor::get(&carrier, "foo.bar"), Some("value")); + } +} diff --git a/opentelemetry-sdk/src/propagation/mod.rs b/opentelemetry-sdk/src/propagation/mod.rs index 6cfb5d07bd..9c53f4756f 100644 --- a/opentelemetry-sdk/src/propagation/mod.rs +++ b/opentelemetry-sdk/src/propagation/mod.rs @@ -1,6 +1,8 @@ //! OpenTelemetry Propagators mod baggage; +mod env_vars_carrier; mod trace_context; pub use baggage::BaggagePropagator; +pub use env_vars_carrier::EnvVarsCarrier; pub use trace_context::TraceContextPropagator; From f15c91f42d38f74bcb4563a97c48d7c57a71f0de Mon Sep 17 00:00:00 2001 From: jsurany Date: Thu, 21 May 2026 17:58:40 -0400 Subject: [PATCH 02/13] fix bad docstring in EnvVarsCarrier docs --- opentelemetry-sdk/src/propagation/env_vars_carrier.rs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/opentelemetry-sdk/src/propagation/env_vars_carrier.rs b/opentelemetry-sdk/src/propagation/env_vars_carrier.rs index 5fedb4837a..278897b6b2 100644 --- a/opentelemetry-sdk/src/propagation/env_vars_carrier.rs +++ b/opentelemetry-sdk/src/propagation/env_vars_carrier.rs @@ -23,7 +23,7 @@ use std::collections::HashMap; /// let mut carrier = EnvVarsCarrier::new(); /// /// // Looks for the normalized "FOO" value in the carrier mapping -/// let val = carrier.get("foo") +/// let val = carrier.get("foo"); /// /// // Sets the value for normalized "FOO" to "bar", does NOT set env vars /// carrier.set("foo", String::from("bar")); From 8f07458a52f916161fa04591df30d7cf0c5b9cb7 Mon Sep 17 00:00:00 2001 From: jsurany Date: Fri, 22 May 2026 09:56:03 -0400 Subject: [PATCH 03/13] add environment variable for testing, add new/default testing for EnvVarsCarrier --- .cargo/config.toml | 3 +++ .../src/propagation/env_vars_carrier.rs | 14 ++++++++++++++ 2 files changed, 17 insertions(+) diff --git a/.cargo/config.toml b/.cargo/config.toml index 18a3cb3183..5c154cef23 100644 --- a/.cargo/config.toml +++ b/.cargo/config.toml @@ -1,3 +1,6 @@ [resolver] # https://doc.rust-lang.org/cargo/reference/config.html#resolverincompatible-rust-versions incompatible-rust-versions = "fallback" + +[env] +ENV_VAR_CARRIER_TEST_VAR = "test" diff --git a/opentelemetry-sdk/src/propagation/env_vars_carrier.rs b/opentelemetry-sdk/src/propagation/env_vars_carrier.rs index 278897b6b2..79d51dd591 100644 --- a/opentelemetry-sdk/src/propagation/env_vars_carrier.rs +++ b/opentelemetry-sdk/src/propagation/env_vars_carrier.rs @@ -132,4 +132,18 @@ mod tests { assert_eq!(Extractor::get(&carrier, "foo.bar"), Some("value")); } + + #[test] + fn test_env_vars_carrier_new() { + let carrier = EnvVarsCarrier::new(); + let entry = carrier.get("ENV_VAR_CARRIER_TEST_VAR").unwrap(); + assert_eq!(entry, "test"); + } + + #[test] + fn test_env_vars_carrier_default() { + let carrier: EnvVarsCarrier = Default::default(); + let entry = carrier.get("ENV_VAR_CARRIER_TEST_VAR").unwrap(); + assert_eq!(entry, "test"); + } } From d29ab9b1350924c9f27d8ea6ca0828ffa06995b2 Mon Sep 17 00:00:00 2001 From: jsurany Date: Fri, 22 May 2026 10:09:16 -0400 Subject: [PATCH 04/13] add coverage for EnvVarsCarrier::keys() method --- opentelemetry-sdk/src/propagation/env_vars_carrier.rs | 2 ++ 1 file changed, 2 insertions(+) diff --git a/opentelemetry-sdk/src/propagation/env_vars_carrier.rs b/opentelemetry-sdk/src/propagation/env_vars_carrier.rs index 79d51dd591..c78eea8465 100644 --- a/opentelemetry-sdk/src/propagation/env_vars_carrier.rs +++ b/opentelemetry-sdk/src/propagation/env_vars_carrier.rs @@ -121,6 +121,8 @@ mod tests { let entry = carrier.map.get("_1FOO_BAR_BAZ").unwrap(); assert_eq!(entry, "bar"); + + assert_eq!(carrier.keys(), vec!["_1FOO_BAR_BAZ"]); } #[test] From 006e8b00e53ed793fbd8c829966f30dee83fb1be Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?tonghuaroot=20=28=E7=AB=A5=E8=AF=9D=29?= Date: Sat, 23 May 2026 01:56:07 +0800 Subject: [PATCH 05/13] Merge commit from fork * fix(sdk): enforce W3C Baggage limits in BaggagePropagator extract path BaggagePropagator::extract_with_context parsed each list-member of an inbound `baggage` header in full -- splitting on `;`, decoding the percent-encoded key/value, allocating Vec for property segments, and constructing a KeyValueMetadata -- before passing the result to Baggage::insert_with_metadata. The storage-side limit (MAX_KEY_VALUE_PAIRS=64, MAX_LEN_OF_ALL_PAIRS=8192) silently dropped entries once full, but per-entry allocation work continued for every attacker-supplied member. This change applies the W3C Baggage limits at the propagator boundary: - Reject the header when its byte length exceeds MAX_BAGGAGE_LENGTH (8192). The header is logged at warn level and an empty Context is returned. - Cap the outer comma-split iterator with .take(MAX_BAGGAGE_ITEMS) so parsing stops after 64 list-members regardless of header content. The injection path is unchanged; the existing Baggage storage type already enforces both limits when entries are added. Two unit tests cover the new behavior: - extract_drops_header_exceeding_max_length verifies that an oversized header produces an empty Baggage. - extract_caps_entry_count_at_max_baggage_items verifies that more than 64 syntactically valid entries are truncated to 64. References: - W3C Baggage limits: https://www.w3.org/TR/baggage/#limits Signed-off-by: tonghuaroot * Clarify CHANGELOG: 8192-byte over-limit drops, 64-entry over-limit truncates Address reviewer feedback that 'dropped at the propagator boundary' was ambiguous: the over-length header path drops the whole header, while the over-count path truncates to the first 64 list members. Split the wording into two clauses so the behavior is unambiguous. Signed-off-by: tonghuaroot --------- Signed-off-by: tonghuaroot --- opentelemetry-sdk/CHANGELOG.md | 8 + opentelemetry-sdk/src/propagation/baggage.rs | 163 ++++++++++++++----- 2 files changed, 126 insertions(+), 45 deletions(-) diff --git a/opentelemetry-sdk/CHANGELOG.md b/opentelemetry-sdk/CHANGELOG.md index 99d1b37d2e..f5618b8deb 100644 --- a/opentelemetry-sdk/CHANGELOG.md +++ b/opentelemetry-sdk/CHANGELOG.md @@ -2,6 +2,14 @@ ## vNext +- `BaggagePropagator` now enforces the W3C Baggage maximum header length + (8192 bytes) and maximum list-member count (64) when extracting an inbound + `baggage` header. Headers exceeding 8192 bytes are dropped at the + propagator boundary; headers with more than 64 list members are + truncated to the first 64 entries. The change keeps the propagator from + parsing attacker-controlled input beyond the W3C limits instead of doing + per-entry parse, decode, and allocation work only to discard the excess + on `Baggage` insert. See https://www.w3.org/TR/baggage/#limits. - Reverted the `SimpleSpanProcessor` telemetry suppression added in 0.32.0 (see #3494), which caused a `RefCell already borrowed` panic when a span was started and dropped inside a `get_active_span` (or `Context::map_current`) diff --git a/opentelemetry-sdk/src/propagation/baggage.rs b/opentelemetry-sdk/src/propagation/baggage.rs index 78bb81a3ab..5c30731627 100644 --- a/opentelemetry-sdk/src/propagation/baggage.rs +++ b/opentelemetry-sdk/src/propagation/baggage.rs @@ -11,6 +11,11 @@ use std::sync::OnceLock; static BAGGAGE_HEADER: &str = "baggage"; const FRAGMENT: &AsciiSet = &CONTROLS.add(b' ').add(b'"').add(b';').add(b',').add(b'='); +// W3C Baggage specification limits. +// See https://www.w3.org/TR/baggage/#limits +const MAX_BAGGAGE_LENGTH: usize = 8192; +const MAX_BAGGAGE_ITEMS: usize = 64; + // TODO Replace this with LazyLock once it is stable. static BAGGAGE_FIELDS: OnceLock<[String; 1]> = OnceLock::new(); #[inline] @@ -101,56 +106,74 @@ impl TextMapPropagator for BaggagePropagator { /// Extracts a `Context` with baggage values from a `Extractor`. fn extract_with_context(&self, cx: &Context, extractor: &dyn Extractor) -> Context { if let Some(header_value) = extractor.get(BAGGAGE_HEADER) { - let baggage = header_value.split(',').filter_map(|context_value| { - if let Some((name_and_value, props)) = context_value - .split(';') - .collect::>() - .split_first() - { - let mut iter = name_and_value.split('='); - if let (Some(name), Some(value)) = (iter.next(), iter.next()) { - let decode_name = percent_decode_str(name).decode_utf8(); - let decode_value = percent_decode_str(value).decode_utf8(); - - if let (Ok(name), Ok(value)) = (decode_name, decode_value) { - // Here we don't store the first ; into baggage since it should be treated - // as separator rather part of metadata - let decoded_props = props - .iter() - .flat_map(|prop| percent_decode_str(prop).decode_utf8()) - .map(|prop| prop.trim().to_string()) - .collect::>() - .join(";"); // join with ; because we deleted all ; when calling split above - - Some(KeyValueMetadata::new( - name.trim().to_owned(), - value.trim().to_string(), - decoded_props.as_str(), - )) + // Enforce the W3C Baggage maximum header length up-front so a + // single oversize header cannot drive per-entry allocation work + // before the entries are dropped on insert. See + // https://www.w3.org/TR/baggage/#limits. + if header_value.len() > MAX_BAGGAGE_LENGTH { + otel_warn!( + name: "BaggagePropagator.Extract.HeaderTooLarge", + message = "Baggage header exceeds W3C maximum length and was dropped", + header_bytes = header_value.len() as i64, + limit_bytes = MAX_BAGGAGE_LENGTH as i64, + ); + return cx.clone(); + } + + let baggage = + header_value + .split(',') + .take(MAX_BAGGAGE_ITEMS) + .filter_map(|context_value| { + if let Some((name_and_value, props)) = context_value + .split(';') + .collect::>() + .split_first() + { + let mut iter = name_and_value.split('='); + if let (Some(name), Some(value)) = (iter.next(), iter.next()) { + let decode_name = percent_decode_str(name).decode_utf8(); + let decode_value = percent_decode_str(value).decode_utf8(); + + if let (Ok(name), Ok(value)) = (decode_name, decode_value) { + // Here we don't store the first ; into baggage since it should be treated + // as separator rather part of metadata + let decoded_props = props + .iter() + .flat_map(|prop| percent_decode_str(prop).decode_utf8()) + .map(|prop| prop.trim().to_string()) + .collect::>() + .join(";"); // join with ; because we deleted all ; when calling split above + + Some(KeyValueMetadata::new( + name.trim().to_owned(), + value.trim().to_string(), + decoded_props.as_str(), + )) + } else { + otel_warn!( + name: "BaggagePropagator.Extract.InvalidUTF8", + message = "Invalid UTF8 string in key values", + baggage_header = header_value, + ); + None + } + } else { + otel_warn!( + name: "BaggagePropagator.Extract.InvalidKeyValueFormat", + message = "Invalid baggage key-value format", + baggage_header = header_value, + ); + None + } } else { otel_warn!( - name: "BaggagePropagator.Extract.InvalidUTF8", - message = "Invalid UTF8 string in key values", - baggage_header = header_value, - ); - None - } - } else { - otel_warn!( - name: "BaggagePropagator.Extract.InvalidKeyValueFormat", - message = "Invalid baggage key-value format", - baggage_header = header_value, - ); - None - } - } else { - otel_warn!( name: "BaggagePropagator.Extract.InvalidFormat", message = "Invalid baggage format", baggage_header = header_value); - None - } - }); + None + } + }); cx.with_baggage(baggage) } else { cx.clone() @@ -320,4 +343,54 @@ mod tests { } } } + + #[test] + fn extract_drops_header_exceeding_max_length() { + let propagator = BaggagePropagator::new(); + + // Build a syntactically valid header longer than the W3C 8192-byte + // limit. The header must be rejected outright; no entries are added. + let mut header = String::with_capacity(MAX_BAGGAGE_LENGTH + 1024); + let mut i = 0u32; + while header.len() < MAX_BAGGAGE_LENGTH + 256 { + if !header.is_empty() { + header.push(','); + } + header.push_str(&format!("k{i}=v{i}")); + i += 1; + } + assert!(header.len() > MAX_BAGGAGE_LENGTH); + + let mut extractor: HashMap = HashMap::new(); + extractor.insert(BAGGAGE_HEADER.to_string(), header); + + let context = propagator.extract(&extractor); + assert_eq!(context.baggage().len(), 0); + } + + #[test] + fn extract_caps_entry_count_at_max_baggage_items() { + let propagator = BaggagePropagator::new(); + + // Build a header whose total size stays under MAX_BAGGAGE_LENGTH but + // contains more than MAX_BAGGAGE_ITEMS entries. + let mut header = String::new(); + let entries = MAX_BAGGAGE_ITEMS + 32; + for i in 0..entries { + if i > 0 { + header.push(','); + } + header.push_str(&format!("k{i:03}=v")); + } + assert!(header.len() <= MAX_BAGGAGE_LENGTH); + + let mut extractor: HashMap = HashMap::new(); + extractor.insert(BAGGAGE_HEADER.to_string(), header); + + let context = propagator.extract(&extractor); + // The storage type also caps at MAX_KEY_VALUE_PAIRS=64; the propagator + // additionally stops parsing at MAX_BAGGAGE_ITEMS so allocator work + // does not grow with attacker-supplied entry counts. + assert_eq!(context.baggage().len(), MAX_BAGGAGE_ITEMS); + } } From 09ae03285a7fb65214d492da51e4e0b79ed5c31c Mon Sep 17 00:00:00 2001 From: Lalit Kumar Bhasin Date: Mon, 25 May 2026 17:28:38 -0700 Subject: [PATCH 06/13] chore: prepare opentelemetry_sdk 0.32.1 release (#3522) --- opentelemetry-sdk/CHANGELOG.md | 4 ++++ opentelemetry-sdk/Cargo.toml | 2 +- 2 files changed, 5 insertions(+), 1 deletion(-) diff --git a/opentelemetry-sdk/CHANGELOG.md b/opentelemetry-sdk/CHANGELOG.md index f5618b8deb..4f91ae7f75 100644 --- a/opentelemetry-sdk/CHANGELOG.md +++ b/opentelemetry-sdk/CHANGELOG.md @@ -2,6 +2,10 @@ ## vNext +## 0.32.1 + +Released 2026-May-23 + - `BaggagePropagator` now enforces the W3C Baggage maximum header length (8192 bytes) and maximum list-member count (64) when extracting an inbound `baggage` header. Headers exceeding 8192 bytes are dropped at the diff --git a/opentelemetry-sdk/Cargo.toml b/opentelemetry-sdk/Cargo.toml index bfd05ff404..1a10758afa 100644 --- a/opentelemetry-sdk/Cargo.toml +++ b/opentelemetry-sdk/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "opentelemetry_sdk" -version = "0.32.0" +version = "0.32.1" description = "The SDK for the OpenTelemetry metrics collection and distributed tracing framework" homepage = "https://github.com/open-telemetry/opentelemetry-rust/tree/main/opentelemetry-sdk" repository = "https://github.com/open-telemetry/opentelemetry-rust/tree/main/opentelemetry-sdk" From f340a8e65da85e328987ba80b982b36377350eb9 Mon Sep 17 00:00:00 2001 From: jsurany Date: Tue, 26 May 2026 11:45:57 -0400 Subject: [PATCH 07/13] pr feedback, rename functions, redefine default, move module --- opentelemetry-sdk/src/propagation/mod.rs | 2 -- .../src/propagation/env_vars_carrier.rs | 30 ++++++++++--------- opentelemetry/src/propagation/mod.rs | 2 ++ 3 files changed, 18 insertions(+), 16 deletions(-) rename {opentelemetry-sdk => opentelemetry}/src/propagation/env_vars_carrier.rs (86%) diff --git a/opentelemetry-sdk/src/propagation/mod.rs b/opentelemetry-sdk/src/propagation/mod.rs index 9c53f4756f..6cfb5d07bd 100644 --- a/opentelemetry-sdk/src/propagation/mod.rs +++ b/opentelemetry-sdk/src/propagation/mod.rs @@ -1,8 +1,6 @@ //! OpenTelemetry Propagators mod baggage; -mod env_vars_carrier; mod trace_context; pub use baggage::BaggagePropagator; -pub use env_vars_carrier::EnvVarsCarrier; pub use trace_context::TraceContextPropagator; diff --git a/opentelemetry-sdk/src/propagation/env_vars_carrier.rs b/opentelemetry/src/propagation/env_vars_carrier.rs similarity index 86% rename from opentelemetry-sdk/src/propagation/env_vars_carrier.rs rename to opentelemetry/src/propagation/env_vars_carrier.rs index c78eea8465..c1019e8649 100644 --- a/opentelemetry-sdk/src/propagation/env_vars_carrier.rs +++ b/opentelemetry/src/propagation/env_vars_carrier.rs @@ -1,5 +1,7 @@ -use opentelemetry::propagation::{Extractor, Injector}; use std::collections::HashMap; + +use crate::propagation::{Extractor, Injector}; + /// Propagates name-value pairs via environment variables. /// /// This propagator provides a mechanism for propagating context information @@ -20,7 +22,7 @@ use std::collections::HashMap; /// use opentelemetry_sdk::propagation::EnvVarsCarrier; /// /// // Builds the carrier, fetching the environment into the carrier mapping. -/// let mut carrier = EnvVarsCarrier::new(); +/// let mut carrier = EnvVarsCarrier::from_env(); /// /// // Looks for the normalized "FOO" value in the carrier mapping /// let val = carrier.get("foo"); @@ -36,15 +38,16 @@ pub struct EnvVarsCarrier { } impl Default for EnvVarsCarrier { + /// Create a new empty `EnvVarsCarrier` object. fn default() -> Self { - Self::new() + Self::empty() } } impl EnvVarsCarrier { /// Create a new `EnvVarsCarrier` object, built from environment variables. /// Environment variables are fetched and normalized at construction time. - pub fn new() -> Self { + pub fn from_env() -> Self { let map = std::env::vars().map(|(k, v)| (normalize(&k), v)).collect(); Self { map } @@ -53,7 +56,7 @@ impl EnvVarsCarrier { /// Create a new `EnvVarsCarrier` object, internally empty. Useful for /// testing and for setting up environment mapping for subprocesses from /// scratch. - pub fn new_empty() -> Self { + pub fn empty() -> Self { Self { map: HashMap::new(), } @@ -90,13 +93,13 @@ fn normalize_char(c: char) -> char { } fn normalize(name: &str) -> String { - let mut bytes = name.chars().peekable(); - let needs_prefix = bytes.peek().is_some_and(|b| b.is_ascii_digit()); + let mut chars = name.chars().peekable(); + let needs_prefix = chars.peek().is_some_and(|b| b.is_ascii_digit()); needs_prefix .then_some('_') .into_iter() - .chain(bytes.map(normalize_char)) + .chain(chars.map(normalize_char)) .collect() } @@ -116,7 +119,7 @@ mod tests { #[test] fn test_env_vars_carrier_injector() { - let mut carrier = EnvVarsCarrier::new_empty(); + let mut carrier = EnvVarsCarrier::empty(); carrier.set("1foo.barᎁbaz", "bar".to_string()); let entry = carrier.map.get("_1FOO_BAR_BAZ").unwrap(); @@ -127,7 +130,7 @@ mod tests { #[test] fn test_env_vars_carrier_extractor() { - let mut carrier = EnvVarsCarrier::new_empty(); + let mut carrier = EnvVarsCarrier::empty(); carrier .map .insert("FOO_BAR".to_string(), "value".to_string()); @@ -136,8 +139,8 @@ mod tests { } #[test] - fn test_env_vars_carrier_new() { - let carrier = EnvVarsCarrier::new(); + fn test_env_vars_carrier_from_env() { + let carrier = EnvVarsCarrier::from_env(); let entry = carrier.get("ENV_VAR_CARRIER_TEST_VAR").unwrap(); assert_eq!(entry, "test"); } @@ -145,7 +148,6 @@ mod tests { #[test] fn test_env_vars_carrier_default() { let carrier: EnvVarsCarrier = Default::default(); - let entry = carrier.get("ENV_VAR_CARRIER_TEST_VAR").unwrap(); - assert_eq!(entry, "test"); + assert!(carrier.keys().is_empty()); } } diff --git a/opentelemetry/src/propagation/mod.rs b/opentelemetry/src/propagation/mod.rs index 7848ffe2a8..2c7b420bcd 100644 --- a/opentelemetry/src/propagation/mod.rs +++ b/opentelemetry/src/propagation/mod.rs @@ -21,11 +21,13 @@ use std::collections::HashMap; +mod env_vars_carrier; pub mod composite; pub mod text_map_propagator; pub use composite::TextMapCompositePropagator; pub use text_map_propagator::TextMapPropagator; +pub use env_vars_carrier::EnvVarsCarrier; /// Injector provides an interface for adding fields from an underlying struct like `HashMap` pub trait Injector { From f9a52709694288c6201819534e3a8f2d1cbde963 Mon Sep 17 00:00:00 2001 From: jsurany Date: Tue, 26 May 2026 11:50:38 -0400 Subject: [PATCH 08/13] fix lint --- opentelemetry/src/propagation/mod.rs | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/opentelemetry/src/propagation/mod.rs b/opentelemetry/src/propagation/mod.rs index 2c7b420bcd..3b90d28e8d 100644 --- a/opentelemetry/src/propagation/mod.rs +++ b/opentelemetry/src/propagation/mod.rs @@ -21,13 +21,13 @@ use std::collections::HashMap; -mod env_vars_carrier; pub mod composite; +mod env_vars_carrier; pub mod text_map_propagator; pub use composite::TextMapCompositePropagator; -pub use text_map_propagator::TextMapPropagator; pub use env_vars_carrier::EnvVarsCarrier; +pub use text_map_propagator::TextMapPropagator; /// Injector provides an interface for adding fields from an underlying struct like `HashMap` pub trait Injector { From 2fb4648bad9aab6ab09e6bd38dbf8ce73b04aed4 Mon Sep 17 00:00:00 2001 From: jsurany Date: Wed, 27 May 2026 11:26:40 -0400 Subject: [PATCH 09/13] fix docs --- opentelemetry/src/propagation/env_vars_carrier.rs | 3 +-- 1 file changed, 1 insertion(+), 2 deletions(-) diff --git a/opentelemetry/src/propagation/env_vars_carrier.rs b/opentelemetry/src/propagation/env_vars_carrier.rs index c1019e8649..bdd851ba37 100644 --- a/opentelemetry/src/propagation/env_vars_carrier.rs +++ b/opentelemetry/src/propagation/env_vars_carrier.rs @@ -18,8 +18,7 @@ use crate::propagation::{Extractor, Injector}; /// /// # Examples /// ``` -/// use opentelemetry::propagation::{Extractor, Injector}; -/// use opentelemetry_sdk::propagation::EnvVarsCarrier; +/// use opentelemetry::propagation::{Extractor, Injector, EnvVarsCarrier}; /// /// // Builds the carrier, fetching the environment into the carrier mapping. /// let mut carrier = EnvVarsCarrier::from_env(); From 63b84c79f85065ae18f9ae8f79d9add02474100b Mon Sep 17 00:00:00 2001 From: Cijo Thomas Date: Tue, 26 May 2026 13:35:25 -0700 Subject: [PATCH 10/13] feat(sdk): add otel.sdk.processor.log.processed self-diagnostics metric (#3514) --- docs/design/observability.md | 101 ++++++++++++ .../examples/basic-otlp-http/src/main.rs | 43 +++-- .../examples/basic-otlp/src/main.rs | 43 +++-- .../src/logs/batch_log_processor.rs | 151 ++++++++++++++++++ scripts/test.sh | 3 + 5 files changed, 313 insertions(+), 28 deletions(-) create mode 100644 docs/design/observability.md diff --git a/docs/design/observability.md b/docs/design/observability.md new file mode 100644 index 0000000000..df4f94a659 --- /dev/null +++ b/docs/design/observability.md @@ -0,0 +1,101 @@ +# SDK Self-Diagnostics via Metrics + +Status: +[Development](https://github.com/open-telemetry/opentelemetry-specification/blob/main/specification/document-status.md) + +The OpenTelemetry Rust SDK can emit metrics about its own internal state, +following the [semantic conventions for SDK metrics](https://github.com/open-telemetry/semantic-conventions/blob/main/docs/otel/sdk-metrics.md). + +## Implemented Metrics + +### `otel.sdk.processor.log.processed` + +- **Instrument**: `Counter` +- **Unit**: `{log_record}` +- **Description**: The number of log records for which processing has finished, + either successful or failed. +- **Component**: `BatchLogProcessor` + +**Attributes:** + +| Attribute | Value | +|-----------|-------| +| `otel.component.type` | `batching_log_processor` | +| `otel.component.name` | `batching_log_processor/{id}` (auto-assigned) | +| `error.type` | `queue_full` when dropped due to full queue; `already_shutdown` when emitted after shutdown. Absent on success. | + +The counter is incremented on every `emit()` call: once for successful +enqueue, once with `error.type=queue_full` when dropped due to a full queue, +and once with `error.type=already_shutdown` when emitted after the processor +has been shut down. + +## Feature Gate + +Self-diagnostics metrics require the `experimental_metrics_bound_instruments` +feature on `opentelemetry_sdk`. This feature is not enabled by default. + +Without bound instruments, every `Counter::add()` call would need to resolve +attributes to the internal aggregation state — roughly 50 ns per call on the +`emit()` hot path. With bound instruments, attributes are resolved once at +construction and subsequent `add()` calls are a single atomic increment at +~1.8 ns. Since `emit()` is called for every log record in the application, +this overhead matters. Bound instruments are what make self-diagnostics +practical without measurable performance impact. + +## Provider Initialization Order + +The `BatchLogProcessor` obtains a `Meter` from the global `MeterProvider` +during construction. Rust's `global::meter()` returns a snapshot — it does +**not** retroactively upgrade if the global provider changes later. + +For self-diagnostics to produce real data, the global `MeterProvider` must be +set **before** creating the `LoggerProvider` (and its `BatchLogProcessor`). +The recommended setup order is: + +```rust +// 1. MeterProvider first. Optionally set up a throwaway thread-local fmt +// subscriber so that any internal logs emitted during MeterProvider +// construction still appear on stdout. This is not required — without it +// those few startup-time debug messages are simply lost. +let meter_provider = { + let _guard = tracing::subscriber::set_default( + tracing_subscriber::fmt().with_env_filter("info").finish(), + ); + let mp = init_metrics(); + global::set_meter_provider(mp.clone()); + mp +}; // _guard drops here, removing the throwaway subscriber + +// 2. LoggerProvider second (BatchLogProcessor picks up the real meter) +let logger_provider = init_logs(); + +// 3. Full tracing subscriber (fmt + OTel bridge) +tracing_subscriber::registry() + .with(OpenTelemetryTracingBridge::new(&logger_provider)) + .with(fmt::layer()) + .init(); + +// 4. TracerProvider last (init logs captured by OTel pipeline) +let tracer_provider = init_traces(); +global::set_tracer_provider(tracer_provider.clone()); +``` + +See the OTLP examples (`basic-otlp`, `basic-otlp-http`) for the full pattern +including a temporary thread-local `fmt` subscriber during MeterProvider setup. + +If the `MeterProvider` is set after the `LoggerProvider`, the counter will be +backed by a no-op meter and silently produce nothing. This is harmless but +means no self-diagnostics data. + +## TODO + +- Emit `otel.sdk.processor.log.queue.size` (current queue depth). +- Emit `otel.sdk.processor.log.queue.capacity` (configured max queue size). +- Emit `otel.sdk.exporter.log.exported` from log exporters. +- Add self-diagnostics to `BatchSpanProcessor`. +- Add self-diagnostics to `SimpleLogProcessor`. +- Record logs lost when `shutdown_with_timeout` times out (the background + thread may still hold unfinished exports and queued records). +- Long-term: `global::meter()` currently returns a snapshot that does not + reflect later calls to `set_meter_provider()`. Investigate whether the + global meter can be made to pick up provider changes after the fact. diff --git a/opentelemetry-otlp/examples/basic-otlp-http/src/main.rs b/opentelemetry-otlp/examples/basic-otlp-http/src/main.rs index 8cffe0c6b1..d1f91be3f5 100644 --- a/opentelemetry-otlp/examples/basic-otlp-http/src/main.rs +++ b/opentelemetry-otlp/examples/basic-otlp-http/src/main.rs @@ -67,6 +67,27 @@ fn init_metrics() -> SdkMeterProvider { #[tokio::main] async fn main() -> Result<(), Box> { + // Setup MeterProvider first so that SDK components (e.g., BatchLogProcessor) + // that obtain a Meter from the global MeterProvider during their + // initialization will get a functional meter for self-diagnostics. + // A temporary thread-local fmt subscriber is used to capture any logs + // emitted during MeterProvider initialization to stdout. + // + // Set the global meter provider using a clone of the meter_provider. + // Setting global meter provider is required if other parts of the application + // uses global::meter() or global::meter_with_version() to get a meter. + // Cloning simply creates a new reference to the same meter provider. It is + // important to hold on to the meter_provider here, so as to invoke + // shutdown on it when application ends. + let meter_provider = { + let _guard = tracing::subscriber::set_default( + tracing_subscriber::fmt().with_env_filter("info").finish(), + ); + let meter_provider = init_metrics(); + global::set_meter_provider(meter_provider.clone()); + meter_provider + }; + let logger_provider = init_logs(); // Create a new OpenTelemetryTracingBridge using the above LoggerProvider. @@ -119,15 +140,6 @@ async fn main() -> Result<(), Box> { // shutdown on it when application ends. global::set_tracer_provider(tracer_provider.clone()); - let meter_provider = init_metrics(); - // Set the global meter provider using a clone of the meter_provider. - // Setting global meter provider is required if other parts of the application - // uses global::meter() or global::meter_with_version() to get a meter. - // Cloning simply creates a new reference to the same meter provider. It is - // important to hold on to the meter_provider here, so as to invoke - // shutdown on it when application ends. - global::set_meter_provider(meter_provider.clone()); - let common_scope_attributes = vec![KeyValue::new("scope-key", "scope-value")]; let scope = InstrumentationScope::builder("basic") .with_version("1.0") @@ -166,20 +178,23 @@ async fn main() -> Result<(), Box> { info!(target: "my-target", "hello from {}. My price is {}", "apple", 1.99); - // Collect all shutdown errors + // Collect all shutdown errors. + // Shutdown order: tracer first, then logger, then meter. + // MeterProvider is shut down last because the LoggerProvider's + // BatchLogProcessor may emit self-diagnostic metrics during its shutdown. let mut shutdown_errors = Vec::new(); if let Err(e) = tracer_provider.shutdown() { shutdown_errors.push(format!("tracer provider: {e}")); } - if let Err(e) = meter_provider.shutdown() { - shutdown_errors.push(format!("meter provider: {e}")); - } - if let Err(e) = logger_provider.shutdown() { shutdown_errors.push(format!("logger provider: {e}")); } + if let Err(e) = meter_provider.shutdown() { + shutdown_errors.push(format!("meter provider: {e}")); + } + // Return an error if any shutdown failed if !shutdown_errors.is_empty() { return Err(format!( diff --git a/opentelemetry-otlp/examples/basic-otlp/src/main.rs b/opentelemetry-otlp/examples/basic-otlp/src/main.rs index ed82cd4ded..2af31f452c 100644 --- a/opentelemetry-otlp/examples/basic-otlp/src/main.rs +++ b/opentelemetry-otlp/examples/basic-otlp/src/main.rs @@ -58,6 +58,27 @@ fn init_logs() -> SdkLoggerProvider { #[tokio::main] async fn main() -> Result<(), Box> { + // Setup MeterProvider first so that SDK components (e.g., BatchLogProcessor) + // that obtain a Meter from the global MeterProvider during their + // initialization will get a functional meter for self-diagnostics. + // A temporary thread-local fmt subscriber is used to capture any logs + // emitted during MeterProvider initialization to stdout. + // + // Set the global meter provider using a clone of the meter_provider. + // Setting global meter provider is required if other parts of the application + // uses global::meter() or global::meter_with_version() to get a meter. + // Cloning simply creates a new reference to the same meter provider. It is + // important to hold on to the meter_provider here, so as to invoke + // shutdown on it when application ends. + let meter_provider = { + let _guard = tracing::subscriber::set_default( + tracing_subscriber::fmt().with_env_filter("info").finish(), + ); + let meter_provider = init_metrics(); + global::set_meter_provider(meter_provider.clone()); + meter_provider + }; + let logger_provider = init_logs(); // Create a new OpenTelemetryTracingBridge using the above LoggerProvider. @@ -110,15 +131,6 @@ async fn main() -> Result<(), Box> { // shutdown on it when application ends. global::set_tracer_provider(tracer_provider.clone()); - let meter_provider = init_metrics(); - // Set the global meter provider using a clone of the meter_provider. - // Setting global meter provider is required if other parts of the application - // uses global::meter() or global::meter_with_version() to get a meter. - // Cloning simply creates a new reference to the same meter provider. It is - // important to hold on to the meter_provider here, so as to invoke - // shutdown on it when application ends. - global::set_meter_provider(meter_provider.clone()); - let common_scope_attributes = vec![KeyValue::new("scope-key", "scope-value")]; let scope = InstrumentationScope::builder("basic") .with_version("1.0") @@ -156,20 +168,23 @@ async fn main() -> Result<(), Box> { info!(name: "my-event", target: "my-target", "hello from {}. My price is {}", "apple", 1.99); - // Collect all shutdown errors + // Collect all shutdown errors. + // Shutdown order: tracer first, then logger, then meter. + // MeterProvider is shut down last because the LoggerProvider's + // BatchLogProcessor may emit self-diagnostic metrics during its shutdown. let mut shutdown_errors = Vec::new(); if let Err(e) = tracer_provider.shutdown() { shutdown_errors.push(format!("tracer provider: {e}")); } - if let Err(e) = meter_provider.shutdown() { - shutdown_errors.push(format!("meter provider: {e}")); - } - if let Err(e) = logger_provider.shutdown() { shutdown_errors.push(format!("logger provider: {e}")); } + if let Err(e) = meter_provider.shutdown() { + shutdown_errors.push(format!("meter provider: {e}")); + } + // Return an error if any shutdown failed if !shutdown_errors.is_empty() { return Err(format!( diff --git a/opentelemetry-sdk/src/logs/batch_log_processor.rs b/opentelemetry-sdk/src/logs/batch_log_processor.rs index f4fb3daf2b..ac65a0a958 100644 --- a/opentelemetry-sdk/src/logs/batch_log_processor.rs +++ b/opentelemetry-sdk/src/logs/batch_log_processor.rs @@ -25,6 +25,9 @@ use std::sync::mpsc::{self, RecvTimeoutError, SyncSender}; use opentelemetry::{otel_debug, otel_error, otel_warn, Context, InstrumentationScope}; +#[cfg(feature = "experimental_metrics_bound_instruments")] +use opentelemetry::KeyValue; + use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering}; use std::{cmp::min, env, sync::Mutex}; use std::{ @@ -147,6 +150,17 @@ pub struct BatchLogProcessor { // Track the maximum queue size that was configured for this processor max_queue_size: usize, + + // Self-diagnostics: otel.sdk.processor.log.processed counter. + // Gated behind experimental_metrics_bound_instruments so the hot-path + // `add` is a single atomic increment (~1.8 ns) with no per-call + // attribute resolution. + #[cfg(feature = "experimental_metrics_bound_instruments")] + processed_success: opentelemetry::metrics::BoundCounter, + #[cfg(feature = "experimental_metrics_bound_instruments")] + processed_queue_full: opentelemetry::metrics::BoundCounter, + #[cfg(feature = "experimental_metrics_bound_instruments")] + processed_after_shutdown: opentelemetry::metrics::BoundCounter, } impl Debug for BatchLogProcessor { @@ -166,6 +180,10 @@ impl LogProcessor for BatchLogProcessor { // match for result and handle each separately match result { Ok(_) => { + // Record successful processing in self-diagnostics + #[cfg(feature = "experimental_metrics_bound_instruments")] + self.processed_success.add(1); + // Successfully sent the log record to the data channel. // Increment the current batch size and check if it has reached // the max export batch size. @@ -207,6 +225,10 @@ impl LogProcessor for BatchLogProcessor { } } Err(mpsc::TrySendError::Full(_)) => { + // Record queue-full drop in self-diagnostics + #[cfg(feature = "experimental_metrics_bound_instruments")] + self.processed_queue_full.add(1); + // Increment dropped logs count. The first time we have to drop // a log, emit a warning. if self.dropped_logs_count.fetch_add(1, Ordering::Relaxed) == 0 { @@ -215,6 +237,10 @@ impl LogProcessor for BatchLogProcessor { } } Err(mpsc::TrySendError::Disconnected(_)) => { + // Record after-shutdown drop in self-diagnostics + #[cfg(feature = "experimental_metrics_bound_instruments")] + self.processed_after_shutdown.add(1); + // The following `otel_warn!` may cause an infinite feedback loop of // 'telemetry-induced-telemetry', potentially causing a stack overflow let _guard = Context::enter_telemetry_suppressed_scope(); @@ -292,6 +318,13 @@ impl LogProcessor for BatchLogProcessor { }) .map_err(|err| match err { RecvTimeoutError::Timeout => { + // TODO: When shutdown times out, log records still + // in the queue or mid-export are silently lost. The + // background thread is not joined and may continue + // running. Consider: (1) recording the lost count + // in the self-diagnostics counter, (2) joining the + // thread with a best-effort wait, or (3) signalling + // the thread to abort the current export. otel_error!( name: "BatchLogProcessor.Shutdown.Timeout", message = "BatchLogProcessor shutdown timing out." @@ -497,6 +530,48 @@ impl BatchLogProcessor { }) .expect("Thread spawn failed."); //TODO: Handle thread spawn failure + // Self-diagnostics: create the otel.sdk.processor.log.processed counter + #[cfg(feature = "experimental_metrics_bound_instruments")] + let (processed_success, processed_queue_full, processed_after_shutdown) = { + static INSTANCE_COUNTER: AtomicUsize = AtomicUsize::new(0); + let instance_id = INSTANCE_COUNTER.fetch_add(1, Ordering::Relaxed); + let component_name = format!("batching_log_processor/{instance_id}"); + + let meter = opentelemetry::global::meter("otel.sdk"); + let counter = meter + .u64_counter("otel.sdk.processor.log.processed") + .with_description( + "The number of log records for which the processing has finished, \ + either successful or failed.", + ) + .with_unit("{log_record}") + .build(); + + // Attribute values follow the OTel semantic conventions for SDK metrics: + // https://github.com/open-telemetry/semantic-conventions/blob/main/docs/otel/sdk-metrics.md#metric-otelsdkprocessorlogprocessed + // https://github.com/open-telemetry/semantic-conventions/blob/main/docs/registry/attributes/otel.md#otel-component-attributes + let success_attrs = [ + KeyValue::new("otel.component.type", "batching_log_processor"), + KeyValue::new("otel.component.name", component_name.clone()), + ]; + let queue_full_attrs = [ + KeyValue::new("error.type", "queue_full"), + KeyValue::new("otel.component.type", "batching_log_processor"), + KeyValue::new("otel.component.name", component_name.clone()), + ]; + let after_shutdown_attrs = [ + KeyValue::new("error.type", "already_shutdown"), + KeyValue::new("otel.component.type", "batching_log_processor"), + KeyValue::new("otel.component.name", component_name), + ]; + + ( + counter.bind(&success_attrs), + counter.bind(&queue_full_attrs), + counter.bind(&after_shutdown_attrs), + ) + }; + // Return batch processor with link to worker BatchLogProcessor { logs_sender, @@ -508,6 +583,12 @@ impl BatchLogProcessor { export_log_message_sent: Arc::new(AtomicBool::new(false)), current_batch_size, max_export_batch_size, + #[cfg(feature = "experimental_metrics_bound_instruments")] + processed_success, + #[cfg(feature = "experimental_metrics_bound_instruments")] + processed_queue_full, + #[cfg(feature = "experimental_metrics_bound_instruments")] + processed_after_shutdown, } } @@ -1122,4 +1203,74 @@ mod tests { Consider reducing queue size or increasing thread count/log volume." ); } + + /// Verifies that `otel.sdk.processor.log.processed` counter records + /// successful log processing when `experimental_metrics_bound_instruments` + /// is enabled and a real MeterProvider is set as global before creating + /// the processor. + /// + /// This test is `#[ignore]`d because it calls + /// `global::set_meter_provider()` which mutates process-wide state. + /// CI runs it in isolation via `test.sh`. + #[cfg(feature = "experimental_metrics_bound_instruments")] + #[test] + #[ignore] + fn self_diagnostics_counter_records_success() { + use crate::metrics::data::{AggregatedMetrics, MetricData}; + use crate::metrics::{InMemoryMetricExporter, SdkMeterProvider}; + + // Setup a real MeterProvider and set it as global BEFORE creating the + // BatchLogProcessor, so the processor picks up a real meter. + let metric_exporter = InMemoryMetricExporter::default(); + let meter_provider = SdkMeterProvider::builder() + .with_periodic_exporter(metric_exporter.clone()) + .build(); + opentelemetry::global::set_meter_provider(meter_provider.clone()); + + let log_exporter = InMemoryLogExporter::default(); + let config = BatchConfigBuilder::default() + .with_max_queue_size(256) + .with_max_export_batch_size(64) + .with_scheduled_delay(Duration::from_secs(60)) + .build(); + let processor = BatchLogProcessor::new(log_exporter, config); + + // Emit 10 logs + let instrumentation = InstrumentationScope::default(); + for _ in 0..10 { + let mut record = SdkLogRecord::new(); + processor.emit(&mut record, &instrumentation); + } + + // Force a metrics collection + meter_provider.force_flush().unwrap(); + + // Find the otel.sdk.processor.log.processed metric and sum all data points + let metrics = metric_exporter.get_finished_metrics().unwrap(); + let mut found = false; + let mut total_value: u64 = 0; + for rm in &metrics { + for sm in &rm.scope_metrics { + for metric in &sm.metrics { + if metric.name == "otel.sdk.processor.log.processed" { + found = true; + if let AggregatedMetrics::U64(MetricData::Sum(sum)) = &metric.data { + for dp in sum.data_points() { + total_value += dp.value(); + } + } + } + } + } + } + + assert!(found, "otel.sdk.processor.log.processed metric not found"); + assert_eq!( + total_value, 10, + "Expected 10 processed logs, got {total_value}" + ); + + processor.shutdown().unwrap(); + meter_provider.shutdown().unwrap(); + } } diff --git a/scripts/test.sh b/scripts/test.sh index bd2d184736..741ab7c899 100755 --- a/scripts/test.sh +++ b/scripts/test.sh @@ -30,6 +30,9 @@ cargo test --manifest-path=opentelemetry-sdk/Cargo.toml --all-features trace::ru cargo test --manifest-path=opentelemetry-sdk/Cargo.toml --all-features trace::runtime_tests::test_set_provider_single_thread_tokio -- --ignored --exact cargo test --manifest-path=opentelemetry-sdk/Cargo.toml --all-features trace::runtime_tests::test_set_provider_single_thread_tokio_shutdown -- --ignored --exact +echo "Running ignored tests for opentelemetry-sdk package (self-diagnostics, requires global MeterProvider)" +cargo test --manifest-path=opentelemetry-sdk/Cargo.toml --all-features logs::batch_log_processor::tests::self_diagnostics_counter_records_success -- --ignored --exact + echo "Running ignored tests for opentelemetry-appender-tracing package (global logger tests)" cargo test --manifest-path=opentelemetry-appender-tracing/Cargo.toml --all-features layer::tests::tracing_appender_standalone_with_tracing_log -- --ignored --exact cargo test --manifest-path=opentelemetry-appender-tracing/Cargo.toml --all-features layer::tests::tracing_appender_inside_tracing_context_with_tracing_log -- --ignored --exact From f20ef4339f23b7bd099b3f5e95e7a5707bcfdabb Mon Sep 17 00:00:00 2001 From: chickenbreeder Date: Wed, 27 May 2026 17:31:52 +0200 Subject: [PATCH 11/13] refactor: parametrize tests (#3256) Co-authored-by: Scott Gerring --- Cargo.toml | 2 +- opentelemetry-sdk/src/metrics/mod.rs | 106 +++++++-------------------- 2 files changed, 27 insertions(+), 81 deletions(-) diff --git a/Cargo.toml b/Cargo.toml index 9313dd299d..36622181ae 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -71,7 +71,7 @@ opentelemetry-proto = { path = "opentelemetry-proto", version = "0.32", default- opentelemetry-semantic-conventions = { path = "opentelemetry-semantic-conventions", version = "0.32", default-features = false } opentelemetry-stdout = { path = "opentelemetry-stdout", version = "0.32", default-features = false } percent-encoding = "2.0" -rstest = "0.23.0" +rstest = "0.26.1" schemars = "1.0" sysinfo = "0.32" tempfile = "3.3.0" diff --git a/opentelemetry-sdk/src/metrics/mod.rs b/opentelemetry-sdk/src/metrics/mod.rs index 558515cb87..e9b3565974 100644 --- a/opentelemetry-sdk/src/metrics/mod.rs +++ b/opentelemetry-sdk/src/metrics/mod.rs @@ -143,6 +143,8 @@ mod tests { use std::thread; use std::time::Duration; + use rstest::rstest; + // Run all tests in this mod // cargo test metrics::tests --features=testing,spec_unstable_metrics_views // Note for all tests from this point onwards in this mod: @@ -391,34 +393,15 @@ mod tests { counter_aggregation_overflow_helper_custom_limit(Temporality::Cumulative); } - #[tokio::test(flavor = "multi_thread", worker_threads = 1)] - async fn counter_aggregation_attribute_order_sorted_first_delta() { - // Run this test with stdout enabled to see output. - // cargo test counter_aggregation_attribute_order_sorted_first_delta --features=testing -- --nocapture - counter_aggregation_attribute_order_helper(Temporality::Delta, true); - } - - #[tokio::test(flavor = "multi_thread", worker_threads = 1)] - async fn counter_aggregation_attribute_order_sorted_first_cumulative() { - // Run this test with stdout enabled to see output. - // cargo test counter_aggregation_attribute_order_sorted_first_cumulative --features=testing -- --nocapture - counter_aggregation_attribute_order_helper(Temporality::Cumulative, true); - } - - #[tokio::test(flavor = "multi_thread", worker_threads = 1)] - async fn counter_aggregation_attribute_order_unsorted_first_delta() { + #[rstest] + #[case(Temporality::Delta, true)] + #[case(Temporality::Delta, false)] + #[case(Temporality::Cumulative, true)] + #[case(Temporality::Cumulative, false)] + fn counter_aggregation(#[case] temporality: Temporality, #[case] start_sorted: bool) { // Run this test with stdout enabled to see output. - // cargo test counter_aggregation_attribute_order_unsorted_first_delta --features=testing -- --nocapture - - counter_aggregation_attribute_order_helper(Temporality::Delta, false); - } - - #[tokio::test(flavor = "multi_thread", worker_threads = 1)] - async fn counter_aggregation_attribute_order_unsorted_first_cumulative() { - // Run this test with stdout enabled to see output. - // cargo test counter_aggregation_attribute_order_unsorted_first_cumulative --features=testing -- --nocapture - - counter_aggregation_attribute_order_helper(Temporality::Cumulative, false); + // cargo test counter_aggregation_attribute --features=testing -- --nocapture + counter_aggregation_attribute_order_helper(temporality, start_sorted); } #[tokio::test(flavor = "multi_thread", worker_threads = 1)] @@ -503,60 +486,23 @@ mod tests { observable_gauge_aggregation_helper(Temporality::Cumulative, true); } - #[tokio::test(flavor = "multi_thread", worker_threads = 1)] - async fn observable_counter_aggregation_cumulative_non_zero_increment() { - // Run this test with stdout enabled to see output. - // cargo test observable_counter_aggregation_cumulative_non_zero_increment --features=testing -- --nocapture - observable_counter_aggregation_helper(Temporality::Cumulative, 100, 10, 4, false); - } - - #[tokio::test(flavor = "multi_thread", worker_threads = 1)] - async fn observable_counter_aggregation_cumulative_non_zero_increment_no_attrs() { - // Run this test with stdout enabled to see output. - // cargo test observable_counter_aggregation_cumulative_non_zero_increment_no_attrs --features=testing -- --nocapture - observable_counter_aggregation_helper(Temporality::Cumulative, 100, 10, 4, true); - } - - #[tokio::test(flavor = "multi_thread", worker_threads = 1)] - async fn observable_counter_aggregation_delta_non_zero_increment() { - // Run this test with stdout enabled to see output. - // cargo test observable_counter_aggregation_delta_non_zero_increment --features=testing -- --nocapture - observable_counter_aggregation_helper(Temporality::Delta, 100, 10, 4, false); - } - - #[tokio::test(flavor = "multi_thread", worker_threads = 1)] - async fn observable_counter_aggregation_delta_non_zero_increment_no_attrs() { - // Run this test with stdout enabled to see output. - // cargo test observable_counter_aggregation_delta_non_zero_increment_no_attrs --features=testing -- --nocapture - observable_counter_aggregation_helper(Temporality::Delta, 100, 10, 4, true); - } - - #[tokio::test(flavor = "multi_thread", worker_threads = 1)] - async fn observable_counter_aggregation_cumulative_zero_increment() { - // Run this test with stdout enabled to see output. - // cargo test observable_counter_aggregation_cumulative_zero_increment --features=testing -- --nocapture - observable_counter_aggregation_helper(Temporality::Cumulative, 100, 0, 4, false); - } - - #[tokio::test(flavor = "multi_thread", worker_threads = 1)] - async fn observable_counter_aggregation_cumulative_zero_increment_no_attrs() { - // Run this test with stdout enabled to see output. - // cargo test observable_counter_aggregation_cumulative_zero_increment_no_attrs --features=testing -- --nocapture - observable_counter_aggregation_helper(Temporality::Cumulative, 100, 0, 4, true); - } - - #[tokio::test(flavor = "multi_thread", worker_threads = 1)] - async fn observable_counter_aggregation_delta_zero_increment() { - // Run this test with stdout enabled to see output. - // cargo test observable_counter_aggregation_delta_zero_increment --features=testing -- --nocapture - observable_counter_aggregation_helper(Temporality::Delta, 100, 0, 4, false); - } - - #[tokio::test(flavor = "multi_thread", worker_threads = 1)] - async fn observable_counter_aggregation_delta_zero_increment_no_attrs() { + #[rstest] + #[case(Temporality::Cumulative, 10, false)] + #[case(Temporality::Cumulative, 10, true)] + #[case(Temporality::Delta, 10, false)] + #[case(Temporality::Delta, 10, true)] + #[case(Temporality::Cumulative, 0, false)] + #[case(Temporality::Cumulative, 0, true)] + #[case(Temporality::Delta, 0, false)] + #[case(Temporality::Delta, 0, true)] + fn observable_counter_aggregation( + #[case] temporality: Temporality, + #[case] increment: u64, + #[case] is_empty_attributes: bool, + ) { // Run this test with stdout enabled to see output. - // cargo test observable_counter_aggregation_delta_zero_increment_no_attrs --features=testing -- --nocapture - observable_counter_aggregation_helper(Temporality::Delta, 100, 0, 4, true); + // cargo test observable_counter_aggregation --features=testing -- --nocapture + observable_counter_aggregation_helper(temporality, 100, increment, 4, is_empty_attributes); } fn observable_counter_aggregation_helper( From 983f10c24320be92a338569e7e2d806bcc851d00 Mon Sep 17 00:00:00 2001 From: jsurany Date: Tue, 2 Jun 2026 15:27:06 -0400 Subject: [PATCH 12/13] add new EnvVarsCarrier test, update docs --- opentelemetry/src/propagation/env_vars_carrier.rs | 13 ++++++++++++- 1 file changed, 12 insertions(+), 1 deletion(-) diff --git a/opentelemetry/src/propagation/env_vars_carrier.rs b/opentelemetry/src/propagation/env_vars_carrier.rs index bdd851ba37..a4d8e7ab84 100644 --- a/opentelemetry/src/propagation/env_vars_carrier.rs +++ b/opentelemetry/src/propagation/env_vars_carrier.rs @@ -28,6 +28,9 @@ use crate::propagation::{Extractor, Injector}; /// /// // Sets the value for normalized "FOO" to "bar", does NOT set env vars /// carrier.set("foo", String::from("bar")); +/// +/// // Fetches the list of (normalized) keys +/// let keys = carrier.keys(); // vec!["FOO"] /// ``` /// /// [POSIX.1-2024]: https://pubs.opengroup.org/onlinepubs/9799919799/basedefs/V1_chap08.html @@ -64,7 +67,7 @@ impl EnvVarsCarrier { impl Injector for EnvVarsCarrier { /// Set the value for the normalized key in the carrier mapping. - /// Does NOT set environment variables, but may + /// Does NOT set environment variables fn set(&mut self, key: &str, value: String) { self.map.insert(normalize(key), value); } @@ -137,6 +140,14 @@ mod tests { assert_eq!(Extractor::get(&carrier, "foo.bar"), Some("value")); } + #[test] + fn test_env_vars_carrier_inject_and_extract() { + let mut carrier = EnvVarsCarrier::empty(); + Injector::set(&mut carrier, "foo.bar", "value".to_string()); + assert_eq!(Extractor::get(&carrier, "foo.bar"), Some("value")); + assert_eq!(carrier.keys(), vec!["FOO_BAR"]); + } + #[test] fn test_env_vars_carrier_from_env() { let carrier = EnvVarsCarrier::from_env(); From 714680ee6a1209a92c3c5bd5f31f5191fff818f3 Mon Sep 17 00:00:00 2001 From: jsurany Date: Mon, 15 Jun 2026 12:26:21 -0400 Subject: [PATCH 13/13] update per specification PR#5144 --- .cargo/config.toml | 1 + .../src/propagation/env_vars_carrier.rs | 48 ++++++++++++++++++- 2 files changed, 47 insertions(+), 2 deletions(-) diff --git a/.cargo/config.toml b/.cargo/config.toml index 5c154cef23..14c1fe34ba 100644 --- a/.cargo/config.toml +++ b/.cargo/config.toml @@ -4,3 +4,4 @@ incompatible-rust-versions = "fallback" [env] ENV_VAR_CARRIER_TEST_VAR = "test" +ENV-VAR-CARRIER-TEST-VAR2 = "test2" # should be ignored by EnvVarsCarrier diff --git a/opentelemetry/src/propagation/env_vars_carrier.rs b/opentelemetry/src/propagation/env_vars_carrier.rs index a4d8e7ab84..2b5f57b87e 100644 --- a/opentelemetry/src/propagation/env_vars_carrier.rs +++ b/opentelemetry/src/propagation/env_vars_carrier.rs @@ -21,6 +21,7 @@ use crate::propagation::{Extractor, Injector}; /// use opentelemetry::propagation::{Extractor, Injector, EnvVarsCarrier}; /// /// // Builds the carrier, fetching the environment into the carrier mapping. +/// // Filters any environment variables that are not already normalized. /// let mut carrier = EnvVarsCarrier::from_env(); /// /// // Looks for the normalized "FOO" value in the carrier mapping @@ -50,8 +51,7 @@ impl EnvVarsCarrier { /// Create a new `EnvVarsCarrier` object, built from environment variables. /// Environment variables are fetched and normalized at construction time. pub fn from_env() -> Self { - let map = std::env::vars().map(|(k, v)| (normalize(&k), v)).collect(); - + let map = std::env::vars().filter(|(k, _)| is_normalized(k)).collect(); Self { map } } @@ -105,6 +105,21 @@ fn normalize(name: &str) -> String { .collect() } +fn is_normalized(name: &str) -> bool { + let mut chars = name.chars(); + + if !chars + .next() + .is_some_and(|c| c.is_ascii_uppercase() || c == '_') + { + return false; + } + + chars + .find(|c| !c.is_ascii_uppercase() && !c.is_ascii_digit() && *c != '_') + .is_none() +} + #[cfg(test)] mod tests { use super::*; @@ -119,6 +134,31 @@ mod tests { assert_eq!(normalize("ⵕu⾫tⅭf⼤8"), "_U_T_F_8"); } + #[test] + fn test_is_normalized_prefix() { + assert!(is_normalized("ABC")); + + assert!(!is_normalized("3ABC")); // begins with number + assert!(is_normalized("_3ABC")); + + assert!(!is_normalized("aBC")); // begins with lowercase letter + assert!(is_normalized("ABC")); + + assert!(!is_normalized(".ABC")); // begins with non-alphanumeric/non-underscore character + assert!(is_normalized("_ABC")); + } + + #[test] + fn test_is_normalized_body() { + assert!(is_normalized("HELLO_WORLD")); + + assert!(!is_normalized("foo.bar")); + assert!(!is_normalized("3abc")); + assert!(!is_normalized("a.b.c")); + assert!(!is_normalized("key with spaces")); + assert!(!is_normalized("ⵕu⾫tⅭf⼤8")); + } + #[test] fn test_env_vars_carrier_injector() { let mut carrier = EnvVarsCarrier::empty(); @@ -150,9 +190,13 @@ mod tests { #[test] fn test_env_vars_carrier_from_env() { + // refer to .cargo/config.toml for the environment variable definitions let carrier = EnvVarsCarrier::from_env(); let entry = carrier.get("ENV_VAR_CARRIER_TEST_VAR").unwrap(); assert_eq!(entry, "test"); + + let vars = carrier.keys(); + assert!(!vars.contains(&"ENV-VAR-CARRIER-TEST-VAR2")); } #[test]