1mod log_sink;
11mod otel;
12
13use std::sync::Arc;
14
15use leviath_core::config::{ObservabilityConfig, TelemetryExporterKind};
16use leviath_core::telemetry::TelemetrySink;
17
18pub use log_sink::LogSink;
19pub use otel::OtelSink;
20
21pub type LogLayer = Box<dyn tracing_subscriber::Layer<tracing_subscriber::Registry> + Send + Sync>;
24
25pub struct BuiltTelemetry {
29 pub sink: Arc<dyn TelemetrySink>,
31 pub log_layer: Option<LogLayer>,
34}
35
36pub fn build_sink(cfg: &ObservabilityConfig) -> Option<BuiltTelemetry> {
41 if !cfg.enabled {
42 return None;
43 }
44 match cfg.exporter {
45 TelemetryExporterKind::None => None,
46 TelemetryExporterKind::Stdout => Some(BuiltTelemetry {
47 sink: Arc::new(LogSink),
48 log_layer: None,
49 }),
50 TelemetryExporterKind::Otlp => {
51 let cfg = cfg.clone();
55 let built = std::thread::spawn(move || OtelSink::from_config(&cfg))
56 .join()
57 .expect("exporter construction reports errors rather than panicking");
58 match built {
59 Ok(sink) => {
60 let log_layer = Some(sink.tracing_log_layer());
61 Some(BuiltTelemetry {
62 sink: Arc::new(sink),
63 log_layer,
64 })
65 }
66 Err(err) => {
67 tracing::warn!("telemetry disabled: {err}");
68 None
69 }
70 }
71 }
72 }
73}
74
75#[cfg(test)]
76mod tests {
77 use super::*;
78
79 fn config(enabled: bool, exporter: TelemetryExporterKind) -> ObservabilityConfig {
80 ObservabilityConfig {
81 enabled,
82 exporter,
83 endpoint: None,
84 service_name: None,
85 }
86 }
87
88 #[test]
89 fn disabled_config_builds_no_sink() {
90 assert!(build_sink(&config(false, TelemetryExporterKind::Otlp)).is_none());
91 }
92
93 #[test]
94 fn none_exporter_builds_no_sink() {
95 assert!(build_sink(&config(true, TelemetryExporterKind::None)).is_none());
96 }
97
98 #[test]
99 fn stdout_exporter_builds_the_log_sink_without_a_log_layer() {
100 let built = build_sink(&config(true, TelemetryExporterKind::Stdout)).unwrap();
101 assert!(built.log_layer.is_none());
102 }
103
104 #[tokio::test(flavor = "multi_thread")]
105 async fn otlp_exporter_builds_from_a_runtime_thread_with_a_log_layer() {
106 let built = build_sink(&config(true, TelemetryExporterKind::Otlp)).unwrap();
109 assert!(built.log_layer.is_some());
110 }
111
112 #[test]
113 fn an_unparseable_endpoint_disables_telemetry_with_a_warning() {
114 let cfg = ObservabilityConfig {
115 enabled: true,
116 exporter: TelemetryExporterKind::Otlp,
117 endpoint: Some("not a url at all".to_string()),
118 service_name: None,
119 };
120 assert!(build_sink(&cfg).is_none());
121 }
122}