1use crate::config::{LogConfig, LogFormat, TelemetryConfig, TelemetryProtocol};
8use http::Uri;
9use std::sync::Arc;
10use thiserror::Error;
11use tracing_subscriber::{EnvFilter, fmt, layer::SubscriberExt, reload, util::SubscriberInitExt};
12
13#[cfg(feature = "telemetry-otlp")]
14use opentelemetry::{KeyValue, trace::TracerProvider as _};
15#[cfg(feature = "telemetry-otlp")]
16use opentelemetry_otlp::WithExportConfig as _;
17#[cfg(feature = "telemetry-otlp")]
18use opentelemetry_otlp::WithTonicConfig as _;
19#[cfg(feature = "telemetry-otlp")]
20use opentelemetry_otlp::tonic_types::transport::ClientTlsConfig;
21#[cfg(feature = "telemetry-otlp")]
22use opentelemetry_sdk::{Resource, propagation::TraceContextPropagator, trace::SdkTracerProvider};
23
24#[derive(Clone)]
35pub struct FilterReloadHandle {
36 reload: Arc<dyn Fn(EnvFilter) -> Result<(), String> + Send + Sync>,
37}
38
39impl std::fmt::Debug for FilterReloadHandle {
40 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
41 f.debug_struct("FilterReloadHandle").finish_non_exhaustive()
42 }
43}
44
45impl FilterReloadHandle {
46 fn from_handle<S>(handle: reload::Handle<EnvFilter, S>) -> Self
50 where
51 S: tracing::Subscriber + 'static,
52 {
53 Self {
54 reload: Arc::new(move |filter| {
55 handle.reload(filter).map_err(|error| error.to_string())
56 }),
57 }
58 }
59
60 pub fn apply_directive(&self, directive: &str) -> Result<(), String> {
70 let filter = EnvFilter::try_new(directive).map_err(|error| error.to_string())?;
71 (self.reload)(filter)
72 }
73
74 #[cfg(test)]
79 #[must_use]
80 pub fn accept_all_for_test() -> Self {
81 Self {
82 reload: Arc::new(|_filter| Ok(())),
83 }
84 }
85}
86
87#[derive(Debug, Clone, Copy, PartialEq, Eq)]
89#[non_exhaustive]
90pub enum ResolvedLogFormat {
91 Pretty,
93 Json,
95}
96
97#[derive(Debug, Clone, PartialEq, Eq)]
99pub struct TelemetryRuntime {
100 pub log_format: ResolvedLogFormat,
102 pub trace_export: TraceExport,
104 pub warning: Option<String>,
106}
107
108#[derive(Debug, Clone, PartialEq, Eq)]
110pub enum TraceExport {
111 Disabled,
113 Otlp(OtlpTraceRuntime),
115}
116
117#[derive(Debug, Clone, PartialEq, Eq)]
119pub struct OtlpTraceRuntime {
120 pub endpoint: String,
122 pub protocol: TelemetryProtocol,
124 pub resource: TelemetryResource,
126}
127
128#[derive(Debug, Clone, PartialEq, Eq)]
130pub struct TelemetryResource {
131 pub service_name: String,
133 pub service_namespace: Option<String>,
135 pub service_version: String,
137 pub environment: String,
139}
140
141#[derive(Debug, Error)]
143#[non_exhaustive]
144pub enum TelemetryInitError {
145 #[error("telemetry is enabled but no OTLP endpoint was configured")]
147 MissingEndpoint,
148 #[error("telemetry service_name must not be empty")]
150 EmptyServiceName,
151 #[error("invalid OTLP endpoint {endpoint:?}: {reason}")]
153 InvalidEndpoint {
154 endpoint: String,
156 reason: String,
158 },
159 #[error("telemetry-otlp cargo feature is not enabled")]
161 #[allow(dead_code)]
162 FeatureDisabled,
163 #[error("failed to initialize OTLP exporter: {0}")]
165 #[allow(dead_code)]
166 #[allow(dead_code)]
167 #[allow(dead_code)]
168 #[allow(dead_code)]
169 ExporterInit(String),
170 #[error("failed to initialize tracing subscriber: {0}")]
172 SubscriberInit(String),
173}
174
175#[must_use]
182#[derive(Debug)]
183pub struct TelemetryGuard {
184 #[cfg(feature = "telemetry-otlp")]
185 provider: Option<SdkTracerProvider>,
186 pub log_buffer: Option<crate::log::capture::LogBuffer>,
189 pub filter_reload: Option<FilterReloadHandle>,
198}
199
200impl TelemetryGuard {
201 pub const fn disabled() -> Self {
208 Self {
209 #[cfg(feature = "telemetry-otlp")]
210 provider: None,
211 log_buffer: None,
212 filter_reload: None,
213 }
214 }
215
216 #[cfg(feature = "telemetry-otlp")]
217 const fn with_provider(provider: SdkTracerProvider) -> Self {
218 Self {
219 provider: Some(provider),
220 log_buffer: None,
221 filter_reload: None,
222 }
223 }
224
225 fn with_log_buffer(mut self, buffer: crate::log::capture::LogBuffer) -> Self {
226 self.log_buffer = Some(buffer);
227 self
228 }
229
230 fn with_filter_reload(mut self, handle: FilterReloadHandle) -> Self {
231 self.filter_reload = Some(handle);
232 self
233 }
234}
235
236impl Drop for TelemetryGuard {
237 fn drop(&mut self) {
238 #[cfg(feature = "telemetry-otlp")]
239 if let Some(provider) = self.provider.take() {
240 let _ = provider.shutdown();
241 }
242 }
243}
244
245impl TelemetryRuntime {
246 pub fn from_config(
256 log: &LogConfig,
257 telemetry: &TelemetryConfig,
258 profile: Option<&str>,
259 ) -> Result<Self, TelemetryInitError> {
260 let log_format = resolve_log_format(log, profile);
261 if !telemetry.enabled {
262 return Ok(Self {
263 log_format,
264 trace_export: TraceExport::Disabled,
265 warning: None,
266 });
267 }
268
269 if telemetry.service_name.trim().is_empty() {
270 return strict_or_fallback(
271 log_format,
272 telemetry.strict,
273 TelemetryInitError::EmptyServiceName,
274 );
275 }
276
277 let Some(endpoint) = telemetry
278 .otlp_endpoint
279 .as_deref()
280 .map(str::trim)
281 .filter(|value| !value.is_empty())
282 else {
283 return strict_or_fallback(
284 log_format,
285 telemetry.strict,
286 TelemetryInitError::MissingEndpoint,
287 );
288 };
289
290 if let Err(error) = validate_otlp_endpoint(endpoint) {
291 return strict_or_fallback(log_format, telemetry.strict, error);
292 }
293
294 Ok(Self {
295 log_format,
296 trace_export: TraceExport::Otlp(OtlpTraceRuntime {
297 endpoint: endpoint.to_owned(),
298 protocol: telemetry.protocol,
299 resource: TelemetryResource {
300 service_name: telemetry.service_name.clone(),
301 service_namespace: telemetry.service_namespace.clone(),
302 service_version: telemetry.service_version.clone(),
303 environment: telemetry.environment.clone(),
304 },
305 }),
306 warning: None,
307 })
308 }
309}
310
311pub fn init(
318 log: &LogConfig,
319 telemetry: &TelemetryConfig,
320 profile: Option<&str>,
321) -> Result<TelemetryGuard, TelemetryInitError> {
322 let runtime = TelemetryRuntime::from_config(log, telemetry, profile)?;
323 if let Some(warning) = runtime.warning.as_deref() {
324 eprintln!("Warning: {warning}");
325 }
326
327 let opted_out_defaults =
328 crate::log::filter::normalized_opt_out_defaults(&log.unfilter_parameters);
329 if !opted_out_defaults.is_empty() {
330 eprintln!(
331 "Warning: log.unfilter_parameters opted out built-in sensitive keys: {}",
332 opted_out_defaults.join(", ")
333 );
334 }
335
336 match &runtime.trace_export {
337 TraceExport::Disabled => init_logging_only(log, runtime.log_format),
338 TraceExport::Otlp(otlp) => {
339 init_otlp_runtime(log, runtime.log_format, telemetry.strict, otlp)
340 }
341 }
342}
343
344fn strict_or_fallback(
345 log_format: ResolvedLogFormat,
346 strict: bool,
347 error: TelemetryInitError,
348) -> Result<TelemetryRuntime, TelemetryInitError> {
349 if strict {
350 Err(error)
351 } else {
352 Ok(TelemetryRuntime {
353 log_format,
354 trace_export: TraceExport::Disabled,
355 warning: Some(error.to_string()),
356 })
357 }
358}
359
360fn resolve_log_format(log: &LogConfig, profile: Option<&str>) -> ResolvedLogFormat {
361 match log.format {
362 LogFormat::Pretty => ResolvedLogFormat::Pretty,
363 LogFormat::Json => ResolvedLogFormat::Json,
364 LogFormat::Auto => {
365 if is_production_profile(profile) || is_production_env() {
366 ResolvedLogFormat::Json
367 } else {
368 ResolvedLogFormat::Pretty
369 }
370 }
371 }
372}
373
374fn is_production_profile(profile: Option<&str>) -> bool {
375 profile.is_some_and(|value| {
376 value.eq_ignore_ascii_case("prod") || value.eq_ignore_ascii_case("production")
377 })
378}
379
380fn is_production_env() -> bool {
381 std::env::var("AUTUMN_ENV").is_ok_and(|value| value.eq_ignore_ascii_case("production"))
382}
383
384#[cfg(feature = "telemetry-otlp")]
386fn is_https_endpoint(endpoint: &str) -> bool {
387 endpoint
388 .parse::<Uri>()
389 .ok()
390 .and_then(|uri| uri.scheme_str().map(|scheme| scheme == "https"))
391 .unwrap_or(false)
392}
393
394fn validate_otlp_endpoint(endpoint: &str) -> Result<(), TelemetryInitError> {
395 let uri: Uri = endpoint.parse().map_err(|error: http::uri::InvalidUri| {
396 TelemetryInitError::InvalidEndpoint {
397 endpoint: endpoint.to_owned(),
398 reason: error.to_string(),
399 }
400 })?;
401
402 if uri.scheme().is_none() {
403 return Err(TelemetryInitError::InvalidEndpoint {
404 endpoint: endpoint.to_owned(),
405 reason: "missing URI scheme".to_owned(),
406 });
407 }
408
409 if uri.authority().is_none() {
410 return Err(TelemetryInitError::InvalidEndpoint {
411 endpoint: endpoint.to_owned(),
412 reason: "missing URI authority".to_owned(),
413 });
414 }
415
416 Ok(())
417}
418
419fn build_filter(log: &LogConfig) -> EnvFilter {
420 EnvFilter::try_new(&log.level).unwrap_or_else(|error| {
421 eprintln!(
422 "Warning: invalid log filter {:?}: {error}, falling back to \"info\"",
423 log.level
424 );
425 EnvFilter::new("info")
426 })
427}
428
429fn build_capture_layer(
430 log: &LogConfig,
431) -> Option<(
432 crate::log::capture::LogCaptureLayer,
433 crate::log::capture::LogBuffer,
434)> {
435 if !log.capture.enabled {
436 return None;
437 }
438 let mut filter_parameters = log.filter_parameters.clone();
440 filter_parameters.extend(crate::encryption::registered_encrypted_column_names());
441 let filter =
442 crate::log::filter::ParameterFilter::new(&filter_parameters, &log.unfilter_parameters);
443 let buffer = crate::log::capture::LogBuffer::new(log.capture.capacity, filter);
444 let layer = crate::log::capture::LogCaptureLayer::new(buffer.clone());
445 Some((layer, buffer))
446}
447
448fn init_logging_only(
449 log: &LogConfig,
450 log_format: ResolvedLogFormat,
451) -> Result<TelemetryGuard, TelemetryInitError> {
452 let filter = build_filter(log);
453 let capture = build_capture_layer(log);
454 let capture_layer = capture.as_ref().map(|(layer, _)| layer.clone());
455
456 let reload_handle = match log_format {
459 ResolvedLogFormat::Json => {
460 let (filter_layer, handle) = reload::Layer::new(filter);
461 tracing_subscriber::registry()
462 .with(filter_layer)
463 .with(fmt::layer().json())
464 .with(capture_layer)
465 .try_init()
466 .map_err(|error| TelemetryInitError::SubscriberInit(error.to_string()))?;
467 FilterReloadHandle::from_handle(handle)
468 }
469 ResolvedLogFormat::Pretty => {
470 let (filter_layer, handle) = reload::Layer::new(filter);
471 tracing_subscriber::registry()
472 .with(filter_layer)
473 .with(fmt::layer().pretty())
474 .with(capture_layer)
475 .try_init()
476 .map_err(|error| TelemetryInitError::SubscriberInit(error.to_string()))?;
477 FilterReloadHandle::from_handle(handle)
478 }
479 };
480
481 let guard = TelemetryGuard::disabled().with_filter_reload(reload_handle);
482 if let Some((_, buffer)) = capture {
483 Ok(guard.with_log_buffer(buffer))
484 } else {
485 Ok(guard)
486 }
487}
488
489#[cfg(feature = "telemetry-otlp")]
490fn init_otlp_runtime(
491 log: &LogConfig,
492 log_format: ResolvedLogFormat,
493 strict: bool,
494 otlp: &OtlpTraceRuntime,
495) -> Result<TelemetryGuard, TelemetryInitError> {
496 let provider = match build_tracer_provider(otlp) {
497 Ok(provider) => provider,
498 Err(error) => {
499 if strict {
500 return Err(error);
501 }
502 eprintln!("Warning: {error}");
503 return init_logging_only(log, log_format);
504 }
505 };
506
507 let tracer = provider.tracer("autumn-web");
508 let filter = build_filter(log);
509 let capture = build_capture_layer(log);
510 let capture_layer = capture.as_ref().map(|(layer, _)| layer.clone());
511
512 let reload_handle = match log_format {
515 ResolvedLogFormat::Json => {
516 let (filter_layer, handle) = reload::Layer::new(filter);
517 tracing_subscriber::registry()
518 .with(filter_layer)
519 .with(fmt::layer().json())
520 .with(tracing_opentelemetry::layer().with_tracer(tracer))
521 .with(capture_layer)
522 .try_init()
523 .map_err(|error| TelemetryInitError::SubscriberInit(error.to_string()))?;
524 FilterReloadHandle::from_handle(handle)
525 }
526 ResolvedLogFormat::Pretty => {
527 let (filter_layer, handle) = reload::Layer::new(filter);
528 tracing_subscriber::registry()
529 .with(filter_layer)
530 .with(fmt::layer().pretty())
531 .with(tracing_opentelemetry::layer().with_tracer(tracer))
532 .with(capture_layer)
533 .try_init()
534 .map_err(|error| TelemetryInitError::SubscriberInit(error.to_string()))?;
535 FilterReloadHandle::from_handle(handle)
536 }
537 };
538
539 let guard = TelemetryGuard::with_provider(provider).with_filter_reload(reload_handle);
540 if let Some((_, buffer)) = capture {
541 Ok(guard.with_log_buffer(buffer))
542 } else {
543 Ok(guard)
544 }
545}
546
547#[cfg(not(feature = "telemetry-otlp"))]
548fn init_otlp_runtime(
549 log: &LogConfig,
550 log_format: ResolvedLogFormat,
551 strict: bool,
552 _otlp: &OtlpTraceRuntime,
553) -> Result<TelemetryGuard, TelemetryInitError> {
554 if strict {
555 return Err(TelemetryInitError::FeatureDisabled);
556 }
557
558 eprintln!("Warning: {}", TelemetryInitError::FeatureDisabled);
559 init_logging_only(log, log_format)
560}
561
562#[cfg(feature = "telemetry-otlp")]
563fn build_tracer_provider(otlp: &OtlpTraceRuntime) -> Result<SdkTracerProvider, TelemetryInitError> {
564 let resource = Resource::builder()
565 .with_service_name(otlp.resource.service_name.clone())
566 .with_attributes(build_resource_attributes(&otlp.resource))
567 .build();
568
569 let exporter = match otlp.protocol {
570 TelemetryProtocol::Grpc => {
571 let builder = opentelemetry_otlp::SpanExporter::builder()
572 .with_tonic()
573 .with_endpoint(otlp.endpoint.clone());
574 let builder = if is_https_endpoint(&otlp.endpoint) {
582 builder.with_tls_config(ClientTlsConfig::new().with_enabled_roots())
583 } else {
584 builder
585 };
586 builder.build()
587 }
588 TelemetryProtocol::HttpProtobuf => opentelemetry_otlp::SpanExporter::builder()
589 .with_http()
590 .with_endpoint(otlp.endpoint.clone())
591 .build(),
592 }
593 .map_err(|error| TelemetryInitError::ExporterInit(error.to_string()))?;
594
595 opentelemetry::global::set_text_map_propagator(TraceContextPropagator::new());
600
601 Ok(SdkTracerProvider::builder()
602 .with_resource(resource)
603 .with_batch_exporter(exporter)
604 .build())
605}
606
607#[cfg(feature = "telemetry-otlp")]
608fn build_resource_attributes(resource: &TelemetryResource) -> [KeyValue; 3] {
609 [
610 KeyValue::new(
611 "service.namespace",
612 resource.service_namespace.clone().unwrap_or_default(),
613 ),
614 KeyValue::new("service.version", resource.service_version.clone()),
615 KeyValue::new("deployment.environment", resource.environment.clone()),
616 ]
617}
618
619pub trait TelemetryProvider: Send + Sync + 'static {
658 fn init(
668 &self,
669 log: &LogConfig,
670 telemetry: &TelemetryConfig,
671 profile: Option<&str>,
672 ) -> Result<TelemetryGuard, TelemetryInitError>;
673}
674
675#[derive(Debug, Default, Clone, Copy)]
681pub struct TracingOtlpTelemetryProvider;
682
683impl TracingOtlpTelemetryProvider {
684 #[must_use]
686 pub const fn new() -> Self {
687 Self
688 }
689}
690
691impl TelemetryProvider for TracingOtlpTelemetryProvider {
692 fn init(
693 &self,
694 log: &LogConfig,
695 telemetry: &TelemetryConfig,
696 profile: Option<&str>,
697 ) -> Result<TelemetryGuard, TelemetryInitError> {
698 init(log, telemetry, profile)
699 }
700}
701
702#[cfg(test)]
703mod tests {
704 use super::*;
705
706 struct NoOpTelemetryProvider;
710
711 impl TelemetryProvider for NoOpTelemetryProvider {
712 fn init(
713 &self,
714 _log: &LogConfig,
715 _telemetry: &TelemetryConfig,
716 _profile: Option<&str>,
717 ) -> Result<TelemetryGuard, TelemetryInitError> {
718 Ok(TelemetryGuard::disabled())
719 }
720 }
721
722 #[test]
723 fn telemetry_provider_trait_returns_supplied_guard() {
724 let provider = NoOpTelemetryProvider;
725 let log = LogConfig::default();
726 let telemetry = TelemetryConfig::default();
727 let guard = provider
729 .init(&log, &telemetry, Some("test"))
730 .expect("no-op provider should succeed");
731 drop(guard);
733 }
734 #[cfg(feature = "telemetry-otlp")]
735 #[test]
736 fn is_https_endpoint_detects_scheme() {
737 assert!(is_https_endpoint("https://collector.example.com:4317"));
738 assert!(!is_https_endpoint("http://localhost:4317"));
739 assert!(!is_https_endpoint("localhost:4317"));
740 assert!(!is_https_endpoint("not a uri"));
741 }
742
743 #[test]
744 fn build_capture_layer_returns_none_when_disabled() {
745 let log = LogConfig::default(); assert!(build_capture_layer(&log).is_none());
747 }
748
749 #[test]
750 fn build_capture_layer_returns_layer_and_buffer_when_enabled() {
751 let log = LogConfig {
752 capture: crate::log::capture::LogCaptureConfig {
753 enabled: true,
754 capacity: 50,
755 },
756 ..Default::default()
757 };
758 let result = build_capture_layer(&log);
759 assert!(
760 result.is_some(),
761 "should build layer when capture.enabled = true"
762 );
763 let (_, buffer) = result.unwrap();
764 assert!(buffer.is_empty(), "newly created buffer should be empty");
765 }
766
767 #[test]
768 fn build_filter_falls_back_to_info_on_invalid_level() {
769 let log = LogConfig {
770 level: "this_is_not_a_valid_directive_it_lacks_an_equal_sign_and_is_not_a_level,foo=bar=baz=invalid".to_owned(),
771 ..Default::default()
772 };
773
774 let filter = build_filter(&log);
775 assert_eq!(filter.to_string(), "info");
776 }
777
778 #[cfg(feature = "telemetry-otlp")]
779 #[test]
780 fn build_resource_attributes_populates_otel_semantic_keys() {
781 let resource = TelemetryResource {
782 service_name: "svc".into(),
783 service_namespace: Some("team".into()),
784 service_version: "1.2.3".into(),
785 environment: "staging".into(),
786 };
787 let attrs = build_resource_attributes(&resource);
788 let pairs: std::collections::HashMap<_, _> = attrs
789 .iter()
790 .map(|kv| (kv.key.as_str().to_owned(), kv.value.to_string()))
791 .collect();
792 assert_eq!(
793 pairs.get("service.namespace").map(String::as_str),
794 Some("team")
795 );
796 assert_eq!(
797 pairs.get("service.version").map(String::as_str),
798 Some("1.2.3")
799 );
800 assert_eq!(
801 pairs.get("deployment.environment").map(String::as_str),
802 Some("staging")
803 );
804 }
805
806 #[cfg(feature = "telemetry-otlp")]
807 struct MapExtractor<'a>(&'a std::collections::HashMap<&'static str, &'static str>);
808
809 #[cfg(feature = "telemetry-otlp")]
810 impl opentelemetry::propagation::Extractor for MapExtractor<'_> {
811 fn get(&self, key: &str) -> Option<&str> {
812 self.0.get(key).copied()
813 }
814 fn keys(&self) -> Vec<&str> {
815 self.0.keys().copied().collect()
816 }
817 }
818
819 #[cfg(feature = "telemetry-otlp")]
820 #[tokio::test]
821 async fn build_tracer_provider_installs_w3c_propagator_and_returns_provider() {
822 use opentelemetry::trace::TraceContextExt as _;
823
824 let otlp = OtlpTraceRuntime {
829 endpoint: "http://127.0.0.1:65530".into(),
830 protocol: TelemetryProtocol::Grpc,
831 resource: TelemetryResource {
832 service_name: "unit-test".into(),
833 service_namespace: None,
834 service_version: "0.0.0".into(),
835 environment: "test".into(),
836 },
837 };
838 let provider = build_tracer_provider(&otlp)
839 .expect("tonic exporter + provider build should succeed lazily");
840
841 let headers = std::collections::HashMap::from([(
844 "traceparent",
845 "00-0af7651916cd43dd8448eb211c80319c-b7ad6b7169203331-01",
846 )]);
847 let cx =
848 opentelemetry::global::get_text_map_propagator(|p| p.extract(&MapExtractor(&headers)));
849 assert!(
850 cx.span().span_context().is_valid(),
851 "global propagator should have been installed"
852 );
853
854 let _ = provider.shutdown();
855 }
856
857 #[cfg(feature = "telemetry-otlp")]
858 #[tokio::test]
859 async fn build_tracer_provider_supports_http_protobuf_protocol() {
860 let otlp = OtlpTraceRuntime {
861 endpoint: "http://127.0.0.1:65531".into(),
862 protocol: TelemetryProtocol::HttpProtobuf,
863 resource: TelemetryResource {
864 service_name: "unit-test".into(),
865 service_namespace: None,
866 service_version: "0.0.0".into(),
867 environment: "test".into(),
868 },
869 };
870 let provider =
871 build_tracer_provider(&otlp).expect("http-protobuf exporter should build lazily");
872 let _ = provider.shutdown();
873 }
874}