use std::collections::BTreeMap;
use crate::model::{KeyAgg, KeyMetric, ScenarioNode, Workload, WorkloadPhase};
const WINDOW: &str = "30d";
fn spine_contract() -> Vec<KeyMetric> {
use KeyAgg::*;
vec![
KeyMetric {
column: "count".into(),
agg: Last,
family: "result_success".into(),
},
KeyMetric {
column: "failures".into(),
agg: Last,
family: "result_failure".into(),
},
KeyMetric {
column: "wall".into(),
agg: Span,
family: String::new(),
},
KeyMetric {
column: "p99".into(),
agg: Max,
family: "result_success_p99".into(),
},
]
}
#[derive(Debug, Default)]
struct View {
coords: Vec<String>,
phases: Vec<String>,
defs: Vec<String>,
}
#[derive(Debug, Clone, PartialEq)]
enum Attach {
Spine,
View(String),
Flattened {
via: String,
},
}
struct Walk {
views: BTreeMap<String, View>,
attachments: BTreeMap<String, Vec<Attach>>,
errors: Vec<String>,
}
impl Walk {
fn walk_nodes(
&mut self,
nodes: &[ScenarioNode],
anchor: Option<&str>,
flattened_via: Option<&str>,
) {
for node in nodes {
match node {
ScenarioNode::Phase(name) => {
let attach = match (flattened_via, anchor) {
(Some(via), _) => Attach::Flattened {
via: via.to_string(),
},
(None, Some(v)) => Attach::View(v.to_string()),
(None, None) => Attach::Spine,
};
let entry = self.attachments.entry(name.clone()).or_default();
if !entry.contains(&attach) {
entry.push(attach.clone());
}
if let Attach::View(v) = &attach {
let view = self
.views
.get_mut(v)
.expect("view registered before descent");
if !view.phases.contains(name) {
view.phases.push(name.clone());
}
}
}
ScenarioNode::Comprehension {
comprehension,
children,
anchor: node_anchor,
..
} => {
let coords = comprehension.coordinate_names();
match node_anchor {
Some(view_name) => {
let view = self.views.entry(view_name.clone()).or_default();
let def = describe_comprehension(comprehension);
if !view.defs.contains(&def) {
view.defs.push(def);
}
if view.coords.is_empty() {
view.coords = coords.clone();
} else if view.coords != coords {
self.errors.push(format!(
"anchor view '{view_name}': coordinate label sets \
disagree across anchors ({:?} vs {:?}) — views \
sharing a name must share coordinates",
view.coords, coords
));
}
self.walk_nodes(children, Some(view_name), None);
}
None => {
let via = format!("for {}", coords.join(","));
self.walk_nodes(children, anchor, Some(&via));
}
}
}
ScenarioNode::IncludedScenario { children, .. }
| ScenarioNode::Bindings { children, .. }
| ScenarioNode::DoWhile { children, .. }
| ScenarioNode::DoUntil { children, .. } => {
self.walk_nodes(children, anchor, flattened_via);
}
}
}
}
}
fn known_families(phase: &WorkloadPhase) -> Vec<String> {
let mut out: Vec<String> = Vec::new();
for (name, _) in &phase.metrics {
out.push(name.clone());
}
for op in &phase.ops {
for (name, _) in &op.metrics {
out.push(name.clone());
}
for source in [&op.params, &op.op] {
if let Some(m) = source
.get("poll")
.and_then(|v| v.as_object())
.and_then(|p| p.get("metric_name"))
.and_then(|v| v.as_str())
{
out.push(m.to_string());
}
}
}
if let Some(poll) = &phase.poll
&& let Some(m) = &poll.metric_name
{
out.push(m.clone());
}
const INSTRUMENTS: &[&str] = &[
"result_success",
"result_failure",
"result_total",
"attempt_total",
"attempt_success",
"attempt_failure",
"result_bytes",
"result_elements",
"cycles_total",
"cycles_servicetime",
];
const SUFFIXES: &[&str] = &[
"", "_p50", "_p75", "_p90", "_p95", "_p99", "_p999", "_mean", "_min", "_max", "_stddev",
"_rate", "_count",
];
for i in INSTRUMENTS {
for sfx in SUFFIXES {
out.push(format!("{i}{sfx}"));
}
}
out
}
fn family_known(phase: &WorkloadPhase, family: &str) -> bool {
family.starts_with("recall_") || known_families(phase).iter().any(|f| f == family)
}
fn agg_desc(km: &KeyMetric) -> String {
use KeyAgg::*;
match km.agg {
Span => "span()".to_string(),
Rate => format!("rate({})", km.family),
Delta => format!("delta({})", km.family),
_ => format!("{}({})", format!("{:?}", km.agg).to_lowercase(), km.family),
}
}
fn describe_comprehension(c: &polydat::iteration::comprehension::Comprehension) -> String {
use polydat::iteration::comprehension::Comprehension as C;
match c {
C::Clause { name, source } => format!("{name} in {}", describe_source(source)),
C::Cartesian { children } => children
.iter()
.map(describe_comprehension)
.collect::<Vec<_>>()
.join(", "),
C::Zip { children, .. } => format!(
"zip({})",
children
.iter()
.map(describe_comprehension)
.collect::<Vec<_>>()
.join(", ")
),
C::Union { children } => children
.iter()
.map(describe_comprehension)
.collect::<Vec<_>>()
.join(" | "),
C::Filter { child, predicate } => {
format!("{} if {predicate}", describe_comprehension(child))
}
C::Order {
child,
strategy,
truncation,
seed,
} => {
let mut text = format!("{} ordered {strategy:?}", describe_comprehension(child));
if let Some(n) = truncation {
text.push_str(&format!(" take {n}"));
}
if let Some(s) = seed {
text.push_str(&format!(" seed {s}"));
}
text
}
}
}
fn describe_source(s: &polydat::iteration::comprehension::Source) -> String {
use polydat::iteration::comprehension::Source;
use polydat::iteration::comprehension::source::LiteralValue;
match s {
Source::Literal { values } => values
.iter()
.map(|v| match v {
LiteralValue::Int(i) => i.to_string(),
LiteralValue::UInt(u) => u.to_string(),
LiteralValue::Float(f) => f.to_string(),
LiteralValue::String(st) => st.clone(),
LiteralValue::Bool(b) => b.to_string(),
LiteralValue::Json(j) => j.to_string(),
})
.collect::<Vec<_>>()
.join(","),
Source::IntRange { lo, hi, step } if *step == 1 => format!("{lo}..{hi}"),
Source::IntRange { lo, hi, step } => format!("{lo}..{hi} step {step}"),
Source::Generator { expr, .. } => expr.clone(),
Source::WorkloadParamList { name, .. } => format!("{{{name}}}"),
Source::ContinuousInterval { .. } => "<continuous interval>".to_string(),
Source::Distribution { distribution, .. } => format!("{distribution:?}(…)"),
}
}
fn query_for(metric: &KeyMetric, phase: &str, by_labels: &str) -> String {
use KeyAgg::*;
let f = &metric.family;
let sel = format!("{{phase=\"{phase}\"}}");
let span =
format!("max(sum_over_time(result_success_interval_ns{sel}[{WINDOW}])) by ({by_labels})");
match metric.agg {
Min => format!("min(min_over_time({f}{sel}[{WINDOW}])) by ({by_labels})"),
Max => format!("max(max_over_time({f}{sel}[{WINDOW}])) by ({by_labels})"),
Avg => format!("avg(avg_over_time({f}{sel}[{WINDOW}])) by ({by_labels})"),
Last => format!("max(last_over_time({f}{sel}[{WINDOW}])) by ({by_labels})"),
First => format!("min(first_over_time({f}{sel}[{WINDOW}])) by ({by_labels})"),
Median => format!("avg(median_over_time({f}{sel}[{WINDOW}])) by ({by_labels})"),
Stddev => format!("avg(stddev_over_time({f}{sel}[{WINDOW}])) by ({by_labels})"),
Sum => format!("sum(sum_over_time({f}{sel}[{WINDOW}])) by ({by_labels})"),
Count => format!("sum(count_over_time({f}{sel}[{WINDOW}])) by ({by_labels})"),
Span => span,
Rate => format!("avg(avg_over_time({f}_rate{sel}[{WINDOW}])) by ({by_labels})"),
Delta => format!(
"max(last_over_time({f}{sel}[{WINDOW}])) by ({by_labels}) - \
min(first_over_time({f}{sel}[{WINDOW}])) by ({by_labels})"
),
}
}
pub fn synthesize(workload: &Workload) -> Result<Option<serde_json::Value>, String> {
if !workload.report.groups.is_empty() {
return Ok(None);
}
synthesize_forced(workload).map(Some)
}
pub fn synthesize_forced(workload: &Workload) -> Result<serde_json::Value, String> {
synthesize_forced_for(workload, None)
}
pub fn synthesize_forced_for(
workload: &Workload,
scenario: Option<&str>,
) -> Result<serde_json::Value, String> {
let mut walk = Walk {
views: BTreeMap::new(),
attachments: BTreeMap::new(),
errors: Vec::new(),
};
let mut scenario_names: Vec<&String> = workload.scenarios.keys().collect();
scenario_names.sort_by_key(|n| (*n != "default", (*n).clone()));
if let Some(sel) = scenario
&& let Some(name) = scenario_names.iter().find(|n| n.as_str() == sel).copied()
{
scenario_names = vec![name];
}
for name in scenario_names {
walk.walk_nodes(&workload.scenarios[name], None, None);
}
let Walk {
views,
attachments,
mut errors,
..
} = walk;
for (phase_name, attaches) in &attachments {
let Some(phase) = workload.phases.get(phase_name) else {
continue;
};
for km in &phase.key_metrics {
if km.agg != KeyAgg::Span && !family_known(phase, &km.family) {
errors.push(format!(
"phase '{phase_name}' key_metrics.{}: family '{}' is not \
emitted by this phase (declared metrics, poll timers, \
SRD-91 instruments and their stat suffixes, recall_*)",
km.column, km.family
));
}
}
if !phase.key_metrics.is_empty() {
for a in attaches {
if let Attach::Flattened { via } = a {
errors.push(format!(
"phase '{phase_name}' designates key metrics but its \
activations multiply through a non-anchored sweep \
({via}): one table row would silently aggregate many \
activations. Anchor that sweep (`anchor: <view>`) or \
remove the designations. There are no implied \
aggregates."
));
}
}
}
}
if !errors.is_empty() {
return Err(format!(
"report synthesis: {} well-formedness error(s):\n - {}",
errors.len(),
errors.join("\n - ")
));
}
let mut groups = serde_json::Map::new();
let spine_phases: Vec<&String> = attachments
.iter()
.filter(|(_, a)| a.contains(&Attach::Spine))
.map(|(n, _)| n)
.collect();
let looped_unanchored: Vec<&String> = attachments
.iter()
.filter(|(_, a)| {
a.iter().any(|x| matches!(x, Attach::Flattened { .. }))
&& !a
.iter()
.any(|x| matches!(x, Attach::View(_)) || *x == Attach::Spine)
})
.map(|(n, _)| n)
.collect();
let mut spine = String::new();
spine.push_str(
"text phases_intro as \"Workload phases — outcomes at a glance\":\n \
One row per workload phase, keyed by the `phase` label. Column \
headers carry each value's definition (aggregate over the \
phase's samples). Each anchored sweep in the scenario renders \
as its own view table below, one row per sweep iteration.\n",
);
if !looped_unanchored.is_empty() {
spine.push_str(&format!(
"text phases_unanchored as \"Not tabulated\":\n \
Looped phases with no anchor and no designations — activations \
would aggregate silently, so no rows are synthesized: {}.\n",
looped_unanchored
.iter()
.map(|s| s.as_str())
.collect::<Vec<_>>()
.join(", ")
));
}
if !spine_phases.is_empty() {
spine.push_str("table phases:\n group_by: phase\n");
spine.push_str(" label \"phases — one row per workload phase\"\n");
let spine_sel = format!(
"{{phase=~\"{}\"}}",
spine_phases
.iter()
.map(|s| regex_escape(s))
.collect::<Vec<_>>()
.join("|")
);
for km in spine_contract() {
let q = query_for_selector(&km, &spine_sel, "phase");
spine.push_str(&format!(" query: {}: {}\n", km.column, q));
spine.push_str(&format!(" header {}: {}\n", km.column, agg_desc(&km)));
}
}
groups.insert("phases".to_string(), serde_json::Value::String(spine));
for (view_name, view) in &views {
let by = view.coords.join(",");
let mut table = String::new();
let mut seen_cols: Vec<String> = Vec::new();
let mut phase_names: Vec<&str> = Vec::new();
let contributing = view
.phases
.iter()
.filter(|p| {
workload
.phases
.get(p.as_str())
.is_some_and(|ph| !ph.key_metrics.is_empty())
})
.count();
for phase_name in &view.phases {
let Some(phase) = workload.phases.get(phase_name) else {
continue;
};
if !phase.key_metrics.is_empty() {
phase_names.push(phase_name);
}
for km in &phase.key_metrics {
let col = if seen_cols.contains(&km.column) {
format!("{phase_name}_{}", km.column)
} else {
km.column.clone()
};
seen_cols.push(col.clone());
let q = query_for(km, phase_name, &by);
table.push_str(&format!(" query: {col}: {q}\n"));
let note = if contributing > 1 {
format!("{} @{phase_name}", agg_desc(km))
} else {
agg_desc(km)
};
table.push_str(&format!(" header {col}: {note}\n"));
}
}
let label = format!(
"{view_name} — one row per {by} ({})",
phase_names.join(", ")
);
let defs = view
.defs
.iter()
.map(|d| format!("`for {d}`"))
.collect::<Vec<_>>()
.join(" and ");
let about = format!(
"text {view_name}_about as \"{view_name} — one row per {by}\":\n \
Rows are the iterations of {defs}; the {by} label(s) \
identify each row within the workload. Columns aggregate \
each iteration's activations of: {}. Headers carry each \
column's definition{}.\n",
phase_names.join(", "),
if contributing > 1 {
" and its @phase provenance"
} else {
""
}
);
let body =
format!("{about}table {view_name}:\n group_by: {by}\n label \"{label}\"\n{table}");
groups.insert(format!("view_{view_name}"), serde_json::Value::String(body));
}
Ok(serde_json::Value::Object(groups))
}
fn query_for_selector(metric: &KeyMetric, sel: &str, by_labels: &str) -> String {
use KeyAgg::*;
let f = &metric.family;
match metric.agg {
Span => format!(
"max(sum_over_time(result_success_interval_ns{sel}[{WINDOW}])) by ({by_labels})"
),
Last => format!("max(last_over_time({f}{sel}[{WINDOW}])) by ({by_labels})"),
Max => format!("max(max_over_time({f}{sel}[{WINDOW}])) by ({by_labels})"),
_ => {
let km = KeyMetric {
column: metric.column.clone(),
agg: metric.agg,
family: f.clone(),
};
let q = query_for(&km, "__sel__", by_labels);
q.replace("{phase=\"__sel__\"}", sel)
}
}
}
pub fn synthesize_yaml(workload: &Workload) -> Result<String, String> {
let value = synthesize_forced(workload)?;
let map = value.as_object().expect("synthesize emits a mapping");
let mut out = String::from("report:\n");
for (group, body) in map {
out.push_str(&format!(" {group}: |\n"));
for line in body.as_str().unwrap_or_default().lines() {
out.push_str(&format!(" {line}\n"));
}
}
Ok(out)
}
fn regex_escape(s: &str) -> String {
let mut out = String::with_capacity(s.len());
for c in s.chars() {
if "\\.^$|?*+()[]{}".contains(c) {
out.push('\\');
}
out.push(c);
}
out
}
#[cfg(test)]
mod tests {
use super::*;
fn wl(yaml: &str) -> Workload {
crate::parse::parse_workload(yaml, &std::collections::HashMap::new())
.expect("workload parses")
}
const BASE: &str = r#"
params: { adapter: testkit }
phases:
tick:
key_metrics:
rows: last(result_success)
spd: rate(result_success)
ops: { t: { stmt: "X" } }
plain:
ops: { t: { stmt: "Y" } }
"#;
#[test]
fn synthesizes_spine_and_anchored_view() {
let y = format!(
"{BASE}
scenarios:
default:
- plain
- for: \"k in 1,2\"
anchor: sweep
phases: [tick]
"
);
let v = synthesize(&wl(&y))
.expect("well-formed")
.expect("synthesized");
let m = v.as_object().unwrap();
assert!(m.contains_key("phases"));
let sweep = m.get("view_sweep").unwrap().as_str().unwrap();
assert!(
sweep.contains("group_by: k"),
"anchor coordinate keys the view"
);
assert!(sweep.contains("last_over_time(result_success{phase=\"tick\"}"));
assert!(
sweep.contains("query: spd:"),
"rate designation synthesized"
);
let spine = m.get("phases").unwrap().as_str().unwrap();
assert!(spine.contains("plain"), "un-looped phase rows the spine");
assert!(
!spine.contains("|tick") && !spine.contains("\"tick"),
"anchored phase stays out of the spine selector"
);
crate::report::parse_report(&v).expect("synthesized section parses");
}
#[test]
fn designated_phase_under_unanchored_sweep_is_an_error() {
let y = format!(
"{BASE}
scenarios:
default:
- for: \"k in 1,2\"
phases: [tick]
"
);
let e = synthesize(&wl(&y)).unwrap_err();
assert!(e.contains("non-anchored sweep"), "{e}");
assert!(
e.contains("no implied") || e.contains("Anchor that sweep"),
"{e}"
);
}
#[test]
fn unknown_family_is_an_error_with_provenance() {
let y = "
params: { adapter: testkit }
phases:
tick:
key_metrics: { bogus: avg(no_such_family) }
ops: { t: { stmt: \"X\" } }
scenarios:
default: [tick]
";
let e = synthesize(&wl(y)).unwrap_err();
assert!(e.contains("no_such_family") && e.contains("tick"), "{e}");
}
#[test]
fn explicit_report_block_suppresses_synthesis() {
let y = "
params: { adapter: testkit }
report:
g: |
text t as \"T\":
body
phases:
tick: { ops: { t: { stmt: \"X\" } } }
scenarios:
default: [tick]
";
assert!(synthesize(&wl(y)).expect("ok").is_none());
}
#[test]
fn unqualified_designation_is_a_parse_error_with_vocabulary() {
let y = "
params: { adapter: testkit }
phases:
tick:
key_metrics: { rows: result_success }
ops: { t: { stmt: \"X\" } }
scenarios:
default: [tick]
";
let e = crate::parse::parse_workload(y, &std::collections::HashMap::new())
.expect_err("must reject unqualified aggregate");
assert!(
e.contains("aggregate qualification required") && e.contains("min, max, avg"),
"{e}"
);
}
}