use std::sync::Arc;
use opentelemetry::trace::TracerProvider as _;
use opentelemetry_otlp::WithTonicConfig as _;
use opentelemetry_sdk::logs::SdkLoggerProvider;
use opentelemetry_sdk::metrics::{Aggregation, SdkMeterProvider};
use opentelemetry_sdk::trace::{BatchConfigBuilder, BatchSpanProcessor, SdkTracerProvider};
use tracing_subscriber::layer::SubscriberExt as _;
use tracing_subscriber::util::SubscriberInitExt as _;
use tracing_subscriber::{EnvFilter, Layer as _};
use crate::shared_reader::SharedManualReader;
use crate::trace_id_format::TraceIdFormat;
use crate::{LogFormat, SpanMetadataCleanupLayer, TelemetryArgs};
const OTLP_EXPORTER_ENV_VAR: &str = "OTEL_EXPORTER_OTLP_ENDPOINT";
#[derive(Debug, Clone, PartialEq, Eq)]
struct ResolvedTraceEndpoints {
rerun_authed: Option<String>,
standard: Option<String>,
}
impl ResolvedTraceEndpoints {
fn resolve(
rerun_telemetry_endpoint: &str,
standard_otel_endpoint: &str,
) -> anyhow::Result<Self> {
let rerun_authed = if rerun_telemetry_endpoint.is_empty() {
None
} else {
let reject = || {
anyhow::anyhow!(
"RERUN_TELEMETRY_ENDPOINT={rerun_telemetry_endpoint:?} is not a supported endpoint URL — \
accepted schemes are rerun://, rerun+http://, rerun+https://, http://, https://"
)
};
let (scheme, rest) = rerun_telemetry_endpoint
.split_once("://")
.ok_or_else(reject)?;
let transport_scheme = match scheme {
"rerun" | "rerun+https" | "https" => "https",
"rerun+http" | "http" => "http",
_ => return Err(reject()),
};
Some(format!("{transport_scheme}://{rest}"))
};
let standard =
(!standard_otel_endpoint.is_empty()).then(|| standard_otel_endpoint.to_owned());
Ok(Self {
rerun_authed,
standard,
})
}
fn any(&self) -> bool {
self.rerun_authed.is_some() || self.standard.is_some()
}
fn trace_mode(&self) -> &'static str {
match (self.rerun_authed.is_some(), self.standard.is_some()) {
(true, true) => "rerun-authed+otlp",
(true, false) => "rerun-authed",
(false, true) => "otlp",
(false, false) => "off",
}
}
fn summary(&self) -> String {
match (&self.rerun_authed, &self.standard) {
(Some(rerun), Some(std)) => format!("{rerun} + {std}"),
(Some(url), None) | (None, Some(url)) => url.clone(),
(None, None) => "off".to_owned(),
}
}
}
#[derive(Debug)]
struct AuthRefreshingSpanExporter<P: re_auth::credentials::CredentialsProvider> {
inner: opentelemetry_otlp::SpanExporter,
provider: Arc<P>,
token_cache: Arc<parking_lot::RwLock<String>>,
refresh_failing: std::sync::atomic::AtomicBool,
}
impl<P> opentelemetry_sdk::trace::SpanExporter for AuthRefreshingSpanExporter<P>
where
P: re_auth::credentials::CredentialsProvider + Send + Sync + std::fmt::Debug + 'static,
{
async fn export(
&self,
batch: Vec<opentelemetry_sdk::trace::SpanData>,
) -> opentelemetry_sdk::error::OTelSdkResult {
use std::sync::atomic::Ordering;
match self.provider.get_token().await {
Ok(Some(jwt)) => {
*self.token_cache.write() = jwt.to_string();
self.refresh_failing.store(false, Ordering::Relaxed);
}
Ok(None) => {
self.token_cache.write().clear();
self.refresh_failing.store(false, Ordering::Relaxed);
}
Err(err) => {
if !self.refresh_failing.swap(true, Ordering::Relaxed) {
tracing::warn!(
"Hub auth token refresh failed, continuing with cached token: {err}"
);
}
}
}
self.inner.export(batch).await
}
fn shutdown_with_timeout(
&self,
timeout: std::time::Duration,
) -> opentelemetry_sdk::error::OTelSdkResult {
self.inner.shutdown_with_timeout(timeout)
}
fn force_flush(&self) -> opentelemetry_sdk::error::OTelSdkResult {
self.inner.force_flush()
}
fn set_resource(&mut self, resource: &opentelemetry_sdk::Resource) {
self.inner.set_resource(resource);
}
}
fn build_rerun_authed_span_exporter(
transport_url: &str,
) -> anyhow::Result<AuthRefreshingSpanExporter<re_auth::credentials::CliCredentialsProvider>> {
use re_auth::credentials::CliCredentialsProvider;
build_rerun_authed_span_exporter_with_provider(
transport_url,
Arc::new(CliCredentialsProvider::new()),
)
}
fn build_rerun_authed_span_exporter_with_provider<P>(
transport_url: &str,
provider: Arc<P>,
) -> anyhow::Result<AuthRefreshingSpanExporter<P>>
where
P: re_auth::credentials::CredentialsProvider + Send + Sync + std::fmt::Debug + 'static,
{
let token_cache: Arc<parking_lot::RwLock<String>> =
Arc::new(parking_lot::RwLock::new(String::new()));
let mut endpoint: tonic::transport::Endpoint = transport_url.parse()?;
if transport_url.starts_with("https://") {
endpoint = endpoint.tls_config(
tonic::transport::ClientTlsConfig::new()
.with_enabled_roots()
.assume_http2(true),
)?;
}
let channel = endpoint.connect_lazy();
let token_for_interceptor: Arc<parking_lot::RwLock<String>> = Arc::clone(&token_cache);
let mut version_interceptor = re_grpc_headers::RerunVersionInterceptor::new_client(None, None);
let parse_failing: Arc<std::sync::atomic::AtomicBool> =
Arc::new(std::sync::atomic::AtomicBool::new(false));
let interceptor = move |mut req: tonic::Request<()>| -> tonic::Result<tonic::Request<()>> {
use std::sync::atomic::Ordering;
let token = token_for_interceptor.read().clone();
if !token.is_empty() {
match format!("Bearer {token}").parse() {
Ok(value) => {
req.metadata_mut().insert("authorization", value);
parse_failing.store(false, Ordering::Relaxed);
}
Err(err) => {
if !parse_failing.swap(true, Ordering::Relaxed) {
tracing::warn!(
"Cached Hub auth token failed to parse as an HTTP header value; aborting send: {err}",
);
}
return Err(tonic::Status::internal(
"cached Hub auth token is not a valid HTTP header value",
));
}
}
}
tonic::service::Interceptor::call(&mut version_interceptor, req)
};
let inner = opentelemetry_otlp::SpanExporter::builder()
.with_tonic()
.with_channel(channel)
.with_interceptor(interceptor)
.with_compression(opentelemetry_otlp::Compression::Gzip)
.build()?;
Ok(AuthRefreshingSpanExporter {
inner,
provider,
token_cache,
refresh_failing: std::sync::atomic::AtomicBool::new(false),
})
}
#[derive(Debug, Clone)]
pub struct Telemetry {
logs: Option<SdkLoggerProvider>,
traces: Option<SdkTracerProvider>,
metrics: Option<SdkMeterProvider>,
metrics_reader: Option<Arc<opentelemetry_sdk::metrics::ManualReader>>,
drop_behavior: TelemetryDropBehavior,
}
#[derive(Debug, Clone, Copy, Default)]
pub enum TelemetryDropBehavior {
Flush,
#[default]
Shutdown,
}
static TELEMETRY_ACTIVE: std::sync::atomic::AtomicBool = std::sync::atomic::AtomicBool::new(false);
pub fn is_telemetry_active() -> bool {
TELEMETRY_ACTIVE.load(std::sync::atomic::Ordering::Acquire)
}
#[cfg(test)]
pub(crate) fn set_telemetry_active_for_test(active: bool) {
TELEMETRY_ACTIVE.store(active, std::sync::atomic::Ordering::Release);
}
impl Telemetry {
pub fn flush(&self) {
let Self {
logs,
traces,
metrics,
metrics_reader: _,
drop_behavior: _,
} = self;
if let Some(logs) = logs
&& let Err(err) = logs.force_flush()
{
tracing::error!(%err, "failed to flush otel log provider");
}
if let Some(traces) = traces
&& let Err(err) = traces.force_flush()
{
tracing::error!(%err, "failed to flush otel trace provider");
}
if let Some(metrics) = metrics
&& let Err(err) = metrics.force_flush()
{
tracing::error!(%err, "failed to flush otel metric provider");
}
}
pub fn shutdown(&self) {
self.flush();
let Self {
logs,
traces,
metrics,
metrics_reader: _,
drop_behavior: _,
} = self;
if let Some(logs) = logs
&& let Err(err) = logs.shutdown()
{
tracing::error!(%err, "failed to shutdown otel log provider");
}
if let Some(traces) = traces
&& let Err(err) = traces.shutdown()
{
tracing::error!(%err, "failed to shutdown otel trace provider");
}
if let Some(metrics) = metrics
&& let Err(err) = metrics.shutdown()
{
tracing::error!(%err, "failed to shutdown otel metric provider");
}
}
}
impl Drop for Telemetry {
fn drop(&mut self) {
match self.drop_behavior {
TelemetryDropBehavior::Flush => self.flush(),
TelemetryDropBehavior::Shutdown => self.shutdown(),
}
}
}
impl Telemetry {
#[cfg(feature = "session_id_reader")]
#[must_use = "dropping this will flush and shutdown all telemetry systems"]
pub fn init_with_session_id_reader(
args: TelemetryArgs,
drop_behavior: TelemetryDropBehavior,
reader: crate::SessionIdReader,
) -> anyhow::Result<Self> {
crate::tracing_session::set_session_id_reader(reader);
Self::init(args, drop_behavior)
}
#[must_use = "dropping this will flush and shutdown all telemetry systems"]
pub fn init(args: TelemetryArgs, drop_behavior: TelemetryDropBehavior) -> anyhow::Result<Self> {
let TelemetryArgs {
tracy_enabled,
enabled,
service_name,
attributes,
log_filter,
log_test_output,
log_format,
log_closed_spans,
log_otlp_enabled,
log_endpoint,
trace_filter,
trace_endpoint,
trace_sampler,
trace_sampler_args,
metric_endpoint,
metric_interval,
metrics_listen_address: _, } = args;
let umbrella_endpoint = std::env::var(OTLP_EXPORTER_ENV_VAR)
.ok()
.filter(|s| !s.is_empty());
let resolve_endpoint = |signal: String| -> String {
if !signal.is_empty() {
signal
} else {
umbrella_endpoint.clone().unwrap_or_default()
}
};
let log_endpoint = resolve_endpoint(log_endpoint);
let trace_endpoint = resolve_endpoint(trace_endpoint);
let metric_endpoint = resolve_endpoint(metric_endpoint);
let rerun_telemetry_endpoint =
std::env::var("RERUN_TELEMETRY_ENDPOINT").unwrap_or_default();
let trace_endpoints = if enabled {
ResolvedTraceEndpoints::resolve(&rerun_telemetry_endpoint, &trace_endpoint)?
} else {
ResolvedTraceEndpoints {
rerun_authed: None,
standard: None,
}
};
let trace_mode: &'static str = trace_endpoints.trace_mode();
let traces_summary = trace_endpoints.summary();
let logs_summary: String = if log_otlp_enabled && !log_endpoint.is_empty() {
log_endpoint.clone()
} else {
"off".to_owned()
};
let metrics_summary: String = if metric_endpoint.is_empty() {
"off".to_owned()
} else {
metric_endpoint.clone()
};
let service_name_summary: String = service_name.as_deref().unwrap_or("<unset>").to_owned();
let result: anyhow::Result<Self> = (move || -> anyhow::Result<Self> {
if !enabled {
if tracy_enabled {
cfg_select! {
feature = "tracy" => {
tracing_subscriber::registry()
.with(self::tracy::tracy_layer())
.try_init()?;
}
_ => {
anyhow::bail!(
"`TRACY_ENABLED=true` but the 'tracy' feature flag is not toggled"
);
}
}
}
return Ok(Self {
logs: None,
metrics: None,
traces: None,
metrics_reader: None,
drop_behavior,
});
}
let Some(service_name) = service_name else {
anyhow::bail!(
"either `OTEL_SERVICE_NAME` or `TelemetryArgs::service_name` must be set in order to initialize telemetry"
);
};
#[expect(unsafe_code)]
unsafe {
if !log_endpoint.is_empty() {
std::env::set_var("OTEL_EXPORTER_OTLP_LOGS_ENDPOINT", &log_endpoint);
}
if !metric_endpoint.is_empty() {
std::env::set_var("OTEL_EXPORTER_OTLP_METRICS_ENDPOINT", &metric_endpoint);
}
if let Some(url) = &trace_endpoints.standard {
std::env::set_var("OTEL_EXPORTER_OTLP_TRACES_ENDPOINT", url);
}
std::env::set_var("OTEL_METRIC_EXPORT_INTERVAL", metric_interval);
std::env::set_var("OTEL_RESOURCE_ATTRIBUTES", attributes);
std::env::set_var("OTEL_SERVICE_NAME", &service_name);
std::env::set_var("OTEL_TRACES_SAMPLER", trace_sampler);
std::env::set_var("OTEL_TRACES_SAMPLER_ARG", trace_sampler_args);
}
let create_filter = |base: &str, forced: &str| {
use crate::EnvFilterExt as _;
EnvFilter::new(base)
.add_directive_if_absent(base, "aws_smithy_runtime", forced)?
.add_directive_if_absent(base, "datafusion", forced)?
.add_directive_if_absent(base, "datafusion_optimizer", forced)?
.add_directive_if_absent(base, "h2", forced)?
.add_directive_if_absent(base, "hyper", forced)?
.add_directive_if_absent(base, "hyper_util", forced)?
.add_directive_if_absent(base, "lance", forced)?
.add_directive_if_absent(base, "lance-arrow", forced)?
.add_directive_if_absent(base, "lance-core", forced)?
.add_directive_if_absent(base, "lance-datafusion", forced)?
.add_directive_if_absent(base, "lance-encoding", forced)?
.add_directive_if_absent(base, "lance-file", forced)?
.add_directive_if_absent(base, "lance-index", forced)?
.add_directive_if_absent(base, "lance-io", forced)?
.add_directive_if_absent(base, "lance-linalg", forced)?
.add_directive_if_absent(base, "lance-table", forced)?
.add_directive_if_absent(base, "lance", forced)?
.add_directive_if_absent(base, "opentelemetry-otlp", forced)?
.add_directive_if_absent(base, "opentelemetry", forced)?
.add_directive_if_absent(base, "opentelemetry_sdk", forced)?
.add_directive_if_absent(base, "rustls", forced)?
.add_directive_if_absent(base, "sqlparser", forced)?
.add_directive_if_absent(base, "tonic", forced)?
.add_directive_if_absent(base, "tonic_web", forced)?
.add_directive_if_absent(base, "tower", forced)?
.add_directive_if_absent(base, "tower_http", forced)?
.add_directive_if_absent(base, "tower_web", forced)?
.add_directive_if_absent(base, "typespec_client_core", forced)?
.add_directive_if_absent(base, "lance::index", "off")?
.add_directive_if_absent(base, "lance::io::exec", "off")?
.add_directive_if_absent(base, "lance::execution", "warn")?
.add_directive_if_absent(base, "lance::dataset::scanner", "off")?
.add_directive_if_absent(base, "lance_index", "off")?
.add_directive_if_absent(base, "lance::dataset::builder", "off")?
.add_directive_if_absent(base, "lance_encoding", "off")
};
let layer_logs_and_traces_stdio = {
let layer = tracing_subscriber::fmt::layer()
.with_writer(std::io::stderr)
.with_file(true)
.with_line_number(true)
.with_target(false)
.with_thread_ids(true)
.with_thread_names(true)
.with_span_events(if log_closed_spans {
tracing_subscriber::fmt::format::FmtSpan::CLOSE
} else {
tracing_subscriber::fmt::format::FmtSpan::NONE
});
macro_rules! handle_format {
($format:ident, $is_json:expr) => {{
let layer = layer
.$format()
.map_event_format(|f| TraceIdFormat::new(f, $is_json));
if log_test_output {
layer.with_test_writer().boxed()
} else {
layer.boxed()
}
}};
}
let layer = match log_format {
LogFormat::Pretty => handle_format!(pretty, false),
LogFormat::Compact => handle_format!(compact, false),
LogFormat::Json => handle_format!(json, true),
};
layer.with_filter(create_filter(&log_filter, "warn")?)
};
let (logger_provider, layer_logs_otlp) = if log_otlp_enabled && !log_endpoint.is_empty()
{
use opentelemetry_appender_tracing::layer::OpenTelemetryTracingBridge;
let exporter = opentelemetry_otlp::LogExporter::builder()
.with_tonic() .build()?;
let provider = SdkLoggerProvider::builder()
.with_batch_exporter(exporter)
.build();
let layer = OpenTelemetryTracingBridge::new(&provider).boxed();
(
Some(provider),
Some(layer.with_filter(create_filter(&log_filter, "warn")?)),
)
} else {
(None, None)
};
let (tracer_provider, layer_traces_otlp) = {
let mut builder = SdkTracerProvider::builder();
if trace_endpoints.any() {
let make_batch_config = || {
BatchConfigBuilder::default()
.with_max_queue_size(8192)
.with_max_export_batch_size(2048)
.build()
};
builder = builder
.with_span_processor(crate::tracestate::RerunSessionRootSpanProcessor);
if let Some(transport_url) = &trace_endpoints.rerun_authed {
let exporter = build_rerun_authed_span_exporter(transport_url)?;
builder = builder.with_span_processor(
BatchSpanProcessor::builder(exporter)
.with_batch_config(make_batch_config())
.build(),
);
}
if trace_endpoints.standard.is_some() {
let exporter = opentelemetry_otlp::SpanExporter::builder()
.with_tonic() .with_compression(opentelemetry_otlp::Compression::Gzip) .build()?;
builder = builder.with_span_processor(
BatchSpanProcessor::builder(exporter)
.with_batch_config(make_batch_config())
.build(),
);
}
}
let provider = builder.build();
let propagators: Vec<
Box<dyn opentelemetry::propagation::TextMapPropagator + Send + Sync>,
> = vec![
Box::new(opentelemetry_sdk::propagation::TraceContextPropagator::new()),
Box::new(crate::tracestate::TraceStateEnricher),
];
opentelemetry::global::set_text_map_propagator(
opentelemetry::propagation::TextMapCompositePropagator::new(propagators),
);
opentelemetry::global::set_tracer_provider(provider.clone());
let layer = tracing_opentelemetry::layer()
.with_tracer(provider.tracer(service_name.clone()))
.with_filter(create_filter(&trace_filter, "info")?)
.boxed();
(Some(provider), Some(layer))
};
let (metric_provider, metrics_reader) = {
let mut builder = SdkMeterProvider::builder();
builder =
builder.with_view(|instrument: &opentelemetry_sdk::metrics::Instrument| {
if instrument.kind()
== opentelemetry_sdk::metrics::InstrumentKind::Histogram
{
opentelemetry_sdk::metrics::Stream::builder()
.with_aggregation(Aggregation::Base2ExponentialHistogram {
max_size: 20,
max_scale: 20,
record_min_max: true,
})
.build()
.ok()
} else {
None
}
});
if !metric_endpoint.is_empty() {
if tokio::runtime::Handle::try_current().is_ok() {
let otlp_exporter = opentelemetry_otlp::MetricExporter::builder()
.with_temporality(opentelemetry_sdk::metrics::Temporality::Cumulative)
.with_http()
.build()?;
let reader = opentelemetry_sdk::metrics::periodic_reader_with_async_runtime::PeriodicReader::builder(
otlp_exporter,
opentelemetry_sdk::runtime::Tokio,
)
.build();
builder = builder.with_reader(reader);
} else {
tracing::warn!(
"OTLP metrics endpoint is set but telemetry was initialized outside a Tokio runtime; \
skipping push-based metric export. Metrics are still available via the Prometheus \
scrape listener if one is configured."
);
}
}
let shared_reader =
SharedManualReader::new(opentelemetry_sdk::metrics::Temporality::Cumulative);
let reader_for_telemetry = shared_reader.inner();
builder = builder.with_reader(shared_reader);
let provider = builder.build();
opentelemetry::global::set_meter_provider(provider.clone());
(Some(provider), Some(reader_for_telemetry))
};
if tracy_enabled {
cfg_select! {
feature = "tracy" => {
tracing_subscriber::registry()
.with(layer_logs_otlp)
.with(layer_logs_and_traces_stdio)
.with(layer_traces_otlp)
.with(SpanMetadataCleanupLayer::default())
.with(self::tracy::tracy_layer())
.try_init()?;
}
_ => {
anyhow::bail!(
"`TRACY_ENABLED=true` but the 'tracy' feature flag is not toggled"
);
}
}
} else {
tracing_subscriber::registry()
.with(layer_logs_otlp)
.with(layer_logs_and_traces_stdio)
.with(layer_traces_otlp)
.with(SpanMetadataCleanupLayer::default())
.try_init()?;
}
crate::memory_telemetry::install_memory_use_meters();
TELEMETRY_ACTIVE.store(true, std::sync::atomic::Ordering::Release);
Ok(Self {
drop_behavior,
logs: logger_provider,
traces: tracer_provider,
metrics: metric_provider,
metrics_reader,
})
})();
match result {
Ok(self_) => {
tracing::info!(
enabled,
service = %service_name_summary,
trace_mode,
traces = %traces_summary,
logs = %logs_summary,
metrics = %metrics_summary,
tracy = tracy_enabled,
"Telemetry initialized"
);
#[cfg(feature = "tracy")]
if tracy_enabled && enabled {
tracing::warn!(
"using tracy in addition to standard telemetry stack, consider `TELEMETRY_ENABLED=false`"
);
}
Ok(self_)
}
Err(err) => {
eprintln!(
"Telemetry init failed (enabled={enabled} service={service_name_summary} trace_mode={trace_mode} traces={traces_summary} logs={logs_summary} metrics={metrics_summary} tracy={tracy_enabled}): {err:#}"
);
Err(err)
}
}
}
pub async fn start_metrics_listener(&self, addr: &str) -> anyhow::Result<()> {
let reader = self.metrics_reader.as_ref()
.ok_or_else(|| anyhow::anyhow!(
"Cannot start metrics listener: telemetry was not initialized with metrics support. \
Ensure TELEMETRY_ENABLED=true"
))?;
let reader_for_server = Arc::clone(reader);
crate::metrics_server::start_metrics_server(addr, reader_for_server).await?;
Ok(())
}
}
#[cfg(feature = "tracy")]
mod tracy {
#[derive(Default)]
pub struct TracyConfig(tracing_subscriber::fmt::format::DefaultFields);
impl tracing_tracy::Config for TracyConfig {
type Formatter = tracing_subscriber::fmt::format::DefaultFields;
fn formatter(&self) -> &Self::Formatter {
&self.0
}
fn format_fields_in_zone_name(&self) -> bool {
false
}
}
pub fn tracy_layer() -> tracing_tracy::TracyLayer<TracyConfig> {
tracing_tracy::TracyLayer::new(TracyConfig::default())
}
}
#[cfg(test)]
mod tests {
use super::ResolvedTraceEndpoints;
#[derive(Debug)]
enum Want {
Endpoints {
rerun_authed: Option<&'static str>,
standard: Option<&'static str>,
},
Err,
}
const fn rerun_only(url: &'static str) -> Want {
Want::Endpoints {
rerun_authed: Some(url),
standard: None,
}
}
const fn standard_only(url: &'static str) -> Want {
Want::Endpoints {
rerun_authed: None,
standard: Some(url),
}
}
const fn both(rerun_authed: &'static str, standard: &'static str) -> Want {
Want::Endpoints {
rerun_authed: Some(rerun_authed),
standard: Some(standard),
}
}
const NONE: Want = Want::Endpoints {
rerun_authed: None,
standard: None,
};
#[test]
fn resolve_behavior() {
let cases: &[(&str, &str, Want)] = &[
("", "", NONE),
(
"",
"https://collector:4317",
standard_only("https://collector:4317"),
),
(
"",
"http://localhost:4317",
standard_only("http://localhost:4317"),
),
("", "grpc://collector", standard_only("grpc://collector")),
(
"",
"rerun://api.example.com",
standard_only("rerun://api.example.com"),
),
(
"rerun://api.example.com",
"",
rerun_only("https://api.example.com"),
),
(
"rerun+https://api.example.com:4317",
"",
rerun_only("https://api.example.com:4317"),
),
(
"rerun+http://localhost:4317",
"",
rerun_only("http://localhost:4317"),
),
(
"rerun://host/foo/bar?x=1",
"",
rerun_only("https://host/foo/bar?x=1"),
),
(
"https://api.example.com:4317",
"",
rerun_only("https://api.example.com:4317"),
),
(
"http://localhost:4317",
"",
rerun_only("http://localhost:4317"),
),
("ftp://collector", "", Want::Err),
("grpc://collector", "", Want::Err),
("garbage", "", Want::Err),
("api.example.com", "", Want::Err),
("RERUN://host", "", Want::Err), ("Rerun+Https://host", "", Want::Err),
("HTTPS://host", "", Want::Err),
("rerun:/host", "", Want::Err),
("rerun", "", Want::Err),
("ftp://bad", "https://otel:4317", Want::Err), (
"rerun://hub",
"https://collector",
both("https://hub", "https://collector"),
),
(
"http://hub",
"https://collector",
both("http://hub", "https://collector"),
),
(
"rerun+http://hub:4317",
"https://collector:4317",
both("http://hub:4317", "https://collector:4317"),
),
];
for (rerun, otel, want) in cases {
let got = ResolvedTraceEndpoints::resolve(rerun, otel);
let matches = match (&got, want) {
(Err(_), Want::Err) => true,
(
Ok(endpoints),
Want::Endpoints {
rerun_authed,
standard,
},
) => {
endpoints.rerun_authed.as_deref() == *rerun_authed
&& endpoints.standard.as_deref() == *standard
}
_ => false,
};
assert!(
matches,
"resolve({rerun:?}, {otel:?})\n got: {got:?}\n expected: {want:?}",
);
}
}
const TEST_JWT: &str = "eyJhbGciOiJIUzI1NiIsInR5cCI6IkpXVCJ9.e30.sig";
#[test]
fn build_rejects_invalid_transport_url() {
use re_auth::credentials::StaticCredentialsProvider;
let jwt = re_auth::Jwt::try_from(TEST_JWT.to_owned()).unwrap();
let provider = super::Arc::new(StaticCredentialsProvider::new(jwt));
let result =
super::build_rerun_authed_span_exporter_with_provider("not a url at all", provider);
assert!(result.is_err(), "expected Err for malformed URL");
}
#[tokio::test(flavor = "multi_thread")]
async fn authed_exporter_sends_bearer_metadata() {
use std::time::Duration;
use opentelemetry::trace::{Tracer as _, TracerProvider as _};
use opentelemetry_sdk::trace::{BatchSpanProcessor, SdkTracerProvider};
use re_auth::credentials::StaticCredentialsProvider;
use re_test_mocks::otlp::MockOtlpCollector;
let collector = MockOtlpCollector::spawn().await;
let jwt = re_auth::Jwt::try_from(TEST_JWT.to_owned()).unwrap();
let provider = super::Arc::new(StaticCredentialsProvider::new(jwt));
let exporter =
super::build_rerun_authed_span_exporter_with_provider(&collector.endpoint(), provider)
.unwrap();
let tracer_provider = SdkTracerProvider::builder()
.with_span_processor(BatchSpanProcessor::builder(exporter).build())
.build();
let tracer = tracer_provider.tracer("test");
{
let span = tracer.start("authed_test_span");
drop(span);
}
tracer_provider.force_flush().ok();
let received = collector
.wait_for(|_| true, Duration::from_secs(10))
.await
.expect("collector should receive at least one span");
let auth = received
.metadata
.get("authorization")
.expect("authorization metadata missing")
.to_str()
.expect("authorization should be ASCII");
assert_eq!(auth, format!("Bearer {TEST_JWT}"));
}
}