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

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions Cargo.lock

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

2 changes: 2 additions & 0 deletions crates/temper-observe/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -26,3 +26,5 @@ sha2 = { workspace = true }

[dev-dependencies]
tokio = { workspace = true }
# The export tests ask the `log` bridge whether logging is enabled.
log = "0.4"
190 changes: 113 additions & 77 deletions crates/temper-observe/src/otel.rs
Original file line number Diff line number Diff line change
Expand Up @@ -18,36 +18,46 @@
//! | `RUST_LOG` | Log level filter (default: `info`) |
//! | `TEMPER_TRACE_QUEUE_SIZE` | Max buffered spans before drop (default: 2048, range: 128–32768) |
//! | `TEMPER_LOG_QUEUE_SIZE` | Max buffered log records before drop (default: 2048, range: 128–32768) |
//! | `OTEL_SERVICE_NAME` | Service name to export under (default: the name built into the binary) |
//! | `OTEL_RESOURCE_ATTRIBUTES` | Extra resource attributes as `key=value,...` (default: none); attributes the server computes itself win |
//! | `OTEL_TRACES_EXPORTER`, `OTEL_METRICS_EXPORTER`, `OTEL_LOGS_EXPORTER` | `none` switches that signal's export off (default: `otlp`) |
//! | `OTEL_TRACES_SAMPLER` | `always_on`, `always_off`, `traceidratio`, `parentbased_always_on`, `parentbased_always_off` or `parentbased_traceidratio` (default: `parentbased_always_on`) |
//! | `OTEL_TRACES_SAMPLER_ARG` | Ratio from 0 to 1 for the two ratio samplers (default: 1) |
//!
//! A value that is not supported never stops the server: it is logged once at
//! startup as a warning and the default applies.

use std::sync::OnceLock;
use std::time::Duration;

use opentelemetry::KeyValue;
use opentelemetry::trace::TracerProvider as _;
use opentelemetry_appender_tracing::layer::OpenTelemetryTracingBridge;
use opentelemetry_otlp::{LogExporter, MetricExporter, SpanExporter, WithExportConfig};
use opentelemetry_sdk::Resource;
use opentelemetry_sdk::logs::{
BatchConfigBuilder as LogBatchConfigBuilder, BatchLogProcessor, SdkLoggerProvider,
};
use opentelemetry_sdk::metrics::{PeriodicReader, SdkMeterProvider};
use opentelemetry_sdk::trace::{
BatchConfigBuilder as SpanBatchConfigBuilder, BatchSpanProcessor, Sampler, SdkTracerProvider,
BatchConfigBuilder as SpanBatchConfigBuilder, BatchSpanProcessor, SdkTracerProvider,
};
use tracing_subscriber::EnvFilter;
use tracing_subscriber::prelude::*;

mod config;
mod log_time;
mod sampler;
mod settings;

use config::{
parse_otlp_headers, read_non_empty_env, resolve_deployment_environment, resolve_otel_config,
resolve_service_version,
};
use log_time::EventTimeLogProcessor;
use sampler::{
DISPATCH_BACKGROUND_SAMPLE_RATE_DEFAULT, NameBasedSampler, TraceSamplerConfig,
WASM_AUXILIARY_SAMPLE_RATE_DEFAULT, record_trace_sampler_config,
};
use settings::{ComputedAttributes, ExportSettings};

const OTEL_EXPORTER_BUILD_RETRY_ATTEMPTS: usize = 3;
const OTEL_EXPORTER_RETRY_BASE_DELAY_MS: u64 = 250;
Expand Down Expand Up @@ -168,13 +178,14 @@ pub fn init_observability(service_name: &str) -> Option<OtelGuard> {
config.logfire_token.is_some(),
);

