use crate::config::telemetry::{
hive::HiveTelemetryConfig,
tracing::{BatchProcessorConfig, OtlpProtocol, TracingExporterConfig},
TelemetryConfig,
};
use datadog_opentelemetry::{configuration::Config as DatadogConfig, DatadogTracingBuilder};
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 use hive_trace_context::record_graphql_document;
pub mod compatibility;
pub mod control;
pub mod hive_console_exporter;
pub(crate) mod hive_trace_context;
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;
enum TraceProviderBuilder {
OpenTelemetry(TracerProviderBuilder),
Datadog(Box<DatadogTracingBuilder>),
}
impl TraceProviderBuilder {
fn with_span_processor(self, processor: impl SpanProcessor + 'static) -> Self {
match self {
Self::OpenTelemetry(builder) => {
Self::OpenTelemetry(builder.with_span_processor(processor))
}
Self::Datadog(builder) => {
Self::Datadog(Box::new((*builder).with_span_processor(processor)))
}
}
}
fn with_span_limits(
self,
config: &crate::config::telemetry::tracing::TracingCollectConfig,
) -> Self {
match self {
Self::OpenTelemetry(builder) => Self::OpenTelemetry(
builder
.with_max_events_per_span(config.max_events_per_span)
.with_max_attributes_per_span(config.max_attributes_per_span)
.with_max_attributes_per_event(config.max_attributes_per_event)
.with_max_attributes_per_link(config.max_attributes_per_link),
),
Self::Datadog(builder) => Self::Datadog(Box::new(
(*builder)
.with_max_events_per_span(config.max_events_per_span)
.with_max_attributes_per_span(config.max_attributes_per_span)
.with_max_attributes_per_event(config.max_attributes_per_event)
.with_max_attributes_per_link(config.max_attributes_per_link),
)),
}
}
fn finish(self) -> SdkTracerProvider {
match self {
Self::OpenTelemetry(builder) => builder.build(),
Self::Datadog(builder) => (*builder).init_local().0,
}
}
}
pub(super) fn build_trace_provider<I>(
config: &TelemetryConfig,
id_generator: I,
resource: Resource,
) -> Result<SdkTracerProvider, TelemetryError>
where
I: IdGenerator + 'static,
{
let mut datadog_configs =
config
.tracing
.exporters
.iter()
.filter_map(|exporter| match exporter {
TracingExporterConfig::Datadog(config) if config.enabled => Some(config.as_ref()),
_ => None,
});
let datadog_config = datadog_configs.next();
if datadog_configs.next().is_some() {
return Err(TelemetryError::TracesExporterSetup(
"only one enabled Datadog exporter may be configured".to_string(),
));
}
let builder = if let Some(datadog) = datadog_config {
let mut datadog_config = DatadogConfig::builder();
if let Some(endpoint) = &datadog.endpoint {
let endpoint = resolve_value_or_expression(endpoint, "Datadog Agent endpoint")?;
validate_datadog_agent_url(&endpoint)?;
datadog_config.set_trace_agent_url(endpoint);
}
datadog_config.set_trace_sample_rate(config.tracing.collect.sampling);
TraceProviderBuilder::Datadog(Box::new(
datadog_opentelemetry::tracing()
.with_config(datadog_config.build())
.with_resource(resource.clone()),
))
} else {
let base_sampler = Sampler::TraceIdRatioBased(config.tracing.collect.sampling);
let mut builder = TracerProviderBuilder::default()
.with_id_generator(id_generator)
.with_resource(resource.clone());
builder = if config.tracing.collect.parent_based_sampler {
builder.with_sampler(Sampler::ParentBased(Box::new(base_sampler)))
} else {
builder.with_sampler(base_sampler)
};
TraceProviderBuilder::OpenTelemetry(builder)
}
.with_span_limits(&config.tracing.collect);
Ok(setup_exporters(config, resource, builder)?.finish())
}
fn validate_datadog_agent_url(endpoint: &str) -> Result<(), TelemetryError> {
if let Some(path) = endpoint.strip_prefix("unix://") {
return if path.is_empty() {
Err(TelemetryError::TracesExporterSetup(
"Datadog Agent endpoint must include a Unix socket path".to_string(),
))
} else {
Ok(())
};
}
if let Some(path) = endpoint.strip_prefix("windows:") {
return if path.is_empty() {
Err(TelemetryError::TracesExporterSetup(
"Datadog Agent endpoint must include a Windows named pipe path".to_string(),
))
} else {
Ok(())
};
}
let uri = endpoint.parse::<http::Uri>().map_err(|error| {
TelemetryError::TracesExporterSetup(format!(
"invalid Datadog Agent endpoint '{endpoint}': {error}"
))
})?;
match uri.scheme_str() {
Some("http" | "https") if uri.authority().is_some() => Ok(()),
Some("http" | "https") => Err(TelemetryError::TracesExporterSetup(format!(
"Datadog Agent endpoint must be absolute: '{endpoint}'"
))),
Some(scheme) => Err(TelemetryError::TracesExporterSetup(format!(
"unsupported Datadog Agent endpoint scheme '{scheme}'; expected http, https, unix, or windows"
))),
None => Err(TelemetryError::TracesExporterSetup(format!(
"Datadog Agent endpoint must include a supported scheme: '{endpoint}'"
))),
}
}
fn setup_exporters(
config: &TelemetryConfig,
resource: Resource,
mut tracer_provider_builder: TraceProviderBuilder,
) -> Result<TraceProviderBuilder, 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,
),
));
}
TracingExporterConfig::Datadog(_) => {}
}
}
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: TraceProviderBuilder,
) -> Result<TraceProviderBuilder, 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))
}
#[cfg(test)]
mod tests {
use super::validate_datadog_agent_url;
#[test]
fn validates_datadog_agent_urls() {
for endpoint in [
"http://localhost:8126",
"https://agent.example.com",
"unix:///var/run/datadog/apm.socket",
r"windows:\\.\pipe\datadog-apm",
] {
assert!(validate_datadog_agent_url(endpoint).is_ok(), "{endpoint}");
}
for endpoint in [
"not a URL",
"http://bad host",
"unix://",
"windows:",
"ftp://localhost:8126",
] {
assert!(validate_datadog_agent_url(endpoint).is_err(), "{endpoint}");
}
}
}
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(()),
}
}