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
}