use crate::{
Error, OtlpMetrics,
data::{self, EncodedEvent, EncodedPayload, EncodedScopeItems, RawEncoder},
internal_metrics::InternalMetrics,
};
use emit::Filter;
use emit_batcher::BatchError;
use std::{cmp, collections::HashMap, fmt, future::Future, sync::Arc, time::Duration};
use self::{
http::{HttpConnection, HttpVersion},
imp::Handle,
};
#[cfg(not(all(
target_arch = "wasm32",
target_vendor = "unknown",
target_os = "unknown"
)))]
#[path = "client/tokio.rs"]
mod imp;
#[cfg(all(
feature = "web",
target_arch = "wasm32",
target_vendor = "unknown",
target_os = "unknown"
))]
#[path = "client/web.rs"]
mod imp;
#[cfg(all(
not(feature = "web"),
target_arch = "wasm32",
target_vendor = "unknown",
target_os = "unknown"
))]
#[path = "client/stub.rs"]
mod imp;
mod channel;
mod http;
mod logs;
mod metrics;
mod traces;
pub(crate) use self::channel::{Channel, ChannelEvent};
pub use self::{logs::*, metrics::*, traces::*};
const DEFAULT_MAX_REQUEST_SIZE_BYTES: usize = 1024 * 1024; const DEFAULT_REQUEST_TIMEOUT: Duration = Duration::from_secs(30);
const DEFAULT_CHANNEL_SIZE_EVENTS: usize = 10_000;
pub struct Otlp {
inner: Option<OtlpInner>,
metrics: Arc<InternalMetrics>,
}
struct OtlpInner {
metrics: Arc<InternalMetrics>,
otlp_logs: Option<emit_batcher::Sender<Channel>>,
otlp_traces: Option<emit_batcher::Sender<Channel>>,
otlp_metrics: Option<emit_batcher::Sender<Channel>>,
#[allow(dead_code)]
handle: Option<Handle>,
}
struct SignalWorker<S, E, R> {
pub(crate) transport: Arc<OtlpTransport<S, E, R>>,
pub(crate) receiver: emit_batcher::Receiver<Channel>,
pub(crate) metrics: Arc<InternalMetrics>,
}
impl Otlp {
pub fn builder() -> OtlpBuilder {
OtlpBuilder::new()
}
pub fn metric_source(&self) -> OtlpMetrics {
OtlpMetrics {
logs_channel_metrics: self
.inner
.as_ref()
.and_then(|inner| inner.otlp_logs.as_ref())
.map(|signal| signal.metric_source()),
traces_channel_metrics: self
.inner
.as_ref()
.and_then(|inner| inner.otlp_traces.as_ref())
.map(|signal| signal.metric_source()),
metrics_channel_metrics: self
.inner
.as_ref()
.and_then(|inner| inner.otlp_metrics.as_ref())
.map(|signal| signal.metric_source()),
metrics: self.metrics.clone(),
}
}
}
#[must_use = "call `.spawn()` to complete the builder"]
pub struct OtlpBuilder {
resource: Option<Resource>,
otlp_logs: Option<OtlpLogsBuilder>,
otlp_traces: Option<OtlpTracesBuilder>,
otlp_metrics: Option<OtlpMetricsBuilder>,
}
impl OtlpBuilder {
pub fn new() -> Self {
OtlpBuilder {
resource: None,
otlp_logs: None,
otlp_traces: None,
otlp_metrics: None,
}
}
pub fn logs(mut self, builder: OtlpLogsBuilder) -> Self {
self.otlp_logs = Some(builder);
self
}
pub fn traces(mut self, builder: OtlpTracesBuilder) -> Self {
self.otlp_traces = Some(builder);
self
}
pub fn metrics(mut self, builder: OtlpMetricsBuilder) -> Self {
self.otlp_metrics = Some(builder);
self
}
pub fn resource(mut self, attributes: impl emit::Props) -> Self {
let mut resource = Resource {
attributes: HashMap::new(),
};
let _ = attributes.for_each(|k, v| {
resource.attributes.insert(k.to_owned(), v.to_owned());
std::ops::ControlFlow::Continue(())
});
self.resource = Some(resource);
self
}
pub fn spawn(self) -> Otlp {
let metrics = Arc::new(InternalMetrics::default());
let inner = match self.try_spawn_inner(metrics.clone()) {
Ok(inner) => Some(inner),
Err(err) => {
emit::error!(
rt: emit::runtime::internal(),
"OTLP configuration is invalid; no events will be written: {err}"
);
metrics.configuration_failed.increment();
None
}
};
Otlp { metrics, inner }
}
fn try_spawn_inner(self, metrics: Arc<InternalMetrics>) -> Result<OtlpInner, Error> {
let (otlp_logs, worker_logs) = match self.otlp_logs {
Some(builder) => {
let transport = Arc::new(builder.build(metrics.clone(), self.resource.as_ref())?);
let (sender, receiver) = emit_batcher::bounded(DEFAULT_CHANNEL_SIZE_EVENTS);
let worker = SignalWorker {
transport,
receiver,
metrics: metrics.clone(),
};
(Some(sender), Some(worker))
}
None => (None, None),
};
let (otlp_traces, worker_traces) = match self.otlp_traces {
Some(builder) => {
let transport = Arc::new(builder.build(metrics.clone(), self.resource.as_ref())?);
let (sender, receiver) = emit_batcher::bounded(DEFAULT_CHANNEL_SIZE_EVENTS);
let worker = SignalWorker {
transport,
receiver,
metrics: metrics.clone(),
};
(Some(sender), Some(worker))
}
None => (None, None),
};
let (otlp_metrics, worker_metrics) = match self.otlp_metrics {
Some(builder) => {
let transport = Arc::new(builder.build(metrics.clone(), self.resource.as_ref())?);
let (sender, receiver) = emit_batcher::bounded(DEFAULT_CHANNEL_SIZE_EVENTS);
let worker = SignalWorker {
transport,
receiver,
metrics: metrics.clone(),
};
(Some(sender), Some(worker))
}
None => (None, None),
};
Self::try_spawn_inner_imp(
otlp_logs,
worker_logs,
otlp_traces,
worker_traces,
otlp_metrics,
worker_metrics,
metrics,
)
}
}
pub struct OtlpTransportBuilder {
protocol: Protocol,
url_base: String,
allow_compression: bool,
url_path: Option<&'static str>,
headers: Vec<(String, String)>,
}
impl OtlpTransportBuilder {
pub fn http(dst: impl Into<String>) -> Self {
OtlpTransportBuilder {
protocol: Protocol::Http,
allow_compression: true,
url_base: dst.into(),
url_path: None,
headers: Vec::new(),
}
}
pub fn grpc(dst: impl Into<String>) -> Self {
OtlpTransportBuilder {
protocol: Protocol::Grpc,
allow_compression: true,
url_base: dst.into(),
url_path: None,
headers: Vec::new(),
}
}
pub fn headers<K: Into<String>, V: Into<String>>(
mut self,
headers: impl IntoIterator<Item = (K, V)>,
) -> Self {
self.headers = headers
.into_iter()
.map(|(k, v)| (k.into(), v.into()))
.collect();
self
}
#[cfg(feature = "gzip")]
pub fn allow_compression(mut self, allow: bool) -> Self {
self.allow_compression = allow;
self
}
fn build<E, R>(
self,
metrics: Arc<InternalMetrics>,
event_encoder: ClientEventEncoder<E>,
resource: Option<EncodedPayload>,
request_encoder: ClientRequestEncoder<R>,
) -> Result<OtlpTransport<HttpConnection, E, R>, Error> {
let mut url = self.url_base;
if let Some(path) = self.url_path {
crate::push_path(&mut url, path);
}
let request_sender = match self.protocol {
Protocol::Http => HttpConnection::new(
HttpVersion::Http1,
metrics.clone(),
url,
self.allow_compression,
self.headers,
|req| Ok(req),
move |res| {
let metrics = metrics.clone();
async move {
let status = res.http_status();
if status >= 200 && status < 300 {
metrics.http_batch_sent.increment();
Ok(())
} else {
metrics.http_batch_failed.increment();
Err(Error::msg(format_args!(
"OTLP HTTP server responded {status}"
)))
}
}
},
)?,
Protocol::Grpc => HttpConnection::new(
HttpVersion::Http2,
metrics.clone(),
url,
self.allow_compression,
self.headers,
|req| {
let content_type_header = match req.content_type_header() {
"application/x-protobuf" => "application/grpc+proto",
content_type => {
return Err(Error::msg(format_args!(
"unsupported content type '{content_type}'"
)));
}
};
let len = (u32::try_from(req.content_payload_len()).unwrap()).to_be_bytes();
Ok(
if let Some(compression) = req.content_encoding_header() {
req.with_content_encoding_header(None)
.with_content_type_header(content_type_header)
.with_headers(match compression {
"gzip" => &[("grpc-encoding", "gzip")],
compression => {
return Err(Error::msg(format_args!(
"unsupported compression '{compression}'"
)));
}
})
.with_content_frame([1, len[0], len[1], len[2], len[3]])
}
else {
req.with_content_type_header(content_type_header)
.with_content_frame([0, len[0], len[1], len[2], len[3]])
},
)
},
move |res| {
let metrics = metrics.clone();
async move {
let mut status = 0;
let mut msg = String::new();
res.stream_trailers(|k, v| match k {
"grpc-status" => {
status = v.parse().unwrap_or(0);
}
"grpc-message" => {
msg = v.into();
}
_ => {}
})
.await?;
if status == 0 {
metrics.grpc_batch_sent.increment();
Ok(())
}
else {
metrics.grpc_batch_failed.increment();
if msg.len() > 0 {
Err(Error::msg(format_args!(
"OTLP gRPC server responded {status} {msg}"
)))
} else {
Err(Error::msg(format_args!(
"OTLP gRPC server responded {status}"
)))
}
}
}
},
)?,
};
Ok(OtlpTransport {
event_encoder,
request_sender,
resource,
request_encoder,
})
}
}
pub(crate) struct OtlpTransport<S, E, R> {
event_encoder: ClientEventEncoder<E>,
request_sender: S,
resource: Option<EncodedPayload>,
request_encoder: ClientRequestEncoder<R>,
}
impl<S: ClientRequestSender, E: data::EventEncoder, R: data::RequestEncoder>
OtlpTransport<S, E, R>
{
pub(crate) async fn send(
&self,
channel: Channel,
metrics: &InternalMetrics,
) -> Result<(), BatchError<Channel>> {
let event_encoder = &self.event_encoder;
channel::batch(
channel,
DEFAULT_MAX_REQUEST_SIZE_BYTES,
metrics,
|event| event_encoder.encode_event(event.get()),
|batch| {
let batch = batch.clone();
async move {
#[emit::span(rt: emit::runtime::internal(), guard: span, "send OTLP batch of {batch_size} events to {uri}", batch_size: batch.total_items(), uri: request_sender.uri())]
async fn send_batch<S: ClientRequestSender, R: data::RequestEncoder>(
request_sender: &S,
resource: &Option<EncodedPayload>,
request_encoder: &ClientRequestEncoder<R>,
batch: &EncodedScopeItems,
) -> Result<(), BatchError<()>> {
let uri = request_sender.uri();
let batch_size = batch.total_items();
match request_sender
.send(
request_encoder.encode_request(resource.as_ref(), &batch)?,
DEFAULT_REQUEST_TIMEOUT,
)
.await
{
Ok(res) => {
span.complete_with(emit::span::completion::from_fn(|evt| {
emit::debug!(
rt: emit::runtime::internal(),
evt,
"OTLP batch of {batch_size} events to {uri}",
batch_size,
)
}));
res
}
Err(err) => {
span.complete_with(emit::span::completion::from_fn(|evt| {
emit::warn!(
rt: emit::runtime::internal(),
evt,
"OTLP batch of {batch_size} events to {uri} failed: {err}",
batch_size,
err,
)
}));
return Err(BatchError::retry(err, ()));
}
};
Ok(())
}
send_batch(
&self.request_sender,
&self.resource,
&self.request_encoder,
&batch,
)
.await
}
},
)
.await
}
}
impl Otlp {
pub async fn flush(&self, timeout: Duration) -> bool {
if let Some(ref inner) = self.inner {
inner.flush(timeout).await
} else {
true
}
}
}
impl emit::Emitter for Otlp {
fn emit<E: emit::event::ToEvent>(&self, evt: E) {
self.inner.emit(evt)
}
fn blocking_flush(&self, timeout: Duration) -> bool {
self.inner.blocking_flush(timeout)
}
}
impl OtlpInner {
fn timeout_per_signal(&self, timeout: Duration) -> Duration {
let logs = self.otlp_logs.as_ref().map(|_| 1).unwrap_or(0);
let traces = self.otlp_traces.as_ref().map(|_| 1).unwrap_or(0);
let metrics = self.otlp_metrics.as_ref().map(|_| 1).unwrap_or(0);
let budget = cmp::max(1, logs + traces + metrics);
timeout / budget
}
async fn flush(&self, timeout: Duration) -> bool {
let timeout = self.timeout_per_signal(timeout);
let mut fully_flushed = true;
if let Some(ref sender) = self.otlp_logs {
if !imp::flush(sender, timeout).await {
fully_flushed = false;
}
}
if let Some(ref sender) = self.otlp_traces {
if !imp::flush(sender, timeout).await {
fully_flushed = false;
}
}
if let Some(ref sender) = self.otlp_metrics {
if !imp::flush(sender, timeout).await {
fully_flushed = false;
}
}
fully_flushed
}
#[allow(dead_code)]
fn take_handle(&mut self) -> Handle {
self.handle.take().expect("handle already taken")
}
}
impl emit::Emitter for OtlpInner {
fn emit<E: emit::event::ToEvent>(&self, evt: E) {
let evt = evt.to_event();
if let Some(ref sender) = self.otlp_metrics {
if emit::kind::is_metric_filter().matches(&evt) {
sender.send(ChannelEvent::from_evt(evt));
return;
}
}
if let Some(ref sender) = self.otlp_traces {
if emit::kind::is_span_filter().matches(&evt) {
sender.send(ChannelEvent::from_evt(evt));
return;
}
}
if let Some(ref sender) = self.otlp_logs {
sender.send(ChannelEvent::from_evt(evt));
return;
}
self.metrics.event_discarded.increment();
}
fn blocking_flush(&self, timeout: Duration) -> bool {
let timeout = self.timeout_per_signal(timeout);
let mut fully_flushed = true;
if let Some(ref sender) = self.otlp_logs {
if !emit_batcher::blocking_flush(sender, timeout) {
fully_flushed = false;
}
}
if let Some(ref sender) = self.otlp_traces {
if !emit_batcher::blocking_flush(sender, timeout) {
fully_flushed = false;
}
}
if let Some(ref sender) = self.otlp_metrics {
if !emit_batcher::blocking_flush(sender, timeout) {
fully_flushed = false;
}
}
fully_flushed
}
}
struct Resource {
attributes: HashMap<emit::Str<'static>, emit::value::OwnedValue>,
}
enum Protocol {
Http,
Grpc,
}
#[derive(Debug, Clone, Copy)]
pub(crate) enum Encoding {
Proto,
Json,
}
impl Encoding {
pub fn of(buf: &EncodedPayload) -> Self {
match buf {
EncodedPayload::Proto(_) => Encoding::Proto,
EncodedPayload::Json(_) => Encoding::Json,
}
}
}
pub(crate) struct ClientEventEncoder<E> {
encoding: Encoding,
encoder: E,
}
impl<E> ClientEventEncoder<E> {
pub fn new(encoding: Encoding, encoder: E) -> Self {
ClientEventEncoder { encoding, encoder }
}
}
impl<E: data::EventEncoder> ClientEventEncoder<E> {
pub fn encode_event(&self, evt: &emit::event::Event<impl emit::Props>) -> Option<EncodedEvent> {
match self.encoding {
Encoding::Proto => self.encoder.encode_event::<data::Proto>(evt),
Encoding::Json => self.encoder.encode_event::<data::Json>(evt),
}
}
}
pub(crate) trait ClientRequestSender {
fn uri(&self) -> &(impl fmt::Display + 'static);
fn send(
&self,
body: EncodedPayload,
timeout: Duration,
) -> impl Future<Output = Result<(), Error>>;
}
pub(crate) struct ClientRequestEncoder<R> {
encoding: Encoding,
encoder: R,
}
impl<R> ClientRequestEncoder<R> {
pub fn new(encoding: Encoding, encoder: R) -> Self {
ClientRequestEncoder { encoding, encoder }
}
}
impl<R: data::RequestEncoder> ClientRequestEncoder<R> {
pub fn encode_request(
&self,
resource: Option<&EncodedPayload>,
items: &EncodedScopeItems,
) -> Result<EncodedPayload, BatchError<()>> {
match self.encoding {
Encoding::Proto => self
.encoder
.encode_request::<data::Proto>(resource, items)
.map_err(BatchError::no_retry),
Encoding::Json => self
.encoder
.encode_request::<data::Json>(resource, items)
.map_err(BatchError::no_retry),
}
}
}
fn encode_resource(encoding: Encoding, resource: &Resource) -> EncodedPayload {
let attributes = data::PropsResourceAttributes(&resource.attributes);
let resource = data::Resource {
attributes: &attributes,
};
match encoding {
Encoding::Proto => data::Proto::encode(&resource),
Encoding::Json => data::Json::encode(&resource),
}
}
#[cfg(test)]
mod tests {
#[test]
#[cfg(not(all(
target_arch = "wasm32",
target_vendor = "unknown",
target_os = "unknown"
)))]
fn otlp_empty_closes_bg_thread_on_drop() {
use super::*;
let mut otlp = Otlp::builder().spawn();
let handle = {
let mut inner = otlp.inner.take().unwrap();
inner.take_handle()
};
drop(otlp);
handle.join().unwrap();
}
#[test]
#[cfg(not(all(
target_arch = "wasm32",
target_vendor = "unknown",
target_os = "unknown"
)))]
fn otlp_non_empty_closes_bg_thread_on_drop() {
use super::*;
let mut otlp = Otlp::builder()
.logs(OtlpLogsBuilder::proto(OtlpTransportBuilder::http(
"http://localhost:4319",
)))
.spawn();
let handle = {
let mut inner = otlp.inner.take().unwrap();
inner.take_handle()
};
drop(otlp);
handle.join().unwrap();
}
}