From 74a931322aef8e0a519d65c67200e40dda75c1c5 Mon Sep 17 00:00:00 2001 From: mdecimus Date: Wed, 4 Dec 2024 09:32:50 +0100 Subject: [PATCH] Tracing: Include all events in OTEL traces + Include spanId in webhooks --- crates/common/src/telemetry/tracers/otel.rs | 36 ++++++++++++++------- crates/common/src/telemetry/webhooks/mod.rs | 2 +- 2 files changed, 26 insertions(+), 12 deletions(-) diff --git a/crates/common/src/telemetry/tracers/otel.rs b/crates/common/src/telemetry/tracers/otel.rs index 88cb1cdd..7591350e 100644 --- a/crates/common/src/telemetry/tracers/otel.rs +++ b/crates/common/src/telemetry/tracers/otel.rs @@ -9,6 +9,7 @@ use std::{ time::{Duration, Instant, SystemTime, UNIX_EPOCH}, }; +use ahash::AHashMap; use mail_parser::DateTime; use opentelemetry::{ logs::{AnyValue, Severity}, @@ -25,6 +26,8 @@ use trc::{ipc::subscriber::SubscriberBuilder, Event, EventDetails, Level, Teleme use crate::{config::telemetry::OtelTracer, telemetry::LONG_SLUMBER}; +const MAX_EVENTS: usize = 2048; + pub(crate) fn spawn_otel_tracer(builder: SubscriberBuilder, mut otel: OtelTracer) { let (_, mut rx) = builder.register(); tokio::spawn(async move { @@ -45,6 +48,8 @@ pub(crate) fn spawn_otel_tracer(builder: SubscriberBuilder, mut otel: OtelTracer let mut pending_logs = Vec::new(); let mut pending_spans = Vec::new(); + let mut active_spans = AHashMap::new(); + loop { // Wait for the next event or timeout let event_or_timeout = tokio::time::timeout(wakeup_time, rx.recv()).await; @@ -52,20 +57,29 @@ pub(crate) fn spawn_otel_tracer(builder: SubscriberBuilder, mut otel: OtelTracer match event_or_timeout { Ok(Some(events)) => { for event in events { - if otel.span_exporter_enable && event.inner.typ.is_span_end() { - if let Some(start_span) = event.inner.span.as_ref() { - pending_spans.push(build_span_data( - start_span, - &event, - [&event].into_iter(), - &instrumentation, - )); - } - } - if otel.log_exporter_enable { pending_logs.push(build_log_record(&event)); } + + if otel.span_exporter_enable { + if let Some(span) = event.inner.span.as_ref() { + let span_id = span.span_id().unwrap(); + if !event.inner.typ.is_span_end() { + let events = + active_spans.entry(span_id).or_insert_with(Vec::new); + if events.len() < MAX_EVENTS { + events.push(event); + } + } else if let Some(events) = active_spans.remove(&span_id) { + pending_spans.push(build_span_data( + span, + &event, + events.iter().chain(std::iter::once(&event)), + &instrumentation, + )); + } + } + } } } Ok(None) => { diff --git a/crates/common/src/telemetry/webhooks/mod.rs b/crates/common/src/telemetry/webhooks/mod.rs index edf6c290..6207fd6c 100644 --- a/crates/common/src/telemetry/webhooks/mod.rs +++ b/crates/common/src/telemetry/webhooks/mod.rs @@ -110,7 +110,7 @@ fn spawn_webhook_handler( tokio::spawn(async move { in_flight.store(true, Ordering::Relaxed); let wrapper = EventWrapper { - events: JsonEventSerializer::new(events).with_id(), + events: JsonEventSerializer::new(events).with_id().with_spans(), }; if let Err(err) = post_webhook_events(&settings, &wrapper).await {