hive-router 0.2.0

GraphQL router for Federation, part of the Hive platform
//! This module builds the `SdkTracerProvider` from config and attaches the appropriate
//! span processors/exporters.
//!
//! Standard OTLP and stdout exporters use the SDK's `BatchSpanProcessor`,
//! while Hive tracing routes through a custom pipeline:
//! -> `TraceBatchSpanProcessor` buffers spans per trace
//! -> `HiveConsoleExporter` normalizes
//! -> OTLP exporter
//!
//! `StandardPipelineExporter` wraps the standard OTLP/stdout pipeline, applying
//! HTTP semantic convention compatibility and pipeline-level redactions/filters
//! without adding overhead to the request hot path.
//!
//! Public helpers like `TracerLayer` and tracing control functions are re-exported here
//! for the rest of the codebase to use.
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 {
    // In order to use non-blocking reqwest client,
    // we need to use BatchSpanProcessor from the span_processor_with_async_runtime module,
    // and also pass a current-thread runtime.
    // Otherwise it will panic and if we switch to blocking reqwest client,
    // then we will break the hive-console export pipeline.
    // Yeah, fun stuff. Very fun. Yeah.
    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(()),
    }
}