use std::collections::BTreeMap;
use serde::{Deserialize, Serialize};
use crate::event::{Envelope, Source};
use crate::graph::NodeStatus;
use crate::journal;
use crate::projection;
pub const TELEMETRY_SCHEMA_VERSION: u32 = 1;
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub struct RunTelemetry {
pub schema_version: u32,
pub run_id: String,
pub wall_ms: u64,
pub buckets: Vec<Bucket>,
pub dispatches: u64,
pub settled_done: u64,
pub no_diff: u64,
pub surfaces_queued: u64,
pub surfaces_read: u64,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub struct Bucket {
pub name: BucketName,
pub ms: u64,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash, Serialize, Deserialize)]
#[serde(rename_all = "kebab-case")]
pub enum BucketName {
Dispatching,
AwaitingPlanner,
AwaitingHuman,
Orchestrating,
}
impl BucketName {
pub const ALL: [Self; 4] = [
Self::Dispatching,
Self::AwaitingPlanner,
Self::AwaitingHuman,
Self::Orchestrating,
];
pub fn as_str(self) -> &'static str {
match self {
Self::Dispatching => "dispatching",
Self::AwaitingPlanner => "awaiting-planner",
Self::AwaitingHuman => "awaiting-human",
Self::Orchestrating => "orchestrating",
}
}
}
pub fn of_run(run: &str, 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 in_flight: u64 = 0;
let mut awaiting_planner = false;
let mut awaiting_human: u64 = 0;
let mut previous = first;
for (ms, event) in &stamps {
let span = ms.saturating_sub(previous);
if span > 0 {
let bucket = if in_flight > 0 {
BucketName::Dispatching
} else if awaiting_planner {
BucketName::AwaitingPlanner
} else if awaiting_human > 0 {
BucketName::AwaitingHuman
} else {
BucketName::Orchestrating
};
*totals.entry(bucket).or_insert(0) += span;
}
previous = *ms;
if event.source != Source::Pipeline {
continue;
}
match journal::PipelineKind::from_wire(&event.kind) {
Some(journal::PipelineKind::NodeDispatched) => in_flight += 1,
Some(journal::PipelineKind::NodeSettled) => {
let status = event
.payload
.get("status")
.and_then(|v| v.as_str())
.and_then(NodeStatus::parse);
if state
.dispatched_at
.contains_key(event.labels.node.as_deref().unwrap_or_default())
{
in_flight = in_flight.saturating_sub(1);
}
if status == Some(NodeStatus::Waiting) {
awaiting_human += 1;
}
}
Some(journal::PipelineKind::HumanAttested) => {
awaiting_human = awaiting_human.saturating_sub(1)
}
Some(journal::PipelineKind::PlannerSurfaced) => {
awaiting_planner = event
.payload
.get("blocking")
.and_then(|v| v.as_bool())
.unwrap_or(false);
}
Some(journal::PipelineKind::PlannerReplied) => awaiting_planner = false,
_ => {}
}
}
let mut buckets: Vec<Bucket> = BucketName::ALL
.into_iter()
.map(|name| Bucket {
name,
ms: totals.get(&name).copied().unwrap_or(0),
})
.collect();
let counted: u64 = buckets.iter().map(|bucket| bucket.ms).sum();
if let Some(residue) = wall_ms.checked_sub(counted) {
if let Some(bucket) = buckets
.iter_mut()
.find(|b| b.name == BucketName::Orchestrating)
{
bucket.ms += residue;
}
} else if let Some(bucket) = buckets
.iter_mut()
.find(|b| b.name == BucketName::Orchestrating)
{
bucket.ms = bucket.ms.saturating_sub(counted - wall_ms);
}
RunTelemetry {
schema_version: TELEMETRY_SCHEMA_VERSION,
run_id: run.to_string(),
wall_ms,
buckets,
dispatches: state.dispatched_at.len() as u64,
settled_done: state
.recorded
.values()
.filter(|status| **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,
}
}
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 {
let share = (bucket.ms * 100)
.checked_div(telemetry.wall_ms)
.unwrap_or(0);
out.push_str(&format!(
" {:<18} {:>10} {share:>3}%\n",
bucket.name.as_str(),
duration(bucket.ms)
));
}
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
}
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::{Labels, ENVELOPE_VERSION};
use crate::plan::{Node, Plan, PLAN_SCHEMA_VERSION};
use serde_json::json;
fn at(
seconds: u64,
kind: journal::PipelineKind,
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: Source::Pipeline,
kind: kind.into(),
labels: Labels {
run_id: Some("demo".into()),
round: Some(1),
node: node.map(str::to_string),
..Labels::default()
},
payload: journal::payload(fields),
artifacts: Vec::new(),
}
}
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()
}],
}
}
#[test]
fn the_buckets_sum_exactly_to_the_wall_clock() {
let events = vec![
at(
0,
journal::PipelineKind::RunStarted,
None,
&[("plan", json!(plan()))],
),
at(
10,
journal::PipelineKind::NodeDispatched,
Some("build"),
&[],
),
at(
70,
journal::PipelineKind::NodeSettled,
Some("build"),
&[("status", json!("done"))],
),
at(100, journal::PipelineKind::RoundFinished, None, &[]),
];
let telemetry = of_run("demo", &events);
assert_eq!(telemetry.wall_ms, 100_000);
let summed: u64 = telemetry.buckets.iter().map(|b| b.ms).sum();
assert_eq!(summed, telemetry.wall_ms, "{:?}", telemetry.buckets);
let bucket = |name: BucketName| {
telemetry
.buckets
.iter()
.find(|b| b.name == name)
.unwrap_or_else(|| panic!("a {} bucket", name.as_str()))
.ms
};
assert_eq!(bucket(BucketName::Dispatching), 60_000);
assert_eq!(bucket(BucketName::Orchestrating), 40_000);
assert_eq!(telemetry.dispatches, 1);
assert_eq!(telemetry.settled_done, 1);
}
#[test]
fn waiting_on_a_person_and_on_the_planner_are_different_buckets() {
let events = vec![
at(
0,
journal::PipelineKind::RunStarted,
None,
&[("plan", json!(plan()))],
),
at(
10,
journal::PipelineKind::NodeSettled,
Some("approve"),
&[("status", json!("waiting"))],
),
at(
40,
journal::PipelineKind::HumanAttested,
None,
&[("ref", json!("approve"))],
),
at(
50,
journal::PipelineKind::PlannerSurfaced,
None,
&[("blocking", json!(true))],
),
at(90, journal::PipelineKind::PlannerReplied, None, &[]),
];
let telemetry = of_run("demo", &events);
let bucket = |name: BucketName| {
telemetry
.buckets
.iter()
.find(|b| b.name == name)
.unwrap_or_else(|| panic!("a {} bucket", name.as_str()))
.ms
};
assert_eq!(bucket(BucketName::AwaitingHuman), 30_000);
assert_eq!(bucket(BucketName::AwaitingPlanner), 40_000);
assert_eq!(
telemetry.buckets.iter().map(|b| b.ms).sum::<u64>(),
telemetry.wall_ms
);
}
#[test]
fn a_non_blocking_surface_does_not_park_the_run_on_the_planner() {
let events = vec![
at(
0,
journal::PipelineKind::RunStarted,
None,
&[("plan", json!(plan()))],
),
at(
10,
journal::PipelineKind::PlannerSurfaced,
None,
&[("blocking", json!(false))],
),
at(50, journal::PipelineKind::RoundFinished, None, &[]),
];
let telemetry = of_run("demo", &events);
let awaiting = telemetry
.buckets
.iter()
.find(|b| b.name == BucketName::AwaitingPlanner)
.expect("the bucket")
.ms;
assert_eq!(awaiting, 0, "a heartbeat parked the run");
}
#[test]
fn an_empty_run_has_a_zero_wall_clock_and_still_balances() {
let telemetry = of_run("demo", &[]);
assert_eq!(telemetry.wall_ms, 0);
assert_eq!(telemetry.buckets.iter().map(|b| b.ms).sum::<u64>(), 0);
assert!(render_breakdown(&telemetry).contains("WALL 0s"));
}
#[test]
fn a_bucket_serialises_as_the_word_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
);
}
}
#[test]
fn a_clock_that_moved_backwards_still_leaves_the_buckets_summing_to_wall() {
let mut events = vec![
at(
0,
journal::PipelineKind::RunStarted,
None,
&[("plan", json!(plan()))],
),
at(
60,
journal::PipelineKind::NodeDispatched,
Some("build"),
&[],
),
at(30, journal::PipelineKind::RoundFinished, None, &[]),
];
events.reverse();
let telemetry = of_run("demo", &events);
assert_eq!(
telemetry.buckets.iter().map(|b| b.ms).sum::<u64>(),
telemetry.wall_ms,
"{:?}",
telemetry.buckets
);
}
#[test]
fn the_breakdown_names_every_bucket_and_its_share() {
let events = vec![
at(
0,
journal::PipelineKind::RunStarted,
None,
&[("plan", json!(plan()))],
),
at(
10,
journal::PipelineKind::NodeDispatched,
Some("build"),
&[],
),
at(
20,
journal::PipelineKind::NodeSettled,
Some("build"),
&[("status", json!("done"))],
),
];
let rendered = render_breakdown(&of_run("demo", &events));
for name in BucketName::ALL {
assert!(
rendered.contains(name.as_str()),
"{rendered} omits {}",
name.as_str()
);
}
assert!(rendered.contains('%'), "{rendered}");
}
#[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");
}
}