emit_otlp 2.22.1

Emit diagnostic events to an OpenTelemetry-compatible collector.
Documentation
use std::{future::Future, pin::Pin, sync::Arc, time::Duration};

use futures_util::{StreamExt, stream::FuturesUnordered};

use crate::{
    Error,
    client::http::HttpConnection,
    client::{Channel, OtlpBuilder, OtlpInner, SignalWorker},
    data::{
        logs::{LogsEventEncoder, LogsRequestEncoder},
        metrics::{MetricsEventEncoder, MetricsRequestEncoder},
        traces::{TracesEventEncoder, TracesRequestEncoder},
    },
    internal_metrics::InternalMetrics,
};

pub(super) type Handle = std::thread::JoinHandle<()>;

impl OtlpBuilder {
    pub(super) fn try_spawn_inner_imp(
        otlp_logs: Option<emit_batcher::Sender<Channel>>,
        worker_logs: Option<SignalWorker<HttpConnection, LogsEventEncoder, LogsRequestEncoder>>,
        otlp_traces: Option<emit_batcher::Sender<Channel>>,
        worker_traces: Option<
            SignalWorker<HttpConnection, TracesEventEncoder, TracesRequestEncoder>,
        >,
        otlp_metrics: Option<emit_batcher::Sender<Channel>>,
        worker_metrics: Option<
            SignalWorker<HttpConnection, MetricsEventEncoder, MetricsRequestEncoder>,
        >,
        metrics: Arc<InternalMetrics>,
    ) -> Result<OtlpInner, Error> {
        let receive = {
            async move {
                let processors =
                    FuturesUnordered::<Pin<Box<dyn Future<Output = ()> + Send + 'static>>>::new();

                if let Some(worker) = worker_logs {
                    processors.push(Box::pin(emit_batcher::tokio::exec(worker.receiver, {
                        move |batch| {
                            let transport = worker.transport.clone();
                            let metrics = worker.metrics.clone();

                            async move { transport.send(batch, &metrics).await }
                        }
                    })));
                }

                if let Some(worker) = worker_traces {
                    processors.push(Box::pin(emit_batcher::tokio::exec(worker.receiver, {
                        move |batch| {
                            let transport = worker.transport.clone();
                            let metrics = worker.metrics.clone();

                            async move { transport.send(batch, &metrics).await }
                        }
                    })));
                }

                if let Some(worker) = worker_metrics {
                    processors.push(Box::pin(emit_batcher::tokio::exec(worker.receiver, {
                        move |batch| {
                            let transport = worker.transport.clone();
                            let metrics = worker.metrics.clone();

                            async move { transport.send(batch, &metrics).await }
                        }
                    })));
                }

                let _ = processors.into_future().await;
            }
        };

        let handle = std::thread::Builder::new()
            .name("emit_otlp_worker".into())
            .spawn(move || {
                tokio::runtime::Builder::new_current_thread()
                    .enable_all()
                    .build()
                    .unwrap()
                    .block_on(receive);
            })
            .map_err(|e| Error::new("failed to spawn background worker", e))?;

        Ok(OtlpInner {
            otlp_logs,
            otlp_traces,
            otlp_metrics,
            metrics,
            handle: Some(handle),
        })
    }
}

pub(crate) async fn flush(sender: &emit_batcher::Sender<Channel>, timeout: Duration) -> bool {
    emit_batcher::tokio::flush(sender, timeout).await
}