use std::env;
use core::{fmt, time};
use std::borrow::Cow;
use opentelemetry_sdk::error::OTelSdkError;
use opentelemetry_sdk::logs::SdkLoggerProvider;
use opentelemetry_sdk::trace::SdkTracerProvider;
#[cfg(feature = "grpc")]
fn create_metadata_map(headers: &[(String, String)]) -> tonic::metadata::MetadataMap {
use tonic::metadata::{MetadataMap, MetadataKey};
let mut result = MetadataMap::with_capacity(headers.len());
for (key, value) in headers.iter() {
let meta_key = match MetadataKey::from_bytes(key.as_bytes()) {
Ok(meta) => meta,
Err(error) => panic!("Header '{key}' is not valid ASCII value: {error}"),
};
match value.parse() {
Ok(value) => {
result.append(meta_key, value);
}
Err(error) => panic!("Header '{key}' has invalid value: {error}"),
}
}
result
}
#[cfg(all(feature = "datadog", any(feature = "metrics", feature = "tracing-metrics")))]
#[cold]
#[inline(never)]
fn unsupported_datadog_feature() -> ! {
panic!("Attempt to use 'datadog' while it doesn't support metrics functionality")
}
#[cfg(not(feature = "datadog"))]
#[cold]
#[inline(never)]
fn missing_datadog_feature() -> ! {
panic!("Attempt to use 'datadog' when corresponding feature is not enabled")
}
#[cfg(not(feature = "grpc"))]
#[cold]
#[inline(never)]
fn missing_grpc_feature() -> ! {
panic!("Attempt to use 'grpc' when corresponding feature is not enabled")
}
#[cfg(not(feature = "http"))]
#[cold]
#[inline(never)]
fn missing_http_feature() -> ! {
panic!("Attempt to use 'http' when corresponding feature is not enabled")
}
#[derive(Clone)]
#[repr(transparent)]
pub struct Attributes(pub(crate) opentelemetry_sdk::Resource);
impl Attributes {
#[inline]
pub fn builder() -> AttributesBuilder {
AttributesBuilder::new()
}
#[inline]
pub fn builder_env() -> AttributesBuilder {
AttributesBuilder::new_env()
}
pub fn from_env() -> Option<Attributes> {
let builder = AttributesBuilder::new_env();
if builder.is_mutated {
Some(builder.finish())
} else {
None
}
}
}
impl fmt::Debug for Attributes {
#[inline]
fn fmt(&self, fmt: &mut fmt::Formatter<'_>) -> fmt::Result {
fmt::Debug::fmt(&self.0, fmt)
}
}
pub struct AttributesBuilder {
is_mutated: bool,
inner: opentelemetry_sdk::resource::ResourceBuilder
}
impl AttributesBuilder {
#[inline]
pub fn new() -> Self {
Self {
is_mutated: false,
inner: opentelemetry_sdk::resource::Resource::builder()
}
}
pub fn new_env() -> Self {
let mut builder = Self::new();
if let Ok(service_name) = env::var("OTEL_SERVICE_NAME") {
builder = builder.with_service_name(service_name)
}
if let Ok(attrs) = env::var("OTEL_RESOURCE_ATTRIBUTES") {
for key_value in attrs.split(',') {
let mut key_value_iter = key_value.trim().splitn(2, '=');
let key = key_value_iter.next().unwrap();
match key_value_iter.next() {
Some(value) => {
builder = builder.with_attr(key.to_owned(), value.to_owned());
},
None => continue
}
}
}
builder
}
#[inline]
fn with_service_name(mut self, value: impl Into<opentelemetry::Value>) -> Self {
self.is_mutated = true;
self.inner = self.inner.with_service_name(value);
self
}
#[inline]
pub fn with_attr(mut self, key: impl Into<Cow<'static, str>>, value: impl Into<opentelemetry::Value>) -> Self {
self.is_mutated = true;
self.inner = self.inner.with_attribute(opentelemetry::KeyValue::new(key.into(), value.into()));
self
}
#[inline]
pub fn with_env_attr(self, key: impl Into<Cow<'static, str>>, env_name: &str) -> Self {
if let Ok(value) = std::env::var(env_name) {
self.with_attr(key, value)
} else {
self
}
}
#[inline]
pub fn finish(self) -> Attributes {
Attributes(self.inner.build())
}
#[inline]
pub fn finish_if_set(self) -> Option<Attributes> {
if self.is_mutated {
Some(self.finish())
} else {
None
}
}
}
#[derive(Default)]
pub struct ShutdownError {
logs: Option<OTelSdkError>,
trace: Option<OTelSdkError>,
#[cfg(any(feature = "metrics", feature = "tracing-metrics"))]
metrics: Option<OTelSdkError>
}
impl fmt::Debug for ShutdownError {
fn fmt(&self, fmt: &mut fmt::Formatter<'_>) -> fmt::Result {
let mut fmt = fmt.debug_struct("OtlpShutdownError");
if let Some(logs) = self.logs.as_ref() {
fmt.field("logs", logs);
}
if let Some(trace) = self.trace.as_ref() {
fmt.field("trace", trace);
}
#[cfg(any(feature = "metrics", feature = "tracing-metrics"))]
if let Some(metrics) = self.metrics.as_ref() {
fmt.field("metrics", metrics);
}
fmt.finish()
}
}
impl fmt::Display for ShutdownError {
fn fmt(&self, fmt: &mut fmt::Formatter<'_>) -> fmt::Result {
fmt.write_str("Failed to shutdown Otlp:")?;
if let Some(logs) = self.logs.as_ref() {
fmt.write_fmt(format_args!(" logs={logs}"))?
}
if let Some(trace) = self.trace.as_ref() {
fmt.write_fmt(format_args!(" trace={trace}"))?
}
#[cfg(any(feature = "metrics", feature = "tracing-metrics"))]
if let Some(metrics) = self.metrics.as_ref() {
fmt.write_fmt(format_args!(" metrics={metrics}"))?
}
Ok(())
}
}
impl std::error::Error for ShutdownError {}
pub struct Otlp {
logs: Option<SdkLoggerProvider>,
trace: Option<SdkTracerProvider>,
#[cfg(any(feature = "metrics", feature = "tracing-metrics"))]
metrics: Option<opentelemetry_sdk::metrics::SdkMeterProvider>
}
impl Otlp {
#[inline]
const fn new() -> Self {
Self {
logs: None,
trace: None,
#[cfg(any(feature = "metrics", feature = "tracing-metrics"))]
metrics: None,
}
}
#[inline]
pub const fn builder() -> Builder {
Builder::new()
}
pub fn flush(&self) -> Result<(), ShutdownError> {
let mut is_error = false;
let mut errors = ShutdownError::default();
if let Some(logs) = self.logs.as_ref() {
if let Err(error) = logs.force_flush() {
is_error = true;
errors.logs = Some(error);
}
}
if let Some(trace) = self.trace.as_ref() {
if let Err(error) = trace.force_flush() {
is_error = true;
errors.trace = Some(error);
}
}
#[cfg(any(feature = "metrics", feature = "tracing-metrics"))]
if let Some(metrics) = self.metrics.as_ref() {
if let Err(error) = metrics.force_flush() {
is_error = true;
errors.metrics = Some(error);
}
}
if is_error {
Err(errors)
} else {
Ok(())
}
}
pub fn shutdown(&mut self, limit: Option<time::Duration>) -> Result<(), ShutdownError> {
fn get_default_timeout(key: &str) -> time::Duration {
for key in [key, "OTEL_EXPORTER_OTLP_TIMEOUT"] {
if let Ok(value) = env::var(key) {
if let Ok(value) = value.parse() {
return time::Duration::from_secs(value)
}
}
}
time::Duration::from_secs(10)
}
let mut is_error = false;
let mut errors = ShutdownError::default();
if let Some(logs) = self.logs.take() {
let limit = limit.unwrap_or_else(|| get_default_timeout("OTEL_EXPORTER_OTLP_LOGS_TIMEOUT"));
if let Err(error) = logs.shutdown_with_timeout(limit) {
is_error = true;
errors.logs = Some(error);
}
}
if let Some(trace) = self.trace.take() {
let limit = limit.unwrap_or_else(|| get_default_timeout("OTEL_EXPORTER_OTLP_TRACES_TIMEOUT"));
if let Err(error) = trace.shutdown_with_timeout(limit) {
is_error = true;
errors.trace = Some(error);
}
}
#[cfg(any(feature = "metrics", feature = "tracing-metrics"))]
if let Some(metrics) = self.metrics.take() {
let limit = limit.unwrap_or_else(|| get_default_timeout("OTEL_EXPORTER_OTLP_METRICS_TIMEOUT"));
if let Err(error) = metrics.shutdown_with_timeout(limit) {
is_error = true;
errors.metrics = Some(error);
}
}
if is_error {
Err(errors)
} else {
Ok(())
}
}
}
impl fmt::Debug for Otlp {
#[inline]
fn fmt(&self, fmt: &mut fmt::Formatter<'_>) -> fmt::Result {
let mut fmt = fmt.debug_struct("Otlp");
if let Some(logs) = self.logs.as_ref() {
fmt.field("logs", logs);
}
if let Some(trace) = self.trace.as_ref() {
fmt.field("trace", trace);
}
#[cfg(any(feature = "metrics", feature = "tracing-metrics"))]
if let Some(metrics) = self.metrics.as_ref() {
fmt.field("metrics", metrics);
}
fmt.finish()
}
}
impl Drop for Otlp {
#[inline(always)]
fn drop(&mut self) {
let _ = self.shutdown(None);
}
}
#[derive(Copy, Clone, PartialEq, Eq, Debug)]
pub enum Protocol {
Grpc,
HttpBinary,
HttpJson,
DatadogAgent,
}
impl Protocol {
pub fn select_default() -> Option<Self> {
if cfg!(feature = "http") {
Some(Self::HttpBinary)
} else if cfg!(feature = "grpc") {
Some(Self::Grpc)
} else if cfg!(feature = "datadog") {
Some(Self::DatadogAgent)
} else {
None
}
}
pub fn from_env() -> Option<Self> {
match env::var("OTEL_EXPORTER_OTLP_PROTOCOL") {
Ok(protocol) => {
if protocol == "grpc" {
return Some(Self::Grpc);
}
if protocol == "http/protobuf" {
return Some(Self::HttpBinary);
}
if protocol == "http/json" {
return Some(Self::HttpJson);
}
if protocol == "datadog" {
return Some(Self::DatadogAgent);
}
None
},
Err(_) => None,
}
}
#[allow(unused)]
#[inline]
const fn into_otel(self) -> opentelemetry_otlp::Protocol {
match self {
#[cfg(feature = "grpc")]
Self::Grpc => opentelemetry_otlp::Protocol::Grpc,
#[cfg(feature = "http")]
Self::HttpJson => opentelemetry_otlp::Protocol::HttpJson,
#[cfg(feature = "http")]
Self::HttpBinary => opentelemetry_otlp::Protocol::HttpBinary,
_ => unreachable!(),
}
}
}
pub struct Destination<'a> {
pub protocol: Protocol,
pub url: Cow<'a, str>,
pub attributes: Option<&'a Attributes>
}
impl Destination<'_> {
#[cfg_attr(not(all(feature = "grpc", feature = "http", feature = "datadog")), allow(unused))]
fn get_service_attrs(&self) -> Option<Attributes> {
match self.attributes {
Some(attrs) => Some(attrs.clone()),
None => Attributes::from_env(),
}
}
}
impl Destination<'static> {
pub fn from_env() -> Self {
let protocol = Protocol::from_env().or_else(Protocol::select_default).expect("Unable to determine Destination::protocol");
let url = match env::var("OTEL_EXPORTER_OTLP_ENDPOINT") {
Ok(url) => url.into(),
Err(_) => match protocol {
Protocol::Grpc => "http://localhost:4317".into(),
Protocol::HttpBinary | Protocol::HttpJson => "http://localhost:4318".into(),
Protocol::DatadogAgent => "http://localhost:8126".into()
}
};
Self {
protocol,
url,
attributes: None,
}
}
}
macro_rules! declare_trace_limits {
({$($name:ident,)+}) => {
struct SpanLimits {
$(
$name: u32,
)+
}
impl SpanLimits {
const DEFAULT: u32 = 128;
#[inline(always)]
const fn new() -> Self {
Self {
$(
$name: Self::DEFAULT,
)+
}
}
#[allow(unused)]
#[inline(always)]
fn apply_to(&self, mut builder: opentelemetry_sdk::trace::TracerProviderBuilder) -> opentelemetry_sdk::trace::TracerProviderBuilder {
$(
if self.$name != Self::DEFAULT {
builder = builder.$name(self.$name);
}
)+
builder
}
}
};
}
declare_trace_limits!({
with_max_events_per_span,
with_max_attributes_per_span,
with_max_links_per_span,
with_max_attributes_per_link,
with_max_attributes_per_event,
});
#[allow(unused)]
#[derive(Clone, Debug)]
pub(crate) struct ParentBasedSampler<T> {
pub(crate) sampler: T,
}
impl<T: opentelemetry_sdk::trace::ShouldSample + Clone + 'static> opentelemetry_sdk::trace::ShouldSample for ParentBasedSampler<T> {
#[inline(always)]
fn should_sample(&self, parent_context: Option<&opentelemetry::Context>, trace_id: opentelemetry::TraceId, name: &str, span_kind: &opentelemetry::trace::SpanKind, attributes: &[opentelemetry::KeyValue], links: &[opentelemetry::trace::Link]) -> opentelemetry_sdk::trace::SamplingResult {
use opentelemetry::trace::TraceContextExt;
use opentelemetry_sdk::trace::SamplingDecision;
if let Some(active_parent) = parent_context.filter(|ctx| ctx.has_active_span()) {
let parent_span = active_parent.span();
let parent_span_context = parent_span.span_context();
let decision = if parent_span_context.is_sampled() {
SamplingDecision::RecordAndSample
} else {
SamplingDecision::Drop
};
opentelemetry_sdk::trace::SamplingResult {
decision,
attributes: Vec::new(),
trace_state: parent_span_context.trace_state().clone(),
}
} else {
self.sampler.should_sample(parent_context, trace_id, name, span_kind, attributes, links)
}
}
}
#[allow(unused)]
#[derive(Copy, Clone, Debug)]
pub(crate) struct AlwaysOnSampler;
impl opentelemetry_sdk::trace::ShouldSample for AlwaysOnSampler {
#[inline(always)]
fn should_sample(&self, parent_context: Option<&opentelemetry::Context>, _: opentelemetry::TraceId, _: &str, _: &opentelemetry::trace::SpanKind, _: &[opentelemetry::KeyValue], _: &[opentelemetry::trace::Link]) -> opentelemetry_sdk::trace::SamplingResult {
use opentelemetry::trace::TraceContextExt;
opentelemetry_sdk::trace::SamplingResult {
decision: opentelemetry_sdk::trace::SamplingDecision::RecordAndSample,
attributes: Vec::new(),
trace_state: match parent_context {
Some(ctx) => ctx.span().span_context().trace_state().clone(),
None => opentelemetry::trace::TraceState::default(),
},
}
}
}
#[allow(unused)]
#[derive(Copy, Clone, Debug)]
pub(crate) struct AlwaysOffSampler;
impl opentelemetry_sdk::trace::ShouldSample for AlwaysOffSampler {
#[inline(always)]
fn should_sample(&self, parent_context: Option<&opentelemetry::Context>, _: opentelemetry::TraceId, _: &str, _: &opentelemetry::trace::SpanKind, _: &[opentelemetry::KeyValue], _: &[opentelemetry::trace::Link]) -> opentelemetry_sdk::trace::SamplingResult {
use opentelemetry::trace::TraceContextExt;
opentelemetry_sdk::trace::SamplingResult {
decision: opentelemetry_sdk::trace::SamplingDecision::Drop,
attributes: Vec::new(),
trace_state: match parent_context {
Some(ctx) => ctx.span().span_context().trace_state().clone(),
None => opentelemetry::trace::TraceState::default(),
},
}
}
}
pub struct TraceSettings {
#[allow(unused)]
name: Cow<'static, str>,
#[allow(unused)]
sample_rate: f64,
#[allow(unused)]
limits: SpanLimits,
#[allow(unused)]
respect_parent: bool,
}
macro_rules! set_trace_limit {
($limits:expr, $name:ident) => {
$limits.$name = $name;
};
}
impl TraceSettings {
pub const fn new(name: Cow<'static, str>, sample_rate: f64) -> Self {
Self {
name,
sample_rate,
limits: SpanLimits::new(),
respect_parent: true,
}
}
pub const fn with_respect_parent_sampling(mut self, value: bool) -> Self {
self.respect_parent = value;
self
}
pub const fn with_max_events_per_span(mut self, with_max_events_per_span: u32) -> Self {
set_trace_limit!(self.limits, with_max_events_per_span);
self
}
pub const fn with_max_attributes_per_span(mut self, with_max_attributes_per_span: u32) -> Self {
set_trace_limit!(self.limits, with_max_attributes_per_span);
self
}
pub const fn with_max_links_per_span(mut self, with_max_links_per_span: u32) -> Self {
set_trace_limit!(self.limits, with_max_links_per_span);
self
}
pub const fn with_max_attributes_per_event(mut self, with_max_attributes_per_event: u32) -> Self {
set_trace_limit!(self.limits, with_max_attributes_per_event);
self
}
pub const fn with_max_attributes_per_link(mut self, with_max_attributes_per_link: u32) -> Self {
set_trace_limit!(self.limits, with_max_attributes_per_link);
self
}
#[inline]
#[cfg(any(feature = "grpc", feature = "http", feature = "datadog"))]
fn create_sampler(&self) -> Box<dyn opentelemetry_sdk::trace::ShouldSample> {
let sample_rate = self.sample_rate.clamp(0.0, 1.0);
if self.respect_parent {
Box::new(ParentBasedSampler {
sampler: opentelemetry_sdk::trace::Sampler::TraceIdRatioBased(sample_rate)
})
} else {
if sample_rate == 0.0 {
Box::new(AlwaysOffSampler)
} else if sample_rate == 1.0 {
Box::new(AlwaysOnSampler)
} else {
Box::new(opentelemetry_sdk::trace::Sampler::TraceIdRatioBased(sample_rate))
}
}
}
}
#[cfg(any(feature = "metrics", feature = "tracing-metrics"))]
pub struct MetricsSettings {
temporality: opentelemetry_sdk::metrics::Temporality,
}
#[cfg(any(feature = "metrics", feature = "tracing-metrics"))]
impl MetricsSettings {
#[inline]
pub const fn new() -> Self {
Self {
temporality: opentelemetry_sdk::metrics::Temporality::Cumulative
}
}
#[inline]
pub const fn with_delta(mut self) -> Self {
self.temporality = opentelemetry_sdk::metrics::Temporality::Delta;
self
}
#[inline]
pub const fn with_low_memory(mut self) -> Self {
self.temporality = opentelemetry_sdk::metrics::Temporality::LowMemory;
self
}
}
#[derive(Copy, Clone)]
pub enum ExportRuntime {
Threaded,
Tokio,
TokioCurrentThrad
}
impl ExportRuntime {
#[inline]
pub fn is_threaded(&self) -> bool {
matches!(self, Self::Threaded)
}
pub fn auto_detect() -> Self {
#[cfg(feature = "rt-tokio")]
if let Ok(handle) = tokio::runtime::Handle::try_current() {
match handle.runtime_flavor() {
tokio::runtime::RuntimeFlavor::CurrentThread => return Self::TokioCurrentThrad,
tokio::runtime::RuntimeFlavor::MultiThread => return Self::Tokio,
_ => (),
}
}
Self::Threaded
}
#[cfg(any(feature = "grpc", feature = "http", feature = "datadog"))]
fn create_logger_exporter<E: opentelemetry_sdk::logs::LogExporter + 'static>(self, exporter: E, config: opentelemetry_sdk::logs::BatchConfig, destination: &Destination<'_>) -> SdkLoggerProvider {
let mut builder = SdkLoggerProvider::builder();
if let Some(attrs) = destination.get_service_attrs() {
builder = builder.with_resource(attrs.0);
}
match self {
Self::Threaded => {
let exporter = opentelemetry_sdk::logs::BatchLogProcessor::builder(exporter).with_batch_config(config).build();
builder.with_log_processor(exporter).build()
},
#[cfg(feature = "rt-tokio")]
Self::Tokio => {
let exporter = opentelemetry_sdk::logs::log_processor_with_async_runtime::BatchLogProcessor::builder(exporter, opentelemetry_sdk::runtime::Tokio).with_batch_config(config).build();
builder.with_log_processor(exporter).build()
},
#[cfg(feature = "rt-tokio")]
Self::TokioCurrentThrad => {
let exporter = opentelemetry_sdk::logs::log_processor_with_async_runtime::BatchLogProcessor::builder(exporter, opentelemetry_sdk::runtime::TokioCurrentThread).with_batch_config(config).build();
builder.with_log_processor(exporter).build()
},
#[cfg(not(feature = "rt-tokio"))]
_ => panic!("rt-tokio feature must be enabled for async runtime"),
}
}
#[cfg(any(feature = "grpc", feature = "http", feature = "datadog"))]
fn create_tracer_exporter<E: opentelemetry_sdk::trace::SpanExporter + 'static>(self, exporter: E, config: opentelemetry_sdk::trace::BatchConfig, destination: &Destination<'_>, settings: &TraceSettings) -> SdkTracerProvider {
let sampler = settings.create_sampler();
let mut builder = SdkTracerProvider::builder().with_id_generator(opentelemetry_sdk::trace::RandomIdGenerator::default()).with_sampler(sampler);
builder = settings.limits.apply_to(builder);
if let Some(attrs) = destination.get_service_attrs() {
builder = builder.with_resource(attrs.0);
}
match self {
Self::Threaded => {
let exporter = opentelemetry_sdk::trace::BatchSpanProcessor::builder(exporter).with_batch_config(config).build();
builder.with_span_processor(exporter).build()
},
#[cfg(feature = "rt-tokio")]
Self::Tokio => {
let exporter = opentelemetry_sdk::trace::span_processor_with_async_runtime::BatchSpanProcessor::builder(exporter, opentelemetry_sdk::runtime::Tokio).with_batch_config(config).build();
builder.with_span_processor(exporter).build()
},
#[cfg(feature = "rt-tokio")]
Self::TokioCurrentThrad => {
let exporter = opentelemetry_sdk::trace::span_processor_with_async_runtime::BatchSpanProcessor::builder(exporter, opentelemetry_sdk::runtime::TokioCurrentThread).with_batch_config(config).build();
builder.with_span_processor(exporter).build()
},
#[cfg(not(feature = "rt-tokio"))]
_ => panic!("rt-tokio feature must be enabled for async runtime"),
}
}
#[cfg(all(feature = "metrics", any(feature = "grpc", feature = "http")))]
fn create_metrics_exporter<E: opentelemetry_sdk::metrics::exporter::PushMetricExporter + 'static>(self, exporter: E, destination: &Destination<'_>, export_interval: time::Duration) -> opentelemetry_sdk::metrics::SdkMeterProvider {
let mut builder = opentelemetry_sdk::metrics::SdkMeterProvider::builder();
if let Some(attrs) = destination.get_service_attrs() {
builder = builder.with_resource(attrs.0.clone());
}
match self {
Self::Threaded => {
let mut reader = opentelemetry_sdk::metrics::PeriodicReader::builder(exporter);
if !export_interval.is_zero() {
reader = reader.with_interval(export_interval);
}
builder.with_reader(reader.build()).build()
},
#[cfg(feature = "rt-tokio")]
Self::Tokio => {
let mut reader = opentelemetry_sdk::metrics::periodic_reader_with_async_runtime::PeriodicReader::builder(exporter, opentelemetry_sdk::runtime::Tokio);
if !export_interval.is_zero() {
reader = reader.with_interval(export_interval);
}
builder.with_reader(reader.build()).build()
},
#[cfg(feature = "rt-tokio")]
Self::TokioCurrentThrad => {
let mut reader = opentelemetry_sdk::metrics::periodic_reader_with_async_runtime::PeriodicReader::builder(exporter, opentelemetry_sdk::runtime::TokioCurrentThread);
if !export_interval.is_zero() {
reader = reader.with_interval(export_interval);
}
builder.with_reader(reader.build()).build()
},
#[cfg(not(feature = "rt-tokio"))]
_ => panic!("rt-tokio feature must be enabled for async runtime"),
}
}
}
#[derive(Copy, Clone, PartialEq, Eq)]
pub struct RetryPolicy {
max_retries: usize,
initial_delay_ms: u64,
max_delay_ms: u64,
jitter_ms: u64,
}
impl RetryPolicy {
pub const fn new() -> Self {
Self {
max_retries: 5,
initial_delay_ms: 500,
max_delay_ms: 30_000,
jitter_ms: 100,
}
}
pub const fn with_max_retries(mut self, max_retries: usize) -> Self {
self.max_retries = max_retries;
self
}
pub const fn with_initial_delay(mut self, initial_delay: time::Duration) -> Self {
self.initial_delay_ms = initial_delay.as_millis() as _;
self
}
pub const fn with_max_delay(mut self, delay: time::Duration) -> Self {
self.max_delay_ms = delay.as_millis() as _;
self
}
pub const fn with_jitter(mut self, jitter: time::Duration) -> Self {
self.jitter_ms = jitter.as_millis() as _;
self
}
}
#[cfg(any(feature = "grpc-retry", feature = "http-retry"))]
impl From<RetryPolicy> for opentelemetry_otlp::retry::RetryPolicy {
#[inline]
fn from(RetryPolicy { max_delay_ms, max_retries, initial_delay_ms, jitter_ms }: RetryPolicy) -> Self {
Self {
max_retries,
initial_delay_ms,
max_delay_ms,
jitter_ms,
}
}
}
impl Default for RetryPolicy {
#[inline]
fn default() -> Self {
Self::new()
}
}
pub struct Builder {
otlp: Otlp,
headers: Vec<(String, String)>,
timeout: time::Duration,
export_interval: time::Duration,
queue_size: usize,
compression: bool,
#[cfg(feature = "http-ureq")]
ureq: Option<crate::ureq::HttpClient>,
runtime: ExportRuntime,
retry: RetryPolicy,
}
impl Builder {
#[inline]
pub const fn new() -> Self {
Self {
otlp: Otlp::new(),
headers: Vec::new(),
timeout: time::Duration::from_secs(0),
export_interval: time::Duration::ZERO,
queue_size: 0,
compression: true,
#[cfg(feature = "http-ureq")]
ureq: None,
runtime: ExportRuntime::Threaded,
retry: RetryPolicy::new(),
}
}
#[cfg(feature = "http-ureq")]
pub fn with_ureq_http_client_shared(mut self, ureq: crate::ureq::HttpClient) -> Self {
self.ureq = Some(ureq);
self
}
#[inline(always)]
#[cfg(feature = "http-ureq")]
pub fn with_ureq_http_client(self) -> Self {
self.with_ureq_http_client_shared(crate::ureq::HttpClient::new())
}
#[inline]
pub fn with_compression(mut self, compression: bool) -> Self {
self.compression = compression;
self
}
#[inline]
pub fn with_timeout(mut self, timeout: time::Duration) -> Self {
self.timeout = timeout;
self
}
#[inline]
pub fn with_header(mut self, key: impl Into<String>, value: impl Into<String>) -> Self {
self.headers.push((key.into(), value.into()));
self
}
#[inline]
pub fn with_interval(mut self, interval: time::Duration) -> Self {
self.export_interval = interval;
self
}
#[inline]
pub fn with_queue_size(mut self, size: usize) -> Self {
self.queue_size = size;
self
}
pub fn with_runtime(mut self, runtime: ExportRuntime) -> Self {
self.runtime = runtime;
self
}
pub fn with_retry(mut self, retry: RetryPolicy) -> Self {
self.retry = retry;
self
}
#[cfg(feature = "grpc")]
fn apply_otel_grpc_config<T: opentelemetry_otlp::WithTonicConfig + opentelemetry_otlp::WithExportConfig>(&self, mut builder: T) -> T {
if cfg!(feature = "grpc-compression") && self.compression {
builder = builder.with_compression(opentelemetry_otlp::Compression::Gzip)
}
if !self.headers.is_empty() {
let headers = create_metadata_map(&self.headers);
builder = builder.with_metadata(headers);
}
#[cfg(feature = "grpc-retry")]
if !self.runtime.is_threaded() {
builder = builder.with_retry_policy(self.retry.into());
}
if self.timeout.is_zero() {
builder
} else {
builder.with_timeout(self.timeout)
}
}
#[cfg(feature = "http")]
fn apply_otel_http_config<T: opentelemetry_otlp::WithHttpConfig + opentelemetry_otlp::WithExportConfig>(&self, mut builder: T) -> T {
if cfg!(feature = "http-compression") && self.compression {
builder = builder.with_compression(opentelemetry_otlp::Compression::Gzip)
}
#[cfg(feature = "http-ureq")]
if let Some(ureq) = self.ureq.as_ref() {
builder = builder.with_http_client(ureq.clone());
}
if !self.headers.is_empty() {
let headers = self.headers.iter().map(|(key, value)| (key.clone(), value.clone())).collect();
builder = builder.with_headers(headers);
}
#[cfg(feature = "http-retry")]
if !self.runtime.is_threaded() {
builder = builder.with_retry_policy(self.retry.into());
}
if self.timeout.is_zero() {
builder
} else {
builder.with_timeout(self.timeout)
}
}
fn create_logs(&mut self, _destination: &Destination<'_>) -> opentelemetry_sdk::logs::SdkLoggerProvider {
if self.otlp.logs.is_some() {
panic!("Logs is already initialized")
}
let mut batch_config = opentelemetry_sdk::logs::BatchConfigBuilder::default();
if !self.export_interval.is_zero() {
batch_config = batch_config.with_scheduled_delay(self.export_interval);
}
if self.queue_size != 0 {
batch_config = batch_config.with_max_queue_size(self.queue_size);
}
let _batch_config = batch_config.build();
let _logs = match _destination.protocol {
#[cfg(feature = "grpc")]
Protocol::Grpc => {
use opentelemetry_otlp::{WithExportConfig};
let mut builder = opentelemetry_otlp::LogExporter::builder().with_tonic().with_endpoint(_destination.url.clone().into_owned());
builder = self.apply_otel_grpc_config(builder);
let exporter = builder.build().expect("Failed to initialize logs grpc exporter");
self.runtime.create_logger_exporter(exporter, _batch_config, &_destination)
},
#[cfg(not(feature = "grpc"))]
Protocol::Grpc => missing_grpc_feature(),
#[cfg(feature = "datadog")]
Protocol::DatadogAgent => {
let attributes = _destination.get_service_attrs();
if let Some(file_path) = _destination.url.strip_prefix("file://") {
self.runtime.create_logger_exporter(crate::datadog::file_exporter(file_path.to_owned().into()).with_attrs(attributes), _batch_config, &_destination)
} else {
self.runtime.create_logger_exporter(crate::datadog::stdout_exporter().with_attrs(attributes), _batch_config, &_destination)
}
}
#[cfg(not(feature = "datadog"))]
Protocol::DatadogAgent => missing_datadog_feature(),
#[cfg(feature = "http")]
http => {
use opentelemetry_otlp::WithExportConfig;
let url = format!("{}/logs", _destination.url.trim_end_matches('/'));
let mut builder = opentelemetry_otlp::LogExporter::builder().with_http().with_protocol(http.into_otel()).with_endpoint(url);
builder = self.apply_otel_http_config(builder);
let exporter = builder.build().expect("Failed to initialize logs http exporter");
self.runtime.create_logger_exporter(exporter, _batch_config, _destination)
},
#[cfg(not(feature = "http"))]
_ => missing_http_feature(),
};
#[cfg(any(feature = "grpc", feature = "http", feature = "datadog"))]
{
let this = self;
this.otlp.logs = Some(_logs.clone());
return _logs;
}
}
pub fn with_logs<S: tracing::Subscriber + for<'a> tracing_subscriber::registry::LookupSpan<'a>>(&mut self, destination: &Destination<'_>) -> impl tracing_subscriber::Layer<S> + Send + Sync + use<S> {
opentelemetry_appender_tracing::layer::OpenTelemetryTracingBridge::new(&self.create_logs(destination))
}
fn create_tracer(&mut self, _destination: &Destination<'_>, _settings: TraceSettings) -> opentelemetry_sdk::trace::SdkTracer {
if self.otlp.trace.is_some() {
panic!("Trace is already initialized")
}
let mut batch_config = opentelemetry_sdk::trace::BatchConfigBuilder::default();
if !self.export_interval.is_zero() {
batch_config = batch_config.with_scheduled_delay(self.export_interval);
}
if self.queue_size != 0 {
batch_config = batch_config.with_max_queue_size(self.queue_size);
}
let _batch_config = batch_config.build();
let _trace = match _destination.protocol {
#[cfg(feature = "grpc")]
Protocol::Grpc => {
use opentelemetry_otlp::{WithExportConfig};
let mut builder = opentelemetry_otlp::SpanExporter::builder().with_tonic().with_endpoint(_destination.url.clone().into_owned());
builder = self.apply_otel_grpc_config(builder);
let exporter = builder.build().expect("Failed to initialize trace grpc exporter");
self.runtime.create_tracer_exporter(exporter, _batch_config, &_destination, &_settings)
},
#[cfg(not(feature = "grpc"))]
Protocol::Grpc => missing_grpc_feature(),
#[cfg(feature = "datadog")]
Protocol::DatadogAgent => {
use crate::datadog::{SERVICE_NAME, SERVICE_VERSION, SERVICE_ENV};
let mut exporter = opentelemetry_datadog::new_pipeline().with_agent_endpoint(_destination.url.clone());
if let Some(attrs) = _destination.get_service_attrs() {
if let Some(service_name) = attrs.0.get(&SERVICE_NAME) {
exporter = exporter.with_service_name(service_name.to_string());
}
if let Some(service_version) = attrs.0.get(&SERVICE_VERSION) {
exporter = exporter.with_version(service_version.to_string());
}
if let Some(service_env) = attrs.0.get(&SERVICE_ENV) {
exporter = exporter.with_env(service_env.to_string());
}
}
#[cfg(feature = "http-ureq")]
if let Some(ureq) = self.ureq.as_ref() {
exporter = exporter.with_http_client(ureq.clone());
}
let exporter = exporter.build_exporter().expect("Failed to initialize datadog exporter");
self.runtime.create_tracer_exporter(exporter, _batch_config, &_destination, &_settings)
},
#[cfg(not(feature = "datadog"))]
Protocol::DatadogAgent => missing_datadog_feature(),
#[cfg(feature = "http")]
http => {
use opentelemetry_otlp::WithExportConfig;
let url = format!("{}/traces", _destination.url.trim_end_matches('/'));
let mut builder = opentelemetry_otlp::SpanExporter::builder().with_http().with_protocol(http.into_otel()).with_endpoint(url);
builder = self.apply_otel_http_config(builder);
let exporter = builder.build().expect("Failed to initialize trace http exporter");
self.runtime.create_tracer_exporter(exporter, _batch_config, &_destination, &_settings)
},
#[cfg(not(feature = "http"))]
_ => missing_http_feature(),
};
#[cfg(any(feature = "grpc", feature = "http", feature = "datadog"))]
{
use opentelemetry::trace::TracerProvider;
let this = self;
let tracer = _trace.tracer(_settings.name);
this.otlp.trace = Some(_trace);
return tracer;
}
}
pub fn with_trace<S: tracing::Subscriber + for<'a> tracing_subscriber::registry::LookupSpan<'a>>(&mut self, destination: &Destination<'_>, settings: TraceSettings) -> impl tracing_subscriber::Layer<S> + use<S> {
tracing_opentelemetry::OpenTelemetryLayer::new(self.create_tracer(destination, settings)).with_level(true)
.with_target(true)
.with_tracked_inactivity(true)
.with_error_events_to_status(true)
.with_error_events_to_exceptions(true)
.with_error_records_to_exceptions(true)
}
#[cfg(any(feature = "metrics", feature = "tracing-metrics"))]
fn create_metrics(&mut self, _destination: &Destination<'_>, _settings: MetricsSettings) -> opentelemetry_sdk::metrics::SdkMeterProvider {
if self.otlp.metrics.is_some() {
panic!("Trace is already initialized")
}
let _metrics = match _destination.protocol {
#[cfg(feature = "grpc")]
Protocol::Grpc => {
use opentelemetry_otlp::{WithExportConfig};
let mut builder = opentelemetry_otlp::MetricExporter::builder().with_tonic().with_endpoint(_destination.url.clone().into_owned()).with_temporality(_settings.temporality);
builder = self.apply_otel_grpc_config(builder);
let exporter = builder.build().expect("Failed to initialize metrics grpc exporter");
self.runtime.create_metrics_exporter(exporter, &_destination, self.export_interval)
},
#[cfg(not(feature = "grpc"))]
Protocol::Grpc => missing_grpc_feature(),
#[cfg(feature = "datadog")]
Protocol::DatadogAgent => unsupported_datadog_feature(),
#[cfg(not(feature = "datadog"))]
Protocol::DatadogAgent => missing_datadog_feature(),
#[cfg(feature = "http")]
http => {
use opentelemetry_otlp::WithExportConfig;
let url = format!("{}/metrics", _destination.url.trim_end_matches('/'));
let mut builder = opentelemetry_otlp::MetricExporter::builder().with_http().with_protocol(http.into_otel()).with_endpoint(url).with_temporality(_settings.temporality);
builder = self.apply_otel_http_config(builder);
let exporter = builder.build().expect("Failed to initialize metrics grpc exporter");
self.runtime.create_metrics_exporter(exporter, &_destination, self.export_interval)
},
#[cfg(not(feature = "http"))]
_ => missing_http_feature(),
};
#[cfg(any(feature = "grpc", feature = "http"))]
{
let this = self;
this.otlp.metrics = Some(_metrics.clone());
return _metrics;
}
}
#[inline(always)]
#[cfg(feature = "metrics")]
pub fn with_metrics(&mut self, destination: &Destination<'_>, settings: MetricsSettings, name: &'static str) -> impl crate::metrics::Recorder + Send + Sync + 'static {
use opentelemetry::metrics::MeterProvider;
let metrics = self.create_metrics(destination, settings);
let meter = metrics.meter(name);
let metrics = metrics_opentelemetry::OpenTelemetryMetrics::new(meter);
metrics_opentelemetry::OpenTelemetryRecorder::new(metrics)
}
#[inline(always)]
#[cfg(feature = "tracing-metrics")]
pub fn with_tracing_metrics<S: tracing::Subscriber + for<'a> tracing_subscriber::registry::LookupSpan<'a>>(&mut self, destination: &Destination<'_>, settings: MetricsSettings) -> impl tracing_subscriber::Layer<S> + Send + Sync + use<S> {
tracing_opentelemetry::MetricsLayer::new(self.create_metrics(destination, settings))
}
#[inline]
pub fn finish(self) -> Otlp {
self.otlp
}
}