diff --git a/Cargo.lock b/Cargo.lock index 23fca4f37..f7c6aaf9f 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -7187,6 +7187,7 @@ name = "temper-observe" version = "0.1.0" dependencies = [ "chrono", + "log", "opentelemetry", "opentelemetry-appender-tracing", "opentelemetry-otlp", diff --git a/crates/temper-observe/Cargo.toml b/crates/temper-observe/Cargo.toml index 956e82641..959b8f40e 100644 --- a/crates/temper-observe/Cargo.toml +++ b/crates/temper-observe/Cargo.toml @@ -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" diff --git a/crates/temper-observe/src/otel.rs b/crates/temper-observe/src/otel.rs index bcef6adbc..8dbd10eb9 100644 --- a/crates/temper-observe/src/otel.rs +++ b/crates/temper-observe/src/otel.rs @@ -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; @@ -168,13 +178,14 @@ pub fn init_observability(service_name: &str) -> Option { 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) @@ -196,6 +207,16 @@ pub fn init_tracing( endpoint: &str, service_name: &str, ) -> Result> { + init_pipeline(endpoint, service_name, &ExportSettings::from_env()) +} + +fn init_pipeline( + endpoint: &str, + service_name: &str, + settings: &ExportSettings, +) -> Result> { + 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)) @@ -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: @@ -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. @@ -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, @@ -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 { diff --git a/crates/temper-observe/src/otel/log_time.rs b/crates/temper-observe/src/otel/log_time.rs new file mode 100644 index 000000000..9f251c09c --- /dev/null +++ b/crates/temper-observe/src/otel/log_time.rs @@ -0,0 +1,103 @@ +//! Event time for exported log records. + +use std::time::SystemTime; + +use opentelemetry::InstrumentationScope; +use opentelemetry::logs::LogRecord as _; +use opentelemetry_sdk::error::OTelSdkResult; +use opentelemetry_sdk::logs::{LogProcessor, SdkLogRecord}; + +/// Gives a log record an event time when it has none. +/// +/// The `tracing` bridge leaves the event time unset and the SDK fills in only +/// the observed time. A backend that dates records by their event time then +/// sees a record from 1970 and can drop it while still answering with +/// success. Records that already carry an event time are left untouched. +/// +/// Register it before the batch processor: processors run in registration +/// order, and the batch processor copies the record it is given. +#[derive(Debug)] +pub(super) struct EventTimeLogProcessor; + +impl LogProcessor for EventTimeLogProcessor { + fn emit(&self, record: &mut SdkLogRecord, _scope: &InstrumentationScope) { + if record.timestamp().is_none() { + // The SDK sets the observed time before it calls any processor. + // determinism-ok: telemetry timestamp, not a simulation variable. + let event_time = record.observed_timestamp().unwrap_or_else(SystemTime::now); + record.set_timestamp(event_time); + } + } + + fn force_flush(&self) -> OTelSdkResult { + Ok(()) + } + + fn shutdown(&self) -> OTelSdkResult { + Ok(()) + } +} + +#[cfg(test)] +mod tests { + use std::sync::{Arc, Mutex}; + use std::time::Duration; + + use opentelemetry::logs::{Logger as _, LoggerProvider as _}; + use opentelemetry_sdk::logs::SdkLoggerProvider; + + use super::*; + + type Times = (Option, Option); + + /// Records the event time and observed time of every record it is given. + #[derive(Debug, Default)] + struct RecordTimes(Arc>>); + + impl LogProcessor for RecordTimes { + fn emit(&self, record: &mut SdkLogRecord, _scope: &InstrumentationScope) { + self.0 + .lock() + .expect("times lock") + .push((record.timestamp(), record.observed_timestamp())); + } + + fn force_flush(&self) -> OTelSdkResult { + Ok(()) + } + + fn shutdown(&self) -> OTelSdkResult { + Ok(()) + } + } + + fn emit_with_event_time(event_time: Option) -> Times { + let seen = Arc::new(Mutex::new(Vec::new())); + let provider = SdkLoggerProvider::builder() + .with_log_processor(EventTimeLogProcessor) + .with_log_processor(RecordTimes(Arc::clone(&seen))) + .build(); + let logger = provider.logger("test"); + let mut record = logger.create_log_record(); + if let Some(event_time) = event_time { + record.set_timestamp(event_time); + } + logger.emit(record); + seen.lock().expect("times lock")[0] + } + + #[test] + fn record_without_event_time_gets_its_observed_time() { + let (event_time, observed_time) = emit_with_event_time(None); + assert!(event_time.is_some(), "event time must be set"); + assert_eq!(event_time, observed_time); + } + + #[test] + fn record_with_event_time_is_left_untouched() { + let original = SystemTime::UNIX_EPOCH + Duration::from_secs(1_700_000_000); + let (event_time, observed_time) = emit_with_event_time(Some(original)); + assert_eq!(event_time, Some(original)); + assert_ne!(observed_time, Some(original)); + } +} diff --git a/crates/temper-observe/src/otel/sampler.rs b/crates/temper-observe/src/otel/sampler.rs index 34488c5bc..2af8e47d4 100644 --- a/crates/temper-observe/src/otel/sampler.rs +++ b/crates/temper-observe/src/otel/sampler.rs @@ -196,11 +196,14 @@ fn trace_id_sample_at(trace_id: TraceId, rate_pct: u8) -> bool { /// Sampler that drops specific span names outright, reduced-samples specific /// prefixes, and delegates every other span decision to the wrapped inner -/// sampler (default: parent-based AlwaysOn). +/// sampler (default: parent-based AlwaysOn; `OTEL_TRACES_SAMPLER` chooses +/// another). /// -/// Implemented here rather than configured via `OTEL_TRACES_SAMPLER_ARG` so -/// the drop rules live in source (grep-able) and survive env-var rewrites by -/// tenant/deploy tooling. +/// The drop rules are implemented here rather than configured through the +/// environment so they live in source (grep-able) and survive env-var +/// rewrites by tenant/deploy tooling. `OTEL_TRACES_SAMPLER` and +/// `OTEL_TRACES_SAMPLER_ARG` choose only the inner sampler; the rules apply +/// around whichever one is chosen. #[derive(Debug, Clone)] pub(super) struct NameBasedSampler { pub(super) inner: Sampler, diff --git a/crates/temper-observe/src/otel/settings.rs b/crates/temper-observe/src/otel/settings.rs new file mode 100644 index 000000000..7ecda7955 --- /dev/null +++ b/crates/temper-observe/src/otel/settings.rs @@ -0,0 +1,275 @@ +//! Export settings read from the standard OpenTelemetry environment variables. +//! +//! Every setting is optional, and with none of them present the export is +//! what it was before they existed. A bad value never stops the server: it +//! becomes a warning that is reported once at startup, and the default +//! applies. + +use std::collections::BTreeMap; + +use opentelemetry::KeyValue; +use opentelemetry_sdk::Resource; +use opentelemetry_sdk::trace::Sampler; + +use super::config::read_non_empty_env; + +const SERVICE_NAME_ENV: &str = "OTEL_SERVICE_NAME"; +const RESOURCE_ATTRIBUTES_ENV: &str = "OTEL_RESOURCE_ATTRIBUTES"; +const TRACES_EXPORTER_ENV: &str = "OTEL_TRACES_EXPORTER"; +const METRICS_EXPORTER_ENV: &str = "OTEL_METRICS_EXPORTER"; +const LOGS_EXPORTER_ENV: &str = "OTEL_LOGS_EXPORTER"; +const TRACES_SAMPLER_ENV: &str = "OTEL_TRACES_SAMPLER"; +const TRACES_SAMPLER_ARG_ENV: &str = "OTEL_TRACES_SAMPLER_ARG"; +/// The sampler used when `OTEL_TRACES_SAMPLER` is unset. +const DEFAULT_SAMPLER: &str = "parentbased_always_on"; + +/// Resource attributes the server computes itself. They win over +/// `OTEL_RESOURCE_ATTRIBUTES`. +#[derive(Clone, Debug)] +pub(super) struct ComputedAttributes { + /// `deployment.environment.name`, when the environment is known. + pub(super) environment: Option, + /// `service.version`, when the version has a variable of its own. + pub(super) version: Option, + /// `runtime-id`, generated once per process. + pub(super) runtime_id: String, +} + +/// Which signals are exported. All of them unless the operator switches one +/// off. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub(super) struct Signals { + pub(super) traces: bool, + pub(super) metrics: bool, + pub(super) logs: bool, +} + +impl Default for Signals { + fn default() -> Self { + Self { + traces: true, + metrics: true, + logs: true, + } + } +} + +impl Signals { + /// The exported signals for the startup log line, for example + /// `traces + metrics + logs`. + pub(super) fn label(self) -> String { + let exported: Vec<&str> = [ + ("traces", self.traces), + ("metrics", self.metrics), + ("logs", self.logs), + ] + .into_iter() + .filter_map(|(signal, on)| on.then_some(signal)) + .collect(); + if exported.is_empty() { + "no signals".to_string() + } else { + exported.join(" + ") + } + } +} + +/// What the operator asked for through the standard variables. +#[derive(Clone, Debug, Default)] +pub(super) struct ExportSettings { + service_name: Option, + resource_attributes: BTreeMap, + signals: Signals, + /// The sampler `OTEL_TRACES_SAMPLER` names, when it is set and supported. + sampler: Option, + warnings: Vec, +} + +impl ExportSettings { + /// Read the settings from the process environment. Called once at + /// startup. + pub(super) fn from_env() -> Self { + Self::from_lookup(read_non_empty_env) + } + + /// Read the settings through `lookup`, which returns a variable's + /// trimmed value, or `None` when it is unset or empty. + pub(super) fn from_lookup(lookup: impl Fn(&str) -> Option) -> Self { + let mut settings = Self { + service_name: lookup(SERVICE_NAME_ENV), + ..Self::default() + }; + if let Some(raw) = lookup(RESOURCE_ATTRIBUTES_ENV) { + settings.read_resource_attributes(&raw); + } + settings.signals = Signals { + traces: settings.read_exporter(TRACES_EXPORTER_ENV, &lookup), + metrics: settings.read_exporter(METRICS_EXPORTER_ENV, &lookup), + logs: settings.read_exporter(LOGS_EXPORTER_ENV, &lookup), + }; + settings.sampler = settings.read_sampler(&lookup); + settings + } + + /// The sampler `OTEL_TRACES_SAMPLER` names. A name that is not supported + /// is reported and the default sampler is used. + fn read_sampler(&mut self, lookup: &impl Fn(&str) -> Option) -> Option { + let name = lookup(TRACES_SAMPLER_ENV)?; + let sampler = match name.to_ascii_lowercase().as_str() { + "always_on" => Sampler::AlwaysOn, + "always_off" => Sampler::AlwaysOff, + "parentbased_always_on" => parent_based(Sampler::AlwaysOn), + "parentbased_always_off" => parent_based(Sampler::AlwaysOff), + "traceidratio" => Sampler::TraceIdRatioBased(self.read_sampler_ratio(lookup)), + "parentbased_traceidratio" => { + parent_based(Sampler::TraceIdRatioBased(self.read_sampler_ratio(lookup))) + } + _ => { + self.warnings.push(format!( + "{TRACES_SAMPLER_ENV}={name} is not supported (expected always_on, \ + always_off, traceidratio, parentbased_always_on, parentbased_always_off \ + or parentbased_traceidratio); using {DEFAULT_SAMPLER}" + )); + return None; + } + }; + Some(sampler) + } + + /// The ratio for the two ratio samplers, from `OTEL_TRACES_SAMPLER_ARG`. + /// Unset means 1, as does a value that is not a number from 0 to 1, + /// which is also reported. + fn read_sampler_ratio(&mut self, lookup: &impl Fn(&str) -> Option) -> f64 { + let Some(raw) = lookup(TRACES_SAMPLER_ARG_ENV) else { + return 1.0; + }; + match raw.parse::() { + Ok(ratio) if (0.0..=1.0).contains(&ratio) => ratio, + _ => { + self.warnings.push(format!( + "{TRACES_SAMPLER_ARG_ENV}={raw} is not a ratio from 0 to 1; using 1" + )); + 1.0 + } + } + } + + /// Whether the signal behind an `OTEL_*_EXPORTER` variable is exported. + /// `none` switches it off; unset or `otlp` exports it. OTLP is the only + /// exporter there is, so any other value is reported and exported as + /// `otlp`. + fn read_exporter(&mut self, variable: &str, lookup: &impl Fn(&str) -> Option) -> bool { + let Some(value) = lookup(variable) else { + return true; + }; + if value.eq_ignore_ascii_case("none") { + return false; + } + if !value.eq_ignore_ascii_case("otlp") { + self.warnings.push(format!( + "{variable}={value} is not supported (expected otlp or none); exporting as otlp" + )); + } + true + } + + /// `OTEL_RESOURCE_ATTRIBUTES` is a comma-separated list of `key=value` + /// pairs. Keys and values are trimmed and otherwise taken as written; a + /// key given twice keeps its last value. + fn read_resource_attributes(&mut self, raw: &str) { + let mut ignored = 0usize; + for entry in raw.split(',').filter(|entry| !entry.trim().is_empty()) { + match entry.split_once('=') { + Some((key, value)) if !key.trim().is_empty() => { + self.resource_attributes + .insert(key.trim().to_string(), value.trim().to_string()); + } + _ => ignored += 1, + } + } + // The entries themselves are not printed: a mistyped variable can + // hold anything. + match ignored { + 0 => {} + 1 => self.warnings.push(format!( + "{RESOURCE_ATTRIBUTES_ENV} has 1 entry that is not key=value; it is ignored" + )), + _ => self.warnings.push(format!( + "{RESOURCE_ATTRIBUTES_ENV} has {ignored} entries that are not key=value; \ + they are ignored" + )), + } + } + + /// The service name to export under: `OTEL_SERVICE_NAME` when it is set, + /// otherwise the name built into the binary. + pub(super) fn service_name<'a>(&'a self, built_in: &'a str) -> &'a str { + self.service_name.as_deref().unwrap_or(built_in) + } + + /// The resource shared by traces, metrics and logs. + /// + /// `OTEL_RESOURCE_ATTRIBUTES` goes in first and what the server computes + /// itself goes in after it, so a deployment that sets none of the new + /// variables keeps the attributes it has today. A `service.name` in + /// `OTEL_RESOURCE_ATTRIBUTES` does not replace the built-in name; only + /// `OTEL_SERVICE_NAME` does. + pub(super) fn resource( + &self, + built_in_service_name: &str, + computed: ComputedAttributes, + ) -> Resource { + let mut attributes = self.resource_attributes.clone(); + attributes.insert( + "service.name".to_string(), + self.service_name(built_in_service_name).to_string(), + ); + if let Some(environment) = computed.environment { + attributes.insert("deployment.environment.name".to_string(), environment); + } + if let Some(version) = computed.version { + attributes.insert("service.version".to_string(), version); + } + attributes.insert("runtime-id".to_string(), computed.runtime_id); + + Resource::builder_empty() + .with_attributes( + attributes + .into_iter() + .map(|(key, value)| KeyValue::new(key, value)), + ) + .build() + } + + /// Which signals to build an exporter for. + pub(super) fn signals(&self) -> Signals { + self.signals + } + + /// The sampler the name-based filter delegates to: the one + /// `OTEL_TRACES_SAMPLER` names, or by default the caller's decision + /// with every root span sampled. + pub(super) fn sampler(&self) -> Sampler { + self.sampler + .clone() + .unwrap_or_else(|| parent_based(Sampler::AlwaysOn)) + } + + /// The sampler, when `OTEL_TRACES_SAMPLER` chose it. + pub(super) fn chosen_sampler(&self) -> Option<&Sampler> { + self.sampler.as_ref() + } + + /// Problems found in the settings, to report once at startup. + pub(super) fn warnings(&self) -> &[String] { + &self.warnings + } +} + +fn parent_based(root: Sampler) -> Sampler { + Sampler::ParentBased(Box::new(root)) +} + +#[cfg(test)] +#[path = "settings_test.rs"] +mod tests; diff --git a/crates/temper-observe/src/otel/settings_test.rs b/crates/temper-observe/src/otel/settings_test.rs new file mode 100644 index 000000000..11af99376 --- /dev/null +++ b/crates/temper-observe/src/otel/settings_test.rs @@ -0,0 +1,240 @@ +use opentelemetry::Key; +use opentelemetry::trace::{ + SamplingDecision, SpanContext, SpanId, SpanKind, TraceContextExt, TraceFlags, TraceId, + TraceState, +}; +use opentelemetry_sdk::trace::ShouldSample; + +use super::*; + +fn settings(variables: &[(&str, &str)]) -> ExportSettings { + ExportSettings::from_lookup(|name| { + variables + .iter() + .find(|(variable, _)| *variable == name) + .map(|(_, value)| value.to_string()) + }) +} + +fn computed() -> ComputedAttributes { + ComputedAttributes { + environment: None, + version: None, + runtime_id: "runtime-1".to_string(), + } +} + +fn attribute(resource: &Resource, key: &'static str) -> Option { + resource + .get(&Key::from_static_str(key)) + .map(|value| value.to_string()) +} + +#[test] +fn exporter_variables_switch_signals_off_one_by_one() { + assert_eq!(settings(&[]).signals(), Signals::default()); + assert_eq!(Signals::default().label(), "traces + metrics + logs"); + + let settings = settings(&[ + (TRACES_EXPORTER_ENV, "otlp"), + (METRICS_EXPORTER_ENV, "None"), + (LOGS_EXPORTER_ENV, "none"), + ]); + let expected = Signals { + traces: true, + metrics: false, + logs: false, + }; + assert_eq!(settings.signals(), expected); + assert_eq!(expected.label(), "traces"); + assert!(settings.warnings().is_empty()); +} + +#[test] +fn unsupported_exporter_is_reported_and_exported_as_otlp() { + let settings = settings(&[(METRICS_EXPORTER_ENV, "prometheus")]); + assert_eq!(settings.signals(), Signals::default()); + let [warning] = settings.warnings() else { + panic!("expected one warning, got {:?}", settings.warnings()); + }; + assert!(warning.contains("OTEL_METRICS_EXPORTER=prometheus is not supported")); +} + +#[test] +fn nothing_set_gives_the_built_in_resource_and_no_warnings() { + let settings = settings(&[]); + let resource = settings.resource("built-in", computed()); + assert_eq!(resource.len(), 2); + assert_eq!( + attribute(&resource, "service.name").as_deref(), + Some("built-in") + ); + assert_eq!( + attribute(&resource, "runtime-id").as_deref(), + Some("runtime-1") + ); + assert!(settings.warnings().is_empty()); +} + +#[test] +fn resource_attributes_are_trimmed_and_the_last_duplicate_wins() { + let settings = settings(&[(RESOURCE_ATTRIBUTES_ENV, " a = 1 ,b=2,a=3,empty=,")]); + let resource = settings.resource("built-in", computed()); + assert_eq!(attribute(&resource, "a").as_deref(), Some("3")); + assert_eq!(attribute(&resource, "b").as_deref(), Some("2")); + assert_eq!(attribute(&resource, "empty").as_deref(), Some("")); + assert!(settings.warnings().is_empty()); +} + +#[test] +fn entries_that_are_not_key_value_are_ignored_with_one_warning() { + let settings = settings(&[(RESOURCE_ATTRIBUTES_ENV, "a=1,token-without-key,=x,b=2")]); + let resource = settings.resource("built-in", computed()); + assert_eq!(attribute(&resource, "a").as_deref(), Some("1")); + assert_eq!(attribute(&resource, "b").as_deref(), Some("2")); + assert_eq!(resource.len(), 4); + let [warning] = settings.warnings() else { + panic!("expected one warning, got {:?}", settings.warnings()); + }; + assert!(warning.contains("OTEL_RESOURCE_ATTRIBUTES has 2 entries")); + assert!(!warning.contains("token-without-key")); +} + +#[test] +fn a_single_malformed_entry_is_reported_in_the_singular() { + let settings = settings(&[(RESOURCE_ATTRIBUTES_ENV, "a=1,oops")]); + assert_eq!( + settings.warnings(), + ["OTEL_RESOURCE_ATTRIBUTES has 1 entry that is not key=value; it is ignored"] + ); +} + +#[test] +fn computed_attributes_win_over_resource_attributes() { + let settings = settings(&[( + RESOURCE_ATTRIBUTES_ENV, + "service.name=x,deployment.environment.name=x,service.version=x,runtime-id=x", + )]); + let resource = settings.resource( + "built-in", + ComputedAttributes { + environment: Some("prod".to_string()), + version: Some("1.0".to_string()), + runtime_id: "runtime-1".to_string(), + }, + ); + assert_eq!( + attribute(&resource, "service.name").as_deref(), + Some("built-in") + ); + assert_eq!( + attribute(&resource, "deployment.environment.name").as_deref(), + Some("prod") + ); + assert_eq!( + attribute(&resource, "service.version").as_deref(), + Some("1.0") + ); + assert_eq!( + attribute(&resource, "runtime-id").as_deref(), + Some("runtime-1") + ); +} + +#[test] +fn service_name_variable_replaces_the_built_in_name() { + let settings = settings(&[ + (SERVICE_NAME_ENV, "orders-eu"), + (RESOURCE_ATTRIBUTES_ENV, "service.name=x"), + ]); + assert_eq!(settings.service_name("built-in"), "orders-eu"); + let resource = settings.resource("built-in", computed()); + assert_eq!( + attribute(&resource, "service.name").as_deref(), + Some("orders-eu") + ); +} + +/// Whether `sampler` samples a root span, a span whose caller sampled, and a +/// span whose caller did not sample. +fn decisions(sampler: &Sampler) -> [bool; 3] { + let trace_id = TraceId::from_hex("0af7651916cd43dd8448eb211c80319c").expect("trace id"); + let caller = |flags: TraceFlags| { + let span_id = SpanId::from_hex("b7ad6b7169203331").expect("span id"); + let context = SpanContext::new(trace_id, span_id, flags, true, TraceState::default()); + opentelemetry::Context::new().with_remote_span_context(context) + }; + let sampled = |parent: Option<&opentelemetry::Context>| { + let result = sampler.should_sample(parent, trace_id, "span", &SpanKind::Server, &[], &[]); + matches!(result.decision, SamplingDecision::RecordAndSample) + }; + [ + sampled(None), + sampled(Some(&caller(TraceFlags::SAMPLED))), + sampled(Some(&caller(TraceFlags::default()))), + ] +} + +#[test] +fn sampler_defaults_to_following_the_caller() { + let settings = settings(&[(TRACES_SAMPLER_ARG_ENV, "0")]); + assert!(settings.chosen_sampler().is_none()); + assert_eq!(decisions(&settings.sampler()), [true, true, false]); + assert!(settings.warnings().is_empty()); +} + +#[test] +fn every_supported_sampler_is_selected_by_name() { + let cases = [ + ("always_on", "", [true, true, true]), + ("ALWAYS_ON", "", [true, true, true]), + ("always_off", "", [false, false, false]), + ("parentbased_always_on", "", [true, true, false]), + ("parentbased_always_off", "", [false, true, false]), + ("traceidratio", "", [true, true, true]), + ("traceidratio", "1", [true, true, true]), + ("traceidratio", "0", [false, false, false]), + ("parentbased_traceidratio", "1.0", [true, true, false]), + ("parentbased_traceidratio", "0.0", [false, true, false]), + ]; + for (name, arg, expected) in cases { + let mut variables = vec![(TRACES_SAMPLER_ENV, name)]; + if !arg.is_empty() { + variables.push((TRACES_SAMPLER_ARG_ENV, arg)); + } + let settings = settings(&variables); + assert!(settings.chosen_sampler().is_some(), "{name}"); + assert_eq!(decisions(&settings.sampler()), expected, "{name} {arg}"); + assert!(settings.warnings().is_empty(), "{name} {arg}"); + } +} + +#[test] +fn unsupported_sampler_is_reported_and_the_default_is_used() { + let settings = settings(&[(TRACES_SAMPLER_ENV, "jaeger_remote")]); + assert!(settings.chosen_sampler().is_none()); + assert_eq!(decisions(&settings.sampler()), [true, true, false]); + let [warning] = settings.warnings() else { + panic!("expected one warning, got {:?}", settings.warnings()); + }; + assert!(warning.contains("OTEL_TRACES_SAMPLER=jaeger_remote is not supported")); + assert!(warning.contains("using parentbased_always_on")); +} + +#[test] +fn ratio_that_is_not_between_0_and_1_is_reported_and_1_is_used() { + for arg in ["half", "1.5", "-0.1", "NaN"] { + let settings = settings(&[ + (TRACES_SAMPLER_ENV, "traceidratio"), + (TRACES_SAMPLER_ARG_ENV, arg), + ]); + assert_eq!(decisions(&settings.sampler()), [true, true, true], "{arg}"); + let [warning] = settings.warnings() else { + panic!("expected one warning, got {:?}", settings.warnings()); + }; + assert!( + warning.contains(&format!("OTEL_TRACES_SAMPLER_ARG={arg} is not a ratio")), + "{warning}" + ); + } +} diff --git a/crates/temper-observe/src/otel/tests.rs b/crates/temper-observe/src/otel/tests.rs index e7edcbfd4..2c8cca197 100644 --- a/crates/temper-observe/src/otel/tests.rs +++ b/crates/temper-observe/src/otel/tests.rs @@ -3,7 +3,7 @@ use super::sampler::{DISPATCH_BACKGROUND_PREFIXES, WASM_AUXILIARY_PREFIXES}; use super::*; use opentelemetry::trace::{SamplingDecision, SpanKind, TraceId}; -use opentelemetry_sdk::trace::ShouldSample; +use opentelemetry_sdk::trace::{Sampler, ShouldSample}; use std::sync::Mutex; static ENV_LOCK: Mutex<()> = Mutex::new(()); diff --git a/crates/temper-observe/tests/otel_export/bad_values.rs b/crates/temper-observe/tests/otel_export/bad_values.rs new file mode 100644 index 000000000..ab8e9ee75 --- /dev/null +++ b/crates/temper-observe/tests/otel_export/bad_values.rs @@ -0,0 +1,137 @@ +//! Values that are not supported, empty values and credentials. + +use crate::*; + +/// The child started, reported `expected_warning` exactly once, on its own +/// log output and as an exported log record, and exported all three signals. +fn assert_one_warning(run: &Run, expected_warning: &str) { + assert_eq!(warnings(run), [expected_warning]); + let exported: Vec = run + .logs() + .into_iter() + .filter(|record| record.severity_text == "WARN" && record.scope == "temper_observe::otel") + .map(|record| record.body) + .collect(); + assert_eq!(exported, [expected_warning]); + for path in SIGNAL_PATHS { + assert!(run.requests_to(path) > 0, "no export reached {path}"); + } +} + +/// An exporter that is not supported is reported once and the signal is +/// exported as by default. +#[test] +fn unsupported_exporter_is_reported_once_and_exported() { + for variable in [ + "OTEL_TRACES_EXPORTER", + "OTEL_METRICS_EXPORTER", + "OTEL_LOGS_EXPORTER", + ] { + let run = child::run(&[(variable, "console")]); + assert_one_warning( + &run, + &format!( + "OTEL export setting: {variable}=console is not supported \ + (expected otlp or none); exporting as otlp" + ), + ); + assert_eq!(span_names(&run), SPANS_FOLLOWING_THE_CALLER, "{variable}"); + } +} + +/// A sampler that is not supported is reported once and the default sampler +/// is used. +#[test] +fn unsupported_sampler_is_reported_once_and_the_default_is_used() { + let run = child::run(&[("OTEL_TRACES_SAMPLER", "jaeger_remote")]); + assert_one_warning( + &run, + "OTEL export setting: OTEL_TRACES_SAMPLER=jaeger_remote is not supported \ + (expected always_on, always_off, traceidratio, parentbased_always_on, \ + parentbased_always_off or parentbased_traceidratio); using parentbased_always_on", + ); + assert_eq!(span_names(&run), SPANS_FOLLOWING_THE_CALLER); +} + +/// A sampler ratio that is not a number from 0 to 1 is reported once and the +/// ratio's default, 1, is used. +#[test] +fn bad_sampler_ratio_is_reported_once_and_one_is_used() { + let run = child::run(&[ + ("OTEL_TRACES_SAMPLER", "traceidratio"), + ("OTEL_TRACES_SAMPLER_ARG", "half"), + ]); + assert_one_warning( + &run, + "OTEL export setting: OTEL_TRACES_SAMPLER_ARG=half is not a ratio from 0 to 1; using 1", + ); + assert_eq!(span_names(&run), EVERY_SPAN); +} + +/// Resource attribute entries that are not `key=value` are reported once, +/// without being printed, and the resource is the default one. +#[test] +fn malformed_resource_attributes_are_reported_once_and_ignored() { + let run = child::run(&[("OTEL_RESOURCE_ATTRIBUTES", "no-equals-sign,=no-key")]); + assert_one_warning( + &run, + "OTEL export setting: OTEL_RESOURCE_ATTRIBUTES has 2 entries that are not key=value; \ + they are ignored", + ); + for (signal, resource) in resources_of_all_signals(&run) { + let keys: Vec<&str> = resource.keys().map(String::as_str).collect(); + assert_eq!(keys, ["runtime-id", "service.name"], "resource of {signal}"); + } +} + +/// An empty value is the same as an unset variable: default behaviour and +/// nothing reported. +#[test] +fn empty_values_are_the_same_as_unset() { + let run = child::run(&[ + ("OTEL_SERVICE_NAME", ""), + ("OTEL_RESOURCE_ATTRIBUTES", " "), + ("OTEL_TRACES_EXPORTER", ""), + ("OTEL_METRICS_EXPORTER", ""), + ("OTEL_LOGS_EXPORTER", ""), + ("OTEL_TRACES_SAMPLER", ""), + ("OTEL_TRACES_SAMPLER_ARG", ""), + ]); + assert_matches_baseline(&run, "the export with every setting empty"); +} + +/// The credential in the OTLP headers reaches the backend on every signal +/// and is never printed or exported, whatever else is reported at startup. +#[test] +fn credentials_are_sent_and_never_logged() { + const SECRET: &str = "s3cr3t-credential"; + let run = child::run(&[ + ( + "OTEL_EXPORTER_OTLP_HEADERS", + "authorization=Bearer s3cr3t-credential", + ), + ("OTEL_RESOURCE_ATTRIBUTES", "team=storage,s3cr3t-credential"), + ("OTEL_SERVICE_NAME", "orders-eu"), + ("OTEL_TRACES_EXPORTER", "console"), + ("OTEL_TRACES_SAMPLER", "traceidratio"), + ("OTEL_TRACES_SAMPLER_ARG", "half"), + ]); + for path in SIGNAL_PATHS { + assert!(run.requests_to(path) > 0, "no export reached {path}"); + } + for request in &run.received { + assert_eq!( + request.header("authorization"), + "Bearer s3cr3t-credential", + "authorization header on {}", + request.path + ); + } + assert_eq!(warnings(&run).len(), 3, "{:?}", warnings(&run)); + assert!(!run.stdout.contains(SECRET), "credential on stdout"); + assert!(!run.stderr.contains(SECRET), "credential on stderr"); + assert!( + !render::render(&run).contains(SECRET), + "credential in the exported telemetry" + ); +} diff --git a/crates/temper-observe/tests/otel_export/child.rs b/crates/temper-observe/tests/otel_export/child.rs new file mode 100644 index 000000000..8a8359b9e --- /dev/null +++ b/crates/temper-observe/tests/otel_export/child.rs @@ -0,0 +1,250 @@ +//! Runs the telemetry setup in a child process and records what it exports. +//! +//! The setup installs process-wide state (the global tracer and meter +//! providers and the `tracing` subscriber) and reads its settings from the +//! environment, so each scenario gets a process of its own: this test binary, +//! started again with only [`entry`] selected. + +use std::process::Command; + +use opentelemetry::trace::{SpanContext, SpanId, TraceContextExt, TraceFlags, TraceId, TraceState}; +use tracing_opentelemetry::OpenTelemetrySpanExt; + +use crate::listener::{Listener, Received}; +use crate::wire::{self, Attributes, LogRecord, Metric, Span}; + +const SCENARIO_ENV: &str = "OTEL_EXPORT_TEST_SCENARIO"; +const WORKLOAD: &str = "workload"; +const LOG_PROBE_THEN_WORKLOAD: &str = "log-probe-then-workload"; + +/// Service name `temper serve` passes to the telemetry setup. +pub const BUILT_IN_SERVICE_NAME: &str = "temper-platform"; + +/// Trace ID of the caller whose request arrives marked "not sampled". +pub const UNSAMPLED_CALLER_TRACE_ID: &str = "0af7651916cd43dd8448eb211c80319c"; +/// Trace ID of the caller whose request arrives marked "sampled". +pub const SAMPLED_CALLER_TRACE_ID: &str = "4bf92f3577b34da6a3ce929d0e0e4736"; +/// The reduced-rate rules keep a trace when the low 64 bits of its ID, modulo +/// 100, are below the rule's rate. 100 % 100 = 0 is kept at every rate above 0. +pub const REDUCED_RATE_KEPT_TRACE_ID: &str = "00000000000000000000000000000064"; +/// 99 % 100 = 99 is dropped at every rate below 100. +pub const REDUCED_RATE_DROPPED_TRACE_ID: &str = "00000000000000000000000000000063"; + +const CALLER_SPAN_ID: &str = "b7ad6b7169203331"; + +/// What one child process exported and printed. +pub struct Run { + pub endpoint: String, + pub received: Vec, + /// The `tracing` log lines, one JSON object per line, among the test + /// harness's own output. + pub stdout: String, + /// What the setup prints before the `tracing` subscriber exists. + pub stderr: String, +} + +impl Run { + pub fn spans(&self) -> Vec { + self.bodies("/v1/traces").flat_map(wire::spans).collect() + } + + pub fn logs(&self) -> Vec { + self.bodies("/v1/logs").flat_map(wire::logs).collect() + } + + /// The metrics of the last export. Every export carries all metrics + /// again, so the earlier ones would only repeat them. + pub fn metrics(&self) -> Vec { + self.bodies("/v1/metrics") + .last() + .map(wire::metrics) + .unwrap_or_default() + } + + /// The resource each signal was exported under: one entry per signal + /// that reached the listener. A signal sent under more than one resource + /// is a failure. + pub fn resources(&self) -> Vec<(&'static str, Attributes)> { + let traces = self.spans().into_iter().map(|span| span.resource); + let metrics = self.metrics().into_iter().map(|metric| metric.resource); + let logs = self.logs().into_iter().map(|record| record.resource); + let mut out = Vec::new(); + for (signal, mut resources) in [ + ("traces", traces.collect::>()), + ("metrics", metrics.collect()), + ("logs", logs.collect()), + ] { + resources.dedup(); + assert!( + resources.len() <= 1, + "{signal} were exported under more than one resource: {resources:?}" + ); + out.extend(resources.into_iter().map(|resource| (signal, resource))); + } + out + } + + /// Number of requests the listener received on `path`. + pub fn requests_to(&self, path: &str) -> usize { + self.bodies(path).count() + } + + /// The log lines the child printed, parsed, in order. + pub fn log_lines(&self) -> Vec> { + self.stdout + .lines() + // The test harness may print the test's name in front of the + // first line. + .filter_map(|line| line.find('{').map(|start| &line[start..])) + .filter_map(|line| serde_json::from_str(line).ok()) + .collect() + } + + fn bodies<'a>(&'a self, path: &'a str) -> impl Iterator { + self.received + .iter() + .filter(move |request| request.path == path) + .map(|request| request.body.as_slice()) + } +} + +/// Start the telemetry setup in a child process with `env` added to a clean +/// environment, run the workload, and collect what reached the listener. +pub fn run(env: &[(&str, &str)]) -> Run { + run_scenario(WORKLOAD, env) +} + +/// Like [`run`], but the child asks the `log` bridge whether logging is +/// enabled, and logs nothing, before it runs the workload. +pub fn run_after_log_probe(env: &[(&str, &str)]) -> Run { + run_scenario(LOG_PROBE_THEN_WORKLOAD, env) +} + +fn run_scenario(scenario: &str, env: &[(&str, &str)]) -> Run { + let listener = Listener::start(); + let mut command = Command::new(std::env::current_exe().expect("path of the test binary")); + command.args(["--exact", "child::entry", "--nocapture"]); + for (name, _) in std::env::vars_os() { + if is_telemetry_variable(&name.to_string_lossy()) { + command.env_remove(name); + } + } + command + .env(SCENARIO_ENV, scenario) + .env("OTLP_ENDPOINT", listener.endpoint()); + for (name, value) in env { + command.env(name, value); + } + + let output = command.output().expect("start the child process"); + let stdout = String::from_utf8_lossy(&output.stdout).into_owned(); + let stderr = String::from_utf8_lossy(&output.stderr).into_owned(); + assert!( + output.status.success(), + "child process for scenario {scenario} failed:\n{stdout}\n{stderr}" + ); + Run { + endpoint: listener.endpoint(), + received: listener.received(), + stdout, + stderr, + } +} + +/// Variables the telemetry setup reads; none may leak in from the machine +/// running the tests. +fn is_telemetry_variable(name: &str) -> bool { + ["OTEL_", "OTLP_", "DD_", "LOGFIRE_", "TEMPER_"] + .iter() + .any(|prefix| name.starts_with(prefix)) + || name == "RUST_LOG" + || name == "FOREGROUND_LOGS" +} + +/// The child side. A no-op unless a parent test started this process. +#[test] +fn entry() { + let Ok(scenario) = std::env::var(SCENARIO_ENV) else { + return; + }; + let guard = temper_observe::otel::init_observability(BUILT_IN_SERVICE_NAME) + .expect("the OTEL pipeline must start when an endpoint is configured"); + if scenario == LOG_PROBE_THEN_WORKLOAD { + probe_log_enabled(); + } + emit_workload(); + guard.shutdown(); +} + +/// Ask the `log` bridge whether logging is enabled without logging anything: +/// once for a level and target that are enabled, once for a target the +/// default filter turns down, and once for a level that is off. +fn probe_log_enabled() { + assert!(log::log_enabled!(log::Level::Info)); + assert!(!log::log_enabled!(target: "hyper", log::Level::Info)); + assert!(!log::log_enabled!(log::Level::Trace)); +} + +/// A span named `$name` whose remote parent is the given caller. +macro_rules! span_from_caller { + ($name:literal, $trace_id:expr, $flags:expr) => {{ + let span = tracing::info_span!($name); + span.set_parent(caller_context($trace_id, $flags)); + span.in_scope(|| {}); + }}; +} + +/// The same telemetry in every scenario: spans, log records and a metric. +fn emit_workload() { + tracing::info_span!("test.request").in_scope(|| { + tracing::info!(detail = "inside", "test log inside a span"); + tracing::info_span!("test.child").in_scope(|| {}); + }); + tracing::warn!("test log outside a span"); + + span_from_caller!( + "test.unsampled_caller", + UNSAMPLED_CALLER_TRACE_ID, + TraceFlags::default() + ); + span_from_caller!( + "test.sampled_caller", + SAMPLED_CALLER_TRACE_ID, + TraceFlags::SAMPLED + ); + + // Names the name-based filter drops outright. + tracing::info_span!("clock_time_get").in_scope(|| {}); + span_from_caller!( + "turso.configured_connection", + SAMPLED_CALLER_TRACE_ID, + TraceFlags::SAMPLED + ); + // A name the name-based filter keeps at a reduced rate, by trace ID. + span_from_caller!( + "wasm:workspace_fs.read", + REDUCED_RATE_KEPT_TRACE_ID, + TraceFlags::SAMPLED + ); + span_from_caller!( + "wasm:workspace_fs.read", + REDUCED_RATE_DROPPED_TRACE_ID, + TraceFlags::SAMPLED + ); + + opentelemetry::global::meter("otel_export_test") + .u64_counter("otel_export_test_counter") + .build() + .add(1, &[]); +} + +fn caller_context(trace_id: &str, flags: TraceFlags) -> opentelemetry::Context { + let caller = SpanContext::new( + TraceId::from_hex(trace_id).expect("valid trace id"), + SpanId::from_hex(CALLER_SPAN_ID).expect("valid span id"), + flags, + true, + TraceState::default(), + ); + opentelemetry::Context::new().with_remote_span_context(caller) +} diff --git a/crates/temper-observe/tests/otel_export/default_export.baseline.txt b/crates/temper-observe/tests/otel_export/default_export.baseline.txt new file mode 100644 index 000000000..fe4fe7037 --- /dev/null +++ b/crates/temper-observe/tests/otel_export/default_export.baseline.txt @@ -0,0 +1,101 @@ +== stderr == +OTEL export configured: endpoint=http://127.0.0.1: source=OTLP_ENDPOINT logfire_auth=false +== log lines == +{"fields":{"endpoint":"http://127.0.0.1:","log_batch":512,"log_queue":2048,"message":"OTEL initialised (traces + metrics + logs)","service_name":"temper-platform","trace_batch":512,"trace_dispatch_background_sample_pct":25,"trace_queue":2048,"trace_wasm_auxiliary_sample_pct":5},"level":"INFO","target":"temper_observe::otel"} +{"fields":{"endpoint":"http://127.0.0.1:","endpoint_source":"OTLP_ENDPOINT","logfire_auth":false,"message":"OTEL export pipeline active","service_name":"temper-platform"},"level":"INFO","target":"temper_observe::otel"} +{"fields":{"detail":"inside","message":"test log inside a span"},"level":"INFO","span":{"name":"test.request"},"target":"otel_export::child"} +{"fields":{"message":"test log outside a span"},"level":"WARN","target":"otel_export::child"} +== requests == +/v1/logs application/x-protobuf +/v1/metrics application/x-protobuf +/v1/traces application/x-protobuf +== traces == +span "test.child" scope=temper trace=trace#1 id=span#1 parent=span#2 kind=1 flags=1 status=0 events=[] + resource runtime-id = "" + resource service.name = "temper-platform" + attribute busy_ns = "" + attribute code.filepath = "crates/temper-observe/tests/otel_export/child.rs" + attribute code.lineno = "" + attribute code.namespace = "otel_export::child" + attribute idle_ns = "" + attribute thread.id = "" + attribute thread.name = "" +span "test.request" scope=temper trace=trace#1 id=span#2 parent=[] kind=1 flags=1 status=0 events=["test log inside a span"] + resource runtime-id = "" + resource service.name = "temper-platform" + attribute busy_ns = "" + attribute code.filepath = "crates/temper-observe/tests/otel_export/child.rs" + attribute code.lineno = "" + attribute code.namespace = "otel_export::child" + attribute idle_ns = "" + attribute thread.id = "" + attribute thread.name = "" +span "test.sampled_caller" scope=temper trace=[4bf92f3577b34da6a3ce929d0e0e4736] id=span#3 parent=span#4 kind=1 flags=1 status=0 events=[] + resource runtime-id = "" + resource service.name = "temper-platform" + attribute busy_ns = "" + attribute code.filepath = "crates/temper-observe/tests/otel_export/child.rs" + attribute code.lineno = "" + attribute code.namespace = "otel_export::child" + attribute idle_ns = "" + attribute thread.id = "" + attribute thread.name = "" +span "wasm:workspace_fs.read" scope=temper trace=[00000000000000000000000000000064] id=span#5 parent=span#4 kind=1 flags=1 status=0 events=[] + resource runtime-id = "" + resource service.name = "temper-platform" + attribute busy_ns = "" + attribute code.filepath = "crates/temper-observe/tests/otel_export/child.rs" + attribute code.lineno = "" + attribute code.namespace = "otel_export::child" + attribute idle_ns = "" + attribute thread.id = "" + attribute thread.name = "" +== logs == +log "test log inside a span" scope=otel_export::child severity=INFO(9) time= observed= trace=[] span=[] + resource runtime-id = "" + resource service.name = "temper-platform" + attribute detail = "inside" +log "test log outside a span" scope=otel_export::child severity=WARN(13) time= observed= trace=[] span=[] + resource runtime-id = "" + resource service.name = "temper-platform" +log "OTEL initialised (traces + metrics + logs)" scope=temper_observe::otel severity=INFO(9) time= observed= trace=[] span=[] + resource runtime-id = "" + resource service.name = "temper-platform" + attribute endpoint = "http://127.0.0.1:" + attribute log_batch = "512" + attribute log_queue = "2048" + attribute service_name = "temper-platform" + attribute trace_batch = "512" + attribute trace_dispatch_background_sample_pct = "25" + attribute trace_queue = "2048" + attribute trace_wasm_auxiliary_sample_pct = "5" +log "OTEL export pipeline active" scope=temper_observe::otel severity=INFO(9) time= observed= trace=[] span=[] + resource runtime-id = "" + resource service.name = "temper-platform" + attribute endpoint = "http://127.0.0.1:" + attribute endpoint_source = "OTLP_ENDPOINT" + attribute logfire_auth = "false" + attribute service_name = "temper-platform" +== metrics == +metric "otel_export_test_counter" scope=otel_export_test kind=sum unit="" description="" + resource runtime-id = "" + resource service.name = "temper-platform" + point {} = 1 +metric "temper_trace_sampler_configured_rules" scope=temper.observe kind=gauge unit="" description="Configured trace sampler rules by rule kind. Values are counts, not rates." + resource runtime-id = "" + resource service.name = "temper-platform" + point {"rule_kind": "drop_exact"} = 2 + point {"rule_kind": "reduced_prefix"} = 3 +metric "temper_trace_sampler_decisions_total" scope=temper.observe kind=sum unit="" description="Trace sampler decisions by bounded rule name and final sampling decision." + resource runtime-id = "" + resource service.name = "temper-platform" + point {"decision": "drop", "rule": "delegate"} = 1 + point {"decision": "drop", "rule": "drop_exact"} = 2 + point {"decision": "record_and_sample", "rule": "delegate"} = 3 + point {"decision": "record_and_sample", "rule": "wasm_auxiliary"} = 1 + point {"decision": "reduced_drop", "rule": "wasm_auxiliary"} = 1 +metric "temper_trace_sampler_reduced_sample_rate_pct" scope=temper.observe kind=gauge unit="%" description="Effective trace sampler keep rate for reduced-prefix rules." + resource runtime-id = "" + resource service.name = "temper-platform" + point {"rule": "dispatch_background"} = 25 + point {"rule": "wasm_auxiliary"} = 5 diff --git a/crates/temper-observe/tests/otel_export/identity.rs b/crates/temper-observe/tests/otel_export/identity.rs new file mode 100644 index 000000000..16301b19e --- /dev/null +++ b/crates/temper-observe/tests/otel_export/identity.rs @@ -0,0 +1,99 @@ +//! Service name and resource attributes from the standard variables. + +use crate::*; + +/// `OTEL_SERVICE_NAME` replaces the built-in service name on every signal. +#[test] +fn service_name_follows_otel_service_name() { + let run = child::run(&[("OTEL_SERVICE_NAME", "orders-eu")]); + for (signal, resource) in resources_of_all_signals(&run) { + assert_eq!( + attribute(&resource, "service.name"), + Some("orders-eu"), + "service.name on {signal}" + ); + } +} + +/// Without `OTEL_SERVICE_NAME` the service name is the built-in one. +#[test] +fn service_name_is_built_in_when_unset() { + let run = child::run(&[]); + for (signal, resource) in resources_of_all_signals(&run) { + assert_eq!( + attribute(&resource, "service.name"), + Some(BUILT_IN_SERVICE_NAME), + "service.name on {signal}" + ); + } +} + +/// `OTEL_RESOURCE_ATTRIBUTES` is merged into the resource of every signal, +/// next to what the server computes itself. +#[test] +fn resource_attributes_are_merged_on_every_signal() { + let run = child::run(&[("OTEL_RESOURCE_ATTRIBUTES", "a=1,b=2")]); + for (signal, resource) in resources_of_all_signals(&run) { + assert_eq!(attribute(&resource, "a"), Some("1"), "a on {signal}"); + assert_eq!(attribute(&resource, "b"), Some("2"), "b on {signal}"); + assert_eq!( + attribute(&resource, "service.name"), + Some(BUILT_IN_SERVICE_NAME), + "service.name on {signal}" + ); + let runtime_id = attribute(&resource, "runtime-id").unwrap_or_default(); + assert!(!runtime_id.is_empty(), "runtime-id missing on {signal}"); + } +} + +/// What the server computes itself wins over `OTEL_RESOURCE_ATTRIBUTES`: +/// the runtime id always, the environment and the version when their own +/// variables are set, and the service name unless `OTEL_SERVICE_NAME` is set. +#[test] +fn computed_resource_attributes_keep_precedence() { + let run = child::run(&[ + ( + "OTEL_RESOURCE_ATTRIBUTES", + "runtime-id=from-attributes,deployment.environment.name=from-attributes,\ + service.version=from-attributes,service.name=from-attributes,team=storage", + ), + ("DD_ENV", "from-env-variable"), + ("DD_VERSION", "from-version-variable"), + ]); + for (signal, resource) in resources_of_all_signals(&run) { + let expect = |key: &str, value: &str| { + assert_eq!(attribute(&resource, key), Some(value), "{key} on {signal}"); + }; + expect("deployment.environment.name", "from-env-variable"); + expect("service.version", "from-version-variable"); + expect("service.name", BUILT_IN_SERVICE_NAME); + expect("team", "storage"); + assert_ne!( + attribute(&resource, "runtime-id"), + Some("from-attributes"), + "runtime-id on {signal}" + ); + } +} + +/// Without their own variables, the environment and the version come from +/// `OTEL_RESOURCE_ATTRIBUTES`, and `OTEL_SERVICE_NAME` wins over a +/// `service.name` given there. +#[test] +fn resource_attributes_supply_what_has_no_variable_of_its_own() { + let run = child::run(&[ + ( + "OTEL_RESOURCE_ATTRIBUTES", + "deployment.environment.name=staging, service.version = 1.2.3 ,service.name=ignored", + ), + ("OTEL_SERVICE_NAME", "orders-eu"), + ]); + for (signal, resource) in resources_of_all_signals(&run) { + let expect = |key: &str, value: &str| { + assert_eq!(attribute(&resource, key), Some(value), "{key} on {signal}"); + }; + expect("deployment.environment.name", "staging"); + expect("service.version", "1.2.3"); + expect("service.name", "orders-eu"); + } +} diff --git a/crates/temper-observe/tests/otel_export/listener.rs b/crates/temper-observe/tests/otel_export/listener.rs new file mode 100644 index 000000000..6b3f2898a --- /dev/null +++ b/crates/temper-observe/tests/otel_export/listener.rs @@ -0,0 +1,103 @@ +//! A local OTLP/HTTP listener that records every request it receives. + +use std::io::{BufRead, BufReader, Read, Write}; +use std::net::{TcpListener, TcpStream}; +use std::sync::{Arc, Mutex}; + +/// One HTTP request as the listener received it. +#[derive(Clone, Debug)] +pub struct Received { + pub path: String, + /// Header names are lower-cased. + pub headers: Vec<(String, String)>, + pub body: Vec, +} + +impl Received { + /// The value of a header, or the empty string when it is absent. + pub fn header(&self, name: &str) -> &str { + self.headers + .iter() + .find(|(header, _)| header == name) + .map(|(_, value)| value.as_str()) + .unwrap_or_default() + } +} + +/// Accepts OTLP/HTTP exports on a local port and answers each with success. +pub struct Listener { + port: u16, + received: Arc>>, +} + +impl Listener { + pub fn start() -> Self { + let socket = TcpListener::bind(("127.0.0.1", 0)).expect("bind the OTLP listener"); + let port = socket.local_addr().expect("listener address").port(); + let received = Arc::new(Mutex::new(Vec::new())); + let sink = Arc::clone(&received); + std::thread::spawn(move || { + for stream in socket.incoming().flatten() { + let sink = Arc::clone(&sink); + std::thread::spawn(move || serve_connection(stream, &sink)); + } + }); + Self { port, received } + } + + pub fn endpoint(&self) -> String { + format!("http://127.0.0.1:{}", self.port) + } + + /// Every request received so far, in arrival order. + pub fn received(&self) -> Vec { + self.received.lock().expect("listener lock").clone() + } +} + +fn serve_connection(stream: TcpStream, sink: &Mutex>) { + let mut reader = BufReader::new(stream); + // A request is recorded before it is answered, so once the exporter has + // its response the request is visible to the test. + while let Some(request) = read_request(&mut reader) { + sink.lock().expect("listener lock").push(request); + let response = + b"HTTP/1.1 200 OK\r\ncontent-type: application/x-protobuf\r\ncontent-length: 0\r\n\r\n"; + if reader.get_mut().write_all(response).is_err() { + return; + } + } +} + +fn read_request(reader: &mut BufReader) -> Option { + let mut request_line = String::new(); + if reader.read_line(&mut request_line).ok()? == 0 { + return None; + } + let path = request_line.split_whitespace().nth(1)?.to_string(); + + let mut headers = Vec::new(); + loop { + let mut line = String::new(); + reader.read_line(&mut line).ok()?; + let line = line.trim_end(); + if line.is_empty() { + break; + } + let (name, value) = line.split_once(':')?; + headers.push((name.trim().to_ascii_lowercase(), value.trim().to_string())); + } + + let mut request = Received { + path, + headers, + body: Vec::new(), + }; + let length = request + .header("content-length") + .parse::() + .expect("the exporter sends a content-length"); + request.body.resize(length, 0); + reader.read_exact(&mut request.body).ok()?; + Some(request) +} diff --git a/crates/temper-observe/tests/otel_export/main.rs b/crates/temper-observe/tests/otel_export/main.rs new file mode 100644 index 000000000..3a9545314 --- /dev/null +++ b/crates/temper-observe/tests/otel_export/main.rs @@ -0,0 +1,157 @@ +//! End-to-end tests of the OTLP export settings. +//! +//! Each test starts the telemetry setup in a child process, pointed at a +//! local OTLP/HTTP listener, and asserts on what the listener received and on +//! what the child printed at startup. + +mod bad_values; +mod child; +mod identity; +mod listener; +mod render; +mod sampling; +mod signals; +mod wire; + +use child::{BUILT_IN_SERVICE_NAME, Run}; +use wire::Attributes; + +const BASELINE_PATH: &str = concat!( + env!("CARGO_MANIFEST_DIR"), + "/tests/otel_export/default_export.baseline.txt" +); + +const SIGNAL_PATHS: [&str; 3] = ["/v1/traces", "/v1/metrics", "/v1/logs"]; + +/// What the workload exports when every span the name-based filter lets +/// through is sampled, whatever its caller said. +const EVERY_SPAN: [&str; 5] = [ + "test.child", + "test.request", + "test.sampled_caller", + "test.unsampled_caller", + "wasm:workspace_fs.read", +]; + +/// What it exports when the caller's decision is followed, as by default. +const SPANS_FOLLOWING_THE_CALLER: [&str; 4] = [ + "test.child", + "test.request", + "test.sampled_caller", + "wasm:workspace_fs.read", +]; + +/// What it exports when only spans with a sampled caller are kept. +const SPANS_WITH_A_SAMPLED_CALLER: [&str; 2] = ["test.sampled_caller", "wasm:workspace_fs.read"]; + +/// The resource of each of the three signals, which must all have exported. +fn resources_of_all_signals(run: &Run) -> Vec<(&'static str, Attributes)> { + let resources = run.resources(); + let signals: Vec<&str> = resources.iter().map(|(signal, _)| *signal).collect(); + assert_eq!(signals, ["traces", "metrics", "logs"]); + resources +} + +fn attribute<'a>(resource: &'a Attributes, key: &str) -> Option<&'a str> { + resource.get(key).map(String::as_str) +} + +/// The messages of the warnings the telemetry setup logged. +fn warnings(run: &Run) -> Vec { + run.log_lines() + .iter() + .filter(|line| line["level"] == "WARN" && line["target"] == "temper_observe::otel") + .map(message) + .collect() +} + +/// The message of every log line the child printed. +fn log_messages(run: &Run) -> Vec { + run.log_lines().iter().map(message).collect() +} + +fn message(line: &serde_json::Map) -> String { + line["fields"]["message"] + .as_str() + .unwrap_or_default() + .to_string() +} + +/// Name and trace ID of every exported span, sorted. +fn exported_spans(run: &Run) -> Vec<(String, String)> { + let mut spans: Vec<(String, String)> = run + .spans() + .into_iter() + .map(|span| (span.name, span.trace_id)) + .collect(); + spans.sort(); + spans +} + +fn span_names(run: &Run) -> Vec { + exported_spans(run) + .into_iter() + .map(|(name, _)| name) + .collect() +} + +/// The startup output, the log lines and the export of `run` are those of +/// the recorded baseline. +fn assert_matches_baseline(run: &Run, what: &str) { + let expected = std::fs::read_to_string(BASELINE_PATH).expect("read the baseline"); + let actual = render::render(run); + assert!( + actual == expected, + "{what} differs from the recorded baseline\n\ + --- recorded ({BASELINE_PATH})\n{expected}\n--- actual\n{actual}" + ); +} + +/// With none of the export settings present, the startup output and the +/// exported telemetry match the recorded baseline. +/// +/// To record the baseline again after an intended change, run this test with +/// `UPDATE_OTEL_EXPORT_BASELINE=1` and review the diff of the baseline file. +#[test] +fn default_export_matches_recorded_baseline() { + let run = child::run(&[]); + for path in SIGNAL_PATHS { + assert!(run.requests_to(path) > 0, "no export reached {path}"); + } + if std::env::var_os("UPDATE_OTEL_EXPORT_BASELINE").is_some() { + std::fs::write(BASELINE_PATH, render::render(&run)).expect("write the baseline"); + return; + } + assert_matches_baseline(&run, "the default export"); +} + +/// Every exported log record carries an event time. A backend that dates +/// records by their event time treats a record without one as very old. +#[test] +fn every_log_record_has_an_event_time() { + let run = child::run(&[]); + let records = run.logs(); + assert!(!records.is_empty(), "no log records were exported"); + for record in records { + assert_ne!( + record.time_unix_nano, 0, + "log record {:?} was exported with an event time of zero", + record.body + ); + assert_eq!( + record.time_unix_nano, record.observed_time_unix_nano, + "log record {:?} has an event time that is not its observed time", + record.body + ); + } +} + +/// Asking "is logging enabled?" through the `log` bridge, with no log line +/// after it, must not cost the next span or log line on that thread: the +/// export and the log lines are the same as without the question. +#[test] +fn log_enabled_probe_does_not_drop_the_next_span() { + let run = child::run_after_log_probe(&[]); + assert_eq!(span_names(&run), SPANS_FOLLOWING_THE_CALLER); + assert_matches_baseline(&run, "the export after a log probe"); +} diff --git a/crates/temper-observe/tests/otel_export/render.rs b/crates/temper-observe/tests/otel_export/render.rs new file mode 100644 index 000000000..aecef1c39 --- /dev/null +++ b/crates/temper-observe/tests/otel_export/render.rs @@ -0,0 +1,195 @@ +//! Renders one run as text, with whatever differs between two runs of the +//! same code masked, so two runs can be compared line by line. + +use std::fmt::Write as _; + +use crate::child::{ + REDUCED_RATE_DROPPED_TRACE_ID, REDUCED_RATE_KEPT_TRACE_ID, Run, SAMPLED_CALLER_TRACE_ID, + UNSAMPLED_CALLER_TRACE_ID, +}; +use crate::wire::Attributes; + +/// Attributes whose values change from run to run, or with how the test +/// harness runs the child (the thread name). +const MASKED_ATTRIBUTES: &[&str] = &[ + "runtime-id", + "thread.id", + "thread.name", + "busy_ns", + "idle_ns", + "code.lineno", +]; + +/// Trace IDs the workload chooses itself; every other trace ID is random. +const CALLER_TRACE_IDS: &[&str] = &[ + UNSAMPLED_CALLER_TRACE_ID, + SAMPLED_CALLER_TRACE_ID, + REDUCED_RATE_KEPT_TRACE_ID, + REDUCED_RATE_DROPPED_TRACE_ID, +]; + +/// The startup output, the log lines and every exported span, log record and +/// metric. +pub fn render(run: &Run) -> String { + let mut ids = Ids::default(); + let mut out = String::new(); + + out.push_str("== stderr ==\n"); + out.push_str(&run.stderr); + + out.push_str("== log lines ==\n"); + for mut line in run.log_lines() { + line.remove("timestamp"); + let _ = writeln!(out, "{}", sorted_json(&serde_json::Value::Object(line))); + } + + out.push_str("== requests ==\n"); + let mut requests: Vec = run + .received + .iter() + .map(|request| format!("{} {}", request.path, request.header("content-type"))) + .collect(); + requests.sort(); + requests.dedup(); + for request in requests { + let _ = writeln!(out, "{request}"); + } + + out.push_str("== traces ==\n"); + // An export groups its items by scope in no fixed order; within a scope + // the order is the order of emission. + let mut spans = run.spans(); + spans.sort_by(|a, b| a.scope.cmp(&b.scope)); + for span in spans { + let _ = writeln!( + out, + "span {:?} scope={} trace={} id={} parent={} kind={} flags={} status={} events={:?}", + span.name, + span.scope, + ids.trace(&span.trace_id), + ids.span(&span.span_id), + ids.span(&span.parent_span_id), + span.kind, + span.flags, + span.status_code, + span.events, + ); + attributes(&mut out, "resource", &span.resource); + attributes(&mut out, "attribute", &span.attributes); + } + + out.push_str("== logs ==\n"); + let mut records = run.logs(); + records.sort_by(|a, b| a.scope.cmp(&b.scope)); + for record in records { + let time = match record.time_unix_nano { + 0 => "", + time if time == record.observed_time_unix_nano => "", + _ => "", + }; + let observed = match record.observed_time_unix_nano { + 0 => "", + _ => "", + }; + let _ = writeln!( + out, + "log {:?} scope={} severity={}({}) time={time} observed={observed} trace={} span={}", + record.body, + record.scope, + record.severity_text, + record.severity_number, + ids.trace(&record.trace_id), + ids.span(&record.span_id), + ); + attributes(&mut out, "resource", &record.resource); + attributes(&mut out, "attribute", &record.attributes); + } + + out.push_str("== metrics ==\n"); + let mut metrics = run.metrics(); + metrics.sort_by(|a, b| (&a.scope, &a.name).cmp(&(&b.scope, &b.name))); + for metric in metrics { + let _ = writeln!( + out, + "metric {:?} scope={} kind={} unit={:?} description={:?}", + metric.name, metric.scope, metric.kind, metric.unit, metric.description, + ); + attributes(&mut out, "resource", &metric.resource); + let mut points = metric.points; + points.sort(); + for (point_attributes, value) in points { + let _ = writeln!(out, " point {point_attributes:?} = {value}"); + } + } + + out.replace(&run.endpoint, "http://127.0.0.1:") +} + +fn attributes(out: &mut String, label: &str, attributes: &Attributes) { + for (key, value) in attributes { + let value = if MASKED_ATTRIBUTES.contains(&key.as_str()) { + "" + } else { + value.as_str() + }; + let _ = writeln!(out, " {label} {key} = {value:?}"); + } +} + +/// Compact JSON with the keys of every object sorted, so the text does not +/// depend on whether `serde_json` keeps keys in insertion order, which other +/// crates in a workspace build can switch on. +fn sorted_json(value: &serde_json::Value) -> String { + match value { + serde_json::Value::Object(fields) => { + let mut fields: Vec<(&String, &serde_json::Value)> = fields.iter().collect(); + fields.sort_by(|a, b| a.0.cmp(b.0)); + let fields: Vec = fields + .into_iter() + .map(|(key, value)| { + let key = serde_json::Value::from(key.as_str()); + format!("{key}:{}", sorted_json(value)) + }) + .collect(); + format!("{{{}}}", fields.join(",")) + } + serde_json::Value::Array(items) => { + let items: Vec = items.iter().map(sorted_json).collect(); + format!("[{}]", items.join(",")) + } + scalar => scalar.to_string(), + } +} + +/// Replaces random IDs by labels numbered in order of first appearance. +#[derive(Default)] +struct Ids { + traces: Vec, + spans: Vec, +} + +impl Ids { + fn trace(&mut self, id: &str) -> String { + if id.is_empty() || CALLER_TRACE_IDS.contains(&id) { + return format!("[{id}]"); + } + format!("trace#{}", position(&mut self.traces, id)) + } + + fn span(&mut self, id: &str) -> String { + if id.is_empty() { + return "[]".to_string(); + } + format!("span#{}", position(&mut self.spans, id)) + } +} + +fn position(seen: &mut Vec, id: &str) -> usize { + match seen.iter().position(|known| known == id) { + Some(index) => index + 1, + None => { + seen.push(id.to_string()); + seen.len() + } + } +} diff --git a/crates/temper-observe/tests/otel_export/sampling.rs b/crates/temper-observe/tests/otel_export/sampling.rs new file mode 100644 index 000000000..ee049f2a0 --- /dev/null +++ b/crates/temper-observe/tests/otel_export/sampling.rs @@ -0,0 +1,81 @@ +//! Choosing the trace sampler with `OTEL_TRACES_SAMPLER`. + +use crate::*; + +/// With `always_on`, a request that arrives marked "not sampled" still +/// produces a span, and the span keeps the caller's trace ID. +#[test] +fn always_on_records_a_request_marked_not_sampled() { + let run = child::run(&[("OTEL_TRACES_SAMPLER", "always_on")]); + let spans = exported_spans(&run); + assert!( + spans.contains(&( + "test.unsampled_caller".to_string(), + child::UNSAMPLED_CALLER_TRACE_ID.to_string() + )), + "no span with the caller's trace ID in {spans:?}" + ); + assert_eq!(span_names(&run), EVERY_SPAN); +} + +/// With the sampler unset, a request marked "not sampled" produces no span. +#[test] +fn default_sampler_follows_a_caller_that_did_not_sample() { + let run = child::run(&[]); + assert_eq!(span_names(&run), SPANS_FOLLOWING_THE_CALLER); +} + +#[test] +fn traceidratio_zero_exports_no_spans() { + let run = child::run(&[ + ("OTEL_TRACES_SAMPLER", "traceidratio"), + ("OTEL_TRACES_SAMPLER_ARG", "0"), + ]); + assert_eq!(span_names(&run), Vec::::new()); +} + +#[test] +fn traceidratio_one_exports_every_span() { + let run = child::run(&[ + ("OTEL_TRACES_SAMPLER", "traceidratio"), + ("OTEL_TRACES_SAMPLER_ARG", "1"), + ]); + assert_eq!(span_names(&run), EVERY_SPAN); +} + +/// Every supported sampler exports what it should, and under every one of +/// them the name-based filter drops what it drops by default: the two names +/// it drops outright, and the reduced-rate name on a trace ID outside the +/// rate. +#[test] +fn every_sampler_keeps_the_name_based_filter() { + let cases: [(&str, &str, &[&str]); 9] = [ + ("always_on", "", &EVERY_SPAN), + ("always_off", "", &[]), + ("parentbased_always_on", "", &SPANS_FOLLOWING_THE_CALLER), + ("parentbased_always_off", "", &SPANS_WITH_A_SAMPLED_CALLER), + ("traceidratio", "1", &EVERY_SPAN), + ("traceidratio", "0", &[]), + ("traceidratio", "", &EVERY_SPAN), + ("parentbased_traceidratio", "1", &SPANS_FOLLOWING_THE_CALLER), + ( + "parentbased_traceidratio", + "0", + &SPANS_WITH_A_SAMPLED_CALLER, + ), + ]; + for (sampler, arg, expected) in cases { + let run = child::run(&[ + ("OTEL_TRACES_SAMPLER", sampler), + ("OTEL_TRACES_SAMPLER_ARG", arg), + ]); + let case = format!("OTEL_TRACES_SAMPLER={sampler} OTEL_TRACES_SAMPLER_ARG={arg}"); + assert_eq!(span_names(&run), expected, "{case}"); + for (name, trace_id) in exported_spans(&run) { + assert_ne!(name, "clock_time_get", "{case}"); + assert_ne!(name, "turso.configured_connection", "{case}"); + assert_ne!(trace_id, child::REDUCED_RATE_DROPPED_TRACE_ID, "{case}"); + } + assert_eq!(warnings(&run), Vec::::new(), "{case}"); + } +} diff --git a/crates/temper-observe/tests/otel_export/signals.rs b/crates/temper-observe/tests/otel_export/signals.rs new file mode 100644 index 000000000..41ef07da1 --- /dev/null +++ b/crates/temper-observe/tests/otel_export/signals.rs @@ -0,0 +1,56 @@ +//! Switching a signal off with `OTEL_*_EXPORTER`. + +use crate::*; + +/// With `variable` set to `none`, nothing reaches `silent_path`, the other +/// two signals are still exported, and the process logs as before. +fn assert_only_one_signal_is_off(variable: &str, silent_path: &str, exported: &str) { + let run = child::run(&[(variable, "none")]); + let messages = log_messages(&run); + for expected in [ + format!("OTEL initialised ({exported})").as_str(), + "test log inside a span", + "test log outside a span", + ] { + assert!( + messages.iter().any(|message| message == expected), + "{variable}=none: no log line {expected:?} in {messages:?}" + ); + } + for path in SIGNAL_PATHS { + let requests = run.requests_to(path); + if path == silent_path { + assert_eq!(requests, 0, "{variable}=none still exported to {path}"); + } else { + assert!(requests > 0, "{variable}=none stopped the export to {path}"); + } + } + assert_eq!(warnings(&run), Vec::::new()); +} + +#[test] +fn traces_exporter_none_switches_only_traces_off() { + assert_only_one_signal_is_off("OTEL_TRACES_EXPORTER", "/v1/traces", "metrics + logs"); +} + +#[test] +fn metrics_exporter_none_switches_only_metrics_off() { + assert_only_one_signal_is_off("OTEL_METRICS_EXPORTER", "/v1/metrics", "traces + logs"); +} + +#[test] +fn logs_exporter_none_switches_only_logs_off() { + assert_only_one_signal_is_off("OTEL_LOGS_EXPORTER", "/v1/logs", "traces + metrics"); +} + +/// `otlp` is the default spelled out: everything is exported as when the +/// variables are unset, and nothing is reported. +#[test] +fn exporters_set_to_otlp_export_as_by_default() { + let run = child::run(&[ + ("OTEL_TRACES_EXPORTER", "otlp"), + ("OTEL_METRICS_EXPORTER", "OTLP"), + ("OTEL_LOGS_EXPORTER", " otlp "), + ]); + assert_matches_baseline(&run, "the export with every exporter set to otlp"); +} diff --git a/crates/temper-observe/tests/otel_export/wire.rs b/crates/temper-observe/tests/otel_export/wire.rs new file mode 100644 index 000000000..45bc8d9ec --- /dev/null +++ b/crates/temper-observe/tests/otel_export/wire.rs @@ -0,0 +1,322 @@ +//! Minimal reader for the OTLP protobuf requests the exporters send. +//! +//! Decodes only the fields the tests look at, straight from the protobuf +//! wire format, so the tests need no protobuf dependency. + +use std::collections::BTreeMap; + +/// Attribute keys mapped to their values rendered as text. +pub type Attributes = BTreeMap; + +/// One exported span with the resource and scope it was sent under. +#[derive(Clone, Debug)] +pub struct Span { + pub resource: Attributes, + pub scope: String, + pub trace_id: String, + pub span_id: String, + pub parent_span_id: String, + pub name: String, + pub kind: u64, + pub flags: u32, + pub attributes: Attributes, + pub events: Vec, + pub status_code: u64, +} + +/// One exported log record with the resource and scope it was sent under. +#[derive(Clone, Debug)] +pub struct LogRecord { + pub resource: Attributes, + pub scope: String, + pub time_unix_nano: u64, + pub observed_time_unix_nano: u64, + pub severity_number: u64, + pub severity_text: String, + pub body: String, + pub attributes: Attributes, + pub trace_id: String, + pub span_id: String, +} + +/// One exported metric with the resource and scope it was sent under. +#[derive(Clone, Debug)] +pub struct Metric { + pub resource: Attributes, + pub scope: String, + pub name: String, + pub description: String, + pub unit: String, + pub kind: &'static str, + /// Data point attributes and the point's value rendered as text. + pub points: Vec<(Attributes, String)>, +} + +/// Decode the body of a request to `/v1/traces`. +pub fn spans(body: &[u8]) -> Vec { + let mut out = Vec::new(); + for resource_spans in Message::parse(body).repeated(1) { + let resource = resource_attributes(&resource_spans); + for scope_spans in resource_spans.repeated(2) { + let scope = scope_spans.message(1).string(1); + for span in scope_spans.repeated(2) { + out.push(Span { + resource: resource.clone(), + scope: scope.clone(), + trace_id: hex(span.bytes(1)), + span_id: hex(span.bytes(2)), + parent_span_id: hex(span.bytes(4)), + name: span.string(5), + kind: span.varint(6), + flags: span.fixed32(16), + attributes: attributes(&span, 9), + events: span.repeated(11).iter().map(|e| e.string(2)).collect(), + status_code: span.message(15).varint(3), + }); + } + } + } + out +} + +/// Decode the body of a request to `/v1/logs`. +pub fn logs(body: &[u8]) -> Vec { + let mut out = Vec::new(); + for resource_logs in Message::parse(body).repeated(1) { + let resource = resource_attributes(&resource_logs); + for scope_logs in resource_logs.repeated(2) { + let scope = scope_logs.message(1).string(1); + for record in scope_logs.repeated(2) { + out.push(LogRecord { + resource: resource.clone(), + scope: scope.clone(), + time_unix_nano: record.fixed64(1), + observed_time_unix_nano: record.fixed64(11), + severity_number: record.varint(2), + severity_text: record.string(3), + body: any_value(&record.message(5)), + attributes: attributes(&record, 6), + trace_id: hex(record.bytes(9)), + span_id: hex(record.bytes(10)), + }); + } + } + } + out +} + +/// Decode the body of a request to `/v1/metrics`. +pub fn metrics(body: &[u8]) -> Vec { + let mut out = Vec::new(); + for resource_metrics in Message::parse(body).repeated(1) { + let resource = resource_attributes(&resource_metrics); + for scope_metrics in resource_metrics.repeated(2) { + let scope = scope_metrics.message(1).string(1); + for metric in scope_metrics.repeated(2) { + // Metric.data is a oneof: gauge = 5, sum = 7, histogram = 9. + let (kind, data) = [(5, "gauge"), (7, "sum"), (9, "histogram")] + .into_iter() + .find(|(field, _)| metric.has(*field)) + .map(|(field, kind)| (kind, metric.message(field))) + .unwrap_or(("other", Message::default())); + out.push(Metric { + resource: resource.clone(), + scope: scope.clone(), + name: metric.string(1), + description: metric.string(2), + unit: metric.string(3), + kind, + points: data + .repeated(1) + .iter() + .map(|point| data_point(kind, point)) + .collect(), + }); + } + } + } + out +} + +fn data_point(kind: &str, point: &Message<'_>) -> (Attributes, String) { + if kind == "histogram" { + // HistogramDataPoint: attributes = 9, count = 2. + return (attributes(point, 9), format!("count={}", point.fixed64(2))); + } + // NumberDataPoint: attributes = 7, as_double = 4, as_int = 6. + let value = if point.has(4) { + f64::from_bits(point.fixed64(4)).to_string() + } else { + (point.fixed64(6) as i64).to_string() + }; + (attributes(point, 7), value) +} + +fn resource_attributes(container: &Message<'_>) -> Attributes { + attributes(&container.message(1), 1) +} + +fn attributes(message: &Message<'_>, field: u32) -> Attributes { + message + .repeated(field) + .iter() + .map(|pair| (pair.string(1), any_value(&pair.message(2)))) + .collect() +} + +/// Render an `AnyValue` as text. +fn any_value(value: &Message<'_>) -> String { + match value.fields.first() { + None => String::new(), + Some((1, Wire::Bytes(text))) => String::from_utf8_lossy(text).into_owned(), + Some((2, Wire::Varint(flag))) => (*flag != 0).to_string(), + Some((3, Wire::Varint(int))) => (*int as i64).to_string(), + Some((4, Wire::Fixed64(bits))) => f64::from_bits(*bits).to_string(), + Some((5, Wire::Bytes(array))) => { + let items: Vec = Message::parse(array) + .repeated(1) + .iter() + .map(any_value) + .collect(); + format!("[{}]", items.join(", ")) + } + Some((6, Wire::Bytes(list))) => { + let pairs: Vec = attributes(&Message::parse(list), 1) + .into_iter() + .map(|(key, value)| format!("{key}={value}")) + .collect(); + format!("{{{}}}", pairs.join(", ")) + } + Some((7, Wire::Bytes(bytes))) => hex(bytes), + Some((field, _)) => panic!("unexpected AnyValue field {field}"), + } +} + +fn hex(bytes: &[u8]) -> String { + bytes.iter().map(|byte| format!("{byte:02x}")).collect() +} + +#[derive(Debug)] +enum Wire<'a> { + Varint(u64), + Fixed64(u64), + Fixed32(u32), + Bytes(&'a [u8]), +} + +/// The fields of one protobuf message, in wire order. +#[derive(Debug, Default)] +struct Message<'a> { + fields: Vec<(u32, Wire<'a>)>, +} + +impl<'a> Message<'a> { + fn parse(mut buf: &'a [u8]) -> Self { + let mut fields = Vec::new(); + while !buf.is_empty() { + let key = read_varint(&mut buf); + let value = match key & 7 { + 0 => Wire::Varint(read_varint(&mut buf)), + 1 => Wire::Fixed64(u64::from_le_bytes(take::<8>(&mut buf))), + 2 => { + let len = read_varint(&mut buf) as usize; + assert!(len <= buf.len(), "truncated protobuf field"); + let (head, rest) = buf.split_at(len); + buf = rest; + Wire::Bytes(head) + } + 5 => Wire::Fixed32(u32::from_le_bytes(take::<4>(&mut buf))), + other => panic!("unsupported protobuf wire type {other}"), + }; + fields.push(((key >> 3) as u32, value)); + } + Self { fields } + } + + fn has(&self, field: u32) -> bool { + self.fields.iter().any(|(number, _)| *number == field) + } + + fn bytes(&self, field: u32) -> &'a [u8] { + self.fields + .iter() + .find_map(|(number, value)| match value { + Wire::Bytes(bytes) if *number == field => Some(*bytes), + _ => None, + }) + .unwrap_or_default() + } + + fn string(&self, field: u32) -> String { + String::from_utf8_lossy(self.bytes(field)).into_owned() + } + + fn message(&self, field: u32) -> Message<'a> { + Message::parse(self.bytes(field)) + } + + fn repeated(&self, field: u32) -> Vec> { + self.fields + .iter() + .filter_map(|(number, value)| match value { + Wire::Bytes(bytes) if *number == field => Some(Message::parse(bytes)), + _ => None, + }) + .collect() + } + + /// A varint field; absent means the protobuf default, zero. + fn varint(&self, field: u32) -> u64 { + self.fields + .iter() + .find_map(|(number, value)| match value { + Wire::Varint(int) if *number == field => Some(*int), + _ => None, + }) + .unwrap_or_default() + } + + /// A fixed64 field; absent means the protobuf default, zero. + fn fixed64(&self, field: u32) -> u64 { + self.fields + .iter() + .find_map(|(number, value)| match value { + Wire::Fixed64(int) if *number == field => Some(*int), + _ => None, + }) + .unwrap_or_default() + } + + /// A fixed32 field; absent means the protobuf default, zero. + fn fixed32(&self, field: u32) -> u32 { + self.fields + .iter() + .find_map(|(number, value)| match value { + Wire::Fixed32(int) if *number == field => Some(*int), + _ => None, + }) + .unwrap_or_default() + } +} + +fn read_varint(buf: &mut &[u8]) -> u64 { + let mut value = 0u64; + let mut shift = 0u32; + loop { + let (&byte, rest) = buf.split_first().expect("truncated protobuf varint"); + *buf = rest; + value |= u64::from(byte & 0x7f) << shift; + if byte & 0x80 == 0 { + return value; + } + shift += 7; + assert!(shift < 64, "protobuf varint too long"); + } +} + +fn take(buf: &mut &[u8]) -> [u8; N] { + assert!(N <= buf.len(), "truncated protobuf fixed-width field"); + let (head, rest) = buf.split_at(N); + *buf = rest; + head.try_into().expect("split_at returned N bytes") +} diff --git a/docs/AGENT_GUIDE.md b/docs/AGENT_GUIDE.md index c9f8beba9..f17238a56 100644 --- a/docs/AGENT_GUIDE.md +++ b/docs/AGENT_GUIDE.md @@ -1288,6 +1288,26 @@ temper serve [--port PORT] [--specs-dir DIR] [--tenant NAME] | `ANTHROPIC_API_KEY` | For agent mode | Claude API key | | `RUST_LOG` | No | Log level (default: `info,temper=debug`) | +### Telemetry export settings + +When an OTLP endpoint is configured, these standard OpenTelemetry variables adjust what `temper serve` exports. All are optional. With none of them set, the export is unchanged. + +| Variable | Default | Effect | +|----------|---------|--------| +| `OTEL_SERVICE_NAME` | `temper-platform` | Service name on traces, metrics and logs. | +| `OTEL_RESOURCE_ATTRIBUTES` | none | Extra resource attributes as `key=value,key=value`, added to traces, metrics and logs. | +| `OTEL_TRACES_EXPORTER` | `otlp` | `none` switches the export of traces off. | +| `OTEL_METRICS_EXPORTER` | `otlp` | `none` switches the export of metrics off. | +| `OTEL_LOGS_EXPORTER` | `otlp` | `none` switches the export of logs off. | +| `OTEL_TRACES_SAMPLER` | `parentbased_always_on` | Which spans are recorded: `always_on`, `always_off`, `traceidratio`, `parentbased_always_on`, `parentbased_always_off` or `parentbased_traceidratio`. | +| `OTEL_TRACES_SAMPLER_ARG` | `1` | Ratio from 0 to 1 for `traceidratio` and `parentbased_traceidratio`. | + +- **Resource attributes the server computes itself win** over `OTEL_RESOURCE_ATTRIBUTES`: `runtime-id` always, `deployment.environment.name` when `DD_ENV` or `LOGFIRE_ENVIRONMENT` is set, and `service.version` when `DD_VERSION` is set. A `service.name` in `OTEL_RESOURCE_ATTRIBUTES` does not replace the service name; use `OTEL_SERVICE_NAME`. Keys and values are trimmed and otherwise taken as written. +- **A signal that is switched off** keeps running inside the process and is not exported. Spans still carry trace context and log lines are still printed. +- **Sampling.** The `parentbased_` samplers follow the caller's decision when a span has a remote parent. `always_on` and `traceidratio` do not, so they record spans for requests the caller marked "not sampled". The name-based filter (span names dropped outright, and prefixes kept at a reduced rate through `TEMPER_TRACE_WASM_AUX_SAMPLE_PCT` and `TEMPER_TRACE_DISPATCH_BACKGROUND_SAMPLE_PCT`) applies around whichever sampler is chosen. +- **Bad values never stop the server.** A value that is not supported is logged once at startup as a warning (`OTEL export setting: ...`) and the default applies. An empty value is the same as an unset variable. +- **Log records carry an event time.** Every exported log record has its event time set (to the time it was observed when it has none), so backends that date records by event time accept them. + **OTEL env var precedence:** The OTEL SDK reads signal-specific env vars (`OTEL_EXPORTER_OTLP_TRACES_ENDPOINT`, `OTEL_EXPORTER_OTLP_METRICS_ENDPOINT`, `OTEL_EXPORTER_OTLP_LOGS_ENDPOINT`) *before* the generic `OTEL_EXPORTER_OTLP_ENDPOINT`. If any signal-specific var is set — even by an unrelated tool — it silently overrides Temper's configured endpoint for that signal. `init_tracing()` clears these vars automatically. See Appendix D for details. ## Appendix D: Infrastructure Pitfalls