use std::future::Future;
use std::sync::Arc;
use std::thread;
use std::time::{Instant, SystemTime, UNIX_EPOCH};
use opentelemetry::metrics::{Counter, Histogram, Meter, MeterProvider as _, UpDownCounter};
use opentelemetry::trace::TracerProvider as _;
use opentelemetry::{KeyValue, global};
use opentelemetry_otlp::WithExportConfig;
use opentelemetry_sdk::Resource;
use opentelemetry_sdk::error::OTelSdkError;
use opentelemetry_sdk::metrics::SdkMeterProvider;
use opentelemetry_sdk::propagation::TraceContextPropagator;
use opentelemetry_sdk::trace::SdkTracerProvider;
use thiserror::Error;
use tracing_subscriber::layer::SubscriberExt;
use tracing_subscriber::util::SubscriberInitExt;
use crate::Headers;
use crate::runtime::{
BlanketLayer, Context, Handler, HandlerResult, HealthProbe, HealthState, Layer, Outgoing,
PublishLayer, PublishNext, PublishPipeline, Settle,
};
pub const PUBLISH_TIME_HEADER: &str = "x-ruststream-published-at";
const SEMCONV_DURATION_BUCKETS: [f64; 14] = [
0.005, 0.01, 0.025, 0.05, 0.075, 0.1, 0.25, 0.5, 0.75, 1.0, 2.5, 5.0, 7.5, 10.0,
];
#[derive(Debug, Error)]
#[non_exhaustive]
pub enum OtelInitError {
#[error("failed to build the OTLP exporter")]
Exporter(#[from] opentelemetry_otlp::ExporterBuildError),
#[error("a tracing subscriber is already installed")]
TracingInit(#[source] Box<dyn std::error::Error + Send + Sync + 'static>),
}
#[derive(Debug, Default)]
#[must_use]
pub struct OtelBuilder {
service_name: Option<String>,
endpoint: Option<String>,
messaging_system: Option<String>,
attributes: Vec<KeyValue>,
tracing_bridge: bool,
stamp_publish_time: bool,
}
impl OtelBuilder {
fn new() -> Self {
Self {
tracing_bridge: true,
stamp_publish_time: true,
..Self::default()
}
}
pub fn service_name(mut self, name: impl Into<String>) -> Self {
self.service_name = Some(name.into());
self
}
pub fn otlp_endpoint(mut self, endpoint: impl Into<String>) -> Self {
self.endpoint = Some(endpoint.into());
self
}
pub fn messaging_system(mut self, system: impl Into<String>) -> Self {
self.messaging_system = Some(system.into());
self
}
pub fn attribute(mut self, key: impl Into<String>, value: impl Into<String>) -> Self {
self.attributes
.push(KeyValue::new(key.into(), value.into()));
self
}
pub const fn tracing_bridge(mut self, enabled: bool) -> Self {
self.tracing_bridge = enabled;
self
}
pub const fn stamp_publish_time(mut self, enabled: bool) -> Self {
self.stamp_publish_time = enabled;
self
}
pub fn init(self) -> Result<Otel, OtelInitError> {
let mut resource = Resource::builder();
if let Some(name) = &self.service_name {
resource = resource.with_service_name(name.clone());
}
let resource = resource.with_attributes(self.attributes.clone()).build();
let mut span_exporter = opentelemetry_otlp::SpanExporter::builder().with_tonic();
let mut metric_exporter = opentelemetry_otlp::MetricExporter::builder().with_tonic();
if let Some(endpoint) = &self.endpoint {
span_exporter = span_exporter.with_endpoint(endpoint.clone());
metric_exporter = metric_exporter.with_endpoint(endpoint.clone());
}
let tracer_provider = SdkTracerProvider::builder()
.with_batch_exporter(span_exporter.build()?)
.with_resource(resource.clone())
.build();
let meter_provider = SdkMeterProvider::builder()
.with_periodic_exporter(metric_exporter.build()?)
.with_resource(resource)
.build();
if self.tracing_bridge {
let bridge =
tracing_opentelemetry::layer().with_tracer(tracer_provider.tracer("ruststream"));
tracing_subscriber::registry()
.with(bridge)
.try_init()
.map_err(|err| OtelInitError::TracingInit(err.into()))?;
}
global::set_text_map_propagator(TraceContextPropagator::new());
global::set_tracer_provider(tracer_provider.clone());
global::set_meter_provider(meter_provider.clone());
Ok(self.attach(tracer_provider, meter_provider))
}
#[must_use]
pub fn attach(
self,
tracer_provider: SdkTracerProvider,
meter_provider: SdkMeterProvider,
) -> Otel {
let meter = meter_provider.meter("ruststream");
Otel {
instruments: Arc::new(Instruments::new(&meter, self.messaging_system.clone())),
stamp_publish_time: self.stamp_publish_time,
meter,
tracer_provider,
meter_provider,
}
}
}
#[derive(Clone)]
pub struct Otel {
instruments: Arc<Instruments>,
stamp_publish_time: bool,
meter: Meter,
tracer_provider: SdkTracerProvider,
meter_provider: SdkMeterProvider,
}
impl std::fmt::Debug for Otel {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("Otel").finish_non_exhaustive()
}
}
impl Otel {
pub fn builder() -> OtelBuilder {
OtelBuilder::new()
}
#[must_use]
pub fn consume_layer(&self) -> OtelConsumeLayer {
OtelConsumeLayer {
instruments: Arc::clone(&self.instruments),
}
}
#[must_use]
pub fn publish_layer(&self) -> OtelPublishLayer {
OtelPublishLayer {
instruments: Arc::clone(&self.instruments),
stamp_publish_time: self.stamp_publish_time,
}
}
pub fn observe_health(&self, probe: HealthProbe) {
let _gauge = self
.meter
.u64_observable_gauge("ruststream.app.state")
.with_description("The service lifecycle state, 0/1 per state attribute.")
.with_callback(move |observer| {
let current = probe.state();
let flags = [
("running", current.is_running()),
(
"shutting_down",
matches!(current, HealthState::ShuttingDown),
),
("stopped", matches!(current, HealthState::Stopped)),
("failed", matches!(current, HealthState::Failed { .. })),
];
for (name, active) in flags {
observer.observe(u64::from(active), &[KeyValue::new("state", name)]);
}
})
.build();
}
pub fn shutdown(self) -> Result<(), OTelSdkError> {
let traces = self.tracer_provider.shutdown();
let metrics = self.meter_provider.shutdown();
traces.and(metrics)
}
}
struct Instruments {
system: Option<KeyValue>,
consumed: Counter<u64>,
process_duration: Histogram<f64>,
processed: Counter<u64>,
in_flight: UpDownCounter<i64>,
queue_time: Histogram<f64>,
decode_failures: Counter<u64>,
panics: Counter<u64>,
sent: Counter<u64>,
operation_duration: Histogram<f64>,
payload_size: Histogram<u64>,
}
impl Instruments {
fn new(meter: &Meter, system: Option<String>) -> Self {
Self {
system: system.map(|s| KeyValue::new("messaging.system", s)),
consumed: meter
.u64_counter("messaging.client.consumed.messages")
.with_unit("{message}")
.with_description("Deliveries received, per handler.")
.build(),
process_duration: meter
.f64_histogram("messaging.process.duration")
.with_unit("s")
.with_boundaries(SEMCONV_DURATION_BUCKETS.to_vec())
.with_description("Handler processing time, per handler.")
.build(),
processed: meter
.u64_counter("ruststream.messages.processed")
.with_unit("{message}")
.with_description("Settled deliveries by outcome, per handler.")
.build(),
in_flight: meter
.i64_up_down_counter("ruststream.messages.in_flight")
.with_unit("{message}")
.with_description("Deliveries currently inside handlers, per handler.")
.build(),
queue_time: meter
.f64_histogram("ruststream.message.queue_time")
.with_unit("s")
.with_description(
"Time from publish to handler start, where the publish stamped its \
timestamp header.",
)
.build(),
decode_failures: meter
.u64_counter("ruststream.messages.decode_failures")
.with_unit("{message}")
.with_description("Deliveries whose payload failed to decode, per handler.")
.build(),
panics: meter
.u64_counter("ruststream.messages.panics")
.with_unit("{message}")
.with_description("Handler invocations that panicked, per handler.")
.build(),
sent: meter
.u64_counter("messaging.client.sent.messages")
.with_unit("{message}")
.with_description("Messages published, split by error.type on failures.")
.build(),
operation_duration: meter
.f64_histogram("messaging.client.operation.duration")
.with_unit("s")
.with_boundaries(SEMCONV_DURATION_BUCKETS.to_vec())
.with_description("Duration of the publish operation.")
.build(),
payload_size: meter
.u64_histogram("ruststream.message.payload.size")
.with_unit("By")
.with_description("Published payload sizes (direction attribute).")
.build(),
}
}
fn attrs(&self, destination: &str, extra: Option<KeyValue>) -> Vec<KeyValue> {
let mut attrs = Vec::with_capacity(3);
attrs.push(KeyValue::new(
"messaging.destination.name",
destination.to_owned(),
));
if let Some(system) = &self.system {
attrs.push(system.clone());
}
if let Some(extra) = extra {
attrs.push(extra);
}
attrs
}
}
const fn outcome_attr(result: HandlerResult) -> &'static str {
match result {
HandlerResult::Ack => "ack",
HandlerResult::Nack { requeue: true } => "nack_requeue",
HandlerResult::Nack { requeue: false } => "nack_drop",
HandlerResult::NackAfter { .. } => "retry_after",
}
}
fn queue_time_seconds(headers: &Headers) -> Option<f64> {
let stamped: u128 = headers.get_str(PUBLISH_TIME_HEADER)?.parse().ok()?;
let now = SystemTime::now()
.duration_since(UNIX_EPOCH)
.ok()?
.as_millis();
let elapsed_millis = now.checked_sub(stamped)?;
#[allow(clippy::cast_precision_loss)] Some(elapsed_millis as f64 / 1000.0)
}
#[derive(Clone)]
pub struct OtelConsumeLayer {
instruments: Arc<Instruments>,
}
impl std::fmt::Debug for OtelConsumeLayer {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("OtelConsumeLayer").finish_non_exhaustive()
}
}
impl<H> Layer<H> for OtelConsumeLayer {
type Handler = OtelConsumeHandler<H>;
fn layer(&self, inner: H) -> Self::Handler {
OtelConsumeHandler {
inner,
instruments: Arc::clone(&self.instruments),
}
}
}
impl BlanketLayer for OtelConsumeLayer {
fn apply<M, C, S, H>(&self, handler: H) -> impl Handler<M, C, S> + 'static
where
M: Send + Sync + 'static,
C: Send + 'static,
S: Send + Sync + 'static,
H: Handler<M, C, S> + 'static,
{
self.layer(handler)
}
}
#[derive(Clone)]
pub struct OtelConsumeHandler<H> {
inner: H,
instruments: Arc<Instruments>,
}
impl<H> std::fmt::Debug for OtelConsumeHandler<H> {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("OtelConsumeHandler").finish_non_exhaustive()
}
}
impl<M, C, S, H> Handler<M, C, S> for OtelConsumeHandler<H>
where
M: Sync,
C: Send,
S: Send + Sync,
H: Handler<M, C, S>,
{
fn handle(&self, msg: &M, ctx: &mut Context<'_, C, S>) -> impl Future<Output = Settle> + Send {
let name = ctx.name().to_owned();
let queue_time = queue_time_seconds(ctx.headers());
async move {
let instruments = &self.instruments;
let attrs = instruments.attrs(&name, None);
instruments.consumed.add(1, &attrs);
if let Some(seconds) = queue_time {
instruments.queue_time.record(seconds, &attrs);
}
instruments.in_flight.add(1, &attrs);
let started = Instant::now();
let settle = {
let _in_flight = InFlightGuard {
instruments,
attrs: &attrs,
};
self.inner.handle(msg, ctx).await
};
if ctx.decode_failed() {
instruments.decode_failures.add(1, &attrs);
}
instruments
.process_duration
.record(started.elapsed().as_secs_f64(), &attrs);
instruments.processed.add(
1,
&instruments.attrs(
&name,
Some(KeyValue::new("outcome", outcome_attr(settle.outcome()))),
),
);
settle
}
}
}
struct InFlightGuard<'a> {
instruments: &'a Instruments,
attrs: &'a [KeyValue],
}
impl Drop for InFlightGuard<'_> {
fn drop(&mut self) {
self.instruments.in_flight.add(-1, self.attrs);
if thread::panicking() {
self.instruments.panics.add(1, self.attrs);
}
}
}
#[derive(Clone)]
pub struct OtelPublishLayer {
instruments: Arc<Instruments>,
stamp_publish_time: bool,
}
impl std::fmt::Debug for OtelPublishLayer {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("OtelPublishLayer").finish_non_exhaustive()
}
}
impl PublishLayer for OtelPublishLayer {
fn on_publish<'a, N: PublishPipeline, P: crate::Publisher>(
&'a self,
out: &'a mut Outgoing<'a>,
next: PublishNext<'a, N, P>,
) -> impl Future<Output = Result<(), Box<dyn std::error::Error + Send + Sync>>> + Send + 'a
{
if self.stamp_publish_time {
if let Ok(now) = SystemTime::now().duration_since(UNIX_EPOCH) {
out.headers_mut()
.insert(PUBLISH_TIME_HEADER, now.as_millis().to_string());
}
}
let name = out.name().to_owned();
#[allow(clippy::cast_possible_truncation)] let payload_bytes = out.payload().len() as u64;
async move {
let instruments = &self.instruments;
let attrs = instruments.attrs(&name, None);
instruments.payload_size.record(
payload_bytes,
&instruments.attrs(&name, Some(KeyValue::new("direction", "publish"))),
);
let started = Instant::now();
let result = next.run(out).await;
instruments
.operation_duration
.record(started.elapsed().as_secs_f64(), &attrs);
let sent_attrs = match &result {
Ok(()) => attrs,
Err(_) => instruments.attrs(&name, Some(KeyValue::new("error.type", "_OTHER"))),
};
instruments.sent.add(1, &sent_attrs);
result
}
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn outcome_attr_names_every_settlement() {
assert_eq!(outcome_attr(HandlerResult::Ack), "ack");
assert_eq!(
outcome_attr(HandlerResult::Nack { requeue: true }),
"nack_requeue"
);
assert_eq!(
outcome_attr(HandlerResult::Nack { requeue: false }),
"nack_drop"
);
assert_eq!(
outcome_attr(HandlerResult::retry_after(std::time::Duration::from_secs(
1
))),
"retry_after"
);
}
#[cfg(feature = "json")]
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn publish_layer_splits_failures_by_error_type() {
use opentelemetry_sdk::metrics::InMemoryMetricExporter;
use crate::codec::JsonCodec;
use crate::runtime::{PublishContext, PublishIdentity, PublishStack, TypedPublisher};
struct Failing;
impl crate::Publisher for Failing {
type Error = std::io::Error;
async fn publish(&self, _msg: crate::OutgoingMessage<'_>) -> Result<(), Self::Error> {
Err(std::io::Error::other("no broker"))
}
}
let exporter = InMemoryMetricExporter::default();
let provider = SdkMeterProvider::builder()
.with_periodic_exporter(exporter.clone())
.build();
let otel = Otel::builder().attach(SdkTracerProvider::builder().build(), provider.clone());
let pipeline = PublishStack::new(otel.publish_layer(), PublishIdentity);
let publisher = TypedPublisher::with_codec(Failing, JsonCodec);
let headers = Headers::new();
let cx = PublishContext::new("orders", &headers, &());
let result = publisher.publish("orders", &7_u32, &pipeline, &cx).await;
assert!(
result.is_err(),
"the failing publisher must surface its error"
);
provider.force_flush().expect("flush failed");
let recorded: Vec<String> = exporter
.get_finished_metrics()
.expect("exporter drained")
.iter()
.flat_map(opentelemetry_sdk::metrics::data::ResourceMetrics::scope_metrics)
.flat_map(opentelemetry_sdk::metrics::data::ScopeMetrics::metrics)
.map(|metric| metric.name().to_owned())
.collect();
assert!(
recorded
.iter()
.any(|name| name == "messaging.client.sent.messages"),
"the failed publish must still count as sent (split by error.type): {recorded:?}",
);
}
#[test]
fn queue_time_ignores_absent_garbage_and_future_stamps() {
let empty = Headers::new();
assert_eq!(queue_time_seconds(&empty), None);
let mut garbage = Headers::new();
garbage.insert(PUBLISH_TIME_HEADER, "not-a-number");
assert_eq!(queue_time_seconds(&garbage), None);
let mut future = Headers::new();
future.insert(PUBLISH_TIME_HEADER, u128::MAX.to_string());
assert_eq!(queue_time_seconds(&future), None);
let mut sane = Headers::new();
let now = SystemTime::now()
.duration_since(UNIX_EPOCH)
.unwrap()
.as_millis();
sane.insert(PUBLISH_TIME_HEADER, (now - 1_500).to_string());
let seconds = queue_time_seconds(&sane).expect("a past stamp must produce a queue time");
assert!(
seconds >= 1.0,
"expected at least a second of lag, got {seconds}"
);
}
}