1use 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
30pub const DEFAULT_ENDPOINT: &str = "http://localhost:4318";
33
34pub const DEFAULT_SERVICE_NAME: &str = "leviath";
36
37pub(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
47pub(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
56pub(crate) fn signal_url(base: &str, signal: &str) -> String {
60 format!("{}/v1/{signal}", base.trim_end_matches('/'))
61}
62
63pub(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
68struct 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
93struct OpenStage {
95 index: usize,
96 entered_ms: i64,
98 span: opentelemetry_sdk::trace::Span,
99}
100
101struct OpenRun {
104 span: opentelemetry_sdk::trace::Span,
105 stage: Option<OpenStage>,
106}
107
108struct 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 lane: LaneInstruments,
120 providers: ProviderInstruments,
121}
122
123struct 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 last: Mutex<LaneLevels>,
133}
134
135#[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 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 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
195struct ProviderInstruments {
200 open: UpDownCounter<i64>,
201 opened: Counter<u64>,
202 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 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 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
288pub 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 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 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 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 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 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 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; };
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 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 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; }
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 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 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 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 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 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 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 h.sink.emit(stage_exited("r1", 3, 5));
954 assert!(h.spans.get_finished_spans().unwrap().is_empty());
955 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 #[test]
1055 fn a_completion_counts_even_without_an_open_span() {
1056 let h = harness();
1057 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 #[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 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 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 #[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 providers.record(&[
1176 down("openrouter", "credits-exhausted"),
1177 down("anthropic", "auth-failed"),
1178 ]);
1179 assert_eq!(open(&providers).len(), 2);
1180
1181 providers.record(&[down("anthropic", "auth-failed")]);
1183 let still = open(&providers);
1184 assert_eq!(still.len(), 1);
1185 assert!(still.contains("anthropic"));
1186
1187 providers.record(&[]);
1189 assert!(open(&providers).is_empty());
1190 }
1191
1192 #[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 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 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 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}