1use std::collections::BTreeMap;
36
37use crate::model::{KeyAgg, KeyMetric, ScenarioNode, Workload, WorkloadPhase};
38
39const WINDOW: &str = "30d";
43
44fn spine_contract() -> Vec<KeyMetric> {
49 use KeyAgg::*;
50 vec![
51 KeyMetric {
52 column: "count".into(),
53 agg: Last,
54 family: "result_success".into(),
55 },
56 KeyMetric {
57 column: "failures".into(),
58 agg: Last,
59 family: "result_failure".into(),
60 },
61 KeyMetric {
62 column: "wall".into(),
63 agg: Span,
64 family: String::new(),
65 },
66 KeyMetric {
67 column: "p99".into(),
68 agg: Max,
69 family: "result_success_p99".into(),
70 },
71 ]
72}
73
74#[derive(Debug, Default)]
76struct View {
77 coords: Vec<String>,
80 phases: Vec<String>,
82 defs: Vec<String>,
86}
87
88#[derive(Debug, Clone, PartialEq)]
90enum Attach {
91 Spine,
92 View(String),
93 Flattened {
97 via: String,
98 },
99}
100
101struct Walk {
102 views: BTreeMap<String, View>,
103 attachments: BTreeMap<String, Vec<Attach>>,
106 errors: Vec<String>,
107}
108
109impl Walk {
110 fn walk_nodes(
111 &mut self,
112 nodes: &[ScenarioNode],
113 anchor: Option<&str>,
114 flattened_via: Option<&str>,
115 ) {
116 for node in nodes {
117 match node {
118 ScenarioNode::Phase(name) => {
119 let attach = match (flattened_via, anchor) {
120 (Some(via), _) => Attach::Flattened {
121 via: via.to_string(),
122 },
123 (None, Some(v)) => Attach::View(v.to_string()),
124 (None, None) => Attach::Spine,
125 };
126 let entry = self.attachments.entry(name.clone()).or_default();
127 if !entry.contains(&attach) {
128 entry.push(attach.clone());
129 }
130 if let Attach::View(v) = &attach {
131 let view = self
132 .views
133 .get_mut(v)
134 .expect("view registered before descent");
135 if !view.phases.contains(name) {
136 view.phases.push(name.clone());
137 }
138 }
139 }
140 ScenarioNode::Comprehension {
141 comprehension,
142 children,
143 anchor: node_anchor,
144 ..
145 } => {
146 let coords = comprehension.coordinate_names();
147 match node_anchor {
148 Some(view_name) => {
149 let view = self.views.entry(view_name.clone()).or_default();
150 let def = describe_comprehension(comprehension);
151 if !view.defs.contains(&def) {
152 view.defs.push(def);
153 }
154 if view.coords.is_empty() {
155 view.coords = coords.clone();
156 } else if view.coords != coords {
157 self.errors.push(format!(
158 "anchor view '{view_name}': coordinate label sets \
159 disagree across anchors ({:?} vs {:?}) — views \
160 sharing a name must share coordinates",
161 view.coords, coords
162 ));
163 }
164 self.walk_nodes(children, Some(view_name), None);
168 }
169 None => {
170 let via = format!("for {}", coords.join(","));
173 self.walk_nodes(children, anchor, Some(&via));
174 }
175 }
176 }
177 ScenarioNode::IncludedScenario { children, .. }
178 | ScenarioNode::Bindings { children, .. }
179 | ScenarioNode::DoWhile { children, .. }
180 | ScenarioNode::DoUntil { children, .. } => {
181 self.walk_nodes(children, anchor, flattened_via);
186 }
187 }
188 }
189 }
190}
191
192fn known_families(phase: &WorkloadPhase) -> Vec<String> {
197 let mut out: Vec<String> = Vec::new();
198 for (name, _) in &phase.metrics {
199 out.push(name.clone());
200 }
201 for op in &phase.ops {
202 for (name, _) in &op.metrics {
203 out.push(name.clone());
204 }
205 for source in [&op.params, &op.op] {
210 if let Some(m) = source
211 .get("poll")
212 .and_then(|v| v.as_object())
213 .and_then(|p| p.get("metric_name"))
214 .and_then(|v| v.as_str())
215 {
216 out.push(m.to_string());
217 }
218 }
219 }
220 if let Some(poll) = &phase.poll
221 && let Some(m) = &poll.metric_name
222 {
223 out.push(m.clone());
224 }
225 const INSTRUMENTS: &[&str] = &[
226 "result_success",
227 "result_failure",
228 "result_total",
229 "attempt_total",
230 "attempt_success",
231 "attempt_failure",
232 "result_bytes",
233 "result_elements",
234 "cycles_total",
235 "cycles_servicetime",
236 ];
237 const SUFFIXES: &[&str] = &[
238 "", "_p50", "_p75", "_p90", "_p95", "_p99", "_p999", "_mean", "_min", "_max", "_stddev",
239 "_rate", "_count",
240 ];
241 for i in INSTRUMENTS {
242 for sfx in SUFFIXES {
243 out.push(format!("{i}{sfx}"));
244 }
245 }
246 out
247}
248
249fn family_known(phase: &WorkloadPhase, family: &str) -> bool {
250 family.starts_with("recall_") || known_families(phase).iter().any(|f| f == family)
251}
252
253fn agg_desc(km: &KeyMetric) -> String {
256 use KeyAgg::*;
257 match km.agg {
258 Span => "span()".to_string(),
259 Rate => format!("rate({})", km.family),
260 Delta => format!("delta({})", km.family),
261 _ => format!("{}({})", format!("{:?}", km.agg).to_lowercase(), km.family),
262 }
263}
264
265fn describe_comprehension(c: &polydat::iteration::comprehension::Comprehension) -> String {
268 use polydat::iteration::comprehension::Comprehension as C;
269 match c {
270 C::Clause { name, source } => format!("{name} in {}", describe_source(source)),
271 C::Cartesian { children } => children
272 .iter()
273 .map(describe_comprehension)
274 .collect::<Vec<_>>()
275 .join(", "),
276 C::Zip { children, .. } => format!(
277 "zip({})",
278 children
279 .iter()
280 .map(describe_comprehension)
281 .collect::<Vec<_>>()
282 .join(", ")
283 ),
284 C::Union { children } => children
285 .iter()
286 .map(describe_comprehension)
287 .collect::<Vec<_>>()
288 .join(" | "),
289 C::Filter { child, predicate } => {
290 format!("{} if {predicate}", describe_comprehension(child))
291 }
292 C::Order {
293 child,
294 strategy,
295 truncation,
296 seed,
297 } => {
298 let mut text = format!("{} ordered {strategy:?}", describe_comprehension(child));
299 if let Some(n) = truncation {
300 text.push_str(&format!(" take {n}"));
301 }
302 if let Some(s) = seed {
303 text.push_str(&format!(" seed {s}"));
304 }
305 text
306 }
307 }
308}
309
310fn describe_source(s: &polydat::iteration::comprehension::Source) -> String {
311 use polydat::iteration::comprehension::Source;
312 use polydat::iteration::comprehension::source::LiteralValue;
313 match s {
314 Source::Literal { values } => values
315 .iter()
316 .map(|v| match v {
317 LiteralValue::Int(i) => i.to_string(),
318 LiteralValue::UInt(u) => u.to_string(),
319 LiteralValue::Float(f) => f.to_string(),
320 LiteralValue::String(st) => st.clone(),
321 LiteralValue::Bool(b) => b.to_string(),
322 LiteralValue::Json(j) => j.to_string(),
323 })
324 .collect::<Vec<_>>()
325 .join(","),
326 Source::IntRange { lo, hi, step } if *step == 1 => format!("{lo}..{hi}"),
327 Source::IntRange { lo, hi, step } => format!("{lo}..{hi} step {step}"),
328 Source::Generator { expr, .. } => expr.clone(),
329 Source::WorkloadParamList { name, .. } => format!("{{{name}}}"),
330 Source::ContinuousInterval { .. } => "<continuous interval>".to_string(),
331 Source::Distribution { distribution, .. } => format!("{distribution:?}(…)"),
332 }
333}
334
335fn query_for(metric: &KeyMetric, phase: &str, by_labels: &str) -> String {
340 use KeyAgg::*;
341 let f = &metric.family;
342 let sel = format!("{{phase=\"{phase}\"}}");
343 let span =
350 format!("max(sum_over_time(result_success_interval_ns{sel}[{WINDOW}])) by ({by_labels})");
351 match metric.agg {
352 Min => format!("min(min_over_time({f}{sel}[{WINDOW}])) by ({by_labels})"),
353 Max => format!("max(max_over_time({f}{sel}[{WINDOW}])) by ({by_labels})"),
354 Avg => format!("avg(avg_over_time({f}{sel}[{WINDOW}])) by ({by_labels})"),
355 Last => format!("max(last_over_time({f}{sel}[{WINDOW}])) by ({by_labels})"),
356 First => format!("min(first_over_time({f}{sel}[{WINDOW}])) by ({by_labels})"),
357 Median => format!("avg(median_over_time({f}{sel}[{WINDOW}])) by ({by_labels})"),
358 Stddev => format!("avg(stddev_over_time({f}{sel}[{WINDOW}])) by ({by_labels})"),
359 Sum => format!("sum(sum_over_time({f}{sel}[{WINDOW}])) by ({by_labels})"),
360 Count => format!("sum(count_over_time({f}{sel}[{WINDOW}])) by ({by_labels})"),
361 Span => span,
362 Rate => format!("avg(avg_over_time({f}_rate{sel}[{WINDOW}])) by ({by_labels})"),
363 Delta => format!(
364 "max(last_over_time({f}{sel}[{WINDOW}])) by ({by_labels}) - \
365 min(first_over_time({f}{sel}[{WINDOW}])) by ({by_labels})"
366 ),
367 }
368}
369
370pub fn synthesize(workload: &Workload) -> Result<Option<serde_json::Value>, String> {
375 if !workload.report.groups.is_empty() {
376 return Ok(None);
377 }
378 synthesize_forced(workload).map(Some)
379}
380
381pub fn synthesize_forced(workload: &Workload) -> Result<serde_json::Value, String> {
387 synthesize_forced_for(workload, None)
388}
389
390pub fn synthesize_forced_for(
397 workload: &Workload,
398 scenario: Option<&str>,
399) -> Result<serde_json::Value, String> {
400 let mut walk = Walk {
401 views: BTreeMap::new(),
402 attachments: BTreeMap::new(),
403 errors: Vec::new(),
404 };
405 let mut scenario_names: Vec<&String> = workload.scenarios.keys().collect();
413 scenario_names.sort_by_key(|n| (*n != "default", (*n).clone()));
414 if let Some(sel) = scenario
415 && let Some(name) = scenario_names.iter().find(|n| n.as_str() == sel).copied()
416 {
417 scenario_names = vec![name];
418 }
419 for name in scenario_names {
420 walk.walk_nodes(&workload.scenarios[name], None, None);
421 }
422 let Walk {
423 views,
424 attachments,
425 mut errors,
426 ..
427 } = walk;
428
429 for (phase_name, attaches) in &attachments {
432 let Some(phase) = workload.phases.get(phase_name) else {
433 continue;
434 };
435 for km in &phase.key_metrics {
436 if km.agg != KeyAgg::Span && !family_known(phase, &km.family) {
437 errors.push(format!(
438 "phase '{phase_name}' key_metrics.{}: family '{}' is not \
439 emitted by this phase (declared metrics, poll timers, \
440 SRD-91 instruments and their stat suffixes, recall_*)",
441 km.column, km.family
442 ));
443 }
444 }
445 if !phase.key_metrics.is_empty() {
446 for a in attaches {
447 if let Attach::Flattened { via } = a {
448 errors.push(format!(
449 "phase '{phase_name}' designates key metrics but its \
450 activations multiply through a non-anchored sweep \
451 ({via}): one table row would silently aggregate many \
452 activations. Anchor that sweep (`anchor: <view>`) or \
453 remove the designations. There are no implied \
454 aggregates."
455 ));
456 }
457 }
458 }
459 }
460 if !errors.is_empty() {
461 return Err(format!(
462 "report synthesis: {} well-formedness error(s):\n - {}",
463 errors.len(),
464 errors.join("\n - ")
465 ));
466 }
467
468 let mut groups = serde_json::Map::new();
469
470 let spine_phases: Vec<&String> = attachments
472 .iter()
473 .filter(|(_, a)| a.contains(&Attach::Spine))
474 .map(|(n, _)| n)
475 .collect();
476 let looped_unanchored: Vec<&String> = attachments
477 .iter()
478 .filter(|(_, a)| {
479 a.iter().any(|x| matches!(x, Attach::Flattened { .. }))
480 && !a
481 .iter()
482 .any(|x| matches!(x, Attach::View(_)) || *x == Attach::Spine)
483 })
484 .map(|(n, _)| n)
485 .collect();
486 let mut spine = String::new();
487 spine.push_str(
488 "text phases_intro as \"Workload phases — outcomes at a glance\":\n \
489 One row per workload phase, keyed by the `phase` label. Column \
490 headers carry each value's definition (aggregate over the \
491 phase's samples). Each anchored sweep in the scenario renders \
492 as its own view table below, one row per sweep iteration.\n",
493 );
494 if !looped_unanchored.is_empty() {
495 spine.push_str(&format!(
496 "text phases_unanchored as \"Not tabulated\":\n \
497 Looped phases with no anchor and no designations — activations \
498 would aggregate silently, so no rows are synthesized: {}.\n",
499 looped_unanchored
500 .iter()
501 .map(|s| s.as_str())
502 .collect::<Vec<_>>()
503 .join(", ")
504 ));
505 }
506 if !spine_phases.is_empty() {
511 spine.push_str("table phases:\n group_by: phase\n");
512 spine.push_str(" label \"phases — one row per workload phase\"\n");
513 let spine_sel = format!(
514 "{{phase=~\"{}\"}}",
515 spine_phases
516 .iter()
517 .map(|s| regex_escape(s))
518 .collect::<Vec<_>>()
519 .join("|")
520 );
521 for km in spine_contract() {
522 let q = query_for_selector(&km, &spine_sel, "phase");
523 spine.push_str(&format!(" query: {}: {}\n", km.column, q));
524 spine.push_str(&format!(" header {}: {}\n", km.column, agg_desc(&km)));
525 }
526 }
527 groups.insert("phases".to_string(), serde_json::Value::String(spine));
528
529 for (view_name, view) in &views {
531 let by = view.coords.join(",");
532 let mut table = String::new();
537 let mut seen_cols: Vec<String> = Vec::new();
538 let mut phase_names: Vec<&str> = Vec::new();
539 let contributing = view
540 .phases
541 .iter()
542 .filter(|p| {
543 workload
544 .phases
545 .get(p.as_str())
546 .is_some_and(|ph| !ph.key_metrics.is_empty())
547 })
548 .count();
549 for phase_name in &view.phases {
550 let Some(phase) = workload.phases.get(phase_name) else {
551 continue;
552 };
553 if !phase.key_metrics.is_empty() {
554 phase_names.push(phase_name);
555 }
556 for km in &phase.key_metrics {
557 let col = if seen_cols.contains(&km.column) {
558 format!("{phase_name}_{}", km.column)
559 } else {
560 km.column.clone()
561 };
562 seen_cols.push(col.clone());
563 let q = query_for(km, phase_name, &by);
564 table.push_str(&format!(" query: {col}: {q}\n"));
565 let note = if contributing > 1 {
566 format!("{} @{phase_name}", agg_desc(km))
567 } else {
568 agg_desc(km)
569 };
570 table.push_str(&format!(" header {col}: {note}\n"));
571 }
572 }
573 let label = format!(
576 "{view_name} — one row per {by} ({})",
577 phase_names.join(", ")
578 );
579 let defs = view
583 .defs
584 .iter()
585 .map(|d| format!("`for {d}`"))
586 .collect::<Vec<_>>()
587 .join(" and ");
588 let about = format!(
589 "text {view_name}_about as \"{view_name} — one row per {by}\":\n \
590 Rows are the iterations of {defs}; the {by} label(s) \
591 identify each row within the workload. Columns aggregate \
592 each iteration's activations of: {}. Headers carry each \
593 column's definition{}.\n",
594 phase_names.join(", "),
595 if contributing > 1 {
596 " and its @phase provenance"
597 } else {
598 ""
599 }
600 );
601 let body =
602 format!("{about}table {view_name}:\n group_by: {by}\n label \"{label}\"\n{table}");
603 groups.insert(format!("view_{view_name}"), serde_json::Value::String(body));
604 }
605
606 Ok(serde_json::Value::Object(groups))
607}
608
609fn query_for_selector(metric: &KeyMetric, sel: &str, by_labels: &str) -> String {
612 use KeyAgg::*;
613 let f = &metric.family;
614 match metric.agg {
615 Span => format!(
616 "max(sum_over_time(result_success_interval_ns{sel}[{WINDOW}])) by ({by_labels})"
617 ),
618 Last => format!("max(last_over_time({f}{sel}[{WINDOW}])) by ({by_labels})"),
619 Max => format!("max(max_over_time({f}{sel}[{WINDOW}])) by ({by_labels})"),
620 _ => {
621 let km = KeyMetric {
624 column: metric.column.clone(),
625 agg: metric.agg,
626 family: f.clone(),
627 };
628 let q = query_for(&km, "__sel__", by_labels);
629 q.replace("{phase=\"__sel__\"}", sel)
630 }
631 }
632}
633
634pub fn synthesize_yaml(workload: &Workload) -> Result<String, String> {
639 let value = synthesize_forced(workload)?;
640 let map = value.as_object().expect("synthesize emits a mapping");
641 let mut out = String::from("report:\n");
642 for (group, body) in map {
643 out.push_str(&format!(" {group}: |\n"));
644 for line in body.as_str().unwrap_or_default().lines() {
645 out.push_str(&format!(" {line}\n"));
646 }
647 }
648 Ok(out)
649}
650
651fn regex_escape(s: &str) -> String {
652 let mut out = String::with_capacity(s.len());
653 for c in s.chars() {
654 if "\\.^$|?*+()[]{}".contains(c) {
655 out.push('\\');
656 }
657 out.push(c);
658 }
659 out
660}
661
662#[cfg(test)]
663mod tests {
664 use super::*;
665
666 fn wl(yaml: &str) -> Workload {
667 crate::parse::parse_workload(yaml, &std::collections::HashMap::new())
668 .expect("workload parses")
669 }
670
671 const BASE: &str = r#"
672params: { adapter: testkit }
673phases:
674 tick:
675 key_metrics:
676 rows: last(result_success)
677 spd: rate(result_success)
678 ops: { t: { stmt: "X" } }
679 plain:
680 ops: { t: { stmt: "Y" } }
681"#;
682
683 #[test]
684 fn synthesizes_spine_and_anchored_view() {
685 let y = format!(
686 "{BASE}
687scenarios:
688 default:
689 - plain
690 - for: \"k in 1,2\"
691 anchor: sweep
692 phases: [tick]
693"
694 );
695 let v = synthesize(&wl(&y))
696 .expect("well-formed")
697 .expect("synthesized");
698 let m = v.as_object().unwrap();
699 assert!(m.contains_key("phases"));
700 let sweep = m.get("view_sweep").unwrap().as_str().unwrap();
701 assert!(
702 sweep.contains("group_by: k"),
703 "anchor coordinate keys the view"
704 );
705 assert!(sweep.contains("last_over_time(result_success{phase=\"tick\"}"));
706 assert!(
707 sweep.contains("query: spd:"),
708 "rate designation synthesized"
709 );
710 let spine = m.get("phases").unwrap().as_str().unwrap();
711 assert!(spine.contains("plain"), "un-looped phase rows the spine");
712 assert!(
713 !spine.contains("|tick") && !spine.contains("\"tick"),
714 "anchored phase stays out of the spine selector"
715 );
716 crate::report::parse_report(&v).expect("synthesized section parses");
719 }
720
721 #[test]
722 fn designated_phase_under_unanchored_sweep_is_an_error() {
723 let y = format!(
724 "{BASE}
725scenarios:
726 default:
727 - for: \"k in 1,2\"
728 phases: [tick]
729"
730 );
731 let e = synthesize(&wl(&y)).unwrap_err();
732 assert!(e.contains("non-anchored sweep"), "{e}");
733 assert!(
734 e.contains("no implied") || e.contains("Anchor that sweep"),
735 "{e}"
736 );
737 }
738
739 #[test]
740 fn unknown_family_is_an_error_with_provenance() {
741 let y = "
742params: { adapter: testkit }
743phases:
744 tick:
745 key_metrics: { bogus: avg(no_such_family) }
746 ops: { t: { stmt: \"X\" } }
747scenarios:
748 default: [tick]
749";
750 let e = synthesize(&wl(y)).unwrap_err();
751 assert!(e.contains("no_such_family") && e.contains("tick"), "{e}");
752 }
753
754 #[test]
755 fn explicit_report_block_suppresses_synthesis() {
756 let y = "
757params: { adapter: testkit }
758report:
759 g: |
760 text t as \"T\":
761 body
762phases:
763 tick: { ops: { t: { stmt: \"X\" } } }
764scenarios:
765 default: [tick]
766";
767 assert!(synthesize(&wl(y)).expect("ok").is_none());
768 }
769
770 #[test]
771 fn unqualified_designation_is_a_parse_error_with_vocabulary() {
772 let y = "
773params: { adapter: testkit }
774phases:
775 tick:
776 key_metrics: { rows: result_success }
777 ops: { t: { stmt: \"X\" } }
778scenarios:
779 default: [tick]
780";
781 let e = crate::parse::parse_workload(y, &std::collections::HashMap::new())
782 .expect_err("must reject unqualified aggregate");
783 assert!(
784 e.contains("aggregate qualification required") && e.contains("min, max, avg"),
785 "{e}"
786 );
787 }
788}