use std::fmt; use serde_json::{Map, Number, Value}; use time::{OffsetDateTime, format_description::well_known::Rfc3339}; use tracing::{Event, Subscriber, field::Visit}; use tracing_subscriber::{ EnvFilter, Layer, filter::filter_fn, fmt::{FmtContext, FormatEvent, FormatFields, MakeWriter, format::Writer}, layer::SubscriberExt, registry::LookupSpan, }; use crate::{ ObservabilityConfig, ObservabilityInitError, RedactionLimits, ServiceIdentity, propagation::current_trace_id, redaction::{redact_value, truncate_string}, schema::LogEnvelope, }; pub fn build_subscriber( config: ObservabilityConfig, writer: W, ) -> Result where W: for<'writer> MakeWriter<'writer> + Send + Sync + 'static, { build_subscriber_with_tracer(config, writer, None) } pub(crate) fn build_subscriber_with_tracer( config: ObservabilityConfig, writer: W, tracer: Option, ) -> Result where W: for<'writer> MakeWriter<'writer> + Send + Sync + 'static, { let (identity, filter, limits) = config.into_parts(); limits.validate()?; let filter = EnvFilter::try_new(filter).map_err(|_| ObservabilityInitError::InvalidFilter)?; let formatter = JsonEventFormatter::new(identity, limits); let fmt_layer = tracing_subscriber::fmt::layer() .with_ansi(false) .event_format(formatter) .with_writer(writer) .with_filter(filter); let otel_layer = tracer.map(|tracer| { tracing_opentelemetry::layer() .with_tracer(tracer) .with_filter(filter_fn(|metadata| { metadata.is_span() && metadata.target() == "crank::trace" })) }); Ok(tracing_subscriber::registry() .with(fmt_layer) .with(otel_layer)) } #[derive(Clone, Debug)] struct JsonEventFormatter { identity: ServiceIdentity, limits: RedactionLimits, } impl JsonEventFormatter { fn new(identity: ServiceIdentity, limits: RedactionLimits) -> Self { Self { identity, limits } } fn envelope(&self, event: &Event<'_>) -> Result { let metadata = event.metadata(); let mut visitor = JsonFieldVisitor::default(); event.record(&mut visitor); let mut raw_fields = visitor.fields; let request_id = take_correlation_id(&mut raw_fields, "request_id") .map(|value| truncate_string(&value, self.limits.max_string_bytes)); let trace_id = take_correlation_id(&mut raw_fields, "trace_id") .or_else(current_trace_id) .map(|value| truncate_string(&value, self.limits.max_string_bytes)); let cleaned = redact_value(&Value::Object(raw_fields), self.limits); let fields = cleaned.as_object().cloned().unwrap_or_default(); let timestamp = OffsetDateTime::now_utc() .format(&Rfc3339) .map_err(|_| fmt::Error)?; Ok(LogEnvelope { timestamp, level: metadata.level().as_str().to_owned(), service: self.identity.service().to_owned(), version: self.identity.version().to_owned(), environment: self.identity.environment().to_owned(), target: truncate_string(metadata.target(), self.limits.max_string_bytes), event: truncate_string(metadata.name(), self.limits.max_string_bytes), request_id, trace_id, fields, }) } fn serialize_bounded(&self, mut envelope: LogEnvelope) -> Result { let line_budget = self.limits.max_event_bytes.saturating_sub(1); let serialized = serde_json::to_string(&envelope).map_err(|_| fmt::Error)?; if serialized.len() <= line_budget { return Ok(serialized); } envelope.fields = Map::from_iter([("truncated".to_owned(), Value::Bool(true))]); let fallback_string_limit = self.limits.max_string_bytes.min(64); envelope.target = truncate_string(&envelope.target, fallback_string_limit); envelope.event = truncate_string(&envelope.event, fallback_string_limit); envelope.request_id = envelope .request_id .map(|value| truncate_string(&value, fallback_string_limit)); envelope.trace_id = envelope .trace_id .map(|value| truncate_string(&value, fallback_string_limit)); let serialized = serde_json::to_string(&envelope).map_err(|_| fmt::Error)?; if serialized.len() <= line_budget { return Ok(serialized); } envelope.request_id = None; envelope.trace_id = None; envelope.target = truncate_string(&envelope.target, 16); envelope.event = truncate_string(&envelope.event, 16); let serialized = serde_json::to_string(&envelope).map_err(|_| fmt::Error)?; (serialized.len() <= line_budget) .then_some(serialized) .ok_or(fmt::Error) } } impl FormatEvent for JsonEventFormatter where S: Subscriber + for<'lookup> LookupSpan<'lookup>, N: for<'writer> FormatFields<'writer> + 'static, { fn format_event( &self, _ctx: &FmtContext<'_, S, N>, mut writer: Writer<'_>, event: &Event<'_>, ) -> fmt::Result { let serialized = self.serialize_bounded(self.envelope(event)?)?; writer.write_str(&serialized)?; writer.write_char('\n') } } #[derive(Default)] struct JsonFieldVisitor { fields: Map, } impl JsonFieldVisitor { fn insert(&mut self, field: &tracing::field::Field, value: Value) { self.fields.insert(field.name().to_owned(), value); } } impl Visit for JsonFieldVisitor { fn record_i64(&mut self, field: &tracing::field::Field, value: i64) { self.insert(field, Value::Number(value.into())); } fn record_u64(&mut self, field: &tracing::field::Field, value: u64) { self.insert(field, Value::Number(value.into())); } fn record_bool(&mut self, field: &tracing::field::Field, value: bool) { self.insert(field, Value::Bool(value)); } fn record_f64(&mut self, field: &tracing::field::Field, value: f64) { let value = Number::from_f64(value) .map(Value::Number) .unwrap_or(Value::Null); self.insert(field, value); } fn record_str(&mut self, field: &tracing::field::Field, value: &str) { self.insert(field, Value::String(value.to_owned())); } fn record_debug(&mut self, field: &tracing::field::Field, value: &dyn fmt::Debug) { let rendered = format!("{value:?}"); let value = if is_correlation_field(field.name()) { Value::String(debug_scalar(&rendered)) } else { match serde_json::from_str(&rendered) { Ok(value @ (Value::Object(_) | Value::Array(_))) => value, _ if field.name() == "message" || is_safe_display_scalar(&rendered) => { Value::String(rendered) } _ => Value::String(crate::REDACTED_MARKER.to_owned()), } }; self.insert(field, value); } } fn take_correlation_id(fields: &mut Map, name: &str) -> Option { let value = fields.remove(name)?; let value = match value { Value::String(value) => value, Value::Number(value) => value.to_string(), Value::Bool(value) => value.to_string(), Value::Null | Value::Array(_) | Value::Object(_) => return None, }; (!value.is_empty()).then_some(value) } fn is_correlation_field(name: &str) -> bool { matches!(name, "request_id" | "trace_id" | "correlation_id") } fn debug_scalar(rendered: &str) -> String { serde_json::from_str::(rendered).unwrap_or_else(|_| rendered.to_owned()) } fn is_safe_display_scalar(rendered: &str) -> bool { !rendered.is_empty() && rendered.bytes().all(|byte| { byte.is_ascii_alphanumeric() || matches!(byte, b'-' | b'_' | b'.' | b':' | b'/' | b'+' | b'@') }) } #[cfg(test)] mod tests { use std::{ io, sync::{Arc, Mutex}, }; use opentelemetry::{global, trace::TracerProvider as _}; use opentelemetry_sdk::{ error::OTelSdkResult, propagation::TraceContextPropagator, trace::{SdkTracerProvider, SpanData, SpanExporter}, }; use tracing::info; use super::build_subscriber_with_tracer; use crate::{ ObservabilityConfig, RedactionLimits, ServiceIdentity, inject_current_trace_context, }; #[test] fn trace_spans_ignore_the_log_level_filter() { global::set_text_map_propagator(TraceContextPropagator::new()); let provider = SdkTracerProvider::builder().build(); let tracer = provider.tracer("trace-filter-test"); let subscriber = build_subscriber_with_tracer(test_config("warn"), io::sink, Some(tracer)) .expect("subscriber must build"); let dispatch = tracing::Dispatch::new(subscriber); let _dispatch_guard = tracing::dispatcher::set_default(&dispatch); let span = tracing::info_span!(target: "crank::trace", "http.request"); let _span_guard = span.enter(); let mut headers = axum::http::HeaderMap::new(); assert!(inject_current_trace_context(&mut headers)); assert!(headers.contains_key("traceparent")); } #[test] fn otel_layer_does_not_export_events() { let exported = Arc::new(Mutex::new(Vec::new())); let provider = SdkTracerProvider::builder() .with_simple_exporter(CapturingExporter(Arc::clone(&exported))) .build(); let tracer = provider.tracer("event-filter-test"); let subscriber = build_subscriber_with_tracer(test_config("info"), io::sink, Some(tracer)) .expect("subscriber must build"); let dispatch = tracing::Dispatch::new(subscriber); tracing::dispatcher::with_default(&dispatch, || { let span = tracing::info_span!(target: "crank::trace", "http.request"); let _span_guard = span.enter(); info!(password = "canary-secret", "sensitive event"); }); provider.force_flush().expect("span must be exported"); let spans = exported.lock().expect("capture lock"); assert_eq!(spans.len(), 1); assert!(spans[0].events.is_empty()); assert!( !format!("{:?}", spans[0]) .as_bytes() .windows(b"canary-secret".len()) .any(|window| window == b"canary-secret") ); } fn test_config(filter: &str) -> ObservabilityConfig { ObservabilityConfig::new( ServiceIdentity::try_new("admin-api", "test", "test").unwrap(), filter, RedactionLimits::default(), ) } #[derive(Clone, Debug)] struct CapturingExporter(Arc>>); impl SpanExporter for CapturingExporter { async fn export(&self, batch: Vec) -> OTelSdkResult { self.0.lock().expect("capture lock").extend(batch); Ok(()) } } }