use std::collections::{BTreeMap, BTreeSet};
use oneagentgraph::event::Role;
use serde::{Deserialize, Serialize};
use serde_json::Value;
use crate::event::{Envelope, Source};
use crate::graph::NodeStatus;
use crate::journal;
use crate::ledger::RunPaths;
use crate::projection;
pub const TELEMETRY_SCHEMA_VERSION: u32 = 2;
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub struct RunTelemetry {
#[serde(deserialize_with = "this_version")]
pub schema_version: u32,
pub run_id: String,
pub wall_ms: u64,
#[serde(deserialize_with = "every_bucket")]
pub buckets: Vec<Bucket>,
pub usage: BTreeMap<Party, Usage>,
pub dispatches: u64,
pub settled_done: u64,
pub no_diff: u64,
pub surfaces_queued: u64,
pub surfaces_read: u64,
}
fn this_version<'de, D: serde::Deserializer<'de>>(reader: D) -> Result<u32, D::Error> {
let found = u32::deserialize(reader)?;
if found != TELEMETRY_SCHEMA_VERSION {
return Err(serde::de::Error::custom(format!(
"telemetry schema_version {found}, and this build reads \
{TELEMETRY_SCHEMA_VERSION}"
)));
}
Ok(found)
}
fn every_bucket<'de, D: serde::Deserializer<'de>>(reader: D) -> Result<Vec<Bucket>, D::Error> {
let found = Vec::<Bucket>::deserialize(reader)?;
let named: Vec<BucketName> = found.iter().map(|bucket| bucket.name).collect();
if named != BucketName::ALL {
return Err(serde::de::Error::custom(format!(
"buckets {:?}, and a telemetry document carries exactly {:?}",
named.iter().map(|name| name.as_str()).collect::<Vec<_>>(),
BucketName::ALL
.iter()
.map(|name| name.as_str())
.collect::<Vec<_>>()
)));
}
Ok(found)
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub struct Bucket {
pub name: BucketName,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub ms: Option<u64>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum BucketName {
Agent,
Judge,
Llmlint,
Gate,
PublicationWait,
LockWait,
Setup,
Scheduling,
}
impl BucketName {
pub const ALL: [Self; 8] = [
Self::Agent,
Self::Judge,
Self::Llmlint,
Self::Gate,
Self::PublicationWait,
Self::LockWait,
Self::Setup,
Self::Scheduling,
];
pub fn as_str(self) -> &'static str {
match self {
Self::Agent => "agent",
Self::Judge => "judge",
Self::Llmlint => "llmlint",
Self::Gate => "gate",
Self::PublicationWait => "publication_wait",
Self::LockWait => "lock_wait",
Self::Setup => "setup",
Self::Scheduling => "scheduling",
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum Party {
Agent,
Judge,
Llmlint,
Total,
}
impl Party {
pub const ALL: [Self; 4] = [Self::Agent, Self::Judge, Self::Llmlint, Self::Total];
pub fn as_str(self) -> &'static str {
match self {
Self::Agent => "agent",
Self::Judge => "judge",
Self::Llmlint => "llmlint",
Self::Total => "total",
}
}
}
#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub struct Usage {
#[serde(default, skip_serializing_if = "Option::is_none")]
pub input: Option<u64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub output: Option<u64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub cache_read: Option<u64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub cache_write: Option<u64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub cost_usd: Option<f64>,
}
impl Usage {
pub fn is_empty(&self) -> bool {
self.input.is_none()
&& self.output.is_none()
&& self.cache_read.is_none()
&& self.cache_write.is_none()
&& self.cost_usd.is_none()
}
pub fn add(&mut self, other: &Self) {
let sum = |into: &mut Option<u64>, value: Option<u64>| {
if let Some(value) = value {
*into = Some(into.unwrap_or(0).saturating_add(value));
}
};
sum(&mut self.input, other.input);
sum(&mut self.output, other.output);
sum(&mut self.cache_read, other.cache_read);
sum(&mut self.cache_write, other.cache_write);
if let Some(cost) = other.cost_usd {
self.cost_usd = Some(self.cost_usd.unwrap_or(0.0) + cost);
}
}
fn of(value: &Value) -> Self {
let count = |names: [&str; 2]| {
names
.iter()
.find_map(|name| value.get(*name).and_then(Value::as_u64))
};
Self {
input: count(["input_tokens", "tokens_in"]),
output: count(["output_tokens", "tokens_out"]),
cache_read: count(["cache_read_tokens", "cache_read"]),
cache_write: count(["cache_write_tokens", "cache_write"]),
cost_usd: ["cost_usd", "cost"]
.iter()
.find_map(|name| value.get(*name).and_then(Value::as_f64)),
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum Phase {
Setup,
LockWait,
Gate,
Publication,
}
impl Phase {
fn of(kind: &str) -> Option<Self> {
match kind {
"session-opened" | "fetch" | "commit-preserved" | "lock-acquired"
| "recovery-attested" => Some(Self::Setup),
"lock-wait" => Some(Self::LockWait),
"gate-started" => Some(Self::Gate),
"gate-verdict" | "push" | "change-opened" | "change-check" | "merge-queued"
| "change-merged" | "merge-completed" | "sync-conflict" => Some(Self::Publication),
_ => None,
}
}
fn bucket(self) -> BucketName {
match self {
Self::Setup => BucketName::Setup,
Self::LockWait => BucketName::LockWait,
Self::Gate => BucketName::Gate,
Self::Publication => BucketName::PublicationWait,
}
}
const PRECEDENCE: [Self; 4] = [Self::LockWait, Self::Gate, Self::Publication, Self::Setup];
}
const SESSION_OPENED: &str = "session-opened";
const SESSION_CLOSED: &str = "session-closed";
const WORKING: [&str; 5] = [
"member-started",
"turn-started",
"turn-activity",
"turn-completed",
crate::report::MEMBER_SETTLED,
];
const ROLE: &str = "role";
fn side_of(event: &Envelope) -> Option<Role> {
serde_json::from_value(event.payload.get(ROLE)?.clone()).ok()
}
pub fn of_run(paths: &RunPaths, events: &[Envelope]) -> RunTelemetry {
let state = projection::fold(events);
let stamps: Vec<(u64, &Envelope)> = events
.iter()
.filter_map(|event| projection::millis_of(&event.ts).map(|ms| (ms, event)))
.collect();
let first = stamps.first().map_or(0, |(ms, _)| *ms);
let last = stamps.last().map_or(first, |(ms, _)| *ms);
let wall_ms = last.saturating_sub(first);
let mut totals: BTreeMap<BucketName, u64> = BTreeMap::new();
let mut dispatched: BTreeSet<String> = BTreeSet::new();
let mut phases: BTreeMap<String, Phase> = BTreeMap::new();
let mut judging: BTreeSet<String> = BTreeSet::new();
let mut judge_measured = false;
let mut previous = first;
for (ms, event) in &stamps {
let span = ms.saturating_sub(previous);
if span > 0 {
*totals
.entry(now(&phases, &judging, &dispatched))
.or_insert(0) += span;
}
previous = *ms;
let whose = event
.labels
.node
.clone()
.unwrap_or_else(|| event.stream.clone());
match event.source {
Source::Vcs => match event.kind.0.as_str() {
SESSION_CLOSED => {
phases.remove(&whose);
}
SESSION_OPENED => {
phases.entry(whose).or_insert(Phase::Setup);
}
kind => {
if let Some(phase) = Phase::of(kind) {
phases.insert(whose, phase);
}
}
},
Source::Agentgraph if WORKING.contains(&event.kind.0.as_str()) => {
phases.remove(&whose);
match side_of(event) {
Some(Role::Judge) => {
judge_measured = true;
judging.insert(whose);
}
Some(Role::Agent) => {
judge_measured = true;
judging.remove(&whose);
}
None if event.kind.0 == crate::report::MEMBER_SETTLED => {
judging.remove(&whose);
}
None => {}
}
}
Source::Pipeline => match journal::PipelineKind::from_wire(&event.kind) {
Some(journal::PipelineKind::NodeDispatched) => {
dispatched.insert(whose);
}
Some(journal::PipelineKind::NodeSettled) => {
dispatched.remove(&whose);
judging.remove(&whose);
phases.remove(&whose);
}
_ => {}
},
Source::Agentgraph => {}
}
}
let measured = |name: BucketName| match name {
BucketName::Llmlint => None,
BucketName::Judge if !judge_measured => None,
_ => Some(totals.get(&name).copied().unwrap_or(0)),
};
let mut buckets: Vec<Bucket> = BucketName::ALL
.into_iter()
.map(|name| Bucket {
name,
ms: measured(name),
})
.collect();
balance(&mut buckets, wall_ms);
RunTelemetry {
schema_version: TELEMETRY_SCHEMA_VERSION,
run_id: paths.run.clone(),
wall_ms,
buckets,
usage: usage_of(paths, events),
dispatches: state.dispatched_at.len() as u64,
settled_done: state
.recorded
.values()
.filter(|recorded| recorded.status() == NodeStatus::Done)
.count() as u64,
no_diff: state
.outcomes
.values()
.filter(|outcome| *outcome == "no-changes")
.count() as u64,
surfaces_queued: state.surfaces_queued,
surfaces_read: state.surfaces_read,
}
}
fn now(
phases: &BTreeMap<String, Phase>,
judging: &BTreeSet<String>,
dispatched: &BTreeSet<String>,
) -> BucketName {
for phase in Phase::PRECEDENCE {
if phases.values().any(|open| *open == phase) {
return phase.bucket();
}
}
if !judging.is_empty() {
return BucketName::Judge;
}
if !dispatched.is_empty() {
return BucketName::Agent;
}
BucketName::Scheduling
}
fn balance(buckets: &mut [Bucket], wall_ms: u64) {
let counted: u64 = buckets.iter().filter_map(|bucket| bucket.ms).sum();
let Some(residue) = buckets
.iter_mut()
.find(|bucket| bucket.name == BucketName::Scheduling)
else {
return;
};
let ms = residue.ms.unwrap_or(0);
residue.ms = Some(match wall_ms.checked_sub(counted) {
Some(under) => ms.saturating_add(under),
None => ms.saturating_sub(counted - wall_ms),
});
}
fn usage_of(paths: &RunPaths, events: &[Envelope]) -> BTreeMap<Party, Usage> {
let mut totals: BTreeMap<Party, Usage> = BTreeMap::new();
let mut fold = |party: Party, usage: &Usage| {
if !usage.is_empty() {
totals.entry(party).or_default().add(usage);
}
};
for event in events
.iter()
.filter(|event| event.source == Source::Agentgraph && event.kind.0 == "turn-completed")
{
if let Some(usage) = event.payload.get("usage").filter(|value| value.is_object()) {
fold(Party::Total, &Usage::of(usage));
}
}
for retained in crate::report::evidence(paths, events) {
let Some(document) = crate::report::read(&retained.kept) else {
continue;
};
let Some(telemetry) = document.get("telemetry") else {
continue;
};
for (party, side) in [(Party::Agent, "agent"), (Party::Judge, "judge")] {
if let Some(usage) = telemetry
.get(side)
.and_then(|side| side.get("usage"))
.filter(|value| value.is_object())
{
fold(party, &Usage::of(usage));
}
}
}
totals
}
pub const UNMEASURED: &str = "not measured";
pub fn render_breakdown(telemetry: &RunTelemetry) -> String {
let mut out = format!(
"{} WALL {}\n",
telemetry.run_id,
duration(telemetry.wall_ms)
);
for bucket in &telemetry.buckets {
match bucket.ms {
None => out.push_str(&format!(
" {:<18} {:>10}\n",
bucket.name.as_str(),
UNMEASURED
)),
Some(ms) => {
let share = (ms * 100).checked_div(telemetry.wall_ms).unwrap_or(0);
out.push_str(&format!(
" {:<18} {:>10} {share:>3}%\n",
bucket.name.as_str(),
duration(ms)
));
}
}
}
for party in Party::ALL {
out.push_str(&format!(" usage {:<12} ", party.as_str()));
match telemetry.usage.get(&party) {
None => out.push_str(&format!("{UNMEASURED}\n")),
Some(usage) => out.push_str(&format!(
"in {} out {} cache r {} w {} ${}\n",
tokens(usage.input),
tokens(usage.output),
tokens(usage.cache_read),
tokens(usage.cache_write),
usage
.cost_usd
.map_or_else(|| UNMEASURED.to_string(), |cost| format!("{cost:.4}"))
)),
}
}
out.push_str(&format!(
" {} dispatch(es), {} done, {} no-diff; {} surface(s) sent, {} read\n",
telemetry.dispatches,
telemetry.settled_done,
telemetry.no_diff,
telemetry.surfaces_queued,
telemetry.surfaces_read
));
out
}
fn tokens(count: Option<u64>) -> String {
count.map_or_else(|| UNMEASURED.to_string(), |count| count.to_string())
}
pub fn duration(ms: u64) -> String {
let seconds = ms / 1_000;
if seconds < 60 {
return format!("{seconds}s");
}
if seconds < 3_600 {
return format!("{}m{:02}s", seconds / 60, seconds % 60);
}
format!("{}h{:02}m", seconds / 3_600, (seconds % 3_600) / 60)
}
#[cfg(test)]
mod tests {
use super::*;
use crate::event::{EventKind, Labels, ENVELOPE_VERSION};
use crate::plan::{Node, Plan, PLAN_SCHEMA_VERSION};
use serde_json::json;
fn stamped(
seconds: u64,
source: Source,
kind: EventKind,
node: Option<&str>,
fields: &[(&str, serde_json::Value)],
) -> Envelope {
Envelope {
v: ENVELOPE_VERSION,
ts: crate::sys::rfc3339_from_millis(1_786_000_000_000 + seconds * 1_000),
stream: "s".into(),
seq: seconds,
source,
kind,
labels: Labels {
run_id: Some("demo".into()),
node: node.map(str::to_string),
..Labels::default()
},
payload: journal::payload(fields),
artifacts: Vec::new(),
}
}
fn at(
seconds: u64,
kind: journal::PipelineKind,
node: Option<&str>,
fields: &[(&str, serde_json::Value)],
) -> Envelope {
stamped(seconds, Source::Pipeline, kind.into(), node, fields)
}
fn session(seconds: u64, kind: &str, node: Option<&str>) -> Envelope {
stamped(seconds, Source::Vcs, EventKind(kind.into()), node, &[])
}
fn turn(
seconds: u64,
kind: &str,
node: Option<&str>,
fields: &[(&str, serde_json::Value)],
) -> Envelope {
stamped(
seconds,
Source::Agentgraph,
EventKind(kind.into()),
node,
fields,
)
}
fn paths() -> RunPaths {
let root = std::env::temp_dir().join(format!(
"onepipeline-telemetry-{}-{:?}",
crate::sys::pid(),
std::thread::current().id()
));
let paths = RunPaths::under(&root, "demo");
paths.create().expect("the run directory");
paths
}
fn plan() -> Plan {
Plan {
schema_version: PLAN_SCHEMA_VERSION,
goal: None,
name: Some("demo".into()),
concurrency: 4,
tasks: vec![Node {
id: "build".into(),
persona: Some("engineer".into()),
task: Some("## What\ndo it".into()),
..Node::default()
}],
}
}
fn started() -> Envelope {
at(
0,
journal::PipelineKind::RunStarted,
None,
&[("plan", json!(plan()))],
)
}
fn summed(telemetry: &RunTelemetry) -> u64 {
telemetry.buckets.iter().filter_map(|b| b.ms).sum()
}
fn bucket_of(telemetry: &RunTelemetry, name: BucketName) -> Option<u64> {
telemetry
.buckets
.iter()
.find(|b| b.name == name)
.unwrap_or_else(|| panic!("a {} bucket", name.as_str()))
.ms
}
#[test]
fn the_buckets_sum_exactly_to_the_wall_clock() {
let events = vec![
started(),
at(
10,
journal::PipelineKind::NodeDispatched,
Some("build"),
&[],
),
at(
70,
journal::PipelineKind::NodeSettled,
Some("build"),
&[("status", json!("done"))],
),
at(100, journal::PipelineKind::RunStopped, None, &[]),
];
let telemetry = of_run(&paths(), &events);
assert_eq!(telemetry.wall_ms, 100_000);
assert_eq!(
summed(&telemetry),
telemetry.wall_ms,
"{:?}",
telemetry.buckets
);
assert_eq!(bucket_of(&telemetry, BucketName::Agent), Some(60_000));
assert_eq!(bucket_of(&telemetry, BucketName::Scheduling), Some(40_000));
assert_eq!(telemetry.dispatches, 1);
assert_eq!(telemetry.settled_done, 1);
}
#[test]
fn gate_time_and_lock_waiting_are_separable_from_agent_time() {
let events = vec![
started(),
at(
10,
journal::PipelineKind::NodeDispatched,
Some("service"),
&[],
),
session(10, "session-opened", Some("service")),
turn(20, "turn-started", Some("service"), &[]),
session(50, "lock-wait", Some("service")),
session(60, "lock-acquired", Some("service")),
session(70, "gate-started", Some("service")),
session(100, "gate-verdict", Some("service")),
session(110, "session-closed", Some("service")),
at(
120,
journal::PipelineKind::NodeSettled,
Some("service"),
&[("status", json!("done"))],
),
];
let telemetry = of_run(&paths(), &events);
assert_eq!(
summed(&telemetry),
telemetry.wall_ms,
"{:?}",
telemetry.buckets
);
assert_eq!(bucket_of(&telemetry, BucketName::Setup), Some(20_000));
assert_eq!(bucket_of(&telemetry, BucketName::Agent), Some(40_000));
assert_eq!(bucket_of(&telemetry, BucketName::LockWait), Some(10_000));
assert_eq!(bucket_of(&telemetry, BucketName::Gate), Some(30_000));
assert_eq!(
bucket_of(&telemetry, BucketName::PublicationWait),
Some(10_000)
);
assert_eq!(bucket_of(&telemetry, BucketName::Scheduling), Some(10_000));
}
#[test]
fn a_second_sessions_opening_does_not_rewind_a_publication_already_under_way() {
let mut opened = session(50, "session-opened", Some("service"));
opened.stream = "z-second-session".into();
let telemetry = of_run(
&paths(),
&[
started(),
at(
10,
journal::PipelineKind::NodeDispatched,
Some("service"),
&[],
),
session(20, "session-opened", Some("service")),
turn(30, "turn-activity", Some("service"), &[]),
session(50, "lock-wait", Some("service")),
opened,
session(90, "lock-acquired", Some("service")),
at(
100,
journal::PipelineKind::NodeSettled,
Some("service"),
&[("status", json!("done"))],
),
],
);
assert_eq!(bucket_of(&telemetry, BucketName::LockWait), Some(40_000));
assert_eq!(summed(&telemetry), telemetry.wall_ms);
}
#[test]
fn a_bucket_nothing_measures_is_absent_rather_than_zero() {
let telemetry = of_run(
&paths(),
&[
started(),
at(
10,
journal::PipelineKind::NodeDispatched,
Some("build"),
&[],
),
turn(20, "turn-activity", Some("build"), &[]),
],
);
assert_eq!(bucket_of(&telemetry, BucketName::Llmlint), None);
assert_eq!(bucket_of(&telemetry, BucketName::Judge), None);
assert_eq!(bucket_of(&telemetry, BucketName::Agent), Some(10_000));
let document = serde_json::to_value(&telemetry).expect("it serialises");
let llmlint = document["buckets"]
.as_array()
.expect("buckets")
.iter()
.find(|bucket| bucket["name"] == "llmlint")
.expect("the llmlint bucket is still named");
assert!(
llmlint.get("ms").is_none(),
"an unmeasured bucket carried a number: {llmlint}"
);
assert!(render_breakdown(&telemetry).contains("llmlint"));
assert!(render_breakdown(&telemetry).contains(UNMEASURED));
}
#[test]
fn a_turn_a_producer_attributes_to_the_judge_is_measured_as_the_judges() {
let telemetry = of_run(
&paths(),
&[
started(),
at(
10,
journal::PipelineKind::NodeDispatched,
Some("build"),
&[],
),
turn(
20,
"turn-started",
Some("build"),
&[("role", json!("judge"))],
),
turn(
50,
"turn-started",
Some("build"),
&[("role", json!("agent"))],
),
at(
60,
journal::PipelineKind::NodeSettled,
Some("build"),
&[("status", json!("done"))],
),
],
);
assert_eq!(bucket_of(&telemetry, BucketName::Judge), Some(30_000));
assert_eq!(bucket_of(&telemetry, BucketName::Agent), Some(20_000));
assert_eq!(summed(&telemetry), telemetry.wall_ms);
}
#[test]
fn an_empty_run_has_a_zero_wall_clock_and_still_balances() {
let telemetry = of_run(&paths(), &[]);
assert_eq!(telemetry.wall_ms, 0);
assert_eq!(summed(&telemetry), 0);
assert!(render_breakdown(&telemetry).contains("WALL 0s"));
}
#[test]
fn a_bucket_and_a_party_serialise_as_the_words_the_breakdown_renders() {
for name in BucketName::ALL {
let json = serde_json::to_string(&name).expect("a bucket name serialises");
assert_eq!(json, format!("\"{}\"", name.as_str()));
assert_eq!(
serde_json::from_str::<BucketName>(&json).expect("it reads back"),
name
);
}
for party in Party::ALL {
let json = serde_json::to_string(&party).expect("a party serialises");
assert_eq!(json, format!("\"{}\"", party.as_str()));
assert_eq!(
serde_json::from_str::<Party>(&json).expect("it reads back"),
party
);
}
}
#[test]
fn a_clock_that_moved_backwards_still_leaves_the_buckets_summing_to_wall() {
let mut events = vec![
started(),
at(
60,
journal::PipelineKind::NodeDispatched,
Some("build"),
&[],
),
at(30, journal::PipelineKind::RunStopped, None, &[]),
];
events.reverse();
let telemetry = of_run(&paths(), &events);
assert_eq!(
summed(&telemetry),
telemetry.wall_ms,
"{:?}",
telemetry.buckets
);
}
#[test]
fn the_breakdown_names_every_bucket_and_every_party() {
let rendered = render_breakdown(&of_run(
&paths(),
&[
started(),
at(
10,
journal::PipelineKind::NodeDispatched,
Some("build"),
&[],
),
at(
20,
journal::PipelineKind::NodeSettled,
Some("build"),
&[("status", json!("done"))],
),
],
));
for name in BucketName::ALL {
assert!(
rendered.contains(name.as_str()),
"{rendered} omits {}",
name.as_str()
);
}
for party in Party::ALL {
assert!(
rendered.contains(&format!("usage {}", party.as_str())),
"{rendered} omits the {} party",
party.as_str()
);
}
assert!(rendered.contains('%'), "{rendered}");
}
#[test]
fn usage_is_totalled_from_the_turns_and_split_by_the_report_they_settled_with() {
let paths = paths();
let stored = paths.report_for("s", 20);
std::fs::create_dir_all(paths.reports_dir()).expect("the run's report storage");
std::fs::write(
&stored,
json!({
"telemetry": {
"agent": {"usage": {
"input_tokens": 1_000, "output_tokens": 300,
"cache_read_tokens": 900, "cost_usd": 0.40,
}},
"judge": {"usage": {"input_tokens": 200, "output_tokens": 40}},
},
})
.to_string(),
)
.expect("a stored report");
let telemetry = of_run(
&paths,
&[
started(),
turn(
10,
"turn-completed",
Some("build"),
&[(
"usage",
json!({
"input_tokens": 1_200, "output_tokens": 340,
"cache_read_tokens": 900, "cache_write_tokens": 120,
"cost_usd": 0.42,
}),
)],
),
turn(
20,
crate::report::MEMBER_SETTLED,
Some("build"),
&[(crate::report::REPORT_PATH, json!("/elsewhere/report.json"))],
),
],
);
let total = &telemetry.usage[&Party::Total];
assert_eq!(total.input, Some(1_200));
assert_eq!(total.output, Some(340));
assert_eq!(total.cache_read, Some(900));
assert_eq!(total.cache_write, Some(120));
assert_eq!(total.cost_usd, Some(0.42));
let agent = &telemetry.usage[&Party::Agent];
assert_eq!(agent.input, Some(1_000));
assert_eq!(agent.cost_usd, Some(0.40));
let judge = &telemetry.usage[&Party::Judge];
assert_eq!(judge.input, Some(200));
assert_eq!(judge.cache_read, None);
assert_eq!(judge.cost_usd, None);
assert!(!telemetry.usage.contains_key(&Party::Llmlint));
let rendered = render_breakdown(&telemetry);
assert!(rendered.contains("in 1200"), "{rendered}");
assert!(rendered.contains("$0.4200"), "{rendered}");
let llmlint = rendered
.lines()
.find(|line| line.trim_start().starts_with("usage llmlint"))
.expect("the llmlint party is still named");
assert!(llmlint.contains(UNMEASURED), "{llmlint}");
std::fs::remove_dir_all(&paths.dir).ok();
}
#[test]
fn usage_is_read_in_either_spelling_the_producer_uses() {
let declared = Usage::of(&json!({
"tokens_in": 10, "tokens_out": 4,
"cache_read": 3, "cache_write": 2, "cost": 0.5,
}));
assert_eq!(declared.input, Some(10));
assert_eq!(declared.output, Some(4));
assert_eq!(declared.cache_read, Some(3));
assert_eq!(declared.cache_write, Some(2));
assert_eq!(declared.cost_usd, Some(0.5));
assert!(Usage::of(&json!({})).is_empty());
}
#[test]
fn a_report_this_host_cannot_read_leaves_the_split_absent() {
let telemetry = of_run(
&paths(),
&[
started(),
turn(
10,
"turn-completed",
Some("build"),
&[("usage", json!({"input_tokens": 5}))],
),
turn(
20,
crate::report::MEMBER_SETTLED,
Some("build"),
&[(
crate::report::REPORT_PATH,
json!("/nowhere/onepipeline/report.json"),
)],
),
],
);
assert_eq!(telemetry.usage[&Party::Total].input, Some(5));
assert!(!telemetry.usage.contains_key(&Party::Agent));
assert!(!telemetry.usage.contains_key(&Party::Judge));
}
#[test]
fn a_usage_total_accumulates_only_the_fields_something_reported() {
let mut total = Usage::default();
assert!(total.is_empty());
total.add(&Usage {
input: Some(10),
cost_usd: Some(0.5),
..Usage::default()
});
total.add(&Usage {
input: Some(5),
output: Some(3),
..Usage::default()
});
assert_eq!(total.input, Some(15));
assert_eq!(total.output, Some(3));
assert_eq!(total.cache_read, None);
assert_eq!(total.cost_usd, Some(0.5));
assert!(!total.is_empty());
}
const GOLDEN: &str = include_str!("../tests/golden/telemetry-v2.json");
fn golden() -> RunTelemetry {
let usage = |input, output, cache: Option<(u64, u64)>, cost| Usage {
input: Some(input),
output: Some(output),
cache_read: cache.map(|(read, _)| read),
cache_write: cache.map(|(_, write)| write),
cost_usd: cost,
};
RunTelemetry {
schema_version: TELEMETRY_SCHEMA_VERSION,
run_id: "golden".into(),
wall_ms: 120_000,
buckets: vec![
Bucket {
name: BucketName::Agent,
ms: Some(40_000),
},
Bucket {
name: BucketName::Judge,
ms: None,
},
Bucket {
name: BucketName::Llmlint,
ms: None,
},
Bucket {
name: BucketName::Gate,
ms: Some(30_000),
},
Bucket {
name: BucketName::PublicationWait,
ms: Some(10_000),
},
Bucket {
name: BucketName::LockWait,
ms: Some(0),
},
Bucket {
name: BucketName::Setup,
ms: Some(20_000),
},
Bucket {
name: BucketName::Scheduling,
ms: Some(20_000),
},
],
usage: BTreeMap::from([
(
Party::Agent,
usage(1_000, 300, Some((900, 120)), Some(0.40)),
),
(Party::Judge, usage(200, 40, None, None)),
(
Party::Total,
usage(1_200, 340, Some((900, 120)), Some(0.42)),
),
]),
dispatches: 1,
settled_done: 1,
no_diff: 0,
surfaces_queued: 2,
surfaces_read: 1,
}
}
#[test]
fn a_schema_2_document_is_the_shape_the_golden_pins() {
let rendered = serde_json::to_string_pretty(&golden()).expect("it serialises");
assert_eq!(
rendered.trim(),
GOLDEN.trim(),
"the telemetry document changed shape. If that was deliberate, bump \
TELEMETRY_SCHEMA_VERSION and update tests/golden/telemetry-v2.json together"
);
assert_eq!(summed(&golden()), golden().wall_ms);
}
#[test]
fn a_schema_2_document_round_trips_through_the_types() {
let value = golden();
let read: RunTelemetry =
serde_json::from_str(GOLDEN).expect("the golden reads back into the types");
assert_eq!(read, value);
let again: RunTelemetry =
serde_json::from_str(&serde_json::to_string(&value).expect("it serialises"))
.expect("it reads back");
assert_eq!(again, value);
}
#[test]
fn an_unmeasured_bucket_omits_its_span_and_a_measured_zero_keeps_it() {
let document: Value =
serde_json::from_str(&serde_json::to_string(&golden()).expect("it serialises"))
.expect("it is JSON");
let bucket = |name: &str| {
document["buckets"]
.as_array()
.expect("buckets")
.iter()
.find(|bucket| bucket["name"] == name)
.unwrap_or_else(|| panic!("a {name} bucket"))
.clone()
};
for absent in ["judge", "llmlint"] {
assert!(
bucket(absent).get("ms").is_none(),
"the {absent} bucket carried a span nothing measured"
);
}
assert_eq!(bucket("lock_wait")["ms"], 0);
let read: RunTelemetry = serde_json::from_value(document).expect("it reads back");
assert_eq!(bucket_of(&read, BucketName::Judge), None);
assert_eq!(bucket_of(&read, BucketName::LockWait), Some(0));
}
#[test]
fn an_unreported_usage_field_is_omitted_and_round_trips_as_absent() {
let document: Value =
serde_json::from_str(&serde_json::to_string(&golden()).expect("it serialises"))
.expect("it is JSON");
let judge = &document["usage"]["judge"];
assert_eq!(judge["input"], 200);
for unreported in ["cache_read", "cache_write", "cost_usd"] {
assert!(
judge.get(unreported).is_none(),
"the judge claimed a {unreported} nothing reported: {judge}"
);
}
assert!(
document["usage"].get("llmlint").is_none(),
"a party nothing reported for is on the wire: {}",
document["usage"]
);
let empty = Usage::default();
let rendered = serde_json::to_string(&empty).expect("it serialises");
assert_eq!(rendered, "{}");
let read: Usage = serde_json::from_str(&rendered).expect("it reads back");
assert_eq!(read, empty);
assert!(read.is_empty());
}
#[test]
fn a_schema_1_document_is_refused_rather_than_read_as_a_newer_one() {
let v1 = json!({
"schema_version": 1,
"run_id": "old",
"wall_ms": 100_000,
"buckets": [
{"name": "dispatching", "ms": 60_000},
{"name": "awaiting-planner", "ms": 0},
{"name": "awaiting-human", "ms": 0},
{"name": "orchestrating", "ms": 40_000},
],
"dispatches": 1,
"settled_done": 1,
"no_diff": 0,
"surfaces_queued": 0,
"surfaces_read": 0,
});
let refusal = serde_json::from_value::<RunTelemetry>(v1.clone())
.expect_err("a schema-1 document was read as a schema-2 one");
assert!(
refusal.to_string().contains("schema_version 1")
&& refusal
.to_string()
.contains(&TELEMETRY_SCHEMA_VERSION.to_string()),
"the refusal names neither the version it met nor the one it reads: {refusal}"
);
let mut relabelled = v1;
relabelled["schema_version"] = json!(TELEMETRY_SCHEMA_VERSION);
let refusal = serde_json::from_value::<RunTelemetry>(relabelled.clone())
.expect_err("a schema-1 body under a schema-2 stamp was read");
assert!(
refusal.to_string().contains("dispatching"),
"the refusal does not name what it could not read: {refusal}"
);
let mut renamed = relabelled;
renamed["buckets"] = json!([{"name": "agent", "ms": 100_000}]);
let refusal = serde_json::from_value::<RunTelemetry>(renamed)
.expect_err("a document carrying one bucket was read");
assert!(
refusal.to_string().contains("exactly"),
"the refusal does not say the set is fixed: {refusal}"
);
}
#[test]
fn a_bucket_set_that_is_not_the_eight_is_refused() {
let document = |buckets: Value| {
let mut value = serde_json::to_value(golden()).expect("it serialises");
value["buckets"] = buckets;
serde_json::from_value::<RunTelemetry>(value)
};
let whole = serde_json::to_value(golden()).expect("it serialises")["buckets"].clone();
assert!(document(whole.clone()).is_ok(), "the eight were refused");
let short: Vec<Value> = whole.as_array().expect("buckets")[..7].to_vec();
assert!(document(json!(short)).is_err(), "a set of seven was read");
let mut twice = whole.as_array().expect("buckets").clone();
twice[1] = twice[0].clone();
assert!(document(json!(twice)).is_err(), "a doubled bucket was read");
let mut shuffled = whole.as_array().expect("buckets").clone();
shuffled.reverse();
assert!(
document(json!(shuffled)).is_err(),
"a shuffled set was read"
);
}
#[test]
fn the_schema_version_and_the_golden_name_the_same_number() {
assert_eq!(TELEMETRY_SCHEMA_VERSION, 2);
assert_eq!(golden().schema_version, TELEMETRY_SCHEMA_VERSION);
let document: Value = serde_json::from_str(GOLDEN).expect("the golden is JSON");
assert_eq!(document["schema_version"], TELEMETRY_SCHEMA_VERSION);
}
#[test]
fn a_duration_reads_in_the_units_its_size_calls_for() {
assert_eq!(duration(0), "0s");
assert_eq!(duration(45_000), "45s");
assert_eq!(duration(125_000), "2m05s");
assert_eq!(duration(7_500_000), "2h05m");
}
}