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}