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 Self::setup_logging(name, &config.otlp_endpoint, config.log_level).await;
22 Self::setup_metrics(name, &config.otlp_endpoint).await;
23
24 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 let _ = tracing_log::LogTracer::init();
48
49 let env_filter =
51 EnvFilter::new(log_level.as_str()).add_directive("mainline=info".parse().unwrap());
52
53 let fmt_layer = fmt::layer().compact().with_line_number(true);
55
56 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 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 let tracer_provider = SdkTracerProvider::builder()
80 .with_resource(Self::create_resource(service_name))
81 .with_batch_exporter(tracing_exporter)
82 .build();
83
84 global::set_tracer_provider(tracer_provider.clone());
87
88 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 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 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 let subscriber = Registry::default().with(stdout_layer).with(otlp_layer);
122
123 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 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 let reader = PeriodicReader::builder(metric_exporter)
156 .with_interval(std::time::Duration::from_secs(30))
157 .build();
158
159 let provider = SdkMeterProvider::builder()
161 .with_resource(Self::create_resource(service_name))
162 .with_reader(reader)
163 .build();
164
165 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}