match init_tracing(&config.endpoint, service_name) {
let settings = ExportSettings::from_env();
match init_pipeline(&config.endpoint, service_name, &settings) {
Ok(guard) => {
tracing::info!(
endpoint = %config.endpoint,
endpoint_source = config.endpoint_source.as_str(),
logfire_auth = config.logfire_token.is_some(),
service_name,
service_name = settings.service_name(service_name),
"OTEL export pipeline active",
);
Some(guard)
Expand All @@ -196,6 +207,16 @@ pub fn init_tracing(
endpoint: &str,
service_name: &str,
) -> Result<OtelGuard, Box<dyn std::error::Error>> {
init_pipeline(endpoint, service_name, &ExportSettings::from_env())
}

fn init_pipeline(
endpoint: &str,
service_name: &str,
settings: &ExportSettings,
) -> Result<OtelGuard, Box<dyn std::error::Error>> {
let service_name = settings.service_name(service_name);

// Build auth headers (Logfire or custom).
let mut headers = read_non_empty_env("OTEL_EXPORTER_OTLP_HEADERS")
.map(|raw| parse_otlp_headers(&raw))
Expand Down Expand Up @@ -249,99 +270,104 @@ pub fn init_tracing(
}
}

let mut resource_attrs = vec![KeyValue::new("service.name", service_name.to_string())];
if let Some(environment) = resolve_deployment_environment() {
resource_attrs.push(KeyValue::new("deployment.environment.name", environment));
}
if let Some(version) = resolve_service_version() {
resource_attrs.push(KeyValue::new("service.version", version));
}
let environment = resolve_deployment_environment();
let version = resolve_service_version();
// ADR-0055: runtime-id enables Datadog Profiler ↔ APM trace stitching.
// Generated once at process start; regenerates only on restart.
// determinism-ok: observability-only identifier, not a simulation variable.
resource_attrs.push(KeyValue::new("runtime-id", runtime_id().to_string()));
let resource = Resource::builder_empty()
.with_attributes(resource_attrs)
.build();

// --- Traces ---
let span_exporter = build_with_retry("trace exporter", || {
SpanExporter::builder()
.with_http()
.with_timeout(Duration::from_secs(10))
.build()
})?;

let trace_queue = queue_size_from_env("TEMPER_TRACE_QUEUE_SIZE", TRACE_BATCH_MAX_QUEUE_SIZE);
let trace_batch_config = SpanBatchConfigBuilder::default()
.with_max_queue_size(trace_queue)
.with_max_export_batch_size(TRACE_BATCH_MAX_EXPORT_BATCH_SIZE)
.with_scheduled_delay(Duration::from_millis(TRACE_BATCH_SCHEDULE_DELAY_MS))
.build();
let runtime_id = runtime_id().to_string();
let computed = ComputedAttributes {
environment,
version,
runtime_id,
};
let resource = settings.resource(service_name, computed);

let trace_batch_processor = BatchSpanProcessor::builder(span_exporter)
.with_batch_config(trace_batch_config)
.build();
// A signal switched off keeps its provider and gets no exporter, so the
// rest of the process sees the same tracer, meter and logger either way.
let signals = settings.signals();

// --- Traces ---
// ADR-0052 hygiene: drop known-noisy span names at ingestion.
let trace_sampler_config = TraceSamplerConfig::from_env();
let sampler = NameBasedSampler {
inner: Sampler::ParentBased(Box::new(Sampler::AlwaysOn)),
inner: settings.sampler(),
config: trace_sampler_config.clone(),
};

let tracer_provider = SdkTracerProvider::builder()
.with_span_processor(trace_batch_processor)
let trace_queue = queue_size_from_env("TEMPER_TRACE_QUEUE_SIZE", TRACE_BATCH_MAX_QUEUE_SIZE);
let mut tracer_builder = SdkTracerProvider::builder()
.with_sampler(sampler.clone())
.with_resource(resource.clone())
.build();
.with_resource(resource.clone());
if signals.traces {
let span_exporter = build_with_retry("trace exporter", || {
SpanExporter::builder()
.with_http()
.with_timeout(Duration::from_secs(10))
.build()
})?;

let trace_batch_config = SpanBatchConfigBuilder::default()
.with_max_queue_size(trace_queue)
.with_max_export_batch_size(TRACE_BATCH_MAX_EXPORT_BATCH_SIZE)
.with_scheduled_delay(Duration::from_millis(TRACE_BATCH_SCHEDULE_DELAY_MS))
.build();

let trace_batch_processor = BatchSpanProcessor::builder(span_exporter)
.with_batch_config(trace_batch_config)
.build();
tracer_builder = tracer_builder.with_span_processor(trace_batch_processor);
}
let tracer_provider = tracer_builder.build();

opentelemetry::global::set_tracer_provider(tracer_provider.clone());

// --- Metrics ---
let metric_exporter = build_with_retry("metric exporter", || {
MetricExporter::builder()
.with_http()
.with_timeout(Duration::from_secs(10))
.build()
})?;

// Export every 30 s so metrics are visible quickly and canary gauges stay fresh.
let metric_reader = PeriodicReader::builder(metric_exporter)
.with_interval(Duration::from_secs(30))
.build();

let meter_provider = SdkMeterProvider::builder()
.with_reader(metric_reader)
.with_resource(resource.clone())
.build();
let mut meter_builder = SdkMeterProvider::builder().with_resource(resource.clone());
if signals.metrics {
let metric_exporter = build_with_retry("metric exporter", || {
MetricExporter::builder()
.with_http()
.with_timeout(Duration::from_secs(10))
.build()
})?;

// Export every 30 s so metrics are visible quickly and canary gauges stay fresh.
let metric_reader = PeriodicReader::builder(metric_exporter)
.with_interval(Duration::from_secs(30))
.build();
meter_builder = meter_builder.with_reader(metric_reader);
}
let meter_provider = meter_builder.build();

opentelemetry::global::set_meter_provider(meter_provider.clone());
record_trace_sampler_config(&trace_sampler_config);

// --- Logs ---
let log_exporter = build_with_retry("log exporter", || {
LogExporter::builder()
.with_http()
.with_timeout(Duration::from_secs(10))
.build()
})?;

let log_queue = queue_size_from_env("TEMPER_LOG_QUEUE_SIZE", LOG_BATCH_MAX_QUEUE_SIZE);
let log_batch_config = LogBatchConfigBuilder::default()
.with_max_queue_size(log_queue)
.with_max_export_batch_size(LOG_BATCH_MAX_EXPORT_BATCH_SIZE)
.with_scheduled_delay(Duration::from_millis(LOG_BATCH_SCHEDULE_DELAY_MS))
.build();

let log_batch_processor = BatchLogProcessor::builder(log_exporter)
.with_batch_config(log_batch_config)
.build();
let mut logger_builder = SdkLoggerProvider::builder().with_resource(resource);
if signals.logs {
let log_exporter = build_with_retry("log exporter", || {
LogExporter::builder()
.with_http()
.with_timeout(Duration::from_secs(10))
.build()
})?;

let logger_provider = SdkLoggerProvider::builder()
.with_log_processor(log_batch_processor)
.with_resource(resource)
.build();
let log_batch_config = LogBatchConfigBuilder::default()
.with_max_queue_size(log_queue)
.with_max_export_batch_size(LOG_BATCH_MAX_EXPORT_BATCH_SIZE)
.with_scheduled_delay(Duration::from_millis(LOG_BATCH_SCHEDULE_DELAY_MS))
.build();

let log_batch_processor = BatchLogProcessor::builder(log_exporter)
.with_batch_config(log_batch_config)
.build();
logger_builder = logger_builder
.with_log_processor(EventTimeLogProcessor)
.with_log_processor(log_batch_processor);
}
let logger_provider = logger_builder.build();

// --- Tracing subscriber ---
// Three layers:
Expand All @@ -365,7 +391,9 @@ pub fn init_tracing(
// Restrict the log bridge to WARN+ to avoid flooding Logfire's /v1/logs
// endpoint with high-volume info events. Traces already capture info-level
// spans via the otel_trace_layer, so no diagnostic value is lost.
let otel_log_layer = OpenTelemetryTracingBridge::new(&logger_provider);
let otel_log_layer = signals
.logs
.then(|| OpenTelemetryTracingBridge::new(&logger_provider));

// `.boxed()` unifies the pretty and JSON fmt layer types so a single
// subscriber chain compiles.
Expand All @@ -391,6 +419,13 @@ pub fn init_tracing(
std::io::Error::other(format!("failed to initialize tracing subscriber: {e}"))
})?;

for warning in settings.warnings() {
tracing::warn!("OTEL export setting: {warning}");
}
if let Some(sampler) = settings.chosen_sampler() {
tracing::info!(?sampler, "OTEL trace sampler set by OTEL_TRACES_SAMPLER");
}

tracing::info!(
endpoint,
service_name,
Expand All @@ -404,7 +439,8 @@ pub fn init_tracing(
.unwrap_or(DISPATCH_BACKGROUND_SAMPLE_RATE_DEFAULT),
log_queue,
log_batch = LOG_BATCH_MAX_EXPORT_BATCH_SIZE,
"OTEL initialised (traces + metrics + logs)"
"OTEL initialised ({})",
signals.label()
);

Ok(OtelGuard {
Expand Down
Loading
Loading