1use std::env;
2
3use opentelemetry::{
4 global::{self, BoxedTracer},
5 trace::TracerProvider,
6 KeyValue,
7};
8use opentelemetry_otlp::ExporterBuildError;
9use opentelemetry_sdk::{
10 metrics::{MeterProviderBuilder, PeriodicReader, SdkMeterProvider},
11 trace::{RandomIdGenerator, Sampler, SdkTracerProvider},
12 Resource,
13};
14use opentelemetry_semantic_conventions::{
15 attribute::{DEPLOYMENT_ENVIRONMENT_NAME, SERVICE_NAME, SERVICE_VERSION},
16 SCHEMA_URL,
17};
18use tracing::Dispatch;
19use tracing::{debug, Level};
20use tracing_core::LevelFilter;
21use tracing_opentelemetry::MetricsLayer;
22use tracing_opentelemetry::OpenTelemetryLayer;
23use tracing_subscriber::fmt;
24use tracing_subscriber::layer::SubscriberExt;
25use tracing_subscriber::prelude::*;
26use tracing_subscriber::Registry;
27use tracing_log::LogTracer;
28
29use crate::prelude::ApplicationData;
30static PHLOW_OTEL_ACTIVE: once_cell::sync::Lazy<bool> =
32 once_cell::sync::Lazy::new(|| match std::env::var("PHLOW_OTEL") {
33 Ok(active) => active.parse::<bool>().unwrap_or(false),
34 Err(_) => false,
35 });
36
37static PHLOW_SPAN_ACTIVE: once_cell::sync::Lazy<Level> =
38 once_cell::sync::Lazy::new(|| match std::env::var("PHLOW_SPAN") {
39 Ok(level) => level.parse::<Level>().unwrap_or(Level::INFO),
40 Err(_) => Level::INFO,
41 });
42
43static PHLOW_LOG: once_cell::sync::Lazy<Level> =
44 once_cell::sync::Lazy::new(|| match std::env::var("PHLOW_LOG") {
45 Ok(level) => level.parse::<Level>().unwrap_or(Level::INFO),
46 Err(_) => Level::INFO,
47 });
48
49fn resource(app_data: ApplicationData) -> Resource {
50 let service_name = env::var("OTEL_SERVICE_NAME")
51 .unwrap_or_else(|_| app_data.name.unwrap_or_else(|| "phlow".to_string()));
52 let service_version = env::var("OTEL_SERVICE_VERSION").unwrap_or_else(|_| {
53 app_data.version.unwrap_or_else(|| {
54 env::var("PHLOW_VERSION")
55 .unwrap_or_else(|_| "".to_string())
56 .to_string()
57 })
58 });
59 let deployment_environment_name =
60 env::var("OTEL_DEPLOYMENT_ENVIRONMENT_NAME").unwrap_or_else(|_| {
61 app_data.environment.unwrap_or_else(|| {
62 env::var("PHLOW_ENV")
63 .unwrap_or_else(|_| "development".to_string())
64 .to_string()
65 })
66 });
67
68 let attributes = vec![
69 KeyValue::new(SCHEMA_URL, "https://opentelemetry.io/schemas/1.4.0"),
70 KeyValue::new(SERVICE_NAME, service_name),
71 KeyValue::new(SERVICE_VERSION, service_version),
72 KeyValue::new(DEPLOYMENT_ENVIRONMENT_NAME, deployment_environment_name),
73 ];
74
75 Resource::builder()
76 .with_schema_url(
77 [
78 KeyValue::new(SERVICE_NAME, env!("CARGO_PKG_NAME")),
79 KeyValue::new(SERVICE_VERSION, env!("CARGO_PKG_VERSION")),
80 KeyValue::new(DEPLOYMENT_ENVIRONMENT_NAME, "develop"),
81 ],
82 SCHEMA_URL,
83 )
84 .with_attributes(attributes)
85 .build()
86}
87
88pub fn get_tracer() -> BoxedTracer {
89 global::tracer("phlow-tracing-otel-subscriber")
90}
91
92fn init_meter_provider(resource: Resource) -> Result<SdkMeterProvider, ExporterBuildError> {
93 let exporter = opentelemetry_otlp::MetricExporter::builder()
94 .with_http()
95 .with_temporality(opentelemetry_sdk::metrics::Temporality::default())
96 .build()?;
97 let reader = PeriodicReader::builder(exporter)
98 .with_interval(std::time::Duration::from_secs(30))
99 .build();
100
101 let stdout_reader =
103 PeriodicReader::builder(opentelemetry_stdout::MetricExporter::default()).build();
104
105 let meter_provider = MeterProviderBuilder::default()
106 .with_resource(resource)
107 .with_reader(reader)
108 .with_reader(stdout_reader)
109 .build();
110
111 global::set_meter_provider(meter_provider.clone());
112
113 Ok(meter_provider)
114}
115
116fn init_tracer_provider(resource: Resource) -> Result<SdkTracerProvider, ExporterBuildError> {
117 let exporter = if env::var("OTEL_EXPORTER_OTLP_PROTOCOL") == Ok("grpc".to_string()) {
118 opentelemetry_otlp::SpanExporter::builder()
119 .with_tonic()
120 .build()?
121 } else {
122 opentelemetry_otlp::SpanExporter::builder()
123 .with_http()
124 .build()?
125 };
126
127 Ok(SdkTracerProvider::builder()
128 .with_sampler(Sampler::ParentBased(Box::new(Sampler::TraceIdRatioBased(
130 1.0,
131 ))))
132 .with_id_generator(RandomIdGenerator::default())
134 .with_resource(resource)
135 .with_batch_exporter(exporter)
136 .build())
137}
138
139pub fn get_log_level() -> Level {
140 *PHLOW_LOG
141}
142
143fn get_span_level() -> Level {
144 *PHLOW_SPAN_ACTIVE
145}
146
147pub fn get_otel_active() -> bool {
148 *PHLOW_OTEL_ACTIVE
149}
150
151pub fn init_tracing_subscriber(app_data: ApplicationData) -> OtelGuard {
152 LogTracer::init().unwrap_or_else(|_| {
154 });
156
157 if !get_otel_active() {
158 let subscriber = Registry::default()
159 .with(fmt::layer().with_filter(LevelFilter::from_level(get_log_level())));
160
161 let dispatch = Dispatch::new(subscriber);
162
163 debug!("PHLOW_OTEL is set to false, using default subscriber");
164
165 return OtelGuard {
166 tracer_provider: None,
167 meter_provider: None,
168 dispatch,
169 };
170 }
171
172 let resource = resource(app_data);
173 let tracer_provider = init_tracer_provider(resource.clone()).ok();
174 let meter_provider = init_meter_provider(resource.clone()).ok();
175
176 if let (Some(tp), Some(mp)) = (&tracer_provider, &meter_provider) {
177 let tracer = tp.tracer("tracing-otel-subscriber");
178
179 let fmt_layer = fmt::layer().with_filter(LevelFilter::from_level(get_log_level()));
180 let otel_layer =
181 OpenTelemetryLayer::new(tracer).with_filter(LevelFilter::from_level(get_span_level()));
182 let metrics_layer = MetricsLayer::new(mp.clone());
183
184 let subscriber = Registry::default()
185 .with(fmt_layer)
186 .with(otel_layer)
187 .with(metrics_layer);
188
189 let dispatch = Dispatch::new(subscriber);
190
191 debug!("OpenTelemetry provider found, using OpenTelemetry subscriber");
192
193 OtelGuard {
194 tracer_provider: Some(tp.clone()),
195 meter_provider: Some(mp.clone()),
196 dispatch,
197 }
198 } else {
199 let subscriber = Registry::default()
200 .with(fmt::layer().with_filter(LevelFilter::from_level(get_log_level())));
201 let dispatch = Dispatch::new(subscriber);
202
203 debug!("No OpenTelemetry provider found, using default subscriber");
204
205 OtelGuard {
206 tracer_provider: None,
207 meter_provider: None,
208 dispatch,
209 }
210 }
211}
212
213pub struct OtelGuard {
214 pub tracer_provider: Option<SdkTracerProvider>,
215 pub meter_provider: Option<SdkMeterProvider>,
216 pub dispatch: Dispatch,
217}
218
219impl Drop for OtelGuard {
220 fn drop(&mut self) {
221 if let Some(tracer_provider) = &self.tracer_provider {
222 if let Err(err) = tracer_provider.shutdown() {
223 eprintln!("{err:?}");
224 }
225 }
226 if let Some(meter_provider) = &self.meter_provider {
227 if let Err(err) = meter_provider.shutdown() {
228 eprintln!("{err:?}");
229 }
230 }
231 }
232}