Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
95 changes: 95 additions & 0 deletions .github/workflows/sdk-shutdown-event-live-check.yml
Original file line number Diff line number Diff line change
@@ -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
24 changes: 24 additions & 0 deletions examples/shutdown-event-live-check/Cargo.toml
Original file line number Diff line number Diff line change
@@ -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 }

Check failure on line 19 in examples/shutdown-event-live-check/Cargo.toml

View workflow job for this annotation

GitHub Actions / cargo-shear

shear/unused_dependency

unused dependency `opentelemetry` (remove this dependency)
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"] }
45 changes: 45 additions & 0 deletions examples/shutdown-event-live-check/src/main.rs
Original file line number Diff line number Diff line change
@@ -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;
}
46 changes: 41 additions & 5 deletions opentelemetry-sdk/src/logs/logger_provider.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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::{
Expand All @@ -23,6 +23,7 @@ fn noop_logger_provider() -> &'static SdkLoggerProvider {
processors: Vec::new(),
is_shutdown: AtomicBool::new(true),
}),
component_name: String::new(),
})
}

Expand All @@ -42,6 +43,7 @@ fn noop_logger_provider() -> &'static SdkLoggerProvider {
/// [`Resource`]: crate::Resource
pub struct SdkLoggerProvider {
inner: Arc<LoggerProviderInner>,
component_name: String,
}

impl opentelemetry::logs::LoggerProvider for SdkLoggerProvider {
Expand Down Expand Up @@ -107,14 +109,38 @@ 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.name" = self.component_name.as_str(),
"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",
"otel.component.name" = self.component_name.as_str(),
"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::<Vec<_>>()
Expand Down Expand Up @@ -295,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!(
Expand Down Expand Up @@ -798,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");
Expand Down Expand Up @@ -839,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
Expand Down
45 changes: 42 additions & 3 deletions opentelemetry-sdk/src/metrics/meter_provider.rs
Original file line number Diff line number Diff line change
@@ -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::{
Expand All @@ -12,7 +12,7 @@ use std::{
},
};

use crate::error::OTelSdkResult;
use crate::error::{OTelSdkError, OTelSdkResult};
use crate::Resource;

use super::{
Expand All @@ -34,6 +34,7 @@ use super::{
#[derive(Clone, Debug)]
pub struct SdkMeterProvider {
inner: Arc<SdkMeterProviderInner>,
component_name: String,
}

#[derive(Debug)]
Expand Down Expand Up @@ -111,11 +112,44 @@ 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."
);
self.inner.shutdown()
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",
"otel.component.name" = self.component_name.as_str(),
"otel.component.type" = "meter_provider",
"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(),
"error.type" = "failed",
"otel.component.type" = "meter_provider",
"otel.component.shutdown.duration" = duration_secs,
);
}
}
result
}

/// shutdown with default timeout
Expand Down Expand Up @@ -401,6 +435,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(
Expand All @@ -411,6 +449,7 @@ impl MeterProviderBuilder {
meters: Default::default(),
shutdown_invoked: AtomicBool::new(false),
}),
component_name: format!("meter_provider/{instance_id}"),
};

otel_debug!(
Expand Down
Loading
Loading