use std::sync::Arc;
use bevy_ecs::prelude::*;
use leviath_core::telemetry::{LogKind, TelemetryEvent, TelemetrySink};
use crate::components::{AgentState, AgentStatus};
use crate::persistence::{RunMetadata, TokenTotals};
use crate::pipeline::{StageCursor, StageIoBuffer, StageJustEntered, StageLedger};
#[derive(Resource, Clone)]
pub struct Telemetry(pub Arc<dyn TelemetrySink>);
#[derive(Debug, Clone, PartialEq)]
pub enum ActivityRecord {
Inference {
provider: String,
model: String,
latency_ms: u64,
prompt_tokens: usize,
completion_tokens: usize,
cached_tokens: usize,
success: bool,
},
ToolCall {
tool_name: String,
batch_latency_ms: u64,
success: bool,
},
Compaction { success: bool },
}
#[derive(Component, Debug, Default)]
pub struct StageActivity(pub Vec<ActivityRecord>);
#[derive(Component, Debug, Clone, Default)]
pub struct TelemetryState {
run_open: bool,
last_stage: Option<(usize, String)>,
}
fn terminal_label(status: &AgentStatus) -> Option<&'static str> {
match status {
AgentStatus::Complete => Some("complete"),
AgentStatus::Error { .. } => Some("error"),
AgentStatus::Cancelled => Some("cancelled"),
AgentStatus::Idle | AgentStatus::Active | AgentStatus::Waiting | AgentStatus::Paused => {
None
}
}
}
fn stage_tokens(ledger: Option<&StageLedger>, index: usize) -> (usize, usize) {
ledger
.and_then(|l| l.0.get(index))
.map_or((0, 0), |rec| (rec.prompt_tokens, rec.completion_tokens))
}
#[allow(clippy::type_complexity)]
pub fn observe_lifecycle(
telemetry: Res<Telemetry>,
mut agents: Query<(
Entity,
&RunMetadata,
&AgentState,
Option<&StageCursor>,
Option<&TokenTotals>,
Option<&StageLedger>,
Option<&StageJustEntered>,
Option<&mut TelemetryState>,
Option<&mut StageActivity>,
Option<&StageIoBuffer>,
Option<&crate::persistence::RunOutcomeFlags>,
)>,
mut commands: Commands,
) {
crate::tick_scope::clear();
for (entity, md, state, cursor, totals, ledger, entered, ts, activity, buffer, flags) in
agents.iter_mut()
{
crate::tick_scope::enter(entity);
let now_ms = chrono::Utc::now().timestamp_millis();
let sink = telemetry.0.as_ref();
let mut ts = ts;
let (mut st, is_new) = match ts.as_deref() {
Some(existing) => (existing.clone(), false),
None => (TelemetryState::default(), true),
};
if is_new {
let recovered = state.iteration > 0 || cursor.is_some_and(|c| c.index > 0);
sink.emit(TelemetryEvent::RunStarted {
run_id: md.run_id.clone(),
agent_name: md.agent_name.clone(),
model: md.model.clone(),
parent_run_id: md.parent_run_id.clone(),
recovered,
at_ms: now_ms,
});
st.run_open = true;
}
if st.run_open {
let boundary = match entered {
Some(marker) => Some((marker.index, marker.name.clone())),
None if st.last_stage.is_none() => {
Some((cursor.map_or(0, |c| c.index), state.current_stage.clone()))
}
None => None,
};
if let Some((index, name)) = boundary {
let same_stage = st.last_stage.as_ref().is_some_and(|(i, _)| *i == index);
if !same_stage {
if let Some((prev_index, prev_name)) = st.last_stage.take() {
let (prompt, completion) = stage_tokens(ledger, prev_index);
sink.emit(TelemetryEvent::StageExited {
run_id: md.run_id.clone(),
stage_index: prev_index,
stage_name: prev_name,
prompt_tokens: prompt,
completion_tokens: completion,
at_ms: now_ms,
});
}
sink.emit(TelemetryEvent::StageEntered {
run_id: md.run_id.clone(),
stage_index: index,
stage_name: name.clone(),
at_ms: now_ms,
});
st.last_stage = Some((index, name));
}
}
if let Some(mut activity) = activity {
let (_, ref stage_name) = *st.last_stage.as_ref().expect("stage set at sighting");
let stage_name = stage_name.clone();
for record in activity.0.drain(..) {
sink.emit(match record {
ActivityRecord::Inference {
provider,
model,
latency_ms,
prompt_tokens,
completion_tokens,
cached_tokens,
success,
} => TelemetryEvent::InferenceCompleted {
run_id: md.run_id.clone(),
stage_name: stage_name.clone(),
provider,
model,
latency_ms,
prompt_tokens,
completion_tokens,
cached_tokens,
success,
},
ActivityRecord::ToolCall {
tool_name,
batch_latency_ms,
success,
} => TelemetryEvent::ToolCallCompleted {
run_id: md.run_id.clone(),
stage_name: stage_name.clone(),
tool_name,
batch_latency_ms,
success,
},
ActivityRecord::Compaction { success } => {
TelemetryEvent::CompactionCompleted {
run_id: md.run_id.clone(),
stage_name: stage_name.clone(),
success,
}
}
});
}
}
if let Some(buffer) = buffer {
for ((idx, line), kind) in buffer
.output
.iter()
.map(|l| (l, LogKind::Output))
.chain(buffer.logs.iter().map(|l| (l, LogKind::Runtime)))
{
sink.emit(TelemetryEvent::Log {
run_id: md.run_id.clone(),
stage_index: *idx,
kind,
line: line.clone(),
});
}
}
if let Some(status) = terminal_label(&state.status) {
let (prev_index, prev_name) = st.last_stage.take().expect("stage set at sighting");
let (prompt, completion) = stage_tokens(ledger, prev_index);
sink.emit(TelemetryEvent::StageExited {
run_id: md.run_id.clone(),
stage_index: prev_index,
stage_name: prev_name,
prompt_tokens: prompt,
completion_tokens: completion,
at_ms: now_ms,
});
let totals = totals.copied().unwrap_or_default();
sink.emit(TelemetryEvent::RunCompleted {
run_id: md.run_id.clone(),
status: status.to_string(),
prompt_tokens: totals.prompt_tokens,
completion_tokens: totals.completion_tokens,
tool_calls: totals.tool_calls,
empty_output: flags
.is_some_and(|f| crate::persistence::is_empty_output(&state.status, &f.0)),
at_ms: now_ms,
});
st.run_open = false;
}
}
if is_new {
commands
.entity(entity)
.insert((st, StageActivity::default()));
} else {
*ts.as_deref_mut().expect("state exists when not new") = st;
}
}
}
#[cfg(test)]
mod tests;