Skip to main content

nexus_common/
stack.rs

1use crate::db::{Neo4jConnector, RedisConnector};
2use crate::types::DynError;
3use crate::{Level, StackConfig};
4use opentelemetry::{global, KeyValue};
5use opentelemetry_appender_tracing::layer::OpenTelemetryTracingBridge;
6use opentelemetry_otlp::{LogExporter, MetricExporter, SpanExporter, WithExportConfig};
7use opentelemetry_sdk::logs::SdkLoggerProvider;
8use opentelemetry_sdk::metrics::{PeriodicReader, SdkMeterProvider};
9use opentelemetry_sdk::trace::SdkTracerProvider;
10use opentelemetry_sdk::Resource;
11use std::time::Duration;
12use tracing::{error, info};
13use tracing_subscriber::{fmt, EnvFilter, Layer};
14use tracing_subscriber::{layer::SubscriberExt, Registry};
15
16pub struct StackManager {}
17
18impl StackManager {
19    pub async fn setup(name: &str, config: &StackConfig) -> Result<(), DynError> {
20        // Initialize logging and metrics
21        Self::setup_logging(name, &config.otlp_endpoint, config.log_level).await;
22        Self::setup_metrics(name, &config.otlp_endpoint).await;
23
24        // Initialize Redis and Neo4j
25        RedisConnector::init(&config.db.redis).await?;
26        Neo4jConnector::init(&config.db.neo4j).await?;
27        Ok(())
28    }
29
30    async fn setup_logging(service_name: &str, otel_endpoint: &Option<String>, log_level: Level) {
31        match otel_endpoint {
32            None => Self::setup_local_logging(log_level),
33            Some(endpoint) => {
34                match Self::setup_otlp_logging(service_name, endpoint, log_level).await {
35                    Ok(()) => info!(
36                        "OpenTelemetry Logging initialized for {} service",
37                        service_name
38                    ),
39                    Err(e) => error!("Failed to initialize OpenTelemetry Logging: {:?}", e),
40                }
41            }
42        }
43    }
44
45    fn setup_local_logging(log_level: Level) {
46        // Enable log-to-tracing bridge so that `log`-based crates (e.g., neo4rs) emit through our `tracing` subscriber
47        let _ = tracing_log::LogTracer::init();
48
49        // Build an env‐based filter
50        let env_filter =
51            EnvFilter::new(log_level.as_str()).add_directive("mainline=info".parse().unwrap());
52
53        // Create a formatting layer
54        let fmt_layer = fmt::layer().compact().with_line_number(true);
55
56        // Compose the subscriber
57        let subscriber = Registry::default().with(env_filter).with(fmt_layer);
58
59        if tracing::subscriber::set_global_default(subscriber).is_ok() {
60            tracing::info!("Local application logging initialized");
61        }
62    }
63
64    async fn setup_otlp_logging(
65        service_name: &str,
66        otel_endpoint: &String,
67        log_level: Level,
68    ) -> Result<(), Box<dyn std::error::Error>> {
69        // TODO: Add local tracer, https://github.com/pubky/pubky-nexus/issues/356
70        // Set up OpenTelemetry Tracer (Spans)
71        let tracing_exporter = SpanExporter::builder()
72            .with_tonic()
73            .with_endpoint(otel_endpoint.clone())
74            .with_timeout(Duration::from_secs(3))
75            .build()
76            .map_err(|e| format!("OTLP Tracing Exporter Error: {e}"))?;
77
78        // Collects spans in memory and sends them in batches
79        let tracer_provider = SdkTracerProvider::builder()
80            .with_resource(Self::create_resource(service_name))
81            .with_batch_exporter(tracing_exporter)
82            .build();
83
84        // Registers OpenTelemetry as the global tracing provider
85        // Ensures that all spans created in the app are processed and exported to an OTLP backend (signoz or jaeger)
86        global::set_tracer_provider(tracer_provider.clone());
87
88        // Set up OpenTelemetry Logging
89        let logging_exporter = LogExporter::builder()
90            .with_tonic()
91            .with_endpoint(otel_endpoint.clone())
92            .with_timeout(Duration::from_secs(3))
93            .build()
94            .map_err(|e| format!("OTLP Logging Exporter Error: {e}"))?;
95
96        let logging_provider = SdkLoggerProvider::builder()
97            .with_resource(Self::create_resource(service_name))
98            .with_batch_exporter(logging_exporter)
99            .build();
100
101        // Apply log filters for verbosity control
102        // This ensures only relevant logs are sent to OpenTelemetry, reducing unnecessary data transmission
103        let otlp_layer = OpenTelemetryTracingBridge::new(&logging_provider).with_filter(
104            EnvFilter::new(log_level.as_str())
105                .add_directive("opentelemetry=error".parse().unwrap())
106                .add_directive("h2=error".parse().unwrap())
107                .add_directive("tower=info".parse().unwrap())
108                .add_directive("mainline=info".parse().unwrap()),
109        );
110
111        // Configure the stdout logging layer
112        let stdout_layer = fmt::layer().compact().with_line_number(true).with_filter(
113            EnvFilter::new(log_level.as_str())
114                .add_directive("opentelemetry=error".parse().unwrap())
115                .add_directive("h2=error".parse().unwrap())
116                .add_directive("tower=info".parse().unwrap())
117                .add_directive("mainline=info".parse().unwrap()),
118        );
119
120        // Creates a tracing subscriber
121        let subscriber = Registry::default().with(stdout_layer).with(otlp_layer);
122
123        // Registers a global tracing subscriber that captures logs
124        // TODO: If multiple services run in the same process, only the first call to
125        // `tracing::subscriber::set_global_default(...)` succeeds. That means the logs from
126        // all services will be emitted under the first service's `service_name`
127        // This happens because tracing subscribers and OTEL logger providers are global.
128        // To fix this, use per-task `Dispatch` with isolated subscribers, or run services
129        // in separate processes.
130        if tracing::subscriber::set_global_default(subscriber).is_ok() {
131            info!(
132                "OpenTelemetry endpoint listening on (OTLP_ENDPOINT={})",
133                otel_endpoint
134            );
135        } else {
136            error!("Failed to initialize OpenTelemetry Logging: Already set globally!");
137        }
138
139        Ok(())
140    }
141
142    async fn setup_metrics(service_name: &str, otel_endpoint: &Option<String>) {
143        match otel_endpoint {
144            None => info!("Metrics collection is disabled. No metrics will be exported"),
145            Some(endpoint) => {
146                // Configure the exporter to collect and send metrics to an OTLP
147                let metric_exporter = MetricExporter::builder()
148                    .with_tonic()
149                    .with_endpoint(endpoint.clone())
150                    .with_timeout(Duration::from_secs(3))
151                    .build()
152                    .expect("Failed to create OTLP metric exporter");
153
154                // Create a periodic metrics reader that collects and exports metrics at a fixed interval
155                let reader = PeriodicReader::builder(metric_exporter)
156                    .with_interval(std::time::Duration::from_secs(30))
157                    .build();
158
159                // Createa Meter Provider, which is responsible for managing and exporting metrics
160                let provider = SdkMeterProvider::builder()
161                    .with_resource(Self::create_resource(service_name))
162                    .with_reader(reader)
163                    .build();
164
165                // Register globally the metrics
166                global::set_meter_provider(provider);
167
168                info!(
169                    "OpenTelemetry Metrics initialized for {} service",
170                    service_name
171                );
172            }
173        }
174    }
175
176    fn create_resource(service_name: &str) -> Resource {
177        Resource::builder_empty()
178            .with_attribute(KeyValue::new("service.name", String::from(service_name)))
179            .build()
180    }
181}