use crate::config::telemetry::{
hive::HiveTelemetryConfig,
tracing::{BatchProcessorConfig, OtlpProtocol, TracingExporterConfig},
TelemetryConfig,
};
use opentelemetry_otlp::{
Protocol, SpanExporter, WithExportConfig, WithHttpConfig, WithTonicConfig,
};
#[cfg(not(feature = "noop_otlp_exporter"))]
use opentelemetry_sdk::error::OTelSdkError;
use opentelemetry_sdk::{
error::OTelSdkResult,
runtime,
trace::{
self, span_processor_with_async_runtime, BatchConfigBuilder, IdGenerator, Sampler,
SdkTracerProvider, SpanData, SpanProcessor, TracerProviderBuilder,
},
Resource,
};
use std::{collections::HashMap, sync::Mutex, time::Duration};
#[cfg(feature = "noop_otlp_exporter")]
use self::noop_exporter::NoopExporter;
use self::standard_pipeline_exporter::StandardPipelineExporter;
use crate::telemetry::{
error::TelemetryError,
traces::hive_console_exporter::HiveConsoleExporter,
utils::{build_metadata, build_tls_config, resolve_string_map, resolve_value_or_expression},
};
pub use control::{disabled_span, is_level_enabled, set_tracing_enabled};
pub mod compatibility;
pub mod control;
pub mod hive_console_exporter;
mod noop_exporter;
pub mod spans;
pub mod standard_pipeline_exporter;
pub mod trace_batch_span_processor;
use crate::telemetry::traces::trace_batch_span_processor::TraceBatchSpanProcessor;
pub(super) fn build_trace_provider<I>(
config: &TelemetryConfig,
id_generator: I,
resource: Resource,
) -> Result<SdkTracerProvider, TelemetryError>
where
I: IdGenerator + 'static,
{
let base_sampler = Sampler::TraceIdRatioBased(config.tracing.collect.sampling);
let mut builder = TracerProviderBuilder::default()
.with_id_generator(id_generator)
.with_resource(resource.clone());
if config.tracing.collect.parent_based_sampler {
builder = builder.with_sampler(Sampler::ParentBased(Box::new(base_sampler)));
} else {
builder = builder.with_sampler(base_sampler);
}
builder = builder
.with_max_events_per_span(config.tracing.collect.max_events_per_span)
.with_max_attributes_per_span(config.tracing.collect.max_attributes_per_span)
.with_max_attributes_per_event(config.tracing.collect.max_attributes_per_event)
.with_max_attributes_per_link(config.tracing.collect.max_attributes_per_link);
Ok(setup_exporters(config, resource, builder)?.build())
}
fn setup_exporters(
config: &TelemetryConfig,
resource: Resource,
mut tracer_provider_builder: TracerProviderBuilder,
) -> Result<TracerProviderBuilder, TelemetryError> {
let sem_conv_mode = &config.tracing.instrumentation.spans.mode;
for exporter_config in &config.tracing.exporters {
match exporter_config {
TracingExporterConfig::Otlp(otlp_config) => {
if !otlp_config.enabled {
continue;
}
ensure_single_protocol_config(
"OTLP exporter",
&otlp_config.protocol,
otlp_config.http.is_some(),
otlp_config.grpc.is_some(),
)?;
let endpoint = resolve_value_or_expression(&otlp_config.endpoint, "OTLP endpoint")?;
let exporter = match &otlp_config.protocol {
OtlpProtocol::Grpc => {
let metadata = otlp_config
.grpc
.as_ref()
.map(|grpc_config| {
resolve_string_map(&grpc_config.metadata, "OTLP grpc metadata key")
})
.transpose()?
.unwrap_or_default();
SpanExporter::builder()
.with_tonic()
.with_endpoint(endpoint)
.with_timeout(otlp_config.batch_processor.max_export_timeout)
.with_tls_config(build_tls_config(
otlp_config.grpc.as_ref().map(|g| &g.tls),
)?)
.with_metadata(build_metadata(metadata)?)
.build()
}
OtlpProtocol::Http => {
let headers = otlp_config
.http
.as_ref()
.map(|http_config| {
resolve_string_map(&http_config.headers, "OTLP http header key")
})
.transpose()?
.unwrap_or_default();
SpanExporter::builder()
.with_http()
.with_endpoint(endpoint)
.with_timeout(otlp_config.batch_processor.max_export_timeout)
.with_headers(headers)
.with_protocol(Protocol::HttpBinary)
.build()
}
}
.map_err(|e| TelemetryError::TracesExporterSetup(e.to_string()))?;
#[cfg(feature = "noop_otlp_exporter")]
let exporter = {
let _ = exporter;
NoopExporter::new()
};
tracer_provider_builder =
tracer_provider_builder.with_span_processor(build_batched_span_processor(
&otlp_config.batch_processor,
&resource,
StandardPipelineExporter::new(exporter, sem_conv_mode),
));
}
TracingExporterConfig::Stdout(stdout_config) => {
if !stdout_config.enabled {
continue;
}
tracer_provider_builder =
tracer_provider_builder.with_span_processor(build_batched_span_processor(
&stdout_config.batch_processor,
&resource,
StandardPipelineExporter::new(
opentelemetry_stdout::SpanExporter::default(),
sem_conv_mode,
),
));
}
}
}
if let Some(hive_config) = &config.hive {
if hive_config.tracing.enabled {
tracer_provider_builder =
setup_hive_exporter(hive_config, &resource, tracer_provider_builder)?;
}
}
Ok(tracer_provider_builder)
}
fn build_batched_span_processor(
config: &BatchProcessorConfig,
resource: &Resource,
exporter: impl trace::SpanExporter + 'static,
) -> impl SpanProcessor {
let mut processor = span_processor_with_async_runtime::BatchSpanProcessor::builder(
exporter,
runtime::TokioCurrentThread,
)
.with_batch_config(
BatchConfigBuilder::default()
.with_max_concurrent_exports(config.max_concurrent_exports as usize)
.with_max_export_batch_size(config.max_export_batch_size as usize)
.with_max_export_timeout(config.max_export_timeout)
.with_max_queue_size(config.max_queue_size as usize)
.with_scheduled_delay(config.scheduled_delay)
.build(),
)
.build();
processor.set_resource(resource);
processor
}
struct TargetedHiveExporter {
endpoint: String,
token: String,
timeout: Duration,
resource: Mutex<Option<Resource>>,
}
impl std::fmt::Debug for TargetedHiveExporter {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
let masked: String = self.token.chars().take(5).collect::<String>() + "***";
f.debug_struct("TargetedHiveExporter")
.field("endpoint", &self.endpoint)
.field("token", &masked)
.field("timeout", &self.timeout)
.finish_non_exhaustive()
}
}
impl TargetedHiveExporter {
fn target_by_trace(batch: &[SpanData]) -> HashMap<opentelemetry::TraceId, String> {
batch
.iter()
.filter_map(|span| {
span.attributes
.iter()
.find(|attribute| attribute.key.as_str() == "hive.target")
.and_then(|attribute| match &attribute.value {
opentelemetry::Value::String(target) => {
Some((span.span_context.trace_id(), target.as_str().to_string()))
}
_ => None,
})
})
.collect()
}
}
impl trace::SpanExporter for TargetedHiveExporter {
async fn export(&self, batch: Vec<SpanData>) -> OTelSdkResult {
let targets = Self::target_by_trace(&batch);
let mut partitions: HashMap<String, Vec<SpanData>> = HashMap::new();
for span in batch {
if let Some(target) = targets.get(&span.span_context.trace_id()) {
partitions.entry(target.clone()).or_default().push(span);
}
}
for (target, spans) in partitions {
#[cfg(not(feature = "noop_otlp_exporter"))]
let exporter = SpanExporter::builder()
.with_http()
.with_endpoint(self.endpoint.clone())
.with_timeout(self.timeout)
.with_headers(HashMap::from([
(
"authorization".to_string(),
format!("Bearer {}", self.token),
),
("x-hive-target-ref".to_string(), target),
]))
.with_protocol(Protocol::HttpBinary)
.build()
.map_err(|error| OTelSdkError::InternalFailure(error.to_string()))?;
#[cfg(feature = "noop_otlp_exporter")]
let exporter = {
let _ = (&self.endpoint, &self.token, self.timeout, target);
NoopExporter::new()
};
let mut exporter = HiveConsoleExporter::new(exporter);
if let Some(resource) = self.resource.lock().unwrap().as_ref() {
trace::SpanExporter::set_resource(&mut exporter, resource);
}
trace::SpanExporter::export(&exporter, spans).await?;
}
Ok(())
}
fn set_resource(&mut self, resource: &Resource) {
*self.resource.lock().unwrap() = Some(resource.clone());
}
}
fn setup_hive_exporter(
config: &HiveTelemetryConfig,
resource: &Resource,
tracer_provider_builder: TracerProviderBuilder,
) -> Result<TracerProviderBuilder, TelemetryError> {
let endpoint = resolve_value_or_expression(&config.tracing.endpoint, "Hive Tracing endpoint")?;
let token = match &config.token {
Some(t) => resolve_value_or_expression(t, "Hive Telemetry token")?,
None => {
return Err(TelemetryError::TracesExporterSetup(
"Hive Tracing token is required but not provided".to_string(),
))
}
};
let hive_exporter = TargetedHiveExporter {
endpoint,
token,
timeout: config.tracing.batch_processor.max_export_timeout,
resource: Mutex::new(None),
};
let mut trace_batching_processor =
TraceBatchSpanProcessor::new(hive_exporter, &config.tracing.batch_processor)?;
trace_batching_processor.set_resource(resource);
Ok(tracer_provider_builder.with_span_processor(trace_batching_processor))
}
fn ensure_single_protocol_config(
name: &str,
protocol: &OtlpProtocol,
http_present: bool,
grpc_present: bool,
) -> Result<(), TelemetryError> {
match protocol {
OtlpProtocol::Grpc if http_present => Err(TelemetryError::TracesExporterSetup(format!(
"{name} http configuration found while protocol is set to gRPC"
))),
OtlpProtocol::Http if grpc_present => Err(TelemetryError::TracesExporterSetup(format!(
"{name} grpc configuration found while protocol is set to HTTP"
))),
_ => Ok(()),
}
}