Skip to main content

leviath_telemetry/
otel.rs

1//! The OpenTelemetry-backed sink: events in, OTLP spans/metrics/logs out.
2//!
3//! Span lifecycle in a polling engine: the observer can't hold RAII span
4//! guards across ticks, so this sink keeps the live `agent.run` and
5//! `agent.stage` span handles in a map keyed by run id - opened on
6//! `RunStarted`/`StageEntered`, ended on `StageExited`/`RunCompleted`. Leaf
7//! spans (`agent.inference`, `agent.tool_call`, `agent.compaction`) are
8//! emitted retroactively at completion with explicit start/end times,
9//! parented to the open stage. A daemon crash loses whatever was open;
10//! recovered runs start a fresh trace tagged `leviath.recovered = true`.
11
12use std::collections::{HashMap, HashSet};
13use std::sync::Mutex;
14use std::time::{Duration, SystemTime, UNIX_EPOCH};
15
16use leviath_core::config::ObservabilityConfig;
17use leviath_core::telemetry::{LaneHealth, LogKind, ProviderHealth, TelemetryEvent, TelemetrySink};
18use opentelemetry::logs::{AnyValue, LogRecord, Logger, LoggerProvider, Severity};
19use opentelemetry::metrics::{Counter, Histogram, Meter, MeterProvider, UpDownCounter};
20use opentelemetry::trace::{
21    Span, SpanBuilder, SpanContext, TraceContextExt, Tracer, TracerProvider,
22};
23use opentelemetry::{Context, KeyValue};
24use opentelemetry_sdk::Resource;
25use opentelemetry_sdk::logs::SdkLoggerProvider;
26use opentelemetry_sdk::metrics::SdkMeterProvider;
27use opentelemetry_sdk::trace::SdkTracerProvider;
28use opentelemetry_sdk::trace::Tracer as SdkTracer;
29
30/// The default OTLP **HTTP** endpoint. 4318 is the HTTP/protobuf port; a
31/// collector's 4317 gRPC listener will not answer these requests.
32pub const DEFAULT_ENDPOINT: &str = "http://localhost:4318";
33
34/// The default `service.name` resource attribute.
35pub const DEFAULT_SERVICE_NAME: &str = "leviath";
36
37/// The endpoint to export to: config wins, then the standard
38/// `OTEL_EXPORTER_OTLP_ENDPOINT`, then [`DEFAULT_ENDPOINT`] - the same
39/// file-over-env precedence the provider keys use.
40pub(crate) fn resolve_endpoint(cfg: &ObservabilityConfig) -> String {
41    cfg.endpoint
42        .clone()
43        .or_else(|| std::env::var("OTEL_EXPORTER_OTLP_ENDPOINT").ok())
44        .unwrap_or_else(|| DEFAULT_ENDPOINT.to_string())
45}
46
47/// The `service.name` to report: config, then `OTEL_SERVICE_NAME`, then
48/// [`DEFAULT_SERVICE_NAME`].
49pub(crate) fn resolve_service_name(cfg: &ObservabilityConfig) -> String {
50    cfg.service_name
51        .clone()
52        .or_else(|| std::env::var("OTEL_SERVICE_NAME").ok())
53        .unwrap_or_else(|| DEFAULT_SERVICE_NAME.to_string())
54}
55
56/// The per-signal OTLP HTTP URL. `with_endpoint` on an HTTP exporter takes
57/// the full path (unlike the env var, which names the base), so the signal
58/// suffix is appended here.
59pub(crate) fn signal_url(base: &str, signal: &str) -> String {
60    format!("{}/v1/{signal}", base.trim_end_matches('/'))
61}
62
63/// Milliseconds-since-epoch as the `SystemTime` the span API wants.
64pub(crate) fn ms_to_time(at_ms: i64) -> SystemTime {
65    UNIX_EPOCH + Duration::from_millis(u64::try_from(at_ms).unwrap_or(0))
66}
67
68/// The log bridge behind a hand-rolled target filter.
69///
70/// `Layer::with_filter` would be the obvious spelling, but a `Filtered` layer
71/// only works when it was part of the subscriber at construction time - this
72/// one is swapped into an initially-empty reload slot later, where the
73/// missing `FilterId` registration panics. The bridge only reacts to events,
74/// so gating `on_event` is complete filtering for it.
75struct TargetFilteredBridge<L> {
76    inner: L,
77    targets: tracing_subscriber::filter::Targets,
78}
79
80impl<S, L> tracing_subscriber::Layer<S> for TargetFilteredBridge<L>
81where
82    S: tracing::Subscriber,
83    L: tracing_subscriber::Layer<S>,
84{
85    fn on_event(&self, event: &tracing::Event<'_>, ctx: tracing_subscriber::layer::Context<'_, S>) {
86        let meta = event.metadata();
87        if self.targets.would_enable(meta.target(), meta.level()) {
88            self.inner.on_event(event, ctx);
89        }
90    }
91}
92
93/// A stage span held open between `StageEntered` and `StageExited`.
94struct OpenStage {
95    index: usize,
96    /// When the stage was entered - the other edge of the duration histogram.
97    entered_ms: i64,
98    span: opentelemetry_sdk::trace::Span,
99}
100
101/// A run span (and its current stage) held open between `RunStarted` and
102/// `RunCompleted`.
103struct OpenRun {
104    span: opentelemetry_sdk::trace::Span,
105    stage: Option<OpenStage>,
106}
107
108/// The metric instruments, built once.
109struct Instruments {
110    active: UpDownCounter<i64>,
111    tokens: Counter<u64>,
112    tool_calls: Counter<u64>,
113    stage_duration: Histogram<f64>,
114    inference_latency: Histogram<f64>,
115    runs: Counter<u64>,
116    /// Tool-lane occupancy, re-stated on every health sample. Up-down counters
117    /// rather than plain counters because these go both ways, and the sink adds
118    /// the delta from the previous sample so a scrape reads the current value.
119    lane: LaneInstruments,
120    providers: ProviderInstruments,
121}
122
123/// The daemon-wide instruments, kept together because they share the
124/// last-sample bookkeeping that turns a level into a delta.
125struct LaneInstruments {
126    tools_busy: UpDownCounter<i64>,
127    tools_queued: UpDownCounter<i64>,
128    tools_parked: UpDownCounter<i64>,
129    dead_cycles: Counter<u64>,
130    relief: Counter<u64>,
131    /// The previous sample's levels, so each one can be reported as a delta.
132    last: Mutex<LaneLevels>,
133}
134
135/// The gauge-shaped parts of [`LaneHealth`], as last reported.
136#[derive(Default, Clone, Copy)]
137struct LaneLevels {
138    tools_busy: i64,
139    tools_queued: i64,
140    tools_parked: i64,
141    dead_cycles: u32,
142}
143
144impl LaneInstruments {
145    fn new(meter: &Meter) -> Self {
146        Self {
147            tools_busy: meter
148                .i64_up_down_counter("leviath.tool_lane.busy")
149                .with_description("Tool batches currently holding lane capacity")
150                .build(),
151            tools_queued: meter
152                .i64_up_down_counter("leviath.tool_lane.queued")
153                .with_description("Tool batches waiting for lane capacity")
154                .build(),
155            tools_parked: meter
156                .i64_up_down_counter("leviath.tool_lane.parked")
157                .with_description("Tool batches parked on an unbounded wait, holding no capacity")
158                .build(),
159            dead_cycles: meter
160                .u64_counter("leviath.scheduler.dead_cycles.total")
161                .with_description("Re-drives that found the lanes full and no run moving")
162                .build(),
163            relief: meter
164                .u64_counter("leviath.tool_lane.relief.total")
165                .with_description("Extra tool-lane capacity handed out to break a wedge")
166                .build(),
167            last: Mutex::new(LaneLevels::default()),
168        }
169    }
170
171    /// Report one sample, converting the levels into the deltas an up-down
172    /// counter wants.
173    fn record(&self, health: &LaneHealth) {
174        let now = LaneLevels {
175            tools_busy: health.tools_busy as i64,
176            tools_queued: health.tools_queued as i64,
177            tools_parked: health.tools_parked as i64,
178            dead_cycles: health.dead_cycles,
179        };
180        let mut last = leviath_core::sync::lock(&self.last);
181        self.tools_busy.add(now.tools_busy - last.tools_busy, &[]);
182        self.tools_queued
183            .add(now.tools_queued - last.tools_queued, &[]);
184        self.tools_parked
185            .add(now.tools_parked - last.tools_parked, &[]);
186        // A streak that grew is that many more dead cycles; one that reset is
187        // not negative progress, it is simply nothing to add.
188        self.dead_cycles
189            .add(now.dead_cycles.saturating_sub(last.dead_cycles) as u64, &[]);
190        self.relief.add(health.relief_granted as u64, &[]);
191        *last = now;
192    }
193}
194
195/// Providers taken out of service by their circuit breaker (issue #201).
196///
197/// Per-provider levels rather than one total, because "which provider is down"
198/// is the whole question an operator has; a bare count answers none of it.
199struct ProviderInstruments {
200    open: UpDownCounter<i64>,
201    opened: Counter<u64>,
202    /// Which providers were reported open last sample, so each can be turned
203    /// into a delta the same way [`LaneInstruments`] does.
204    last: Mutex<HashSet<String>>,
205}
206
207impl ProviderInstruments {
208    fn new(meter: &Meter) -> Self {
209        Self {
210            open: meter
211                .i64_up_down_counter("leviath.provider.circuit.open")
212                .with_description("Providers currently out of service, by provider")
213                .build(),
214            opened: meter
215                .u64_counter("leviath.provider.circuit.opened.total")
216                .with_description("Times a provider's circuit has opened, by provider and reason")
217                .build(),
218            last: Mutex::new(HashSet::new()),
219        }
220    }
221
222    fn record(&self, down: &[ProviderHealth]) {
223        let now: HashSet<String> = down.iter().map(|p| p.provider.clone()).collect();
224        let mut last = leviath_core::sync::lock(&self.last);
225        for p in down {
226            let attrs = [KeyValue::new("leviath.provider", p.provider.clone())];
227            if !last.contains(&p.provider) {
228                self.open.add(1, &attrs);
229                self.opened.add(
230                    1,
231                    &[
232                        KeyValue::new("leviath.provider", p.provider.clone()),
233                        KeyValue::new("leviath.reason", p.reason.clone()),
234                    ],
235                );
236            }
237        }
238        // Anything that was open and is not any more came back.
239        for provider in last.difference(&now) {
240            self.open
241                .add(-1, &[KeyValue::new("leviath.provider", provider.clone())]);
242        }
243        *last = now;
244    }
245}
246
247impl Instruments {
248    fn new(meter: &Meter) -> Self {
249        Self {
250            active: meter
251                .i64_up_down_counter("leviath.agents.active")
252                .with_description("Currently running agent runs")
253                .build(),
254            tokens: meter
255                .u64_counter("leviath.tokens.total")
256                .with_description("Cumulative tokens by provider, model, and kind")
257                .build(),
258            tool_calls: meter
259                .u64_counter("leviath.tool_calls.total")
260                .with_description("Tool calls by tool name and outcome")
261                .build(),
262            stage_duration: meter
263                .f64_histogram("leviath.stage_duration")
264                .with_description("Wall-clock seconds per stage")
265                .with_unit("s")
266                .build(),
267            inference_latency: meter
268                .f64_histogram("leviath.inference_latency")
269                .with_description("Per-call inference latency by provider")
270                .with_unit("s")
271                .build(),
272            // Every finished run, tagged with how it ended and whether it
273            // produced anything. One counter rather than a separate "empty
274            // runs" one: a bare count of empty runs cannot be normalized,
275            // whereas this divides into a rate.
276            runs: meter
277                .u64_counter("leviath.runs.total")
278                .with_description(
279                    "Finished runs by terminal status and whether they produced output",
280                )
281                .build(),
282            lane: LaneInstruments::new(meter),
283            providers: ProviderInstruments::new(meter),
284        }
285    }
286}
287
288/// [`TelemetrySink`] that exports to OpenTelemetry providers.
289pub struct OtelSink {
290    tracer_provider: SdkTracerProvider,
291    meter_provider: SdkMeterProvider,
292    logger_provider: SdkLoggerProvider,
293    tracer: SdkTracer,
294    logger: opentelemetry_sdk::logs::SdkLogger,
295    instruments: Instruments,
296    open: Mutex<HashMap<String, OpenRun>>,
297}
298
299impl OtelSink {
300    /// Wrap already-built providers. This is the seam the tests use (with the
301    /// SDK's in-memory exporters); [`OtelSink::from_config`] is the OTLP path.
302    pub fn new(
303        tracer_provider: SdkTracerProvider,
304        meter_provider: SdkMeterProvider,
305        logger_provider: SdkLoggerProvider,
306    ) -> Self {
307        let tracer = tracer_provider.tracer("leviath");
308        let logger = logger_provider.logger("leviath");
309        let instruments = Instruments::new(&meter_provider.meter("leviath"));
310        Self {
311            tracer_provider,
312            meter_provider,
313            logger_provider,
314            tracer,
315            logger,
316            instruments,
317            open: Mutex::new(HashMap::new()),
318        }
319    }
320
321    /// Build the OTLP HTTP/protobuf export pipeline from config + OTEL env
322    /// fallbacks. Constructs a blocking HTTP client - call from a plain
323    /// thread, not a tokio runtime thread ([`crate::build_sink`] does this).
324    pub fn from_config(cfg: &ObservabilityConfig) -> Result<Self, String> {
325        use opentelemetry_otlp::WithExportConfig;
326        let endpoint = resolve_endpoint(cfg);
327        let resource = Resource::builder()
328            .with_service_name(resolve_service_name(cfg))
329            .build();
330        let spans = opentelemetry_otlp::SpanExporter::builder()
331            .with_http()
332            .with_endpoint(signal_url(&endpoint, "traces"))
333            .build();
334        let metrics = opentelemetry_otlp::MetricExporter::builder()
335            .with_http()
336            .with_endpoint(signal_url(&endpoint, "metrics"))
337            .build();
338        let logs = opentelemetry_otlp::LogExporter::builder()
339            .with_http()
340            .with_endpoint(signal_url(&endpoint, "logs"))
341            .build();
342        // The three exporters share endpoint and config, so a build failure
343        // hits all of them the same way; one error path reports whichever
344        // surfaced first (per-exporter early returns would be branches only
345        // that exporter's failure reaches).
346        let (spans, metrics, logs) = match (spans, metrics, logs) {
347            (Ok(spans), Ok(metrics), Ok(logs)) => (spans, metrics, logs),
348            (spans, metrics, logs) => {
349                let err = [spans.err(), metrics.err(), logs.err()]
350                    .into_iter()
351                    .flatten()
352                    .next()
353                    .expect("the non-Ok arm has at least one error");
354                return Err(format!("building the OTLP exporters: {err}"));
355            }
356        };
357        Ok(Self::new(
358            SdkTracerProvider::builder()
359                .with_resource(resource.clone())
360                .with_batch_exporter(spans)
361                .build(),
362            SdkMeterProvider::builder()
363                .with_resource(resource.clone())
364                .with_periodic_exporter(metrics)
365                .build(),
366            SdkLoggerProvider::builder()
367                .with_resource(resource)
368                .with_batch_exporter(logs)
369                .build(),
370        ))
371    }
372
373    /// A `tracing-subscriber` layer that forwards the process's own `tracing`
374    /// events into this sink's OTLP logs pipeline (daemon-level log export).
375    ///
376    /// These records correlate by resource attributes only - the daemon's
377    /// tracing events fire outside any run's spans, so they carry no trace
378    /// ids; the per-run [`TelemetryEvent::Log`] records are the correlated
379    /// ones. Filtered to INFO, with the OTel/HTTP stack's own targets
380    /// silenced: an export failure that logged through this bridge would
381    /// otherwise generate more exports. The filtering is done inside
382    /// `TargetFilteredBridge` rather than `Layer::with_filter` because a
383    /// `Filtered` layer swapped into a reload slot after subscriber
384    /// construction has no registered `FilterId` and panics.
385    pub fn tracing_log_layer(&self) -> crate::LogLayer {
386        use tracing_subscriber::filter::LevelFilter;
387        let bridge = opentelemetry_appender_tracing::layer::OpenTelemetryTracingBridge::new(
388            &self.logger_provider,
389        );
390        let targets = tracing_subscriber::filter::Targets::new()
391            .with_default(LevelFilter::INFO)
392            .with_target("opentelemetry", LevelFilter::OFF)
393            .with_target("opentelemetry_sdk", LevelFilter::OFF)
394            .with_target("opentelemetry-otlp", LevelFilter::OFF)
395            .with_target("hyper", LevelFilter::OFF)
396            .with_target("reqwest", LevelFilter::OFF)
397            .with_target("h2", LevelFilter::OFF);
398        Box::new(TargetFilteredBridge {
399            inner: bridge,
400            targets,
401        })
402    }
403
404    /// Start a span at `at_ms`, optionally as a child of `parent`.
405    fn start_span(
406        &self,
407        name: &'static str,
408        at_ms: i64,
409        attributes: Vec<KeyValue>,
410        parent: Option<&SpanContext>,
411    ) -> opentelemetry_sdk::trace::Span {
412        let builder = SpanBuilder::from_name(name)
413            .with_start_time(ms_to_time(at_ms))
414            .with_attributes(attributes);
415        let cx = match parent {
416            Some(parent) => Context::new().with_remote_span_context(parent.clone()),
417            None => Context::new(),
418        };
419        self.tracer.build_with_context(builder, &cx)
420    }
421
422    /// Emit a completed leaf span under the run's open stage (or, if no stage
423    /// is open, under the run itself), spanning the last `duration_ms`.
424    fn emit_leaf(&self, run_id: &str, name: &'static str, duration_ms: u64, attrs: Vec<KeyValue>) {
425        let mut open = leviath_core::sync::lock(&self.open);
426        let Some(run) = open.get_mut(run_id) else {
427            return; // events for a run this sink never saw start
428        };
429        let parent = match run.stage.as_ref() {
430            Some(stage) => stage.span.span_context().clone(),
431            None => run.span.span_context().clone(),
432        };
433        // The event fires at completion, so the span covers the last
434        // `duration_ms` ending now.
435        let end = SystemTime::now();
436        let start = end
437            .checked_sub(Duration::from_millis(duration_ms))
438            .unwrap_or(UNIX_EPOCH);
439        let builder = SpanBuilder::from_name(name)
440            .with_start_time(start)
441            .with_attributes(attrs);
442        let cx = Context::new().with_remote_span_context(parent);
443        let mut span = self.tracer.build_with_context(builder, &cx);
444        span.end_with_timestamp(end);
445    }
446
447    /// End the run's open stage span (if any) at `at_ms`.
448    fn close_stage(run: &mut OpenRun, at_ms: i64) {
449        if let Some(mut stage) = run.stage.take() {
450            stage.span.end_with_timestamp(ms_to_time(at_ms));
451        }
452    }
453}
454
455impl TelemetrySink for OtelSink {
456    fn emit(&self, event: TelemetryEvent) {
457        match event {
458            TelemetryEvent::RunStarted {
459                run_id,
460                agent_name,
461                model,
462                parent_run_id,
463                recovered,
464                at_ms,
465            } => {
466                let mut attrs = vec![
467                    KeyValue::new("leviath.run.id", run_id.clone()),
468                    KeyValue::new("leviath.agent.name", agent_name),
469                    KeyValue::new("leviath.recovered", recovered),
470                ];
471                if let Some(model) = model {
472                    attrs.push(KeyValue::new("leviath.model", model));
473                }
474                if let Some(parent) = parent_run_id {
475                    attrs.push(KeyValue::new("leviath.parent_run.id", parent));
476                }
477                let span = self.start_span("agent.run", at_ms, attrs, None);
478                self.instruments.active.add(1, &[]);
479                leviath_core::sync::lock(&self.open).insert(run_id, OpenRun { span, stage: None });
480            }
481            TelemetryEvent::StageEntered {
482                run_id,
483                stage_index,
484                stage_name,
485                at_ms,
486            } => {
487                let mut open = leviath_core::sync::lock(&self.open);
488                let Some(run) = open.get_mut(&run_id) else {
489                    return;
490                };
491                let parent = run.span.span_context().clone();
492                let span = self.start_span(
493                    "agent.stage",
494                    at_ms,
495                    vec![
496                        KeyValue::new("leviath.stage.name", stage_name),
497                        KeyValue::new("leviath.stage.index", stage_index as i64),
498                    ],
499                    Some(&parent),
500                );
501                run.stage = Some(OpenStage {
502                    index: stage_index,
503                    entered_ms: at_ms,
504                    span,
505                });
506            }
507            TelemetryEvent::StageExited {
508                run_id,
509                stage_index,
510                stage_name,
511                prompt_tokens,
512                completion_tokens,
513                at_ms,
514            } => {
515                let mut open = leviath_core::sync::lock(&self.open);
516                let Some(run) = open.get_mut(&run_id) else {
517                    return;
518                };
519                let Some(stage) = run.stage.as_mut() else {
520                    return;
521                };
522                if stage.index != stage_index {
523                    return; // stale exit for a stage this sink isn't holding
524                }
525                stage
526                    .span
527                    .set_attribute(KeyValue::new("leviath.tokens.prompt", prompt_tokens as i64));
528                stage.span.set_attribute(KeyValue::new(
529                    "leviath.tokens.completion",
530                    completion_tokens as i64,
531                ));
532                let duration_s = (at_ms - stage.entered_ms).max(0) as f64 / 1000.0;
533                Self::close_stage(run, at_ms);
534                self.instruments.stage_duration.record(
535                    duration_s,
536                    &[KeyValue::new("leviath.stage.name", stage_name)],
537                );
538            }
539            TelemetryEvent::InferenceCompleted {
540                run_id,
541                stage_name: _,
542                provider,
543                model,
544                latency_ms,
545                prompt_tokens,
546                completion_tokens,
547                cached_tokens,
548                success,
549            } => {
550                let token_attrs = [
551                    KeyValue::new("leviath.provider", provider.clone()),
552                    KeyValue::new("leviath.model", model.clone()),
553                ];
554                for (kind, count) in [
555                    ("prompt", prompt_tokens),
556                    ("completion", completion_tokens),
557                    ("cached", cached_tokens),
558                ] {
559                    if count > 0 {
560                        let mut attrs = token_attrs.to_vec();
561                        attrs.push(KeyValue::new("leviath.tokens.kind", kind));
562                        self.instruments.tokens.add(count as u64, &attrs);
563                    }
564                }
565                self.instruments.inference_latency.record(
566                    latency_ms as f64 / 1000.0,
567                    &[KeyValue::new("leviath.provider", provider.clone())],
568                );
569                self.emit_leaf(
570                    &run_id,
571                    "agent.inference",
572                    latency_ms,
573                    vec![
574                        KeyValue::new("leviath.provider", provider),
575                        KeyValue::new("leviath.model", model),
576                        KeyValue::new("leviath.tokens.prompt", prompt_tokens as i64),
577                        KeyValue::new("leviath.tokens.completion", completion_tokens as i64),
578                        KeyValue::new("leviath.tokens.cached", cached_tokens as i64),
579                        KeyValue::new("leviath.success", success),
580                    ],
581                );
582            }
583            TelemetryEvent::ToolCallCompleted {
584                run_id,
585                stage_name: _,
586                tool_name,
587                batch_latency_ms,
588                success,
589            } => {
590                self.instruments.tool_calls.add(
591                    1,
592                    &[
593                        KeyValue::new("leviath.tool.name", tool_name.clone()),
594                        KeyValue::new("leviath.outcome", if success { "ok" } else { "error" }),
595                    ],
596                );
597                self.emit_leaf(
598                    &run_id,
599                    "agent.tool_call",
600                    batch_latency_ms,
601                    vec![
602                        KeyValue::new("leviath.tool.name", tool_name),
603                        KeyValue::new("leviath.success", success),
604                        KeyValue::new("leviath.batch_latency_ms", batch_latency_ms as i64),
605                    ],
606                );
607            }
608            TelemetryEvent::CompactionCompleted {
609                run_id,
610                stage_name: _,
611                success,
612            } => {
613                self.emit_leaf(
614                    &run_id,
615                    "agent.compaction",
616                    0,
617                    vec![KeyValue::new("leviath.success", success)],
618                );
619            }
620            TelemetryEvent::RunCompleted {
621                run_id,
622                status,
623                prompt_tokens,
624                completion_tokens,
625                tool_calls,
626                empty_output,
627                at_ms,
628            } => {
629                // Counted before the span lookup: the metric answers "how many
630                // runs finished, and how many had nothing to show for it",
631                // which must not depend on whether this process happens to
632                // hold the run's open span.
633                self.instruments.runs.add(
634                    1,
635                    &[
636                        KeyValue::new("leviath.status", status.clone()),
637                        KeyValue::new("leviath.empty_output", empty_output),
638                    ],
639                );
640                let mut open = leviath_core::sync::lock(&self.open);
641                let Some(mut run) = open.remove(&run_id) else {
642                    return;
643                };
644                drop(open);
645                // The observer closes the stage first; be defensive anyway so
646                // a crash-path completion still ends cleanly.
647                Self::close_stage(&mut run, at_ms);
648                run.span
649                    .set_attribute(KeyValue::new("leviath.status", status));
650                run.span
651                    .set_attribute(KeyValue::new("leviath.tokens.prompt", prompt_tokens as i64));
652                run.span.set_attribute(KeyValue::new(
653                    "leviath.tokens.completion",
654                    completion_tokens as i64,
655                ));
656                run.span
657                    .set_attribute(KeyValue::new("leviath.tool_calls", tool_calls as i64));
658                run.span
659                    .set_attribute(KeyValue::new("leviath.empty_output", empty_output));
660                run.span.end_with_timestamp(ms_to_time(at_ms));
661                self.instruments.active.add(-1, &[]);
662            }
663            TelemetryEvent::Log {
664                run_id,
665                stage_index,
666                kind,
667                line,
668            } => {
669                let open = leviath_core::sync::lock(&self.open);
670                let Some(run) = open.get(&run_id) else {
671                    return;
672                };
673                let span_context = match run.stage.as_ref() {
674                    Some(stage) => stage.span.span_context().clone(),
675                    None => run.span.span_context().clone(),
676                };
677                let mut record = self.logger.create_log_record();
678                record.set_timestamp(SystemTime::now());
679                record.set_severity_number(Severity::Info);
680                record.set_body(AnyValue::from(line));
681                record.set_trace_context(
682                    span_context.trace_id(),
683                    span_context.span_id(),
684                    Some(span_context.trace_flags()),
685                );
686                record.add_attribute("leviath.run.id", run_id);
687                record.add_attribute("leviath.stage.index", stage_index as i64);
688                record.add_attribute(
689                    "leviath.log.kind",
690                    match kind {
691                        LogKind::Output => "output",
692                        LogKind::Runtime => "runtime",
693                    },
694                );
695                self.logger.emit(record);
696            }
697        }
698    }
699
700    fn observe_lanes(&self, health: LaneHealth) {
701        self.instruments.lane.record(&health);
702    }
703
704    fn observe_providers(&self, down: &[ProviderHealth]) {
705        self.instruments.providers.record(down);
706    }
707
708    fn force_flush(&self) {
709        let _ = self.tracer_provider.force_flush();
710        let _ = self.meter_provider.force_flush();
711        let _ = self.logger_provider.force_flush();
712    }
713}
714
715#[cfg(test)]
716mod tests {
717    use super::*;
718    use opentelemetry_sdk::logs::InMemoryLogExporter;
719    use opentelemetry_sdk::metrics::InMemoryMetricExporter;
720    use opentelemetry_sdk::trace::InMemorySpanExporter;
721    use tracing_subscriber::layer::SubscriberExt;
722
723    #[test]
724    fn tracing_log_layer_forwards_app_events_and_filters_noise() {
725        let h = harness();
726        let subscriber = tracing_subscriber::registry().with(h.sink.tracing_log_layer());
727        tracing::subscriber::with_default(subscriber, || {
728            tracing::info!(target: "leviath::daemon", "exported line");
729            tracing::info!(target: "hyper", "http-stack noise");
730            tracing::debug!(target: "leviath::daemon", "below the info floor");
731        });
732        let logs = h.logs.get_emitted_logs().unwrap();
733        assert_eq!(logs.len(), 1, "{logs:?}");
734        assert!(format!("{:?}", logs[0].record.body()).contains("exported line"));
735    }
736
737    struct Harness {
738        sink: OtelSink,
739        spans: InMemorySpanExporter,
740        metrics: InMemoryMetricExporter,
741        logs: InMemoryLogExporter,
742    }
743
744    fn harness() -> Harness {
745        let spans = InMemorySpanExporter::default();
746        let metrics = InMemoryMetricExporter::default();
747        let logs = InMemoryLogExporter::default();
748        let sink = OtelSink::new(
749            SdkTracerProvider::builder()
750                .with_simple_exporter(spans.clone())
751                .build(),
752            SdkMeterProvider::builder()
753                .with_periodic_exporter(metrics.clone())
754                .build(),
755            SdkLoggerProvider::builder()
756                .with_simple_exporter(logs.clone())
757                .build(),
758        );
759        Harness {
760            sink,
761            spans,
762            metrics,
763            logs,
764        }
765    }
766
767    fn run_started(run_id: &str, at_ms: i64) -> TelemetryEvent {
768        TelemetryEvent::RunStarted {
769            run_id: run_id.to_string(),
770            agent_name: "coder".to_string(),
771            model: Some("mock/m".to_string()),
772            parent_run_id: Some("r0".to_string()),
773            recovered: false,
774            at_ms,
775        }
776    }
777
778    fn stage_entered(run_id: &str, index: usize, at_ms: i64) -> TelemetryEvent {
779        TelemetryEvent::StageEntered {
780            run_id: run_id.to_string(),
781            stage_index: index,
782            stage_name: format!("stage{index}"),
783            at_ms,
784        }
785    }
786
787    fn stage_exited(run_id: &str, index: usize, at_ms: i64) -> TelemetryEvent {
788        TelemetryEvent::StageExited {
789            run_id: run_id.to_string(),
790            stage_index: index,
791            stage_name: format!("stage{index}"),
792            prompt_tokens: 10,
793            completion_tokens: 4,
794            at_ms,
795        }
796    }
797
798    fn run_completed(run_id: &str, at_ms: i64) -> TelemetryEvent {
799        TelemetryEvent::RunCompleted {
800            run_id: run_id.to_string(),
801            status: "complete".to_string(),
802            prompt_tokens: 10,
803            completion_tokens: 4,
804            tool_calls: 1,
805            empty_output: false,
806            at_ms,
807        }
808    }
809
810    #[test]
811    fn run_and_stage_spans_nest_with_explicit_times() {
812        let h = harness();
813        h.sink.emit(run_started("r1", 1_000));
814        h.sink.emit(stage_entered("r1", 0, 1_000));
815        h.sink.emit(stage_exited("r1", 0, 3_500));
816        h.sink.emit(run_completed("r1", 4_000));
817
818        let spans = h.spans.get_finished_spans().unwrap();
819        assert_eq!(spans.len(), 2, "{spans:?}");
820        let stage = &spans[0];
821        let run = &spans[1];
822        assert_eq!(stage.name, "agent.stage");
823        assert_eq!(run.name, "agent.run");
824        // The stage nests under the run, in the same trace.
825        assert_eq!(stage.parent_span_id, run.span_context.span_id());
826        assert_eq!(stage.span_context.trace_id(), run.span_context.trace_id());
827        // Explicit timestamps from the events, not the wall clock.
828        assert_eq!(run.start_time, ms_to_time(1_000));
829        assert_eq!(run.end_time, ms_to_time(4_000));
830        assert_eq!(stage.start_time, ms_to_time(1_000));
831        assert_eq!(stage.end_time, ms_to_time(3_500));
832        // Final status and totals landed on the run span.
833        assert!(
834            run.attributes
835                .iter()
836                .any(|kv| kv.key.as_str() == "leviath.status")
837        );
838    }
839
840    #[test]
841    fn inference_span_nests_under_the_open_stage() {
842        let h = harness();
843        h.sink.emit(run_started("r1", 0));
844        h.sink.emit(stage_entered("r1", 0, 0));
845        h.sink.emit(TelemetryEvent::InferenceCompleted {
846            run_id: "r1".to_string(),
847            stage_name: "stage0".to_string(),
848            provider: "anthropic".to_string(),
849            model: "m".to_string(),
850            latency_ms: 250,
851            prompt_tokens: 10,
852            completion_tokens: 4,
853            cached_tokens: 2,
854            success: true,
855        });
856        let spans = h.spans.get_finished_spans().unwrap();
857        assert_eq!(spans.len(), 1);
858        let inference = &spans[0];
859        assert_eq!(inference.name, "agent.inference");
860        // The stage span is still open; harvest its id by closing everything.
861        h.sink.emit(stage_exited("r1", 0, 1));
862        h.sink.emit(run_completed("r1", 1));
863        let spans = h.spans.get_finished_spans().unwrap();
864        let stage = spans.iter().find(|s| s.name == "agent.stage").unwrap();
865        assert_eq!(inference.parent_span_id, stage.span_context.span_id());
866        assert!(
867            inference
868                .attributes
869                .iter()
870                .any(|kv| kv.key.as_str() == "leviath.provider")
871        );
872    }
873
874    #[test]
875    fn leaf_spans_between_stages_nest_under_the_run() {
876        let h = harness();
877        h.sink.emit(run_started("r1", 0));
878        h.sink.emit(stage_entered("r1", 0, 0));
879        h.sink.emit(stage_exited("r1", 0, 1));
880        h.sink.emit(TelemetryEvent::ToolCallCompleted {
881            run_id: "r1".to_string(),
882            stage_name: "stage0".to_string(),
883            tool_name: "read_file".to_string(),
884            batch_latency_ms: 30,
885            success: false,
886        });
887        h.sink.emit(TelemetryEvent::CompactionCompleted {
888            run_id: "r1".to_string(),
889            stage_name: "stage0".to_string(),
890            success: true,
891        });
892        h.sink.emit(run_completed("r1", 2));
893        let spans = h.spans.get_finished_spans().unwrap();
894        let run = spans.iter().find(|s| s.name == "agent.run").unwrap();
895        let tool = spans.iter().find(|s| s.name == "agent.tool_call").unwrap();
896        let compaction = spans.iter().find(|s| s.name == "agent.compaction").unwrap();
897        assert_eq!(tool.parent_span_id, run.span_context.span_id());
898        assert_eq!(compaction.parent_span_id, run.span_context.span_id());
899    }
900
901    #[test]
902    fn a_run_without_model_or_parent_omits_those_attributes() {
903        let h = harness();
904        h.sink.emit(TelemetryEvent::RunStarted {
905            run_id: "r1".to_string(),
906            agent_name: "coder".to_string(),
907            model: None,
908            parent_run_id: None,
909            recovered: false,
910            at_ms: 0,
911        });
912        h.sink.emit(run_completed("r1", 1));
913        let spans = h.spans.get_finished_spans().unwrap();
914        let run = spans.iter().find(|s| s.name == "agent.run").unwrap();
915        let keys: Vec<&str> = run.attributes.iter().map(|kv| kv.key.as_str()).collect();
916        assert!(!keys.contains(&"leviath.model"), "{keys:?}");
917        assert!(!keys.contains(&"leviath.parent_run.id"), "{keys:?}");
918    }
919
920    #[test]
921    fn events_for_an_unknown_run_are_dropped() {
922        let h = harness();
923        h.sink.emit(stage_entered("ghost", 0, 0));
924        h.sink.emit(stage_exited("ghost", 0, 0));
925        h.sink.emit(TelemetryEvent::InferenceCompleted {
926            run_id: "ghost".to_string(),
927            stage_name: "s".to_string(),
928            provider: "p".to_string(),
929            model: "m".to_string(),
930            latency_ms: 1,
931            prompt_tokens: 0,
932            completion_tokens: 0,
933            cached_tokens: 0,
934            success: true,
935        });
936        h.sink.emit(TelemetryEvent::Log {
937            run_id: "ghost".to_string(),
938            stage_index: 0,
939            kind: LogKind::Runtime,
940            line: "x".to_string(),
941        });
942        h.sink.emit(run_completed("ghost", 0));
943        assert!(h.spans.get_finished_spans().unwrap().is_empty());
944        assert!(h.logs.get_emitted_logs().unwrap().is_empty());
945    }
946
947    #[test]
948    fn a_stale_stage_exit_is_ignored() {
949        let h = harness();
950        h.sink.emit(run_started("r1", 0));
951        h.sink.emit(stage_entered("r1", 1, 0));
952        // Wrong index: not the stage the sink holds open.
953        h.sink.emit(stage_exited("r1", 3, 5));
954        assert!(h.spans.get_finished_spans().unwrap().is_empty());
955        // No stage open at all after a real exit: a second exit is a no-op.
956        h.sink.emit(stage_exited("r1", 1, 6));
957        h.sink.emit(stage_exited("r1", 1, 7));
958        assert_eq!(h.spans.get_finished_spans().unwrap().len(), 1);
959    }
960
961    #[test]
962    fn run_completed_closes_a_still_open_stage() {
963        let h = harness();
964        h.sink.emit(run_started("r1", 0));
965        h.sink.emit(stage_entered("r1", 0, 0));
966        h.sink.emit(run_completed("r1", 9));
967        let spans = h.spans.get_finished_spans().unwrap();
968        assert_eq!(spans.len(), 2);
969        assert_eq!(spans[0].name, "agent.stage");
970        assert_eq!(spans[0].end_time, ms_to_time(9));
971    }
972
973    #[test]
974    fn log_records_carry_the_open_span_context() {
975        let h = harness();
976        h.sink.emit(run_started("r1", 0));
977        h.sink.emit(stage_entered("r1", 0, 0));
978        h.sink.emit(TelemetryEvent::Log {
979            run_id: "r1".to_string(),
980            stage_index: 0,
981            kind: LogKind::Output,
982            line: "hello".to_string(),
983        });
984        h.sink.emit(stage_exited("r1", 0, 1));
985        h.sink.emit(TelemetryEvent::Log {
986            run_id: "r1".to_string(),
987            stage_index: 0,
988            kind: LogKind::Runtime,
989            line: "between stages".to_string(),
990        });
991        h.sink.emit(run_completed("r1", 2));
992
993        let logs = h.logs.get_emitted_logs().unwrap();
994        assert_eq!(logs.len(), 2);
995        let spans = h.spans.get_finished_spans().unwrap();
996        let stage = spans.iter().find(|s| s.name == "agent.stage").unwrap();
997        let run = spans.iter().find(|s| s.name == "agent.run").unwrap();
998        let first = logs[0].record.trace_context().unwrap();
999        assert_eq!(first.trace_id, run.span_context.trace_id());
1000        assert_eq!(first.span_id, stage.span_context.span_id());
1001        let second = logs[1].record.trace_context().unwrap();
1002        assert_eq!(second.span_id, run.span_context.span_id());
1003    }
1004
1005    #[test]
1006    fn metrics_flow_from_the_event_stream() {
1007        let h = harness();
1008        h.sink.emit(run_started("r1", 0));
1009        h.sink.emit(stage_entered("r1", 0, 0));
1010        h.sink.emit(TelemetryEvent::InferenceCompleted {
1011            run_id: "r1".to_string(),
1012            stage_name: "stage0".to_string(),
1013            provider: "anthropic".to_string(),
1014            model: "m".to_string(),
1015            latency_ms: 250,
1016            prompt_tokens: 10,
1017            completion_tokens: 4,
1018            cached_tokens: 0,
1019            success: true,
1020        });
1021        h.sink.emit(TelemetryEvent::ToolCallCompleted {
1022            run_id: "r1".to_string(),
1023            stage_name: "stage0".to_string(),
1024            tool_name: "read_file".to_string(),
1025            batch_latency_ms: 30,
1026            success: true,
1027        });
1028        h.sink.emit(stage_exited("r1", 0, 2_000));
1029        h.sink.emit(run_completed("r1", 2_000));
1030        h.sink.force_flush();
1031
1032        let exported = h.metrics.get_finished_metrics().unwrap();
1033        let names: Vec<String> = exported
1034            .iter()
1035            .flat_map(|rm| rm.scope_metrics())
1036            .flat_map(|sm| sm.metrics())
1037            .map(|m| m.name().to_string())
1038            .collect();
1039        for expected in [
1040            "leviath.agents.active",
1041            "leviath.tokens.total",
1042            "leviath.tool_calls.total",
1043            "leviath.stage_duration",
1044            "leviath.inference_latency",
1045            "leviath.runs.total",
1046        ] {
1047            assert!(names.contains(&expected.to_string()), "{names:?}");
1048        }
1049    }
1050
1051    /// The empty-run counter has to survive a completion whose span this
1052    /// process never opened, or a daemon restart would silently under-count
1053    /// exactly the runs the metric exists to find (issue #192).
1054    #[test]
1055    fn a_completion_counts_even_without_an_open_span() {
1056        let h = harness();
1057        // No `RunStarted` for this id, so the span map has nothing to remove.
1058        h.sink.emit(TelemetryEvent::RunCompleted {
1059            run_id: "ghost".to_string(),
1060            status: "complete".to_string(),
1061            prompt_tokens: 0,
1062            completion_tokens: 0,
1063            tool_calls: 0,
1064            empty_output: true,
1065            at_ms: 10,
1066        });
1067        h.sink.force_flush();
1068
1069        let exported = h.metrics.get_finished_metrics().unwrap();
1070        assert!(
1071            exported
1072                .iter()
1073                .flat_map(|rm| rm.scope_metrics())
1074                .flat_map(|sm| sm.metrics())
1075                .any(|m| m.name() == "leviath.runs.total"),
1076            "a completion with no open span still counts"
1077        );
1078    }
1079
1080    /// The daemon-wide instruments. Lane occupancy goes up and down, so the sink
1081    /// turns each sample's level into a delta; a dead-cycle streak only ever
1082    /// grows, and resetting to zero adds nothing rather than going backwards.
1083    #[test]
1084    fn lane_health_samples_export_as_deltas() {
1085        let h = harness();
1086        h.sink.observe_lanes(LaneHealth {
1087            tools_busy: 3,
1088            tools_queued: 5,
1089            tools_parked: 1,
1090            tools_workers: 8,
1091            dead_cycles: 2,
1092            ..Default::default()
1093        });
1094        // Busier, one more dead cycle, and some relief handed out.
1095        h.sink.observe_lanes(LaneHealth {
1096            agents_active: 6,
1097            agents_waiting: 2,
1098            tools_busy: 8,
1099            tools_queued: 2,
1100            tools_parked: 4,
1101            tools_workers: 8,
1102            dead_cycles: 3,
1103            relief_granted: 2,
1104        });
1105        // The wedge cleared: the streak resets, which must not subtract.
1106        h.sink.observe_lanes(LaneHealth {
1107            tools_workers: 8,
1108            ..Default::default()
1109        });
1110        h.sink.force_flush();
1111
1112        let exported = h.metrics.get_finished_metrics().unwrap();
1113        let names: Vec<String> = exported
1114            .iter()
1115            .flat_map(|rm| rm.scope_metrics())
1116            .flat_map(|sm| sm.metrics())
1117            .map(|m| m.name().to_string())
1118            .collect();
1119        for expected in [
1120            "leviath.tool_lane.busy",
1121            "leviath.tool_lane.queued",
1122            "leviath.tool_lane.parked",
1123            "leviath.scheduler.dead_cycles.total",
1124            "leviath.tool_lane.relief.total",
1125        ] {
1126            assert!(names.contains(&expected.to_string()), "{names:?}");
1127        }
1128    }
1129
1130    fn down(provider: &str, reason: &str) -> ProviderHealth {
1131        ProviderHealth {
1132            provider: provider.to_string(),
1133            reason: reason.to_string(),
1134            consecutive_failures: 3,
1135            retry_in_secs: 240,
1136        }
1137    }
1138
1139    #[test]
1140    fn provider_circuit_samples_export() {
1141        let h = harness();
1142        h.sink
1143            .observe_providers(&[down("openrouter", "credits-exhausted")]);
1144        h.sink.observe_providers(&[]);
1145        h.sink.force_flush();
1146
1147        let exported = h.metrics.get_finished_metrics().unwrap();
1148        let names: Vec<String> = exported
1149            .iter()
1150            .flat_map(|rm| rm.scope_metrics())
1151            .flat_map(|sm| sm.metrics())
1152            .map(|m| m.name().to_string())
1153            .collect();
1154        for expected in [
1155            "leviath.provider.circuit.open",
1156            "leviath.provider.circuit.opened.total",
1157        ] {
1158            assert!(names.contains(&expected.to_string()), "{names:?}");
1159        }
1160    }
1161
1162    /// The level arithmetic on its own. The gauge must go back down when a
1163    /// provider recovers, or a topped-up account reads as still broken.
1164    #[test]
1165    fn a_provider_that_recovers_brings_the_gauge_back_down() {
1166        let meter = SdkMeterProvider::builder().build().meter("test");
1167        let providers = ProviderInstruments::new(&meter);
1168        let open = |p: &ProviderInstruments| p.last.lock().unwrap().clone();
1169
1170        providers.record(&[down("openrouter", "credits-exhausted")]);
1171        assert_eq!(open(&providers).len(), 1);
1172
1173        // Still down, plus a second one. Re-reporting the first must not count
1174        // it as newly opened.
1175        providers.record(&[
1176            down("openrouter", "credits-exhausted"),
1177            down("anthropic", "auth-failed"),
1178        ]);
1179        assert_eq!(open(&providers).len(), 2);
1180
1181        // One recovers.
1182        providers.record(&[down("anthropic", "auth-failed")]);
1183        let still = open(&providers);
1184        assert_eq!(still.len(), 1);
1185        assert!(still.contains("anthropic"));
1186
1187        // Everything recovers.
1188        providers.record(&[]);
1189        assert!(open(&providers).is_empty());
1190    }
1191
1192    /// The delta arithmetic on its own, where the numbers are readable.
1193    #[test]
1194    fn lane_levels_become_deltas_and_streaks_only_climb() {
1195        let meter = SdkMeterProvider::builder().build().meter("test");
1196        let lane = LaneInstruments::new(&meter);
1197        let levels = |lane: &LaneInstruments| *lane.last.lock().unwrap();
1198
1199        lane.record(&LaneHealth {
1200            tools_busy: 4,
1201            tools_queued: 2,
1202            tools_parked: 1,
1203            dead_cycles: 5,
1204            ..Default::default()
1205        });
1206        let after = levels(&lane);
1207        assert_eq!((after.tools_busy, after.tools_queued), (4, 2));
1208        assert_eq!(after.dead_cycles, 5);
1209
1210        // A quiet sample takes the levels back down and the streak to zero.
1211        lane.record(&LaneHealth::default());
1212        let after = levels(&lane);
1213        assert_eq!(
1214            (after.tools_busy, after.tools_queued, after.tools_parked),
1215            (0, 0, 0)
1216        );
1217        assert_eq!(after.dead_cycles, 0, "the streak reset");
1218    }
1219
1220    #[test]
1221    fn endpoint_and_service_resolution_prefer_config_over_env() {
1222        temp_env::with_vars(
1223            [
1224                ("OTEL_EXPORTER_OTLP_ENDPOINT", Some("http://env:4318")),
1225                ("OTEL_SERVICE_NAME", Some("env-name")),
1226            ],
1227            || {
1228                let mut cfg = ObservabilityConfig::default();
1229                assert_eq!(resolve_endpoint(&cfg), "http://env:4318");
1230                assert_eq!(resolve_service_name(&cfg), "env-name");
1231                cfg.endpoint = Some("http://file:4318".to_string());
1232                cfg.service_name = Some("file-name".to_string());
1233                assert_eq!(resolve_endpoint(&cfg), "http://file:4318");
1234                assert_eq!(resolve_service_name(&cfg), "file-name");
1235            },
1236        );
1237        temp_env::with_vars(
1238            [
1239                ("OTEL_EXPORTER_OTLP_ENDPOINT", None::<&str>),
1240                ("OTEL_SERVICE_NAME", None),
1241            ],
1242            || {
1243                let cfg = ObservabilityConfig::default();
1244                assert_eq!(resolve_endpoint(&cfg), DEFAULT_ENDPOINT);
1245                assert_eq!(resolve_service_name(&cfg), DEFAULT_SERVICE_NAME);
1246            },
1247        );
1248    }
1249
1250    #[test]
1251    fn signal_urls_append_paths_without_doubling_slashes() {
1252        assert_eq!(
1253            signal_url("http://localhost:4318", "traces"),
1254            "http://localhost:4318/v1/traces"
1255        );
1256        assert_eq!(
1257            signal_url("http://localhost:4318/", "logs"),
1258            "http://localhost:4318/v1/logs"
1259        );
1260    }
1261
1262    #[test]
1263    fn ms_to_time_clamps_negative_to_epoch() {
1264        assert_eq!(ms_to_time(-5), UNIX_EPOCH);
1265        assert_eq!(ms_to_time(1_000), UNIX_EPOCH + Duration::from_secs(1));
1266    }
1267
1268    #[test]
1269    fn otlp_pipeline_exports_to_a_live_http_endpoint() {
1270        use std::io::{Read, Write};
1271        use std::sync::Arc;
1272        use std::sync::atomic::{AtomicUsize, Ordering};
1273
1274        // A minimal OTLP receiver: answer every POST with 200.
1275        let listener = std::net::TcpListener::bind("127.0.0.1:0").unwrap();
1276        let addr = listener.local_addr().unwrap();
1277        let hits = Arc::new(AtomicUsize::new(0));
1278        let hits_seen = hits.clone();
1279        std::thread::spawn(move || {
1280            // `flatten` skips accept errors; the thread dies with the process.
1281            for mut stream in listener.incoming().flatten() {
1282                let mut buf = [0u8; 65536];
1283                let _ = stream.read(&mut buf);
1284                hits_seen.fetch_add(1, Ordering::SeqCst);
1285                let _ = stream.write_all(b"HTTP/1.1 200 OK\r\ncontent-length: 0\r\n\r\n");
1286            }
1287        });
1288
1289        let cfg = ObservabilityConfig {
1290            enabled: true,
1291            exporter: leviath_core::config::TelemetryExporterKind::Otlp,
1292            endpoint: Some(format!("http://{addr}")),
1293            service_name: Some("leviath-test".to_string()),
1294        };
1295        let sink = OtelSink::from_config(&cfg).unwrap();
1296        sink.emit(run_started("r1", 0));
1297        sink.emit(run_completed("r1", 1));
1298        sink.force_flush();
1299        assert!(
1300            hits.load(Ordering::SeqCst) > 0,
1301            "the exporter never reached the endpoint"
1302        );
1303    }
1304}