Skip to main content

leviath_core/
telemetry.rs

1//! The telemetry seam: pure-data lifecycle events and the sink they flow into.
2//!
3//! The runtime's observability system translates ECS state changes into
4//! [`TelemetryEvent`] values and hands them to whatever [`TelemetrySink`] the
5//! host installed. The events carry plain data only - no SDK types - so the
6//! runtime never depends on an exporter, and tests can assert on the exact
7//! event stream with [`MemorySink`]. The OpenTelemetry-backed sink lives in
8//! `leviath-telemetry`; a host that installs nothing gets [`NoopSink`].
9
10/// What kind of per-run log line a [`TelemetryEvent::Log`] carries.
11///
12/// Mirrors the two per-stage files the persistence layer writes: `output.log`
13/// (the model's own text) and `logs.log` (tool results, token counts, errors).
14#[derive(Debug, Clone, Copy, PartialEq, Eq)]
15pub enum LogKind {
16    /// A line of assistant output (`output.log`).
17    Output,
18    /// A runtime log line - tool results, token counts, errors (`logs.log`).
19    Runtime,
20}
21
22/// One observable moment in an agent run's life.
23///
24/// Timestamps are milliseconds since the Unix epoch (`at_ms`) so a sink can
25/// reconstruct span boundaries without sub-second drift; durations are
26/// measured wall-clock milliseconds at the point the work actually ran.
27#[derive(Debug, Clone, PartialEq)]
28pub enum TelemetryEvent {
29    /// An agent run became visible to the observer.
30    RunStarted {
31        /// The run this is about; the correlation key for every later event.
32        run_id: String,
33        /// The blueprint's `[agent] name`.
34        agent_name: String,
35        /// The run-level model hint from spawn metadata, if one was recorded.
36        model: Option<String>,
37        /// Present when this run is a sub-agent of another run.
38        parent_run_id: Option<String>,
39        /// True when the run was reloaded from disk rather than freshly
40        /// spawned - its earlier life was traced (if at all) by a previous
41        /// daemon process, so this trace starts mid-run.
42        recovered: bool,
43        /// When it happened, in milliseconds since the Unix epoch.
44        at_ms: i64,
45    },
46    /// The run entered a stage (including the first).
47    StageEntered {
48        /// The run this is about.
49        run_id: String,
50        /// Zero-based position of the stage in the blueprint's stage list.
51        stage_index: usize,
52        /// The stage's name, matching its key under `[stages]`.
53        stage_name: String,
54        /// When it happened, in milliseconds since the Unix epoch.
55        at_ms: i64,
56    },
57    /// The run left a stage; token counts are the stage's own totals.
58    StageExited {
59        /// The run this is about.
60        run_id: String,
61        /// Zero-based position of the stage in the blueprint's stage list.
62        stage_index: usize,
63        /// The stage's name, matching its key under `[stages]`.
64        stage_name: String,
65        /// Input tokens billed across this stage's whole life, so a revisited
66        /// stage reports the accumulated figure rather than the last visit's.
67        prompt_tokens: usize,
68        /// Output tokens billed across this stage's whole life.
69        completion_tokens: usize,
70        /// When it happened, in milliseconds since the Unix epoch.
71        at_ms: i64,
72    },
73    /// One inference call finished (successfully or not).
74    InferenceCompleted {
75        /// The run this is about.
76        run_id: String,
77        /// The stage the call was made from.
78        stage_name: String,
79        /// Which provider served it, after fallback resolution - so this is the
80        /// one that answered, not the one the blueprint listed first.
81        provider: String,
82        /// The model identifier sent on the wire.
83        model: String,
84        /// Wall-clock time of the provider call, including retries.
85        latency_ms: u64,
86        /// Input tokens this single call billed.
87        prompt_tokens: usize,
88        /// Output tokens this single call billed.
89        completion_tokens: usize,
90        /// Input tokens served from the provider's prompt cache, already counted
91        /// within `prompt_tokens` rather than in addition to it.
92        cached_tokens: usize,
93        /// Whether a response came back. A refusal or a tool call is still a
94        /// success; only a failed call is not.
95        success: bool,
96        /// What this one call cost, when it could be priced at all.
97        ///
98        /// `None` where the model has no reported cost and no known rate. Not
99        /// zero: a zero sums into a total that looks complete and is short by an
100        /// unknown amount, which is the failure this whole accounting exists to
101        /// avoid.
102        cost_usd: Option<f64>,
103    },
104    /// One tool call finished.
105    ToolCallCompleted {
106        /// The run this is about.
107        run_id: String,
108        /// The stage the call was made from.
109        stage_name: String,
110        /// The tool as the model named it, before alias resolution.
111        tool_name: String,
112        /// Wall-clock time of the batch the call ran in. Tool calls execute
113        /// in batches and the executor reports one duration per batch, so
114        /// every call in a batch carries the same figure.
115        batch_latency_ms: u64,
116        /// Derived from the `[error] ` result-text convention every executor
117        /// uses; a heuristic, not a structured status.
118        success: bool,
119    },
120    /// A context compaction finished.
121    CompactionCompleted {
122        /// The run this is about.
123        run_id: String,
124        /// The stage whose context was compacted.
125        stage_name: String,
126        /// Whether the compaction produced a usable summary. A failure leaves
127        /// the region as it was rather than emptying it.
128        success: bool,
129    },
130    /// The run reached a terminal status; totals are run-wide.
131    RunCompleted {
132        /// The run this is about.
133        run_id: String,
134        /// The terminal status label: `complete`, `error`, or `cancelled`.
135        status: String,
136        /// Input tokens billed across the whole run.
137        prompt_tokens: usize,
138        /// Output tokens billed across the whole run.
139        completion_tokens: usize,
140        /// How many tool calls the run made, across every stage.
141        tool_calls: usize,
142        /// Whether the run stopped having modified nothing, when its blueprint
143        /// gave it a way to. `complete` says the pipeline reached the end, not
144        /// that it achieved anything; this is the difference.
145        empty_output: bool,
146        /// When it happened, in milliseconds since the Unix epoch.
147        at_ms: i64,
148    },
149    /// One per-run log line, as also written to the stage's log files.
150    Log {
151        /// The run this is about. An OTLP sink uses it to look up the run's open
152        /// span and stamp the record with its trace context.
153        run_id: String,
154        /// Which stage's log files the line also went to.
155        stage_index: usize,
156        /// Which log the line belongs in.
157        kind: LogKind,
158        /// The line itself, without a trailing newline.
159        line: String,
160    },
161}
162
163impl TelemetryEvent {
164    /// A stable short name for this event's variant.
165    ///
166    /// Exists so a test can `assert_eq!(event.kind(), "run_started")` rather
167    /// than `assert!(matches!(event, ...))` - the `matches!` non-matching arm
168    /// is a region only a *failing* assertion ever reaches, which reads as
169    /// uncovered under the workspace's 100% gate. Useful in its own right for
170    /// structured logging, where the kind is the field worth indexing on.
171    #[must_use]
172    pub fn kind(&self) -> &'static str {
173        match self {
174            Self::RunStarted { .. } => "run_started",
175            Self::StageEntered { .. } => "stage_entered",
176            Self::StageExited { .. } => "stage_exited",
177            Self::InferenceCompleted { .. } => "inference_completed",
178            Self::ToolCallCompleted { .. } => "tool_call_completed",
179            Self::CompactionCompleted { .. } => "compaction_completed",
180            Self::RunCompleted { .. } => "run_completed",
181            Self::Log { .. } => "log",
182        }
183    }
184
185    /// The run this event belongs to.
186    pub fn run_id(&self) -> &str {
187        match self {
188            Self::RunStarted { run_id, .. }
189            | Self::StageEntered { run_id, .. }
190            | Self::StageExited { run_id, .. }
191            | Self::InferenceCompleted { run_id, .. }
192            | Self::ToolCallCompleted { run_id, .. }
193            | Self::CompactionCompleted { run_id, .. }
194            | Self::RunCompleted { run_id, .. }
195            | Self::Log { run_id, .. } => run_id,
196        }
197    }
198}
199
200/// How the daemon as a whole is doing, sampled once per safety re-drive.
201///
202/// Deliberately not a [`TelemetryEvent`]: every variant of that enum belongs to
203/// one run, and this belongs to none of them. The distinction is the point. A
204/// daemon whose lanes are full and whose runs have all stopped moving emits no
205/// per-run telemetry at all, precisely because nothing is happening, so that
206/// silence is indistinguishable from an idle night without this.
207#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
208pub struct LaneHealth {
209    /// Agents doing work, or ready to.
210    pub agents_active: usize,
211    /// Agents blocked on input, a child, or a prompt.
212    pub agents_waiting: usize,
213    /// Tool batches holding lane capacity and running.
214    pub tools_busy: usize,
215    /// Tool batches waiting for lane capacity.
216    pub tools_queued: usize,
217    /// Tool batches parked on an unbounded wait, holding no capacity.
218    pub tools_parked: usize,
219    /// The tool lane's concurrency cap, including any relief granted.
220    pub tools_workers: usize,
221    /// Consecutive re-drives that found a lane at capacity and no run moving.
222    pub dead_cycles: u32,
223    /// Extra tool-lane capacity handed out on this sample, if any.
224    pub relief_granted: usize,
225}
226
227/// One provider the daemon has stopped sending work to, sampled alongside
228/// [`LaneHealth`].
229///
230/// Daemon-wide for the same reason as `LaneHealth`: a provider out of credits
231/// belongs to no single run, and the runs it kills emit nothing useful because
232/// they die before doing anything.
233///
234/// `reason` is the label rather than the runtime's own enum: this crate sits
235/// below `leviath-providers`, so the type that names it is not in scope here.
236#[derive(Debug, Clone, PartialEq, Eq)]
237pub struct ProviderHealth {
238    /// The provider taken out of service.
239    pub provider: String,
240    /// Why, as a stable lowercase label (`credits-exhausted`, `auth-failed`).
241    pub reason: String,
242    /// Consecutive failures accumulated against it.
243    pub consecutive_failures: u32,
244    /// Seconds until it is probed again.
245    pub retry_in_secs: u64,
246}
247
248/// Where telemetry events go.
249///
250/// Implementations must tolerate being called from the engine's tick loop:
251/// `emit` should hand off or record cheaply, never block on network I/O.
252pub trait TelemetrySink: Send + Sync {
253    /// Record one event.
254    fn emit(&self, event: TelemetryEvent);
255
256    /// Record one daemon-wide health sample. Default: ignore it, so a sink that
257    /// only cares about runs needs no changes.
258    fn observe_lanes(&self, _health: LaneHealth) {}
259
260    /// Record which providers are currently out of service, sampled on the same
261    /// re-drive tick as [`TelemetrySink::observe_lanes`]. Default: ignore it,
262    /// so an existing sink keeps compiling unchanged.
263    fn observe_providers(&self, _down: &[ProviderHealth]) {}
264
265    /// Flush any buffered export before shutdown. Default: nothing buffered.
266    fn force_flush(&self) {}
267}
268
269/// The sink used when no telemetry backend is installed: drops everything.
270pub struct NoopSink;
271
272impl TelemetrySink for NoopSink {
273    fn emit(&self, _event: TelemetryEvent) {}
274}
275
276/// A sink that records every event in memory, for tests to assert on.
277#[derive(Default)]
278pub struct MemorySink {
279    events: std::sync::Mutex<Vec<TelemetryEvent>>,
280    lanes: std::sync::Mutex<Vec<LaneHealth>>,
281    providers: std::sync::Mutex<Vec<Vec<ProviderHealth>>>,
282    flushes: std::sync::atomic::AtomicUsize,
283}
284
285impl MemorySink {
286    /// A snapshot of everything emitted so far, in order.
287    pub fn events(&self) -> Vec<TelemetryEvent> {
288        crate::sync::lock(&self.events).clone()
289    }
290
291    /// Every lane-health sample recorded so far, in order.
292    pub fn lane_samples(&self) -> Vec<LaneHealth> {
293        crate::sync::lock(&self.lanes).clone()
294    }
295
296    /// Every provider-health sample recorded so far, in order.
297    pub fn provider_samples(&self) -> Vec<Vec<ProviderHealth>> {
298        crate::sync::lock(&self.providers).clone()
299    }
300
301    /// How many times `force_flush` was called.
302    pub fn flush_count(&self) -> usize {
303        self.flushes.load(std::sync::atomic::Ordering::SeqCst)
304    }
305}
306
307impl TelemetrySink for MemorySink {
308    fn emit(&self, event: TelemetryEvent) {
309        crate::sync::lock(&self.events).push(event);
310    }
311
312    fn observe_lanes(&self, health: LaneHealth) {
313        crate::sync::lock(&self.lanes).push(health);
314    }
315
316    fn observe_providers(&self, down: &[ProviderHealth]) {
317        crate::sync::lock(&self.providers).push(down.to_vec());
318    }
319
320    fn force_flush(&self) {
321        self.flushes
322            .fetch_add(1, std::sync::atomic::Ordering::SeqCst);
323    }
324}
325
326#[cfg(test)]
327mod tests {
328    use super::*;
329
330    fn run_started(run_id: &str) -> TelemetryEvent {
331        TelemetryEvent::RunStarted {
332            run_id: run_id.to_string(),
333            agent_name: "coder".to_string(),
334            model: Some("claude-sonnet-5".to_string()),
335            parent_run_id: None,
336            recovered: false,
337            at_ms: 1_000,
338        }
339    }
340
341    #[test]
342    fn memory_sink_records_events_in_order() {
343        let sink = MemorySink::default();
344        sink.emit(run_started("r1"));
345        sink.emit(TelemetryEvent::StageEntered {
346            run_id: "r1".to_string(),
347            stage_index: 0,
348            stage_name: "plan".to_string(),
349            at_ms: 1_001,
350        });
351        let events = sink.events();
352        assert_eq!(events.len(), 2);
353        assert_eq!(events[0].kind(), "run_started");
354        assert_eq!(events[1].kind(), "stage_entered");
355    }
356
357    #[test]
358    fn memory_sink_records_lane_samples_and_the_noop_ignores_them() {
359        let sink = MemorySink::default();
360        assert!(sink.lane_samples().is_empty());
361        sink.observe_lanes(LaneHealth {
362            dead_cycles: 3,
363            ..Default::default()
364        });
365        assert_eq!(sink.lane_samples().len(), 1);
366        assert_eq!(sink.lane_samples()[0].dead_cycles, 3);
367
368        // The default trait body: a sink that only cares about runs drops it.
369        NoopSink.observe_lanes(LaneHealth::default());
370    }
371
372    #[test]
373    fn memory_sink_records_provider_samples_and_the_noop_ignores_them() {
374        let sink = MemorySink::default();
375        assert!(sink.provider_samples().is_empty());
376        let down = ProviderHealth {
377            provider: "openrouter".to_string(),
378            reason: "credits-exhausted".to_string(),
379            consecutive_failures: 3,
380            retry_in_secs: 240,
381        };
382        sink.observe_providers(std::slice::from_ref(&down));
383        // The empty sample matters too: it is how a collector sees a provider
384        // come back, not just go away.
385        sink.observe_providers(&[]);
386        assert_eq!(sink.provider_samples().len(), 2);
387        assert_eq!(sink.provider_samples()[0], vec![down]);
388        assert!(sink.provider_samples()[1].is_empty());
389
390        NoopSink.observe_providers(&[]);
391    }
392
393    #[test]
394    fn memory_sink_counts_flushes() {
395        let sink = MemorySink::default();
396        assert_eq!(sink.flush_count(), 0);
397        sink.force_flush();
398        sink.force_flush();
399        assert_eq!(sink.flush_count(), 2);
400    }
401
402    #[test]
403    fn noop_sink_accepts_events_and_default_flush() {
404        let sink = NoopSink;
405        sink.emit(run_started("r1"));
406        // The trait's default force_flush is a no-op; exercise it through the
407        // trait object the runtime actually holds.
408        let boxed: Box<dyn TelemetrySink> = Box::new(NoopSink);
409        boxed.force_flush();
410    }
411
412    #[test]
413    fn run_id_reaches_every_variant() {
414        let events = [
415            run_started("r1"),
416            TelemetryEvent::StageEntered {
417                run_id: "r1".to_string(),
418                stage_index: 0,
419                stage_name: "plan".to_string(),
420                at_ms: 0,
421            },
422            TelemetryEvent::StageExited {
423                run_id: "r1".to_string(),
424                stage_index: 0,
425                stage_name: "plan".to_string(),
426                prompt_tokens: 10,
427                completion_tokens: 5,
428                at_ms: 0,
429            },
430            TelemetryEvent::InferenceCompleted {
431                run_id: "r1".to_string(),
432                stage_name: "plan".to_string(),
433                provider: "anthropic".to_string(),
434                model: "claude-sonnet-5".to_string(),
435                latency_ms: 120,
436                prompt_tokens: 10,
437                completion_tokens: 5,
438                cached_tokens: 0,
439                success: true,
440                cost_usd: Some(0.01),
441            },
442            TelemetryEvent::ToolCallCompleted {
443                run_id: "r1".to_string(),
444                stage_name: "build".to_string(),
445                tool_name: "read_file".to_string(),
446                batch_latency_ms: 8,
447                success: true,
448            },
449            TelemetryEvent::CompactionCompleted {
450                run_id: "r1".to_string(),
451                stage_name: "build".to_string(),
452                success: true,
453            },
454            TelemetryEvent::RunCompleted {
455                run_id: "r1".to_string(),
456                status: "complete".to_string(),
457                prompt_tokens: 10,
458                completion_tokens: 5,
459                tool_calls: 1,
460                empty_output: false,
461                at_ms: 0,
462            },
463            TelemetryEvent::Log {
464                run_id: "r1".to_string(),
465                stage_index: 0,
466                kind: LogKind::Runtime,
467                line: "[Tokens: 10 in, 5 out]".to_string(),
468            },
469        ];
470        let kinds: Vec<&str> = events.iter().map(TelemetryEvent::kind).collect();
471        assert_eq!(
472            kinds,
473            [
474                "run_started",
475                "stage_entered",
476                "stage_exited",
477                "inference_completed",
478                "tool_call_completed",
479                "compaction_completed",
480                "run_completed",
481                "log",
482            ]
483        );
484        for event in &events {
485            assert_eq!(event.run_id(), "r1");
486        }
487    }
488
489    #[test]
490    fn event_clone_debug_and_eq() {
491        let event = run_started("r1");
492        let cloned = event.clone();
493        assert_eq!(event, cloned);
494        assert!(format!("{event:?}").contains("RunStarted"));
495        assert_ne!(LogKind::Output, LogKind::Runtime);
496        assert!(format!("{:?}", LogKind::Output).contains("Output"));
497    }
498}