From 0a2e42c718889fcedc5c33d1b04c284fb202a263 Mon Sep 17 00:00:00 2001 From: cijothomas Date: Wed, 10 Jun 2026 15:11:58 -0700 Subject: [PATCH 01/11] POC: emit otel.sdk.component.shutdown event from BatchLogProcessor and OTLP gRPC log exporter --- opentelemetry-otlp/src/exporter/tonic/logs.rs | 81 ++++++++++++- .../src/logs/batch_log_processor.rs | 112 ++++++++++++++++-- 2 files changed, 180 insertions(+), 13 deletions(-) diff --git a/opentelemetry-otlp/src/exporter/tonic/logs.rs b/opentelemetry-otlp/src/exporter/tonic/logs.rs index 6704a2a2c5..64b6882412 100644 --- a/opentelemetry-otlp/src/exporter/tonic/logs.rs +++ b/opentelemetry-otlp/src/exporter/tonic/logs.rs @@ -5,8 +5,10 @@ use opentelemetry_proto::tonic::collector::logs::v1::{ }; use opentelemetry_sdk::error::{OTelSdkError, OTelSdkResult}; use opentelemetry_sdk::logs::{LogBatch, LogExporter}; +use std::sync::atomic::{AtomicUsize, Ordering}; use std::sync::{Arc, Mutex}; use std::time; +use std::time::Instant; use tonic::{codegen::CompressionEncoding, service::Interceptor, transport::Channel, Request}; use opentelemetry_proto::transform::logs::tonic::group_logs_by_resource_and_scope; @@ -23,6 +25,7 @@ pub(crate) struct TonicLogsClient { #[allow(dead_code)] // would be removed once we support set_resource for metrics. resource: opentelemetry_proto::transform::common::tonic::ResourceAttributesWithSchema, + component_name: String, } struct ClientInner { @@ -52,6 +55,12 @@ impl TonicLogsClient { otel_debug!(name: "TonicsLogsClientBuilt"); + let component_name = { + static INSTANCE_COUNTER: AtomicUsize = AtomicUsize::new(0); + let instance_id = INSTANCE_COUNTER.fetch_add(1, Ordering::Relaxed); + format!("otlp_grpc_log_exporter/{instance_id}") + }; + TonicLogsClient { inner: Mutex::new(Some(ClientInner { client, @@ -64,6 +73,7 @@ impl TonicLogsClient { jitter_ms: 100, }), resource: Default::default(), + component_name, } } } @@ -143,11 +153,25 @@ impl LogExporter for TonicLogsClient { } fn shutdown_with_timeout(&self, _timeout: time::Duration) -> OTelSdkResult { - self.inner + let shutdown_start = Instant::now(); + let was_first_shutdown = match self + .inner .lock() - .map_err(|e| OTelSdkError::InternalFailure(format!("Failed to acquire lock: {e}")))? - .take(); + .map_err(|e| OTelSdkError::InternalFailure(format!("Failed to acquire lock: {e}"))) + { + Ok(mut guard) => guard.take().is_some(), + Err(e) => { + let duration_secs = shutdown_start.elapsed().as_secs_f64(); + self.emit_shutdown_event("failed", duration_secs); + return Err(e); + } + }; + let duration_secs = shutdown_start.elapsed().as_secs_f64(); + if was_first_shutdown { + self.emit_shutdown_event("success", duration_secs); + } + // Idempotent re-invocation: emit no event per spec. Ok(()) } @@ -155,3 +179,54 @@ impl LogExporter for TonicLogsClient { self.resource = resource.into(); } } + +impl TonicLogsClient { + // POC: emit otel.sdk.component.shutdown for this exporter. See + // BatchLogProcessor::emit_shutdown_event for the rationale on using + // opentelemetry::_private (tracing) directly instead of the + // otel_info!/otel_warn! macros. + fn emit_shutdown_event(&self, result_str: &'static str, duration_secs: f64) { + // Exporter has no internal queue; shutdown.dropped is always 0. + let shutdown_dropped: u64 = 0; + // No lifetime drop counter on this exporter (no queue). + let lifetime_dropped: u64 = 0; + + if result_str == "success" { + #[cfg(feature = "internal-logs")] + opentelemetry::_private::info!( + name: "otel.sdk.component.shutdown", + target: env!("CARGO_PKG_NAME"), + name = "otel.sdk.component.shutdown", + "otel.component.type" = "otlp_grpc_log_exporter", + "otel.component.name" = self.component_name.as_str(), + "otel.component.shutdown.result" = result_str, + "otel.component.dropped" = lifetime_dropped, + "otel.component.shutdown.dropped" = shutdown_dropped, + "otel.component.shutdown.duration" = duration_secs, + ); + } else { + #[cfg(feature = "internal-logs")] + opentelemetry::_private::warn!( + name: "otel.sdk.component.shutdown", + target: env!("CARGO_PKG_NAME"), + name = "otel.sdk.component.shutdown", + "otel.component.type" = "otlp_grpc_log_exporter", + "otel.component.name" = self.component_name.as_str(), + "otel.component.shutdown.result" = result_str, + "otel.component.dropped" = lifetime_dropped, + "otel.component.shutdown.dropped" = shutdown_dropped, + "otel.component.shutdown.duration" = duration_secs, + ); + } + + #[cfg(not(feature = "internal-logs"))] + { + let _ = ( + result_str, + duration_secs, + shutdown_dropped, + lifetime_dropped, + ); + } + } +} diff --git a/opentelemetry-sdk/src/logs/batch_log_processor.rs b/opentelemetry-sdk/src/logs/batch_log_processor.rs index ac65a0a958..6b5e0133c2 100644 --- a/opentelemetry-sdk/src/logs/batch_log_processor.rs +++ b/opentelemetry-sdk/src/logs/batch_log_processor.rs @@ -151,6 +151,12 @@ pub struct BatchLogProcessor { // Track the maximum queue size that was configured for this processor max_queue_size: usize, + // Stable component identity used for self-observability events/metrics. + // Format: "batching_log_processor/". Always present, independent of + // the experimental_metrics_bound_instruments feature flag because the + // shutdown event (otel.sdk.component.shutdown) needs it unconditionally. + component_name: String, + // 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 @@ -303,6 +309,23 @@ impl LogProcessor for BatchLogProcessor { ); } + let shutdown_start = Instant::now(); + let result = self.shutdown_inner(timeout); + let duration_secs = shutdown_start.elapsed().as_secs_f64(); + self.emit_shutdown_event(&result, dropped_logs, duration_secs); + result + } + + fn set_resource(&mut self, resource: &Resource) { + let resource = Arc::new(resource.clone()); + let _ = self + .message_sender + .try_send(BatchMessage::SetResource(resource)); + } +} + +impl BatchLogProcessor { + fn shutdown_inner(&self, timeout: Duration) -> OTelSdkResult { let (sender, receiver) = mpsc::sync_channel(1); match self.message_sender.try_send(BatchMessage::Shutdown(sender)) { Ok(_) => { @@ -361,11 +384,74 @@ impl LogProcessor for BatchLogProcessor { } } - fn set_resource(&mut self, resource: &Resource) { - let resource = Arc::new(resource.clone()); - let _ = self - .message_sender - .try_send(BatchMessage::SetResource(resource)); + // POC: emit the otel.sdk.component.shutdown event for this processor. + // Per spec PR open-telemetry/semantic-conventions#3723 the event is INFO + // on success and WARN on any non-success result. We bypass the + // otel_info!/otel_warn! macros because those restrict attribute keys to + // plain identifiers; the spec requires dotted attribute names like + // `otel.component.type`, which tracing's own macros accept via quoted-key + // syntax. Behaviour is otherwise identical to otel_info!/otel_warn!. + fn emit_shutdown_event( + &self, + result: &OTelSdkResult, + lifetime_dropped: usize, + duration_secs: f64, + ) { + let result_str = match result { + Ok(()) => "success", + Err(OTelSdkError::Timeout(_)) => "timed_out", + Err(OTelSdkError::AlreadyShutdown) => { + // Idempotent re-invocation: per spec a component is shut down + // at most once; subsequent calls are no-ops and emit no event. + return; + } + Err(_) => "failed", + }; + let is_success = result_str == "success"; + + // shutdown.dropped: items lost during shutdown itself. On success the + // spec mandates 0. On timeout/failure we cannot enumerate the queue + // residual from outside the worker thread, so we report 0 as a known + // under-count. SPEC GAP: should the attribute be omitted when unknown? + let shutdown_dropped: u64 = 0; + + if is_success { + #[cfg(feature = "internal-logs")] + opentelemetry::_private::info!( + name: "otel.sdk.component.shutdown", + target: env!("CARGO_PKG_NAME"), + name = "otel.sdk.component.shutdown", + "otel.component.type" = "batching_log_processor", + "otel.component.name" = self.component_name.as_str(), + "otel.component.shutdown.result" = result_str, + "otel.component.dropped" = lifetime_dropped as u64, + "otel.component.shutdown.dropped" = shutdown_dropped, + "otel.component.shutdown.duration" = duration_secs, + ); + } else { + #[cfg(feature = "internal-logs")] + opentelemetry::_private::warn!( + name: "otel.sdk.component.shutdown", + target: env!("CARGO_PKG_NAME"), + name = "otel.sdk.component.shutdown", + "otel.component.type" = "batching_log_processor", + "otel.component.name" = self.component_name.as_str(), + "otel.component.shutdown.result" = result_str, + "otel.component.dropped" = lifetime_dropped as u64, + "otel.component.shutdown.dropped" = shutdown_dropped, + "otel.component.shutdown.duration" = duration_secs, + ); + } + + #[cfg(not(feature = "internal-logs"))] + { + let _ = ( + result_str, + lifetime_dropped, + shutdown_dropped, + duration_secs, + ); + } } } @@ -530,13 +616,18 @@ 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) = { + // Stable component identity used for self-observability events/metrics. + // Allocated outside the metrics feature gate because the shutdown event + // needs it unconditionally. + let component_name = { 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}"); + format!("batching_log_processor/{instance_id}") + }; + // 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) = { let meter = opentelemetry::global::meter("otel.sdk"); let counter = meter .u64_counter("otel.sdk.processor.log.processed") @@ -562,7 +653,7 @@ impl BatchLogProcessor { 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), + KeyValue::new("otel.component.name", component_name.clone()), ]; ( @@ -580,6 +671,7 @@ impl BatchLogProcessor { forceflush_timeout: Duration::from_secs(5), // TODO: make this configurable dropped_logs_count: AtomicUsize::new(0), max_queue_size, + component_name, export_log_message_sent: Arc::new(AtomicBool::new(false)), current_batch_size, max_export_batch_size, From 20a8ea99acb589c7e553129f3be15af66debf6f0 Mon Sep 17 00:00:00 2001 From: cijothomas Date: Wed, 10 Jun 2026 18:17:03 -0700 Subject: [PATCH 02/11] Trim shutdown event to result + duration (mirror spec) --- opentelemetry-otlp/src/exporter/tonic/logs.rs | 16 +----------- .../src/logs/batch_log_processor.rs | 26 +++---------------- 2 files changed, 4 insertions(+), 38 deletions(-) diff --git a/opentelemetry-otlp/src/exporter/tonic/logs.rs b/opentelemetry-otlp/src/exporter/tonic/logs.rs index 64b6882412..adde47668c 100644 --- a/opentelemetry-otlp/src/exporter/tonic/logs.rs +++ b/opentelemetry-otlp/src/exporter/tonic/logs.rs @@ -186,11 +186,6 @@ impl TonicLogsClient { // opentelemetry::_private (tracing) directly instead of the // otel_info!/otel_warn! macros. fn emit_shutdown_event(&self, result_str: &'static str, duration_secs: f64) { - // Exporter has no internal queue; shutdown.dropped is always 0. - let shutdown_dropped: u64 = 0; - // No lifetime drop counter on this exporter (no queue). - let lifetime_dropped: u64 = 0; - if result_str == "success" { #[cfg(feature = "internal-logs")] opentelemetry::_private::info!( @@ -200,8 +195,6 @@ impl TonicLogsClient { "otel.component.type" = "otlp_grpc_log_exporter", "otel.component.name" = self.component_name.as_str(), "otel.component.shutdown.result" = result_str, - "otel.component.dropped" = lifetime_dropped, - "otel.component.shutdown.dropped" = shutdown_dropped, "otel.component.shutdown.duration" = duration_secs, ); } else { @@ -213,20 +206,13 @@ impl TonicLogsClient { "otel.component.type" = "otlp_grpc_log_exporter", "otel.component.name" = self.component_name.as_str(), "otel.component.shutdown.result" = result_str, - "otel.component.dropped" = lifetime_dropped, - "otel.component.shutdown.dropped" = shutdown_dropped, "otel.component.shutdown.duration" = duration_secs, ); } #[cfg(not(feature = "internal-logs"))] { - let _ = ( - result_str, - duration_secs, - shutdown_dropped, - lifetime_dropped, - ); + let _ = (result_str, duration_secs); } } } diff --git a/opentelemetry-sdk/src/logs/batch_log_processor.rs b/opentelemetry-sdk/src/logs/batch_log_processor.rs index 6b5e0133c2..1cc04273c2 100644 --- a/opentelemetry-sdk/src/logs/batch_log_processor.rs +++ b/opentelemetry-sdk/src/logs/batch_log_processor.rs @@ -312,7 +312,7 @@ impl LogProcessor for BatchLogProcessor { let shutdown_start = Instant::now(); let result = self.shutdown_inner(timeout); let duration_secs = shutdown_start.elapsed().as_secs_f64(); - self.emit_shutdown_event(&result, dropped_logs, duration_secs); + self.emit_shutdown_event(&result, duration_secs); result } @@ -391,12 +391,7 @@ impl BatchLogProcessor { // plain identifiers; the spec requires dotted attribute names like // `otel.component.type`, which tracing's own macros accept via quoted-key // syntax. Behaviour is otherwise identical to otel_info!/otel_warn!. - fn emit_shutdown_event( - &self, - result: &OTelSdkResult, - lifetime_dropped: usize, - duration_secs: f64, - ) { + fn emit_shutdown_event(&self, result: &OTelSdkResult, duration_secs: f64) { let result_str = match result { Ok(()) => "success", Err(OTelSdkError::Timeout(_)) => "timed_out", @@ -409,12 +404,6 @@ impl BatchLogProcessor { }; let is_success = result_str == "success"; - // shutdown.dropped: items lost during shutdown itself. On success the - // spec mandates 0. On timeout/failure we cannot enumerate the queue - // residual from outside the worker thread, so we report 0 as a known - // under-count. SPEC GAP: should the attribute be omitted when unknown? - let shutdown_dropped: u64 = 0; - if is_success { #[cfg(feature = "internal-logs")] opentelemetry::_private::info!( @@ -424,8 +413,6 @@ impl BatchLogProcessor { "otel.component.type" = "batching_log_processor", "otel.component.name" = self.component_name.as_str(), "otel.component.shutdown.result" = result_str, - "otel.component.dropped" = lifetime_dropped as u64, - "otel.component.shutdown.dropped" = shutdown_dropped, "otel.component.shutdown.duration" = duration_secs, ); } else { @@ -437,20 +424,13 @@ impl BatchLogProcessor { "otel.component.type" = "batching_log_processor", "otel.component.name" = self.component_name.as_str(), "otel.component.shutdown.result" = result_str, - "otel.component.dropped" = lifetime_dropped as u64, - "otel.component.shutdown.dropped" = shutdown_dropped, "otel.component.shutdown.duration" = duration_secs, ); } #[cfg(not(feature = "internal-logs"))] { - let _ = ( - result_str, - lifetime_dropped, - shutdown_dropped, - duration_secs, - ); + let _ = (result_str, duration_secs); } } } From 51370de035abc415bb634fea87ccfe15c523d79b Mon Sep 17 00:00:00 2001 From: cijothomas Date: Wed, 10 Jun 2026 18:22:32 -0700 Subject: [PATCH 03/11] Drop dead idempotency check in OTLP exporter (processor guarantees single call) --- opentelemetry-otlp/src/exporter/tonic/logs.rs | 28 +++++++++---------- 1 file changed, 13 insertions(+), 15 deletions(-) diff --git a/opentelemetry-otlp/src/exporter/tonic/logs.rs b/opentelemetry-otlp/src/exporter/tonic/logs.rs index adde47668c..83ff0b1975 100644 --- a/opentelemetry-otlp/src/exporter/tonic/logs.rs +++ b/opentelemetry-otlp/src/exporter/tonic/logs.rs @@ -154,25 +154,18 @@ impl LogExporter for TonicLogsClient { fn shutdown_with_timeout(&self, _timeout: time::Duration) -> OTelSdkResult { let shutdown_start = Instant::now(); - let was_first_shutdown = match self + let result = self .inner .lock() .map_err(|e| OTelSdkError::InternalFailure(format!("Failed to acquire lock: {e}"))) - { - Ok(mut guard) => guard.take().is_some(), - Err(e) => { - let duration_secs = shutdown_start.elapsed().as_secs_f64(); - self.emit_shutdown_event("failed", duration_secs); - return Err(e); - } - }; - + .map(|mut guard| { + guard.take(); + }); let duration_secs = shutdown_start.elapsed().as_secs_f64(); - if was_first_shutdown { - self.emit_shutdown_event("success", duration_secs); - } - // Idempotent re-invocation: emit no event per spec. - Ok(()) + + let result_str = if result.is_ok() { "success" } else { "failed" }; + self.emit_shutdown_event(result_str, duration_secs); + result } fn set_resource(&mut self, resource: &opentelemetry_sdk::Resource) { @@ -185,6 +178,11 @@ impl TonicLogsClient { // BatchLogProcessor::emit_shutdown_event for the rationale on using // opentelemetry::_private (tracing) directly instead of the // otel_info!/otel_warn! macros. + // + // The exporter relies on its owning processor to invoke shutdown exactly + // once (the processor holds the only reference and is itself guarded by + // LoggerProvider's idempotent shutdown). We therefore do not defend against + // re-invocation here. fn emit_shutdown_event(&self, result_str: &'static str, duration_secs: f64) { if result_str == "success" { #[cfg(feature = "internal-logs")] From 276016b1b40170aa8c659d1efa7aec214cd52a6d Mon Sep 17 00:00:00 2001 From: cijothomas Date: Wed, 10 Jun 2026 18:26:53 -0700 Subject: [PATCH 04/11] Drop dead AlreadyShutdown branch in BLP shutdown event (provider guarantees single call) --- opentelemetry-sdk/src/logs/batch_log_processor.rs | 9 ++++----- 1 file changed, 4 insertions(+), 5 deletions(-) diff --git a/opentelemetry-sdk/src/logs/batch_log_processor.rs b/opentelemetry-sdk/src/logs/batch_log_processor.rs index 1cc04273c2..ded2500a67 100644 --- a/opentelemetry-sdk/src/logs/batch_log_processor.rs +++ b/opentelemetry-sdk/src/logs/batch_log_processor.rs @@ -391,15 +391,14 @@ impl BatchLogProcessor { // plain identifiers; the spec requires dotted attribute names like // `otel.component.type`, which tracing's own macros accept via quoted-key // syntax. Behaviour is otherwise identical to otel_info!/otel_warn!. + // + // The processor relies on its owning LoggerProvider to invoke shutdown + // exactly once (the provider holds the only reference and guards calls + // with an atomic CAS). We therefore do not defend against re-invocation. fn emit_shutdown_event(&self, result: &OTelSdkResult, duration_secs: f64) { let result_str = match result { Ok(()) => "success", Err(OTelSdkError::Timeout(_)) => "timed_out", - Err(OTelSdkError::AlreadyShutdown) => { - // Idempotent re-invocation: per spec a component is shut down - // at most once; subsequent calls are no-ops and emit no event. - return; - } Err(_) => "failed", }; let is_success = result_str == "success"; From 610451a1029f3d7bbf401dd9e346be8868d49050 Mon Sep 17 00:00:00 2001 From: cijothomas Date: Mon, 22 Jun 2026 22:11:53 +0530 Subject: [PATCH 05/11] Document why name field is required in shutdown event macros --- opentelemetry-otlp/src/exporter/tonic/logs.rs | 1 + opentelemetry-sdk/src/logs/batch_log_processor.rs | 3 +++ 2 files changed, 4 insertions(+) diff --git a/opentelemetry-otlp/src/exporter/tonic/logs.rs b/opentelemetry-otlp/src/exporter/tonic/logs.rs index 83ff0b1975..6b20260160 100644 --- a/opentelemetry-otlp/src/exporter/tonic/logs.rs +++ b/opentelemetry-otlp/src/exporter/tonic/logs.rs @@ -189,6 +189,7 @@ impl TonicLogsClient { opentelemetry::_private::info!( name: "otel.sdk.component.shutdown", target: env!("CARGO_PKG_NAME"), + // `name = ...` required as the first field; see batch_log_processor.rs. name = "otel.sdk.component.shutdown", "otel.component.type" = "otlp_grpc_log_exporter", "otel.component.name" = self.component_name.as_str(), diff --git a/opentelemetry-sdk/src/logs/batch_log_processor.rs b/opentelemetry-sdk/src/logs/batch_log_processor.rs index ded2500a67..f35ea26b73 100644 --- a/opentelemetry-sdk/src/logs/batch_log_processor.rs +++ b/opentelemetry-sdk/src/logs/batch_log_processor.rs @@ -408,6 +408,9 @@ impl BatchLogProcessor { opentelemetry::_private::info!( name: "otel.sdk.component.shutdown", target: env!("CARGO_PKG_NAME"), + // `name = ...` is required as the first field to prevent the + // tracing macro from parsing the subsequent quoted-key fields + // (e.g. "otel.component.type") as a message string literal. name = "otel.sdk.component.shutdown", "otel.component.type" = "batching_log_processor", "otel.component.name" = self.component_name.as_str(), From 0e374512c7ccd2c3e3724c65f98c5cf03b316038 Mon Sep 17 00:00:00 2001 From: cijothomas Date: Tue, 23 Jun 2026 07:04:21 +0530 Subject: [PATCH 06/11] fix: allow dead_code on component_name when feature-gated consumers are off --- opentelemetry-otlp/src/exporter/tonic/logs.rs | 1 + opentelemetry-sdk/src/logs/batch_log_processor.rs | 7 +++++++ 2 files changed, 8 insertions(+) diff --git a/opentelemetry-otlp/src/exporter/tonic/logs.rs b/opentelemetry-otlp/src/exporter/tonic/logs.rs index 6b20260160..5d8bb8293a 100644 --- a/opentelemetry-otlp/src/exporter/tonic/logs.rs +++ b/opentelemetry-otlp/src/exporter/tonic/logs.rs @@ -25,6 +25,7 @@ pub(crate) struct TonicLogsClient { #[allow(dead_code)] // would be removed once we support set_resource for metrics. resource: opentelemetry_proto::transform::common::tonic::ResourceAttributesWithSchema, + #[cfg_attr(not(feature = "internal-logs"), allow(dead_code))] component_name: String, } diff --git a/opentelemetry-sdk/src/logs/batch_log_processor.rs b/opentelemetry-sdk/src/logs/batch_log_processor.rs index f35ea26b73..a75049c2ad 100644 --- a/opentelemetry-sdk/src/logs/batch_log_processor.rs +++ b/opentelemetry-sdk/src/logs/batch_log_processor.rs @@ -155,6 +155,13 @@ pub struct BatchLogProcessor { // Format: "batching_log_processor/". Always present, independent of // the experimental_metrics_bound_instruments feature flag because the // shutdown event (otel.sdk.component.shutdown) needs it unconditionally. + #[cfg_attr( + not(any( + feature = "internal-logs", + feature = "experimental_metrics_bound_instruments" + )), + allow(dead_code) + )] component_name: String, // Self-diagnostics: otel.sdk.processor.log.processed counter. From d92310297a7943d4c7904c741eee79f4ea19036f Mon Sep 17 00:00:00 2001 From: cijothomas Date: Fri, 26 Jun 2026 18:55:19 +0530 Subject: [PATCH 07/11] Pivot POC to provider-level shutdown event - Move emission from BatchLogProcessor + TonicLogsClient to LoggerProvider.shutdown_with_timeout (one event per provider). - Use error.type (absent=success, 'timeout'/'failed' on failure) instead of custom otel.component.shutdown.result enum. - otel.component.type = 'logger_provider'. - Revert batch_log_processor.rs and tonic/logs.rs to main (no per-component emission). --- opentelemetry-otlp/src/exporter/tonic/logs.rs | 71 +----------- .../src/logs/batch_log_processor.rs | 101 ++---------------- opentelemetry-sdk/src/logs/logger_provider.rs | 32 +++++- 3 files changed, 42 insertions(+), 162 deletions(-) diff --git a/opentelemetry-otlp/src/exporter/tonic/logs.rs b/opentelemetry-otlp/src/exporter/tonic/logs.rs index 5d8bb8293a..6704a2a2c5 100644 --- a/opentelemetry-otlp/src/exporter/tonic/logs.rs +++ b/opentelemetry-otlp/src/exporter/tonic/logs.rs @@ -5,10 +5,8 @@ use opentelemetry_proto::tonic::collector::logs::v1::{ }; use opentelemetry_sdk::error::{OTelSdkError, OTelSdkResult}; use opentelemetry_sdk::logs::{LogBatch, LogExporter}; -use std::sync::atomic::{AtomicUsize, Ordering}; use std::sync::{Arc, Mutex}; use std::time; -use std::time::Instant; use tonic::{codegen::CompressionEncoding, service::Interceptor, transport::Channel, Request}; use opentelemetry_proto::transform::logs::tonic::group_logs_by_resource_and_scope; @@ -25,8 +23,6 @@ pub(crate) struct TonicLogsClient { #[allow(dead_code)] // would be removed once we support set_resource for metrics. resource: opentelemetry_proto::transform::common::tonic::ResourceAttributesWithSchema, - #[cfg_attr(not(feature = "internal-logs"), allow(dead_code))] - component_name: String, } struct ClientInner { @@ -56,12 +52,6 @@ impl TonicLogsClient { otel_debug!(name: "TonicsLogsClientBuilt"); - let component_name = { - static INSTANCE_COUNTER: AtomicUsize = AtomicUsize::new(0); - let instance_id = INSTANCE_COUNTER.fetch_add(1, Ordering::Relaxed); - format!("otlp_grpc_log_exporter/{instance_id}") - }; - TonicLogsClient { inner: Mutex::new(Some(ClientInner { client, @@ -74,7 +64,6 @@ impl TonicLogsClient { jitter_ms: 100, }), resource: Default::default(), - component_name, } } } @@ -154,65 +143,15 @@ impl LogExporter for TonicLogsClient { } fn shutdown_with_timeout(&self, _timeout: time::Duration) -> OTelSdkResult { - let shutdown_start = Instant::now(); - let result = self - .inner + self.inner .lock() - .map_err(|e| OTelSdkError::InternalFailure(format!("Failed to acquire lock: {e}"))) - .map(|mut guard| { - guard.take(); - }); - let duration_secs = shutdown_start.elapsed().as_secs_f64(); - - let result_str = if result.is_ok() { "success" } else { "failed" }; - self.emit_shutdown_event(result_str, duration_secs); - result + .map_err(|e| OTelSdkError::InternalFailure(format!("Failed to acquire lock: {e}")))? + .take(); + + Ok(()) } fn set_resource(&mut self, resource: &opentelemetry_sdk::Resource) { self.resource = resource.into(); } } - -impl TonicLogsClient { - // POC: emit otel.sdk.component.shutdown for this exporter. See - // BatchLogProcessor::emit_shutdown_event for the rationale on using - // opentelemetry::_private (tracing) directly instead of the - // otel_info!/otel_warn! macros. - // - // The exporter relies on its owning processor to invoke shutdown exactly - // once (the processor holds the only reference and is itself guarded by - // LoggerProvider's idempotent shutdown). We therefore do not defend against - // re-invocation here. - fn emit_shutdown_event(&self, result_str: &'static str, duration_secs: f64) { - if result_str == "success" { - #[cfg(feature = "internal-logs")] - opentelemetry::_private::info!( - name: "otel.sdk.component.shutdown", - target: env!("CARGO_PKG_NAME"), - // `name = ...` required as the first field; see batch_log_processor.rs. - name = "otel.sdk.component.shutdown", - "otel.component.type" = "otlp_grpc_log_exporter", - "otel.component.name" = self.component_name.as_str(), - "otel.component.shutdown.result" = result_str, - "otel.component.shutdown.duration" = duration_secs, - ); - } else { - #[cfg(feature = "internal-logs")] - opentelemetry::_private::warn!( - name: "otel.sdk.component.shutdown", - target: env!("CARGO_PKG_NAME"), - name = "otel.sdk.component.shutdown", - "otel.component.type" = "otlp_grpc_log_exporter", - "otel.component.name" = self.component_name.as_str(), - "otel.component.shutdown.result" = result_str, - "otel.component.shutdown.duration" = duration_secs, - ); - } - - #[cfg(not(feature = "internal-logs"))] - { - let _ = (result_str, duration_secs); - } - } -} diff --git a/opentelemetry-sdk/src/logs/batch_log_processor.rs b/opentelemetry-sdk/src/logs/batch_log_processor.rs index a75049c2ad..ac65a0a958 100644 --- a/opentelemetry-sdk/src/logs/batch_log_processor.rs +++ b/opentelemetry-sdk/src/logs/batch_log_processor.rs @@ -151,19 +151,6 @@ pub struct BatchLogProcessor { // Track the maximum queue size that was configured for this processor max_queue_size: usize, - // Stable component identity used for self-observability events/metrics. - // Format: "batching_log_processor/". Always present, independent of - // the experimental_metrics_bound_instruments feature flag because the - // shutdown event (otel.sdk.component.shutdown) needs it unconditionally. - #[cfg_attr( - not(any( - feature = "internal-logs", - feature = "experimental_metrics_bound_instruments" - )), - allow(dead_code) - )] - component_name: String, - // 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 @@ -316,23 +303,6 @@ impl LogProcessor for BatchLogProcessor { ); } - let shutdown_start = Instant::now(); - let result = self.shutdown_inner(timeout); - let duration_secs = shutdown_start.elapsed().as_secs_f64(); - self.emit_shutdown_event(&result, duration_secs); - result - } - - fn set_resource(&mut self, resource: &Resource) { - let resource = Arc::new(resource.clone()); - let _ = self - .message_sender - .try_send(BatchMessage::SetResource(resource)); - } -} - -impl BatchLogProcessor { - fn shutdown_inner(&self, timeout: Duration) -> OTelSdkResult { let (sender, receiver) = mpsc::sync_channel(1); match self.message_sender.try_send(BatchMessage::Shutdown(sender)) { Ok(_) => { @@ -391,56 +361,11 @@ impl BatchLogProcessor { } } - // POC: emit the otel.sdk.component.shutdown event for this processor. - // Per spec PR open-telemetry/semantic-conventions#3723 the event is INFO - // on success and WARN on any non-success result. We bypass the - // otel_info!/otel_warn! macros because those restrict attribute keys to - // plain identifiers; the spec requires dotted attribute names like - // `otel.component.type`, which tracing's own macros accept via quoted-key - // syntax. Behaviour is otherwise identical to otel_info!/otel_warn!. - // - // The processor relies on its owning LoggerProvider to invoke shutdown - // exactly once (the provider holds the only reference and guards calls - // with an atomic CAS). We therefore do not defend against re-invocation. - fn emit_shutdown_event(&self, result: &OTelSdkResult, duration_secs: f64) { - let result_str = match result { - Ok(()) => "success", - Err(OTelSdkError::Timeout(_)) => "timed_out", - Err(_) => "failed", - }; - let is_success = result_str == "success"; - - if is_success { - #[cfg(feature = "internal-logs")] - opentelemetry::_private::info!( - name: "otel.sdk.component.shutdown", - target: env!("CARGO_PKG_NAME"), - // `name = ...` is required as the first field to prevent the - // tracing macro from parsing the subsequent quoted-key fields - // (e.g. "otel.component.type") as a message string literal. - name = "otel.sdk.component.shutdown", - "otel.component.type" = "batching_log_processor", - "otel.component.name" = self.component_name.as_str(), - "otel.component.shutdown.result" = result_str, - "otel.component.shutdown.duration" = duration_secs, - ); - } else { - #[cfg(feature = "internal-logs")] - opentelemetry::_private::warn!( - name: "otel.sdk.component.shutdown", - target: env!("CARGO_PKG_NAME"), - name = "otel.sdk.component.shutdown", - "otel.component.type" = "batching_log_processor", - "otel.component.name" = self.component_name.as_str(), - "otel.component.shutdown.result" = result_str, - "otel.component.shutdown.duration" = duration_secs, - ); - } - - #[cfg(not(feature = "internal-logs"))] - { - let _ = (result_str, duration_secs); - } + fn set_resource(&mut self, resource: &Resource) { + let resource = Arc::new(resource.clone()); + let _ = self + .message_sender + .try_send(BatchMessage::SetResource(resource)); } } @@ -605,18 +530,13 @@ impl BatchLogProcessor { }) .expect("Thread spawn failed."); //TODO: Handle thread spawn failure - // Stable component identity used for self-observability events/metrics. - // Allocated outside the metrics feature gate because the shutdown event - // needs it unconditionally. - let component_name = { - static INSTANCE_COUNTER: AtomicUsize = AtomicUsize::new(0); - let instance_id = INSTANCE_COUNTER.fetch_add(1, Ordering::Relaxed); - format!("batching_log_processor/{instance_id}") - }; - // 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") @@ -642,7 +562,7 @@ impl BatchLogProcessor { 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.clone()), + KeyValue::new("otel.component.name", component_name), ]; ( @@ -660,7 +580,6 @@ impl BatchLogProcessor { forceflush_timeout: Duration::from_secs(5), // TODO: make this configurable dropped_logs_count: AtomicUsize::new(0), max_queue_size, - component_name, export_log_message_sent: Arc::new(AtomicBool::new(false)), current_batch_size, max_export_batch_size, diff --git a/opentelemetry-sdk/src/logs/logger_provider.rs b/opentelemetry-sdk/src/logs/logger_provider.rs index e8765b456b..8c2a921549 100644 --- a/opentelemetry-sdk/src/logs/logger_provider.rs +++ b/opentelemetry-sdk/src/logs/logger_provider.rs @@ -2,8 +2,8 @@ use super::{BatchLogProcessor, LogProcessor, SdkLogger, SimpleLogProcessor}; use crate::error::{OTelSdkError, OTelSdkResult}; use crate::logs::LogExporter; use crate::Resource; -use opentelemetry::{otel_debug, otel_info, InstrumentationScope}; -use std::time::Duration; +use opentelemetry::{otel_debug, otel_info, otel_warn, InstrumentationScope}; +use std::time::{Duration, Instant}; use std::{ borrow::Cow, sync::{ @@ -107,14 +107,36 @@ impl SdkLoggerProvider { .compare_exchange(false, true, Ordering::SeqCst, Ordering::SeqCst) .is_ok() { + let shutdown_start = Instant::now(); // propagate the shutdown signal to processors - let result = self.inner.shutdown_with_timeout(timeout); - if result.iter().all(|res| res.is_ok()) { + let results = self.inner.shutdown_with_timeout(timeout); + let duration_secs = shutdown_start.elapsed().as_secs_f64(); + + let success = results.iter().all(|res| res.is_ok()); + if success { + otel_info!(name: "otel.sdk.component.shutdown", + "otel.component.type" = "logger_provider", + "otel.component.shutdown.duration" = duration_secs, + ); Ok(()) } else { + // Classify: timeout takes precedence over failed + let error_type = if results + .iter() + .any(|r| matches!(r, Err(OTelSdkError::Timeout(_)))) + { + "timeout" + } else { + "failed" + }; + otel_warn!(name: "otel.sdk.component.shutdown", + "error.type" = error_type, + "otel.component.type" = "logger_provider", + "otel.component.shutdown.duration" = duration_secs, + ); Err(OTelSdkError::InternalFailure(format!( "Shutdown errors: {:?}", - result + results .into_iter() .filter_map(Result::err) .collect::>() From 3c4c0320808ee669d6f3bca631286c704cc756fe Mon Sep 17 00:00:00 2001 From: cijothomas Date: Fri, 26 Jun 2026 19:05:46 +0530 Subject: [PATCH 08/11] Add otel.component.name to provider shutdown event --- opentelemetry-sdk/src/logs/logger_provider.rs | 14 ++++++++++++++ 1 file changed, 14 insertions(+) diff --git a/opentelemetry-sdk/src/logs/logger_provider.rs b/opentelemetry-sdk/src/logs/logger_provider.rs index 8c2a921549..6350476265 100644 --- a/opentelemetry-sdk/src/logs/logger_provider.rs +++ b/opentelemetry-sdk/src/logs/logger_provider.rs @@ -23,6 +23,7 @@ fn noop_logger_provider() -> &'static SdkLoggerProvider { processors: Vec::new(), is_shutdown: AtomicBool::new(true), }), + component_name: String::new(), }) } @@ -42,6 +43,7 @@ fn noop_logger_provider() -> &'static SdkLoggerProvider { /// [`Resource`]: crate::Resource pub struct SdkLoggerProvider { inner: Arc, + component_name: String, } impl opentelemetry::logs::LoggerProvider for SdkLoggerProvider { @@ -115,6 +117,7 @@ impl SdkLoggerProvider { let success = results.iter().all(|res| res.is_ok()); if success { otel_info!(name: "otel.sdk.component.shutdown", + "otel.component.name" = self.component_name.as_str(), "otel.component.type" = "logger_provider", "otel.component.shutdown.duration" = duration_secs, ); @@ -130,6 +133,7 @@ impl SdkLoggerProvider { "failed" }; otel_warn!(name: "otel.sdk.component.shutdown", + "otel.component.name" = self.component_name.as_str(), "error.type" = error_type, "otel.component.type" = "logger_provider", "otel.component.shutdown.duration" = duration_secs, @@ -317,11 +321,17 @@ impl LoggerProviderBuilder { processor.set_resource(&resource); } + static INSTANCE_COUNTER: std::sync::atomic::AtomicUsize = + std::sync::atomic::AtomicUsize::new(0); + let instance_id = INSTANCE_COUNTER.fetch_add(1, std::sync::atomic::Ordering::Relaxed); + let component_name = format!("logger_provider/{instance_id}"); + let logger_provider = SdkLoggerProvider { inner: Arc::new(LoggerProviderInner { processors, is_shutdown: AtomicBool::new(false), }), + component_name, }; otel_debug!( @@ -820,9 +830,11 @@ mod tests { { let logger_provider1 = SdkLoggerProvider { inner: shared_inner.clone(), + component_name: String::new(), }; let logger_provider2 = SdkLoggerProvider { inner: shared_inner.clone(), + component_name: String::new(), }; let logger1 = logger_provider1.logger("test-logger1"); @@ -861,9 +873,11 @@ mod tests { { let logger_provider1 = SdkLoggerProvider { inner: shared_inner.clone(), + component_name: String::new(), }; let logger_provider2 = SdkLoggerProvider { inner: shared_inner.clone(), + component_name: String::new(), }; // Explicitly shut down the logger provider From df1dda12d6e4b082e7f3c81c768746626cc565f7 Mon Sep 17 00:00:00 2001 From: cijothomas Date: Fri, 26 Jun 2026 19:17:58 +0530 Subject: [PATCH 09/11] Expand shutdown event to all 3 providers (TracerProvider, MeterProvider) Each provider now emits otel.sdk.component.shutdown with: - otel.component.type: tracer_provider / meter_provider / logger_provider - otel.component.name: / - error.type: absent on success, 'timeout'/'failed' on failure - otel.component.shutdown.duration: wall-clock seconds --- .../src/metrics/meter_provider.rs | 30 ++++++++++++++++-- opentelemetry-sdk/src/trace/provider.rs | 31 +++++++++++++++++-- 2 files changed, 57 insertions(+), 4 deletions(-) diff --git a/opentelemetry-sdk/src/metrics/meter_provider.rs b/opentelemetry-sdk/src/metrics/meter_provider.rs index 27f7ca9cde..935f95ba3b 100644 --- a/opentelemetry-sdk/src/metrics/meter_provider.rs +++ b/opentelemetry-sdk/src/metrics/meter_provider.rs @@ -1,7 +1,7 @@ use core::fmt; use opentelemetry::{ metrics::{Meter, MeterProvider}, - otel_debug, otel_error, otel_info, InstrumentationScope, + otel_debug, otel_error, otel_info, otel_warn, InstrumentationScope, }; use std::time::Duration; use std::{ @@ -34,6 +34,7 @@ use super::{ #[derive(Clone, Debug)] pub struct SdkMeterProvider { inner: Arc, + component_name: String, } #[derive(Debug)] @@ -115,7 +116,27 @@ impl SdkMeterProvider { name: "MeterProvider.Shutdown", message = "User initiated shutdown of MeterProvider." ); - self.inner.shutdown() + let shutdown_start = std::time::Instant::now(); + let result = self.inner.shutdown(); + let duration_secs = shutdown_start.elapsed().as_secs_f64(); + match &result { + Ok(()) => { + otel_info!(name: "otel.sdk.component.shutdown", + "otel.component.name" = self.component_name.as_str(), + "otel.component.type" = "meter_provider", + "otel.component.shutdown.duration" = duration_secs, + ); + } + Err(_) => { + otel_warn!(name: "otel.sdk.component.shutdown", + "otel.component.name" = self.component_name.as_str(), + "error.type" = "failed", + "otel.component.type" = "meter_provider", + "otel.component.shutdown.duration" = duration_secs, + ); + } + } + result } /// shutdown with default timeout @@ -401,6 +422,10 @@ impl MeterProviderBuilder { builder = format!("{:?}", &self), ); + static INSTANCE_COUNTER: std::sync::atomic::AtomicUsize = + std::sync::atomic::AtomicUsize::new(0); + let instance_id = INSTANCE_COUNTER.fetch_add(1, std::sync::atomic::Ordering::Relaxed); + let meter_provider = SdkMeterProvider { inner: Arc::new(SdkMeterProviderInner { pipes: Arc::new(Pipelines::new( @@ -411,6 +436,7 @@ impl MeterProviderBuilder { meters: Default::default(), shutdown_invoked: AtomicBool::new(false), }), + component_name: format!("meter_provider/{instance_id}"), }; otel_debug!( diff --git a/opentelemetry-sdk/src/trace/provider.rs b/opentelemetry-sdk/src/trace/provider.rs index b8dc113ca7..1cfbe00db4 100644 --- a/opentelemetry-sdk/src/trace/provider.rs +++ b/opentelemetry-sdk/src/trace/provider.rs @@ -71,7 +71,7 @@ use crate::trace::{ use crate::Resource; use crate::{trace::SpanExporter, trace::SpanProcessor}; use opentelemetry::otel_debug; -use opentelemetry::{otel_info, InstrumentationScope}; +use opentelemetry::{otel_info, otel_warn, InstrumentationScope}; use std::borrow::Cow; use std::sync::atomic::{AtomicBool, Ordering}; use std::sync::{Arc, OnceLock}; @@ -97,6 +97,7 @@ fn noop_tracer_provider() -> &'static SdkTracerProvider { }, is_shutdown: AtomicBool::new(true), }), + component_name: String::new(), } }) } @@ -157,6 +158,7 @@ impl Drop for TracerProviderInner { #[derive(Clone, Debug)] pub struct SdkTracerProvider { inner: Arc, + component_name: String, } impl Default for SdkTracerProvider { @@ -168,8 +170,12 @@ impl Default for SdkTracerProvider { impl SdkTracerProvider { /// Build a new tracer provider pub(crate) fn new(inner: TracerProviderInner) -> Self { + static INSTANCE_COUNTER: std::sync::atomic::AtomicUsize = + std::sync::atomic::AtomicUsize::new(0); + let instance_id = INSTANCE_COUNTER.fetch_add(1, std::sync::atomic::Ordering::Relaxed); SdkTracerProvider { inner: Arc::new(inner), + component_name: format!("tracer_provider/{instance_id}"), } } @@ -250,18 +256,39 @@ impl SdkTracerProvider { .compare_exchange(false, true, Ordering::SeqCst, Ordering::SeqCst) .is_ok() { + let shutdown_start = std::time::Instant::now(); // propagate the shutdown signal to processors let results = self.inner.shutdown_with_timeout(timeout); + let duration_secs = shutdown_start.elapsed().as_secs_f64(); if results.iter().all(|res| res.is_ok()) { + otel_info!(name: "otel.sdk.component.shutdown", + "otel.component.name" = self.component_name.as_str(), + "otel.component.type" = "tracer_provider", + "otel.component.shutdown.duration" = duration_secs, + ); Ok(()) } else { + let error_type = if results + .iter() + .any(|r| matches!(r, Err(OTelSdkError::Timeout(_)))) + { + "timeout" + } else { + "failed" + }; + otel_warn!(name: "otel.sdk.component.shutdown", + "otel.component.name" = self.component_name.as_str(), + "error.type" = error_type, + "otel.component.type" = "tracer_provider", + "otel.component.shutdown.duration" = duration_secs, + ); Err(OTelSdkError::InternalFailure(format!( "Shutdown errors: {:?}", results .into_iter() .filter_map(Result::err) - .collect::>() // Collect only the errors + .collect::>() ))) } } else { From 3862701c528ddcf4412c1066ec1d6255c97b72bb Mon Sep 17 00:00:00 2001 From: cijothomas Date: Fri, 26 Jun 2026 20:16:27 +0530 Subject: [PATCH 10/11] ci: add weaver live-check for otel.sdk.component.shutdown event - In-tree semconv registry (semconv/) declaring the shutdown event with attributes (local until semconv#3723 is released upstream). - Dedicated live-check example that creates a LoggerProvider, calls shutdown(), and lets the event flow to weaver via OTLP. - Workflow using weaver-live-check-{start,stop} composite actions with fail-on: violation. --- .../sdk-shutdown-event-live-check.yml | 95 +++++++++++++++++++ examples/shutdown-event-live-check/Cargo.toml | 24 +++++ .../shutdown-event-live-check/src/main.rs | 45 +++++++++ .../event.otel.sdk.component.shutdown.yaml | 29 ++++++ semconv/manifest.yaml | 12 +++ 5 files changed, 205 insertions(+) create mode 100644 .github/workflows/sdk-shutdown-event-live-check.yml create mode 100644 examples/shutdown-event-live-check/Cargo.toml create mode 100644 examples/shutdown-event-live-check/src/main.rs create mode 100644 semconv/groups/event.otel.sdk.component.shutdown.yaml create mode 100644 semconv/manifest.yaml diff --git a/.github/workflows/sdk-shutdown-event-live-check.yml b/.github/workflows/sdk-shutdown-event-live-check.yml new file mode 100644 index 0000000000..fcc96b5f31 --- /dev/null +++ b/.github/workflows/sdk-shutdown-event-live-check.yml @@ -0,0 +1,95 @@ +name: SDK Shutdown Event Live Check + +# Validates that the otel.sdk.component.shutdown event emitted by SDK +# providers matches the in-tree semconv registry under `semconv/`. +# Companion to #3553 (SDK self-obs metrics live-check). + +env: + CI: true + WEAVER_VERSION: v0.24.0 + +permissions: + contents: read + +on: + workflow_dispatch: + pull_request: + paths: + - 'opentelemetry-sdk/src/logs/logger_provider.rs' + - 'opentelemetry-sdk/src/trace/provider.rs' + - 'opentelemetry-sdk/src/metrics/meter_provider.rs' + - 'semconv/**' + - 'examples/shutdown-event-live-check/**' + - '.github/workflows/sdk-shutdown-event-live-check.yml' + merge_group: + push: + branches: + - main + paths: + - 'opentelemetry-sdk/src/logs/logger_provider.rs' + - 'opentelemetry-sdk/src/trace/provider.rs' + - 'opentelemetry-sdk/src/metrics/meter_provider.rs' + - 'semconv/**' + - 'examples/shutdown-event-live-check/**' + - '.github/workflows/sdk-shutdown-event-live-check.yml' + +concurrency: + group: ${{ github.workflow }}-${{ github.ref }} + cancel-in-progress: true + +jobs: + live-check: + runs-on: ubuntu-latest + timeout-minutes: 10 + steps: + - uses: actions/checkout@de0fac2e4500dabe0009e67214ff5f5447ce83dd # v6.0.2 + + - name: Install stable Rust + run: rustup toolchain install stable --profile minimal --no-self-update + + - name: Cache cargo + uses: actions/cache@v4 + with: + path: | + ~/.cargo/registry + ~/.cargo/git + target + key: shutdown-live-check-${{ runner.os }}-${{ hashFiles('**/Cargo.lock') }} + + - name: Build live-check app + run: cargo build --release -p shutdown-event-live-check + + - name: Setup weaver + uses: open-telemetry/weaver/.github/actions/setup-weaver@v0.24.0 + with: + version: ${{ env.WEAVER_VERSION }} + + - name: Start weaver live-check + id: live-check + uses: open-telemetry/weaver/.github/actions/weaver-live-check-start@v0.24.0 + with: + registry: semconv + + - name: Run live-check app + env: + OTEL_BSP_SCHEDULE_DELAY: '500' + OTEL_EXPORTER_OTLP_ENDPOINT: ${{ steps.live-check.outputs.otlp-grpc-endpoint }} + run: | + set -euo pipefail + ./target/release/shutdown-event-live-check > "${{ steps.live-check.outputs.state-dir }}/app.log" 2>&1 + echo "App exited cleanly" + + - name: Stop weaver live-check + id: weaver-stop + if: always() + uses: open-telemetry/weaver/.github/actions/weaver-live-check-stop@v0.24.0 + with: + fail-on: violation + + - name: Upload live-check report + if: always() + uses: actions/upload-artifact@v4 + with: + name: sdk-shutdown-event-live-check + path: ${{ steps.live-check.outputs.state-dir }} + if-no-files-found: warn diff --git a/examples/shutdown-event-live-check/Cargo.toml b/examples/shutdown-event-live-check/Cargo.toml new file mode 100644 index 0000000000..a3e43a8d43 --- /dev/null +++ b/examples/shutdown-event-live-check/Cargo.toml @@ -0,0 +1,24 @@ +[package] +name = "shutdown-event-live-check" +version = "0.1.0" +edition = "2021" +license = "Apache-2.0" +rust-version = "1.75.0" +publish = false +autobenches = false + +# Purpose-built example for the weaver live-check CI workflow that validates +# the otel.sdk.component.shutdown event against the in-tree semconv registry. + +[[bin]] +name = "shutdown-event-live-check" +path = "src/main.rs" +bench = false + +[dependencies] +opentelemetry = { workspace = true } +opentelemetry_sdk = { workspace = true, features = ["logs", "trace", "metrics"] } +opentelemetry-otlp = { workspace = true, features = ["grpc-tonic", "logs"] } +opentelemetry-appender-tracing = { workspace = true } +tracing-subscriber = { workspace = true, features = ["env-filter", "registry", "std", "fmt"] } +tokio = { workspace = true, features = ["rt-multi-thread", "macros"] } diff --git a/examples/shutdown-event-live-check/src/main.rs b/examples/shutdown-event-live-check/src/main.rs new file mode 100644 index 0000000000..4df4d84057 --- /dev/null +++ b/examples/shutdown-event-live-check/src/main.rs @@ -0,0 +1,45 @@ +// Purpose-built binary for validating the otel.sdk.component.shutdown event +// via weaver live-check. +// +// Creates a LoggerProvider with OTLP exporter (pointed at weaver listener), +// then calls shutdown(). The shutdown event flows through the tracing bridge +// → OTel LoggerProvider → OTLP exporter → weaver for validation. + +use opentelemetry_otlp::LogExporter; +use opentelemetry_sdk::logs::SdkLoggerProvider; +use opentelemetry_sdk::Resource; +use tracing_subscriber::{layer::SubscriberExt, util::SubscriberInitExt}; + +const SERVICE_NAME: &str = "shutdown-event-live-check"; + +#[tokio::main] +async fn main() { + // LoggerProvider with OTLP exporter pointed at weaver (via env var) + let log_exporter = LogExporter::builder().with_tonic().build().unwrap(); + let logger_provider = SdkLoggerProvider::builder() + .with_resource(Resource::builder().with_service_name(SERVICE_NAME).build()) + .with_batch_exporter(log_exporter) + .build(); + + // Wire tracing → OTel so the otel_info! shutdown event (which uses + // tracing internally) flows through this LoggerProvider to weaver. + let otel_layer = + opentelemetry_appender_tracing::layer::OpenTelemetryTracingBridge::new(&logger_provider); + // Only capture opentelemetry crate logs (the shutdown event target is + // the SDK crate name). Filter out h2/tonic/hyper noise. + let filter = tracing_subscriber::EnvFilter::new("opentelemetry=trace,opentelemetry_sdk=trace"); + tracing_subscriber::registry() + .with(filter) + .with(otel_layer) + .init(); + + // Brief pause to let the batch processor's background thread start + tokio::time::sleep(std::time::Duration::from_millis(500)).await; + + // Shutdown: this emits the otel.sdk.component.shutdown event + let _ = logger_provider.shutdown(); + + // Give the OTLP exporter time to flush the shutdown event to weaver + // (the event is emitted synchronously but the BLP's export is async) + tokio::time::sleep(std::time::Duration::from_secs(3)).await; +} diff --git a/semconv/groups/event.otel.sdk.component.shutdown.yaml b/semconv/groups/event.otel.sdk.component.shutdown.yaml new file mode 100644 index 0000000000..2f14a3f034 --- /dev/null +++ b/semconv/groups/event.otel.sdk.component.shutdown.yaml @@ -0,0 +1,29 @@ +groups: + - id: event.otel.sdk.component.shutdown + type: event + name: otel.sdk.component.shutdown + stability: development + brief: > + Emitted when an SDK provider's shutdown attempt has ended, whether by + completing successfully, failing, or timing out. + attributes: + - ref: otel.component.type + requirement_level: required + brief: > + The provider type. + examples: ["logger_provider", "tracer_provider", "meter_provider"] + - ref: otel.component.name + requirement_level: required + - ref: error.type + requirement_level: + conditionally_required: if the shutdown was not successful. + brief: > + Describes the class of error that caused shutdown to fail. + examples: ["timeout", "failed"] + - id: otel.component.shutdown.duration + type: double + stability: development + requirement_level: recommended + brief: > + Wall-clock elapsed time of the shutdown attempt, in seconds. + examples: [0.015, 30.0] diff --git a/semconv/manifest.yaml b/semconv/manifest.yaml new file mode 100644 index 0000000000..b63ca3b9c9 --- /dev/null +++ b/semconv/manifest.yaml @@ -0,0 +1,12 @@ +name: io.opentelemetry.sdk +description: > + In-tree semantic conventions for self-observability events emitted by the + OpenTelemetry Rust SDK. Validated by weaver live-check in CI. + + Once open-telemetry/semantic-conventions#3723 is released, these local + definitions can be replaced by a dependency reference to the upstream + registry. +schema_url: https://opentelemetry.io/schemas/1.42.0 +dependencies: + - schema_url: https://opentelemetry.io/schemas/1.42.0 + registry_path: https://github.com/open-telemetry/semantic-conventions.git@v1.42.0[model] From 7c783b14c947cdc2cc6b3efac1012ef26117f558 Mon Sep 17 00:00:00 2001 From: cijothomas Date: Sat, 27 Jun 2026 06:12:14 +0530 Subject: [PATCH 11/11] Fix MeterProvider idempotency: check shutdown_invoked before emitting Second call to shutdown_with_timeout now returns AlreadyShutdown immediately without emitting an event, matching Logger/TracerProvider. --- opentelemetry-sdk/src/metrics/meter_provider.rs | 15 ++++++++++++++- 1 file changed, 14 insertions(+), 1 deletion(-) diff --git a/opentelemetry-sdk/src/metrics/meter_provider.rs b/opentelemetry-sdk/src/metrics/meter_provider.rs index 935f95ba3b..8fd58a61af 100644 --- a/opentelemetry-sdk/src/metrics/meter_provider.rs +++ b/opentelemetry-sdk/src/metrics/meter_provider.rs @@ -12,7 +12,7 @@ use std::{ }, }; -use crate::error::OTelSdkResult; +use crate::error::{OTelSdkError, OTelSdkResult}; use crate::Resource; use super::{ @@ -112,6 +112,13 @@ impl SdkMeterProvider { /// There is no guaranteed that all telemetry be flushed or all resources have /// been released on error. pub fn shutdown_with_timeout(&self, _timeout: Duration) -> OTelSdkResult { + // Check idempotency at the outer level to avoid emitting multiple events. + // The inner shutdown_invoked swap is the authoritative guard; we read it + // here only to short-circuit the timing + emission path. + if self.inner.shutdown_invoked.load(Ordering::SeqCst) { + return Err(OTelSdkError::AlreadyShutdown); + } + otel_debug!( name: "MeterProvider.Shutdown", message = "User initiated shutdown of MeterProvider." @@ -119,6 +126,9 @@ impl SdkMeterProvider { let shutdown_start = std::time::Instant::now(); let result = self.inner.shutdown(); let duration_secs = shutdown_start.elapsed().as_secs_f64(); + + // Only emit the event if this was the first (actual) shutdown call. + // inner.shutdown() returns AlreadyShutdown on races; don't emit in that case. match &result { Ok(()) => { otel_info!(name: "otel.sdk.component.shutdown", @@ -127,6 +137,9 @@ impl SdkMeterProvider { "otel.component.shutdown.duration" = duration_secs, ); } + Err(OTelSdkError::AlreadyShutdown) => { + // Race: another thread won the swap. Don't emit. + } Err(_) => { otel_warn!(name: "otel.sdk.component.shutdown", "otel.component.name" = self.component_name.as_str(),