Skip to main content

phlow_sdk/
otel.rs

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;
30// otel active
31static 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    // For debugging in development
102    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        // Customize sampling strategy
129        .with_sampler(Sampler::ParentBased(Box::new(Sampler::TraceIdRatioBased(
130            1.0,
131        ))))
132        // If export trace to AWS X-Ray, you can use XrayIdGenerator
133        .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    // Initialize the log-tracing bridge to capture log! macros
153    LogTracer::init().unwrap_or_else(|_| {
154        // Bridge already initialized, ignore
155    });
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}