use std::{
collections::HashMap, env, error, fmt, path::Path, str::FromStr, time::Duration,
};
use futures::{future::BoxFuture, FutureExt};
use once_cell::sync::Lazy;
use opentelemetry::{
global,
propagation::{Extractor, Injector},
sdk::{
export::trace::{ExportResult, SpanData, SpanExporter},
propagation::{
BaggagePropagator, TextMapCompositePropagator, TraceContextPropagator,
},
resource::{
EnvResourceDetector, OsResourceDetector, ProcessResourceDetector,
SdkProvidedResourceDetector,
},
trace::{Config, TracerProvider},
Resource,
},
trace::TracerProvider as _,
KeyValue, Value,
};
use opentelemetry_stackdriver::{StackDriverExporter, YupAuthorizer};
use tokio::{sync::RwLock, task::JoinHandle};
use tracing_opentelemetry::OpenTelemetrySpanExt;
use tracing_subscriber::{fmt::format::FmtSpan, prelude::*, Registry};
use crate::{env_extractor::EnvExtractor, env_injector::EnvInjector, TelemetryConfig};
use super::{debug_exporter::DebugExporter, Error, Result};
trait SpanExt {
fn record_result<T, E>(&self, result: Result<T, E>) -> Result<T, E>
where
E: fmt::Display;
}
impl SpanExt for tracing::Span {
fn record_result<T, E>(&self, result: Result<T, E>) -> Result<T, E>
where
E: fmt::Display,
{
match result {
Ok(value) => {
Ok(value)
}
Err(err) => {
self.record("otel.status_code", "ERROR");
self.record("otel.status_message", err.to_string().as_str());
Err(err)
}
}
}
}
#[derive(Debug)]
struct TracerTypeParseError(String);
impl fmt::Display for TracerTypeParseError {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
write!(f, "unknown tracer type: {:?}", self.0)
}
}
impl error::Error for TracerTypeParseError {}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
enum TracerType {
CloudTrace,
Debug,
}
impl FromStr for TracerType {
type Err = TracerTypeParseError;
fn from_str(s: &str) -> Result<Self, Self::Err> {
match s {
"cloud_trace" => Ok(TracerType::CloudTrace),
"debug" => Ok(TracerType::Debug),
_ => Err(TracerTypeParseError(s.to_owned())),
}
}
}
impl fmt::Display for TracerType {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
match self {
TracerType::CloudTrace => "cloud_trace".fmt(f),
TracerType::Debug => "debug".fmt(f),
}
}
}
fn install_opentracing_globals() {
let propagator = TextMapCompositePropagator::new(vec![
Box::new(BaggagePropagator::new()),
Box::new(TraceContextPropagator::new()),
]);
global::set_text_map_propagator(propagator);
}
#[derive(Debug)]
struct BoxExporter(Box<dyn SpanExporter + 'static>);
impl BoxExporter {
fn new<E>(exporter: E) -> Self
where
E: SpanExporter + 'static,
{
Self(Box::new(exporter))
}
async fn for_tracer_type(
tracer_type: TracerType,
) -> Result<(Self, BoxFuture<'static, ()>)> {
match tracer_type {
TracerType::CloudTrace => {
let credentials_str = env::var("GCLOUD_SERVICE_ACCOUNT_KEY_PATH")
.map_err(|_| {
Error::env_var_not_set("GCLOUD_SERVICE_ACCOUNT_KEY_PATH")
})?;
let credentials_path = Path::new(&credentials_str);
let authenticator = YupAuthorizer::new(credentials_path, None)
.await
.map_err(Error::could_not_configure_tracing)?;
let (exporter, future) = StackDriverExporter::builder()
.build(authenticator)
.await
.map_err(Error::could_not_configure_tracing)?;
Ok((BoxExporter::new(exporter), future.boxed()))
}
TracerType::Debug => {
Ok((BoxExporter::new(DebugExporter), async {}.boxed()))
}
}
}
}
impl SpanExporter for BoxExporter {
fn export(&mut self, batch: Vec<SpanData>) -> BoxFuture<'static, ExportResult> {
self.0.export(batch)
}
}
static CRATE_NAME: &str = env!("CARGO_PKG_NAME");
static TRACER_JOIN_HANDLE: Lazy<RwLock<Option<JoinHandle<()>>>> =
Lazy::new(|| RwLock::new(None));
pub async fn start_tracing(config: &TelemetryConfig) -> Result<()> {
install_opentracing_globals();
let filter = tracing_subscriber::EnvFilter::from_default_env();
let tracer_type = env::var("OPINIONATED_TELEMETRY_TRACER")
.ok()
.map(|t| t.parse())
.transpose()
.map_err(Error::could_not_configure_tracing)?;
if let Some(tracer_type) = tracer_type {
let (exporter, future) = BoxExporter::for_tracer_type(tracer_type).await?;
*TRACER_JOIN_HANDLE.write().await = Some(tokio::spawn(future));
let mut resource = Resource::from_detectors(
Duration::from_secs(0),
vec![
Box::new(SdkProvidedResourceDetector),
Box::<EnvResourceDetector>::default(),
Box::new(OsResourceDetector),
Box::new(ProcessResourceDetector),
],
);
let need_service_name = match resource.get("service.name".into()) {
None => true,
Some(value) if value == Value::String("unknown_service".into()) => true,
_ => false,
};
if need_service_name {
resource = resource.merge(&Resource::new(vec![
KeyValue::new("service.name", config.service_name.clone()),
KeyValue::new("service.version", config.service_version.clone()),
]));
}
let config = Config::default().with_resource(resource);
let provider = TracerProvider::builder()
.with_config(config)
.with_simple_exporter(exporter)
.build();
let tracer = provider.tracer(CRATE_NAME);
global::set_tracer_provider(provider);
let subscriber = Registry::default()
.with(tracing_opentelemetry::layer().with_tracer(tracer))
.with(filter);
tracing::subscriber::set_global_default(subscriber)
.expect("Could not set up global logger");
tracing_log::LogTracer::init().expect("could not hook up `log` to tracing");
} else {
tracing_subscriber::fmt::Subscriber::builder()
.with_writer(std::io::stderr)
.with_span_events(FmtSpan::NEW | FmtSpan::CLOSE)
.with_env_filter(filter)
.finish()
.try_init()
.expect("could not install tracing subscriber");
}
Ok(())
}
pub async fn stop_tracing() {
opentelemetry::global::shutdown_tracer_provider();
if let Some(handle) = TRACER_JOIN_HANDLE.write().await.take() {
handle.await.expect("could not join trace exporter");
}
}
pub trait SetParentFromExtractor: Sized {
fn set_parent_from_extractor(&mut self, extractor: &dyn Extractor);
fn set_parent_from_env(&mut self) {
self.set_parent_from_extractor(&EnvExtractor::from_env());
}
}
impl SetParentFromExtractor for tracing::Span {
fn set_parent_from_extractor(&mut self, extractor: &dyn Extractor) {
global::get_text_map_propagator(|propagator| {
let context = propagator.extract(extractor);
self.set_parent(context);
});
}
}
pub fn set_parent_span_from(extractor: &dyn Extractor) {
let mut span: tracing::Span = tracing::Span::current();
span.set_parent_from_extractor(extractor);
}
pub fn set_parent_span_from_env() {
let mut span: tracing::Span = tracing::Span::current();
span.set_parent_from_env();
}
pub fn inject_current_span_into(injector: &mut dyn Injector) {
global::get_text_map_propagator(|propagator| {
let span: tracing::Span = tracing::Span::current();
let context: opentelemetry::Context = span.context();
propagator.inject_context(&context, injector);
});
}
pub fn current_span_as_env() -> impl Iterator<Item = (String, String)> {
let mut injector = EnvInjector::new();
inject_current_span_into(&mut injector);
injector.into_iter()
}
pub fn current_span_as_headers() -> HashMap<String, String> {
let mut injector = HashMap::new();
inject_current_span_into(&mut injector);
injector
}