mod log_claim;
mod task_output_forwarder;
pub(crate) mod task_run_telemetry;
pub(crate) use log_claim::{LOG_CLAIM_ENV, LogClaim, LogClaimWatcher};
pub(crate) use task_output_forwarder::{ExportStreams, TaskOutputForwarder};
pub(crate) use task_run_telemetry::TaskRunTelemetry;
use crate::config::Settings;
use opentelemetry_otlp::WithExportConfig;
use opentelemetry_sdk::Resource;
use opentelemetry_sdk::logs::SdkLoggerProvider;
use opentelemetry_sdk::runtime;
use opentelemetry_sdk::trace::SdkTracerProvider;
use std::time::Duration;
const DEFAULT_EXPORT_TIMEOUT: Duration = Duration::from_secs(3);
pub(crate) fn traces_enabled() -> bool {
let settings = Settings::get();
settings.otel.enabled
&& !settings.offline()
&& (env_is_set("OTEL_EXPORTER_OTLP_ENDPOINT")
|| env_is_set("OTEL_EXPORTER_OTLP_TRACES_ENDPOINT"))
}
fn env_is_set(key: &str) -> bool {
std::env::var(key).is_ok_and(|v| !v.trim().is_empty())
}
fn http_protocol(signal_var: &str) -> opentelemetry_otlp::Protocol {
select_http_protocol(
std::env::var(signal_var).ok().as_deref(),
std::env::var("OTEL_EXPORTER_OTLP_PROTOCOL").ok().as_deref(),
)
}
fn select_http_protocol(
signal: Option<&str>,
generic: Option<&str>,
) -> opentelemetry_otlp::Protocol {
use opentelemetry_otlp::Protocol;
let parse = |value: Option<&str>| match value?.trim().to_ascii_lowercase().as_str() {
"http/json" => Some(Protocol::HttpJson),
"http/protobuf" => Some(Protocol::HttpBinary),
_ => None,
};
parse(signal)
.or_else(|| parse(generic))
.unwrap_or(Protocol::HttpBinary)
}
fn export_timeout(signal_var: &str) -> Option<Duration> {
(!env_is_set(signal_var) && !env_is_set("OTEL_EXPORTER_OTLP_TIMEOUT"))
.then_some(DEFAULT_EXPORT_TIMEOUT)
}
pub(crate) fn logs_enabled() -> bool {
let settings = Settings::get();
settings.otel.logs
&& !settings.offline()
&& (env_is_set("OTEL_EXPORTER_OTLP_ENDPOINT")
|| env_is_set("OTEL_EXPORTER_OTLP_LOGS_ENDPOINT"))
}
pub(crate) fn build_resource() -> Resource {
let resource = Resource::builder().build();
let service_name = opentelemetry::Key::from_static_str("service.name");
let has_service_name = resource
.get(&service_name)
.is_some_and(|name| !name.as_str().starts_with("unknown_service"));
if has_service_name {
resource
} else {
Resource::builder().with_service_name("mise").build()
}
}
pub(crate) fn build_tracer_provider(resource: Resource) -> Option<SdkTracerProvider> {
let mut builder = opentelemetry_otlp::SpanExporter::builder()
.with_http()
.with_protocol(http_protocol("OTEL_EXPORTER_OTLP_TRACES_PROTOCOL"));
if let Some(timeout) = export_timeout("OTEL_EXPORTER_OTLP_TRACES_TIMEOUT") {
builder = builder.with_timeout(timeout);
}
let exporter = match builder.build() {
Ok(e) => e,
Err(err) => {
debug!("otel: failed to build span exporter: {err}");
return None;
}
};
Some(
SdkTracerProvider::builder()
.with_span_processor(
opentelemetry_sdk::trace::span_processor_with_async_runtime::BatchSpanProcessor::builder(exporter, runtime::Tokio)
.build(),
)
.with_resource(resource)
.build(),
)
}
pub(crate) fn build_logger_provider(resource: Resource) -> Option<SdkLoggerProvider> {
let mut builder = opentelemetry_otlp::LogExporter::builder()
.with_http()
.with_protocol(http_protocol("OTEL_EXPORTER_OTLP_LOGS_PROTOCOL"));
if let Some(timeout) = export_timeout("OTEL_EXPORTER_OTLP_LOGS_TIMEOUT") {
builder = builder.with_timeout(timeout);
}
let exporter = match builder.build() {
Ok(e) => e,
Err(err) => {
debug!("otel: failed to build log exporter: {err}");
return None;
}
};
Some(
SdkLoggerProvider::builder()
.with_log_processor(
opentelemetry_sdk::logs::log_processor_with_async_runtime::BatchLogProcessor::builder(exporter, runtime::Tokio)
.build(),
)
.with_resource(resource)
.build(),
)
}
#[cfg(test)]
mod tests {
use super::*;
use opentelemetry_otlp::Protocol;
#[test]
fn http_protocol_defaults_to_protobuf() {
assert_eq!(select_http_protocol(None, None), Protocol::HttpBinary);
assert_eq!(
select_http_protocol(Some("grpc"), None),
Protocol::HttpBinary
);
}
#[test]
fn http_protocol_prefers_a_valid_signal_value() {
assert_eq!(
select_http_protocol(Some("http/protobuf"), Some("http/json")),
Protocol::HttpBinary
);
assert_eq!(
select_http_protocol(Some(" HTTP/JSON "), None),
Protocol::HttpJson
);
}
#[test]
fn http_protocol_skips_an_invalid_signal_value() {
assert_eq!(
select_http_protocol(Some("typo"), Some("http/json")),
Protocol::HttpJson
);
assert_eq!(
select_http_protocol(Some(""), Some("http/json")),
Protocol::HttpJson
);
}
}