Skip to content

Commit f97deb2

Browse files
committed
feat: added open telemetry support
1 parent 5b82c9c commit f97deb2

3 files changed

Lines changed: 357 additions & 0 deletions

File tree

src/flow_telemetry/errors.rs

Lines changed: 145 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,145 @@
1+
use std::{
2+
backtrace::Backtrace,
3+
error::Error,
4+
sync::atomic::{AtomicBool, Ordering},
5+
};
6+
7+
static CAPTURE_BACKTRACES: AtomicBool = AtomicBool::new(false);
8+
9+
pub const EXCEPTION_TARGET: &str = "telemetry::exception";
10+
pub const SUMMARY_TARGET: &str = "telemetry::error_summary";
11+
12+
pub fn enable_backtraces() {
13+
CAPTURE_BACKTRACES.store(true, Ordering::Relaxed);
14+
}
15+
16+
pub fn record(
17+
category: &'static str,
18+
operation: &'static str,
19+
error: &(dyn Error + 'static),
20+
context: impl AsRef<str>,
21+
) {
22+
let message = error.to_string();
23+
let chain = error_chain(error);
24+
25+
tracing::error!(
26+
target: SUMMARY_TARGET,
27+
error_category = category,
28+
operation,
29+
error_message = message.as_str(),
30+
error_context = context.as_ref(),
31+
"operation failed"
32+
);
33+
34+
record_exception(category, operation, &message, &chain, context.as_ref());
35+
}
36+
37+
pub fn record_message(
38+
category: &'static str,
39+
operation: &'static str,
40+
message: impl AsRef<str>,
41+
context: impl AsRef<str>,
42+
) {
43+
let message = message.as_ref();
44+
tracing::error!(
45+
target: SUMMARY_TARGET,
46+
error_category = category,
47+
operation,
48+
error_message = message,
49+
error_context = context.as_ref(),
50+
"operation failed"
51+
);
52+
53+
record_exception(category, operation, message, "", context.as_ref());
54+
}
55+
56+
pub fn panic(message: &str, location: &str) {
57+
tracing::error!(
58+
target: SUMMARY_TARGET,
59+
error_category = "panic",
60+
operation = "process",
61+
error_message = message,
62+
error_context = location,
63+
"process panicked"
64+
);
65+
66+
record_exception("panic", "process", message, "", location);
67+
}
68+
69+
fn error_chain(error: &(dyn Error + 'static)) -> String {
70+
let mut messages = Vec::new();
71+
let mut current = error.source();
72+
73+
while let Some(source) = current {
74+
messages.push(source.to_string());
75+
current = source.source();
76+
}
77+
78+
messages.join(": ")
79+
}
80+
81+
fn record_exception(
82+
category: &'static str,
83+
operation: &'static str,
84+
message: &str,
85+
chain: &str,
86+
context: &str,
87+
) {
88+
if !CAPTURE_BACKTRACES.load(Ordering::Relaxed) {
89+
return;
90+
}
91+
92+
tracing::error!(
93+
target: EXCEPTION_TARGET,
94+
error_category = category,
95+
operation,
96+
error_message = message,
97+
error_chain = chain,
98+
error_context = context,
99+
exception_stacktrace = %Backtrace::force_capture(),
100+
"exception"
101+
);
102+
}
103+
104+
#[cfg(test)]
105+
mod tests {
106+
use std::{error::Error, fmt};
107+
108+
use super::error_chain;
109+
110+
#[derive(Debug)]
111+
struct TestError {
112+
message: &'static str,
113+
source: Option<Box<TestError>>,
114+
}
115+
116+
impl fmt::Display for TestError {
117+
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
118+
formatter.write_str(self.message)
119+
}
120+
}
121+
122+
impl Error for TestError {
123+
fn source(&self) -> Option<&(dyn Error + 'static)> {
124+
self.source
125+
.as_deref()
126+
.map(|source| source as &(dyn Error + 'static))
127+
}
128+
}
129+
130+
#[test]
131+
fn formats_source_chain_from_nearest_to_root() {
132+
let error = TestError {
133+
message: "outer",
134+
source: Some(Box::new(TestError {
135+
message: "middle",
136+
source: Some(Box::new(TestError {
137+
message: "root",
138+
source: None,
139+
})),
140+
})),
141+
};
142+
143+
assert_eq!(error_chain(&error), "middle: root");
144+
}
145+
}

src/flow_telemetry/mod.rs

Lines changed: 209 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,209 @@
1+
pub mod errors;
2+
use std::error::Error;
3+
4+
use opentelemetry::{KeyValue, global, trace::TracerProvider as _};
5+
use opentelemetry_appender_tracing::layer::OpenTelemetryTracingBridge;
6+
use opentelemetry_otlp::WithExportConfig;
7+
use opentelemetry_sdk::{
8+
Resource,
9+
logs::SdkLoggerProvider,
10+
metrics::{PeriodicReader, SdkMeterProvider},
11+
propagation::TraceContextPropagator,
12+
trace::SdkTracerProvider,
13+
};
14+
use serde::{Deserialize, Serialize};
15+
use tracing_subscriber::{
16+
EnvFilter, Layer, filter::filter_fn, layer::SubscriberExt, util::SubscriberInitExt,
17+
};
18+
19+
#[derive(Clone, Debug, Deserialize, Serialize)]
20+
#[serde(default)]
21+
pub struct OpenTelemetry {
22+
pub enabled: bool,
23+
pub service_name: String,
24+
pub logs_endpoint: Option<String>,
25+
pub metrics_endpoint: Option<String>,
26+
pub traces_endpoint: Option<String>,
27+
}
28+
29+
impl OpenTelemetry {
30+
pub fn logs_endpoint(&self) -> Option<&str> {
31+
non_empty_url(&self.logs_endpoint)
32+
}
33+
34+
pub fn metrics_endpoint(&self) -> Option<&str> {
35+
non_empty_url(&self.metrics_endpoint)
36+
}
37+
38+
pub fn traces_endpoint(&self) -> Option<&str> {
39+
non_empty_url(&self.traces_endpoint)
40+
}
41+
42+
pub fn has_enabled_exporter(&self) -> bool {
43+
self.logs_endpoint().is_some()
44+
|| self.metrics_endpoint().is_some()
45+
|| self.traces_endpoint().is_some()
46+
}
47+
}
48+
49+
impl Default for OpenTelemetry {
50+
fn default() -> Self {
51+
Self {
52+
enabled: false,
53+
service_name: env!("CARGO_PKG_NAME").into(),
54+
logs_endpoint: None,
55+
metrics_endpoint: None,
56+
traces_endpoint: None,
57+
}
58+
}
59+
}
60+
61+
pub struct TelemetrySettings<'a> {
62+
pub environment: &'a str,
63+
pub default_log_level: &'a str,
64+
pub service_version: &'a str,
65+
pub instrumentation_name: &'a str,
66+
pub initialize_metrics: Option<fn()>,
67+
}
68+
69+
pub struct Telemetry {
70+
logger_provider: Option<SdkLoggerProvider>,
71+
meter_provider: Option<SdkMeterProvider>,
72+
tracer_provider: Option<SdkTracerProvider>,
73+
}
74+
75+
impl Telemetry {
76+
pub fn initialize(
77+
config: &OpenTelemetry,
78+
settings: TelemetrySettings<'_>,
79+
) -> Result<Self, Box<dyn Error + Send + Sync>> {
80+
let filter = EnvFilter::try_from_default_env()
81+
.or_else(|_| EnvFilter::try_new(settings.default_log_level))?;
82+
let fmt_layer = tracing_subscriber::fmt::layer()
83+
.compact()
84+
.with_target(false)
85+
.with_filter(filter_fn(|metadata| {
86+
metadata.target() != errors::EXCEPTION_TARGET
87+
}));
88+
89+
if !config.enabled || !config.has_enabled_exporter() {
90+
tracing_subscriber::registry()
91+
.with(filter)
92+
.with(fmt_layer)
93+
.init();
94+
return Ok(Self {
95+
logger_provider: None,
96+
meter_provider: None,
97+
tracer_provider: None,
98+
});
99+
}
100+
101+
let resource = Resource::builder()
102+
.with_service_name(config.service_name.clone())
103+
.with_attributes([
104+
KeyValue::new("service.version", settings.service_version.to_owned()),
105+
KeyValue::new(
106+
"deployment.environment.name",
107+
settings.environment.to_owned(),
108+
),
109+
])
110+
.build();
111+
112+
let tracer_provider = if let Some(endpoint) = config.traces_endpoint() {
113+
let exporter = opentelemetry_otlp::SpanExporter::builder()
114+
.with_tonic()
115+
.with_endpoint(endpoint.to_owned())
116+
.build()?;
117+
let provider = SdkTracerProvider::builder()
118+
.with_resource(resource.clone())
119+
.with_batch_exporter(exporter)
120+
.build();
121+
global::set_text_map_propagator(TraceContextPropagator::new());
122+
Some(provider)
123+
} else {
124+
None
125+
};
126+
127+
let logger_provider = if let Some(endpoint) = config.logs_endpoint() {
128+
let exporter = opentelemetry_otlp::LogExporter::builder()
129+
.with_tonic()
130+
.with_endpoint(endpoint.to_owned())
131+
.build()?;
132+
Some(
133+
SdkLoggerProvider::builder()
134+
.with_resource(resource.clone())
135+
.with_batch_exporter(exporter)
136+
.build(),
137+
)
138+
} else {
139+
None
140+
};
141+
142+
let meter_provider = if let Some(endpoint) = config.metrics_endpoint() {
143+
let exporter = opentelemetry_otlp::MetricExporter::builder()
144+
.with_tonic()
145+
.with_endpoint(endpoint.to_owned())
146+
.build()?;
147+
let reader = PeriodicReader::builder(exporter).build();
148+
let provider = SdkMeterProvider::builder()
149+
.with_resource(resource)
150+
.with_reader(reader)
151+
.build();
152+
global::set_meter_provider(provider.clone());
153+
if let Some(initialize_metrics) = settings.initialize_metrics {
154+
initialize_metrics();
155+
}
156+
Some(provider)
157+
} else {
158+
None
159+
};
160+
161+
if config.logs_endpoint().is_some() || config.traces_endpoint().is_some() {
162+
errors::enable_backtraces();
163+
}
164+
165+
let trace_layer = tracer_provider.as_ref().map(|provider| {
166+
tracing_opentelemetry::layer()
167+
.with_tracer(provider.tracer(settings.instrumentation_name.to_owned()))
168+
.with_error_records_to_exceptions(true)
169+
.with_error_events_to_status(true)
170+
.with_filter(filter_fn(|metadata| {
171+
metadata.target() != errors::SUMMARY_TARGET
172+
}))
173+
});
174+
let log_layer = logger_provider.as_ref().map(|provider| {
175+
OpenTelemetryTracingBridge::new(provider).with_filter(filter_fn(|metadata| {
176+
metadata.target() != errors::SUMMARY_TARGET
177+
}))
178+
});
179+
180+
tracing_subscriber::registry()
181+
.with(filter)
182+
.with(fmt_layer)
183+
.with(trace_layer)
184+
.with(log_layer)
185+
.init();
186+
187+
Ok(Self {
188+
logger_provider,
189+
meter_provider,
190+
tracer_provider,
191+
})
192+
}
193+
194+
pub fn shutdown(self) {
195+
if let Some(provider) = self.logger_provider {
196+
let _ = provider.shutdown();
197+
}
198+
if let Some(provider) = self.meter_provider {
199+
let _ = provider.shutdown();
200+
}
201+
if let Some(provider) = self.tracer_provider {
202+
let _ = provider.shutdown();
203+
}
204+
}
205+
}
206+
207+
fn non_empty_url(url: &Option<String>) -> Option<&str> {
208+
url.as_deref().filter(|value| !value.trim().is_empty())
209+
}

src/lib.rs

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -9,3 +9,6 @@ pub mod flow_health;
99

1010
#[cfg(feature = "flow_service")]
1111
pub mod flow_service;
12+
13+
#[cfg(feature = "flow_telemetry")]
14+
pub mod flow_telemetry;

0 commit comments

Comments
 (0)