faucet_core/observability/
install.rs1use thiserror::Error;
6
7#[derive(Debug, Clone, Default)]
10pub struct ObservabilityConfig {
11 pub prometheus: Option<PrometheusConfig>,
12 pub tracing: Option<TracingConfig>,
13 pub otel: Option<crate::observability::otel::OtelConfig>,
14}
15
16#[cfg(feature = "observability-install")]
20#[derive(Debug, Clone, Copy, PartialEq, Eq)]
21pub(crate) enum MetricsMode {
22 None,
23 PrometheusOnly,
24 OtelOnly,
25 Fanout,
26}
27
28#[cfg(feature = "observability-install")]
29impl MetricsMode {
30 pub(crate) fn select(prometheus: bool, otel_metrics: bool) -> Self {
31 match (prometheus, otel_metrics) {
32 (true, true) => MetricsMode::Fanout,
33 (true, false) => MetricsMode::PrometheusOnly,
34 (false, true) => MetricsMode::OtelOnly,
35 (false, false) => MetricsMode::None,
36 }
37 }
38}
39
40#[derive(Debug, Clone)]
41pub struct PrometheusConfig {
42 pub listen: String,
45 pub buckets: Option<Vec<f64>>,
48}
49
50#[derive(Debug, Clone)]
51pub struct TracingConfig {
52 pub level: String,
54}
55
56#[derive(Debug, Clone, Default)]
59pub struct InstallReport {
60 pub prometheus_listen: Option<String>,
61 pub prometheus_already_installed: bool,
62 pub tracing_already_installed: bool,
63 pub otel_installed: bool,
64 pub otel_signals: Vec<&'static str>,
65}
66
67#[derive(Debug, Error)]
68pub enum InstallError {
69 #[error("failed to bind Prometheus listener at {listen}: {source}")]
70 PrometheusBind {
71 listen: String,
72 #[source]
73 source: std::io::Error,
74 },
75 #[error("failed to install Prometheus recorder: {0}")]
76 PrometheusInstall(String),
77}
78
79#[cfg(feature = "observability-install")]
92pub fn install_observability(cfg: &ObservabilityConfig) -> Result<InstallReport, InstallError> {
93 let mut report = InstallReport::default();
94
95 #[cfg(feature = "otel")]
98 let mut otel_tracer: Option<opentelemetry_sdk::trace::SdkTracerProvider> = None;
99
100 #[cfg(feature = "otel")]
102 let otel_metrics = cfg
103 .otel
104 .as_ref()
105 .map(|o| o.exports(crate::observability::otel::OtelSignal::Metrics))
106 .unwrap_or(false);
107 #[cfg(not(feature = "otel"))]
108 let otel_metrics = false;
109
110 let mode = MetricsMode::select(cfg.prometheus.is_some(), otel_metrics);
111
112 #[cfg(feature = "otel")]
114 let mut otel_meter: Option<opentelemetry_sdk::metrics::SdkMeterProvider> = None;
115
116 match mode {
117 MetricsMode::None => {}
118 MetricsMode::PrometheusOnly => {
119 install_prometheus_only(cfg.prometheus.as_ref().unwrap(), &mut report)?;
120 }
121 #[cfg(feature = "otel")]
122 MetricsMode::OtelOnly => {
123 if let Some(otel) = cfg.otel.as_ref() {
124 match crate::observability::otel::build_meter_provider(otel) {
125 Ok((mp, recorder)) => {
126 if metrics::set_global_recorder(recorder).is_err() {
127 tracing::warn!("metrics recorder already installed; continuing");
128 report.prometheus_already_installed = true;
129 } else {
130 otel_meter = Some(mp);
131 report.otel_signals.push("metrics");
132 }
133 }
134 Err(e) => tracing::warn!("OTLP metrics exporter init failed; skipping: {e}"),
135 }
136 }
137 }
138 #[cfg(feature = "otel")]
139 MetricsMode::Fanout => {
140 install_fanout(
141 cfg.prometheus.as_ref().unwrap(),
142 cfg.otel.as_ref().unwrap(),
143 &mut report,
144 &mut otel_meter,
145 )?;
146 }
147 #[cfg(not(feature = "otel"))]
148 MetricsMode::OtelOnly | MetricsMode::Fanout => {
149 unreachable!("otel metrics mode selected without the otel feature")
150 }
151 }
152
153 if let Some(t) = cfg.tracing.as_ref() {
155 use tracing_subscriber::EnvFilter;
156 use tracing_subscriber::layer::SubscriberExt;
157 use tracing_subscriber::util::SubscriberInitExt;
158
159 let make_filter =
160 || EnvFilter::try_new(&t.level).unwrap_or_else(|_| EnvFilter::new("info"));
161
162 #[cfg_attr(not(feature = "otel"), allow(unused_mut))]
165 let mut installed = false;
166
167 #[cfg(feature = "otel")]
168 {
169 let otel_traces = cfg
170 .otel
171 .as_ref()
172 .map(|o| o.exports(crate::observability::otel::OtelSignal::Traces))
173 .unwrap_or(false);
174 if let (true, Some(otel)) = (otel_traces, cfg.otel.as_ref()) {
175 match crate::observability::otel::build_trace_provider(otel) {
176 Ok(tp) => {
177 use opentelemetry::trace::TracerProvider as _;
178 let tracer = tp.tracer("faucet");
179 crate::observability::otel::install_propagator();
180 let reg = tracing_subscriber::registry()
181 .with(make_filter())
182 .with(tracing_subscriber::fmt::layer())
183 .with(tracing_opentelemetry::layer().with_tracer(tracer))
184 .with(crate::observability::otel::OtelErrorCountLayer);
185 if reg.try_init().is_err() {
186 tracing::warn!("tracing subscriber already installed; continuing");
187 report.tracing_already_installed = true;
188 } else {
191 report.otel_signals.push("traces");
192 otel_tracer = Some(tp);
193 }
194 installed = true;
195 }
196 Err(e) => tracing::warn!("OTLP trace exporter init failed; logs-only: {e}"),
197 }
198 }
199 }
200
201 if !installed {
202 #[cfg(feature = "otel")]
208 let otel_error_layer = cfg
209 .otel
210 .as_ref()
211 .map(|o| {
212 o.exports(crate::observability::otel::OtelSignal::Traces)
213 || o.exports(crate::observability::otel::OtelSignal::Metrics)
214 })
215 .unwrap_or(false)
216 .then_some(crate::observability::otel::OtelErrorCountLayer);
217 #[cfg(not(feature = "otel"))]
218 let otel_error_layer: Option<tracing_subscriber::layer::Identity> = None;
219
220 let reg = tracing_subscriber::registry()
221 .with(make_filter())
222 .with(tracing_subscriber::fmt::layer())
223 .with(otel_error_layer);
224 if reg.try_init().is_err() {
225 tracing::warn!("tracing subscriber already installed; continuing");
229 report.tracing_already_installed = true;
230 }
231 }
232 }
233
234 #[cfg(feature = "otel")]
236 {
237 if otel_tracer.is_some() || otel_meter.is_some() {
238 crate::observability::otel::describe();
239 let _ = crate::observability::otel::set_guard(crate::observability::otel::OtelGuard {
240 tracer: otel_tracer,
241 meter: otel_meter,
242 });
243 report.otel_installed = true;
244 }
245 }
246
247 crate::observability::resilience::describe();
251 crate::observability::cleanup::describe();
252 crate::observability::drift::describe();
253 crate::observability::describe_roundtrip_metrics();
254 register_build_info();
255
256 Ok(report)
257}
258
259#[cfg(feature = "observability-install")]
263fn install_prometheus_only(
264 p: &PrometheusConfig,
265 report: &mut InstallReport,
266) -> Result<(), InstallError> {
267 use metrics_exporter_prometheus::{BuildError, PrometheusBuilder};
268
269 let listen: std::net::SocketAddr =
270 p.listen
271 .parse()
272 .map_err(|e: std::net::AddrParseError| InstallError::PrometheusBind {
273 listen: p.listen.clone(),
274 source: std::io::Error::new(std::io::ErrorKind::InvalidInput, e.to_string()),
275 })?;
276
277 const DEFAULT_BUCKETS: &[f64] = &[
278 0.001, 0.005, 0.01, 0.05, 0.1, 0.5, 1.0, 5.0, 10.0, 30.0, 60.0, 300.0,
279 ];
280 let buckets = p.buckets.as_deref().unwrap_or(DEFAULT_BUCKETS);
281
282 let builder = PrometheusBuilder::new()
283 .with_http_listener(listen)
284 .set_buckets(buckets)
285 .map_err(|e| InstallError::PrometheusInstall(e.to_string()))?;
286
287 match builder.install() {
288 Ok(()) => report.prometheus_listen = Some(p.listen.clone()),
289 Err(e) => match e {
293 BuildError::FailedToSetGlobalRecorder(_) => {
296 tracing::warn!("Prometheus recorder already installed; continuing");
297 report.prometheus_already_installed = true;
298 }
299 BuildError::FailedToCreateHTTPListener(msg) => {
305 return Err(InstallError::PrometheusBind {
306 listen: p.listen.clone(),
307 source: std::io::Error::other(msg),
308 });
309 }
310 other => return Err(InstallError::PrometheusInstall(other.to_string())),
311 },
312 }
313 Ok(())
314}
315
316#[cfg(all(feature = "observability-install", feature = "otel"))]
322fn install_fanout(
323 p: &PrometheusConfig,
324 otel: &crate::observability::otel::OtelConfig,
325 report: &mut InstallReport,
326 otel_meter: &mut Option<opentelemetry_sdk::metrics::SdkMeterProvider>,
327) -> Result<(), InstallError> {
328 use metrics_exporter_prometheus::{BuildError, PrometheusBuilder};
329 use metrics_util::layers::FanoutBuilder;
330
331 let listen: std::net::SocketAddr =
332 p.listen
333 .parse()
334 .map_err(|e: std::net::AddrParseError| InstallError::PrometheusBind {
335 listen: p.listen.clone(),
336 source: std::io::Error::new(std::io::ErrorKind::InvalidInput, e.to_string()),
337 })?;
338 const DEFAULT_BUCKETS: &[f64] = &[
339 0.001, 0.005, 0.01, 0.05, 0.1, 0.5, 1.0, 5.0, 10.0, 30.0, 60.0, 300.0,
340 ];
341 let buckets = p.buckets.as_deref().unwrap_or(DEFAULT_BUCKETS);
342
343 let (prom_recorder, prom_exporter) = PrometheusBuilder::new()
346 .with_http_listener(listen)
347 .set_buckets(buckets)
348 .map_err(|e| InstallError::PrometheusInstall(e.to_string()))?
349 .build()
350 .map_err(|e| match e {
351 BuildError::FailedToCreateHTTPListener(msg) => InstallError::PrometheusBind {
352 listen: p.listen.clone(),
353 source: std::io::Error::other(msg),
354 },
355 other => InstallError::PrometheusInstall(other.to_string()),
356 })?;
357
358 match crate::observability::otel::build_meter_provider(otel) {
359 Ok((mp, otel_recorder)) => {
360 let fanout = FanoutBuilder::default()
361 .add_recorder(prom_recorder)
362 .add_recorder(otel_recorder)
363 .build();
364 if metrics::set_global_recorder(fanout).is_err() {
365 tracing::warn!("metrics recorder already installed; continuing");
366 report.prometheus_already_installed = true;
367 } else {
368 report.prometheus_listen = Some(p.listen.clone());
369 report.otel_signals.push("metrics");
370 *otel_meter = Some(mp);
371 tokio::spawn(prom_exporter);
372 }
373 }
374 Err(e) => {
375 tracing::warn!("OTLP metrics exporter init failed; Prometheus-only: {e}");
378 if metrics::set_global_recorder(prom_recorder).is_err() {
379 report.prometheus_already_installed = true;
380 } else {
381 report.prometheus_listen = Some(p.listen.clone());
382 tokio::spawn(prom_exporter);
383 }
384 }
385 }
386 Ok(())
387}
388
389#[cfg(not(feature = "observability-install"))]
391pub fn install_observability(_cfg: &ObservabilityConfig) -> Result<InstallReport, InstallError> {
392 crate::observability::resilience::describe();
393 crate::observability::cleanup::describe();
394 crate::observability::drift::describe();
395 crate::observability::describe_roundtrip_metrics();
396 crate::observability::otel::describe();
397 register_build_info();
398 Ok(InstallReport::default())
399}
400
401pub fn register_build_info() {
411 metrics::gauge!(
412 "faucet_build_info",
413 "version" => env!("CARGO_PKG_VERSION"),
414 )
415 .set(1.0);
416}
417
418#[cfg(all(test, feature = "observability-install"))]
419mod tests {
420 use super::*;
421 use std::sync::Mutex;
422
423 static LOCK: Mutex<()> = Mutex::new(());
424
425 #[test]
426 fn metrics_mode_selection() {
427 use super::MetricsMode;
428 assert_eq!(MetricsMode::select(true, true), MetricsMode::Fanout);
429 assert_eq!(
430 MetricsMode::select(true, false),
431 MetricsMode::PrometheusOnly
432 );
433 assert_eq!(MetricsMode::select(false, true), MetricsMode::OtelOnly);
434 assert_eq!(MetricsMode::select(false, false), MetricsMode::None);
435 }
436
437 #[test]
438 fn no_config_returns_empty_report() {
439 let _g = LOCK.lock().unwrap_or_else(|e| e.into_inner());
440 let r = install_observability(&ObservabilityConfig::default()).unwrap();
441 assert!(r.prometheus_listen.is_none());
442 assert!(!r.prometheus_already_installed);
443 assert!(!r.tracing_already_installed);
444 }
445
446 #[test]
447 fn malformed_listen_returns_bind_error() {
448 let _g = LOCK.lock().unwrap_or_else(|e| e.into_inner());
449 let cfg = ObservabilityConfig {
450 prometheus: Some(PrometheusConfig {
451 listen: "not-a-socket".into(),
452 buckets: None,
453 }),
454 tracing: None,
455 otel: None,
456 };
457 match install_observability(&cfg) {
458 Err(InstallError::PrometheusBind { .. }) => {}
459 other => panic!("expected PrometheusBind error, got {other:?}"),
460 }
461 }
462
463 #[test]
464 fn register_build_info_is_callable_and_idempotent() {
465 let _g = LOCK.lock().unwrap_or_else(|e| e.into_inner());
466 register_build_info();
469 register_build_info();
470 }
471
472 #[test]
473 fn install_prometheus_and_tracing_returns_ok() {
474 let _g = LOCK.lock().unwrap_or_else(|e| e.into_inner());
480 let cfg = ObservabilityConfig {
481 prometheus: Some(PrometheusConfig {
482 listen: "127.0.0.1:0".into(),
483 buckets: Some(vec![0.01, 0.1, 1.0]),
485 }),
486 tracing: Some(TracingConfig {
487 level: "info".into(),
488 }),
489 otel: None,
490 };
491 let report = install_observability(&cfg).expect("install must return Ok");
492 assert!(
494 report.prometheus_listen.is_some() || report.prometheus_already_installed,
495 "prometheus install must either bind or report already-installed"
496 );
497 }
498
499 #[test]
500 fn install_tracing_with_invalid_directive_falls_back_to_info() {
501 let _g = LOCK.lock().unwrap_or_else(|e| e.into_inner());
504 let cfg = ObservabilityConfig {
505 prometheus: None,
506 tracing: Some(TracingConfig {
507 level: "this is !!! not a valid filter".into(),
509 }),
510 otel: None,
511 };
512 install_observability(&cfg).expect("invalid tracing directive must not fail install");
513 }
514}