nmbrs_workload/model.rs
1// Copyright 2024-2026 Jonathan Shook
2// SPDX-License-Identifier: Apache-2.0
3
4//! Normalized workload model: the canonical ParsedOp representation.
5//!
6//! All YAML shorthand forms normalize to this model. This is what
7//! driver adapters consume.
8
9use serde::{Deserialize, Serialize};
10use std::collections::HashMap;
11
12/// A complete workload definition after normalization.
13#[derive(Debug, Clone, Serialize, Deserialize)]
14pub struct Workload {
15 #[serde(default)]
16 pub description: Option<String>,
17 #[serde(default)]
18 pub scenarios: HashMap<String, Vec<ScenarioStep>>,
19 #[serde(default)]
20 pub ops: Vec<ParsedOp>,
21 /// Workload-level Polydat bindings declared via the top-level
22 /// `bindings:` block. These compile into the workload-root
23 /// kernel directly, separate from per-op bindings — so
24 /// declarations like `cursor row = range(0, 50)` are visible
25 /// to scenario-level comprehensions (e.g.,
26 /// `for xval in all(row)`) without needing to be threaded
27 /// through phase-level ops.
28 #[serde(default)]
29 pub bindings: BindingsDef,
30 /// Resolved workload parameters. These are available as bind points
31 /// in op templates and as constants in Polydat bindings.
32 /// Populated from: workload `params:` defaults, CLI overrides, env vars.
33 #[serde(default)]
34 pub params: HashMap<String, String>,
35 /// Phase definitions. Each phase has its own config and either
36 /// inline ops or tag filters to select from blocks/top-level ops.
37 #[serde(default)]
38 pub phases: HashMap<String, WorkloadPhase>,
39 /// Phase names in YAML definition order. HashMap does not preserve
40 /// insertion order, so this Vec tracks the order phases appeared in
41 /// the workload YAML for deterministic default scenario execution.
42 #[serde(default)]
43 pub phase_order: Vec<String>,
44 /// SRD-83 — workload-shell stop conditions. Declarations distribute
45 /// by their `each:` selector: `each: phase` applies the predicate at
46 /// every phase, while `each: [self, workload]` evaluates it at the
47 /// workload shell itself (reading the `children_*` aggregate of
48 /// child phase outcomes).
49 #[serde(default)]
50 pub stop_when: Vec<StopConditionSpec>,
51 /// Param names declared in the workload YAML `params:` section.
52 /// Used to detect unrecognized CLI params. Does not include
53 /// ad-hoc CLI params.
54 #[serde(default)]
55 pub declared_params: Vec<String>,
56 /// Unified report block (SRD-46): plots and tables under one
57 /// schema with figure enumeration, palette/style cascade, and
58 /// declaration-order rendering. Replaces the separate
59 /// `plot:` and `summary:` blocks (gone, no shim).
60 #[serde(default)]
61 pub report: crate::report::Report,
62 /// Non-fatal warnings emitted by the report-block parser
63 /// (SRD-46). Empty in normal mode; strict mode (SRD-15)
64 /// promotes them to errors. Plumbed up so the runner /
65 /// validator decide how to surface them.
66 #[serde(default, skip_serializing)]
67 pub report_warnings: Vec<String>,
68 /// Non-fatal reference-resolution warnings from the
69 /// `extends:` chain (SRD-85 nearest-first): a target name
70 /// that matched multiple resources resolved to the nearest,
71 /// and the shadowing is surfaced here — never silently.
72 /// Logged by the runner; strict mode promotes to errors.
73 #[serde(default, skip_serializing)]
74 pub resolution_warnings: Vec<String>,
75 /// Fatal scenario-parse errors collected during
76 /// `parse_scenario_nodes` — typically "unknown scenario-
77 /// node key" cases the parser used to silently drop. Per
78 /// the project's "Never Ignore Silently" rule (memory),
79 /// the runner promotes these to hard errors before
80 /// dispatching the workload. Unlike `report_warnings`,
81 /// these are always-fatal regardless of strict mode —
82 /// a malformed scenario-tree node never produces useful
83 /// behavior, so a downstream `phase 'iterate' not found`
84 /// error masks the real bug.
85 #[serde(default, skip_serializing)]
86 pub scenario_parse_errors: Vec<String>,
87 /// Workload-wide default for the per-phase
88 /// [`WorkloadPhase::status_metrics`] field. Phases that don't
89 /// declare their own `status_metrics:` inherit this list.
90 /// Supports glob-style patterns (`recall*`, `latency*`) so a
91 /// single doc-root entry can emphasize a metric family across
92 /// every phase that produces it.
93 ///
94 /// Empty (default) → no metrics tail anywhere; per-phase
95 /// declarations are still honoured.
96 #[serde(default)]
97 pub status_metrics: Vec<String>,
98 /// Resolved `readouts:` block bindings (SRD-63 §5).
99 /// One entry per event slot the workload bound; the
100 /// runtime binder reads this map and dispatches at fire
101 /// time. Empty (default) → all slots fall back to the
102 /// hard-coded built-ins activity.rs uses today.
103 ///
104 /// Each value is a list of literal body strings — one
105 /// per readout invocation in the slot. The body strings
106 /// haven't been parsed against the readout grammar yet;
107 /// that happens at activity-init time once the workload
108 /// kernel is in place. Push 3 ships the data shape
109 /// only; Push 4 wires resolved layered overrides
110 /// (CLI / extends).
111 #[serde(default)]
112 pub readouts: ReadoutsBindings,
113 /// SRD-32a Push 3 — workload-root wrapper composition
114 /// override. When present, every op template in this
115 /// workload uses this innermost-to-outermost order
116 /// instead of the runtime's default tiebreaker order.
117 /// Per-op `wrappers: { order: ... }` shadows this entry
118 /// entirely (no cascading merge).
119 #[serde(default, skip_serializing_if = "Option::is_none")]
120 pub wrappers: Option<WrappersConfig>,
121 /// SRD-108 Part B — names the blueprint this document
122 /// IMPLEMENTS (resolved local-first, then bundled catalog,
123 /// like `extends:` targets). A document carrying this is an
124 /// implementation module: it provides op bodies for the
125 /// blueprint's abstract slots and must carry no phase
126 /// scaffolding of its own.
127 #[serde(default, skip_serializing_if = "Option::is_none")]
128 pub implements: Option<String>,
129 /// SRD-106 Part 3 — `stick_session: true` declares this
130 /// workload's intended usage as iterative re-attachment:
131 /// when the operator passes no explicit session selection
132 /// and `sessions/latest` exists, the run re-attaches to it
133 /// and layers a new execution per SRD-77, announcing the
134 /// re-attachment as the run's first notable event. CLI
135 /// `stick_session=true|false` overrides; `--session new`
136 /// forces a fresh session. Absent → today's fresh-session
137 /// behavior.
138 #[serde(default, skip_serializing_if = "Option::is_none")]
139 pub stick_session: Option<bool>,
140}
141
142impl Workload {
143 /// Unification (2026-05-27): the scenario-tree executor is
144 /// the sole execution path. Workloads that pre-date the
145 /// `phases:`/`scenarios:` shape — `op=...` inline CLI,
146 /// `blocks:` YAML, top-level `ops:` lists — would
147 /// historically run via a separate single-activity branch
148 /// in the runner that bypassed `run_phase`. After
149 /// unification that branch is gone; this method
150 /// synthesizes an implicit `main` phase + `default`
151 /// scenario so the runner has the phased shape to walk.
152 ///
153 /// **Idempotent**: when `phases` is already populated the
154 /// method returns without changes. The synthesized phase
155 /// owns the ops; `Workload::ops` stays populated because
156 /// downstream compile-time inspectors (the workload-root
157 /// kernel, wrapper-cascade resolver, bind-point
158 /// validator) still walk the top-level list.
159 ///
160 /// **Synthesized shape**: phase name `main`, scenario name
161 /// `default` containing `[Phase("main")]`. No `cycles` /
162 /// `concurrency` / `rate` / `errors` overrides — those
163 /// inherit from CLI / workload defaults via the existing
164 /// phased resolution.
165 pub fn synthesize_default_phase(&mut self) {
166 if !self.phases.is_empty() {
167 return;
168 }
169 if self.ops.is_empty() {
170 return;
171 }
172 const SYNTHETIC: &str = "main";
173 // Promote CLI-style `cycles=N` / `concurrency=N` /
174 // `rate=N` from `self.params` onto the synthetic
175 // phase. Pre-unification, the now-deleted single-
176 // activity branch read these directly off CLI params
177 // and set them on `ActivityConfig`; the phased branch
178 // reads them off `WorkloadPhase`. Forwarding here
179 // preserves the legacy contract — `nmbrs run op=...
180 // cycles=20` still runs 20 cycles after unification.
181 //
182 // The `==ops:N` wrap on cycles tells the per-phase
183 // resolver (executor.rs's `phase_cycles` block) to
184 // treat the number as a raw op-iteration count
185 // instead of the standard "N stanzas" multiplication.
186 // The legacy single-activity branch always used the
187 // op-count interpretation; without this, a 2-op
188 // stanza with `cycles=4` would run 8 ops instead of
189 // the historical 4.
190 let cycles = self.params.get("cycles").map(|c| format!("==ops:{c}"));
191 let concurrency = self.params.get("concurrency").cloned();
192 let rate = self.params.get("rate").cloned();
193 // Move GK-syntax workload-root bindings DOWN onto the
194 // synthetic phase. The legacy single-activity branch
195 // compiled workload-root bindings into the SAME kernel
196 // as the op templates, so destructure-target names
197 // (`(device, reading) := ...`) were locally visible.
198 // The phased path puts workload-root bindings on a
199 // separate kernel and exposes only the manifest names
200 // to child kernels — but the manifest lists the
201 // destructure tuple as a single entry, not the
202 // individual targets. Putting the bindings on the
203 // phase kernel preserves the legacy locality.
204 //
205 // The legacy `Map` form (`user_id: Hash(); Mod(...)`)
206 // gets a translation pass at workload-root that the
207 // phase-level parser doesn't apply — so we leave that
208 // form alone. Only `PolydatSource` (native Polydat string form)
209 // moves down. This split matches the two-form parser
210 // contract and avoids re-implementing translation.
211 let bindings = match &self.bindings {
212 BindingsDef::PolydatSource(_) => std::mem::take(&mut self.bindings),
213 BindingsDef::Map(_) => BindingsDef::default(),
214 };
215 let phase = WorkloadPhase {
216 ops: self.ops.clone(),
217 bindings,
218 cycles,
219 concurrency,
220 rate,
221 ..Default::default()
222 };
223 self.phases.insert(SYNTHETIC.to_string(), phase);
224 self.phase_order.push(SYNTHETIC.to_string());
225 // Only seed the default scenario when none was
226 // declared — an operator-authored `scenarios:` block
227 // (even with no phases yet) is honoured as-is.
228 if self.scenarios.is_empty() {
229 self.scenarios.insert(
230 "default".to_string(),
231 vec![ScenarioStep::Phase(SYNTHETIC.to_string())],
232 );
233 }
234 }
235}
236
237/// SRD-32a Push 3 — wrapper-composition override block.
238/// Carries an explicit innermost-to-outermost order list
239/// that the resolver uses in place of its built-in
240/// default-order tiebreaker. The list must be a permutation
241/// of the wrappers the op actually triggers (after
242/// transitive activation); listing a non-triggered wrapper
243/// or omitting a triggered one is a hard error per SRD-32a
244/// §"Workload-level override".
245#[derive(Debug, Clone, Default, Serialize, Deserialize)]
246pub struct WrappersConfig {
247 /// Innermost-to-outermost wrapper-name list. Empty list
248 /// is treated as "no override" (equivalent to leaving
249 /// `wrappers:` off the workload).
250 #[serde(default, skip_serializing_if = "Vec::is_empty")]
251 pub order: Vec<String>,
252}
253
254/// Per-event-slot list of readout body strings declared in
255/// the workload's `readouts:` block. See SRD-63 §5.0 for
256/// the three legal forms.
257///
258/// The lower-case slot keys here mirror the
259/// [`Event::slot_name`] return values
260/// (`on_phase_end`, `on_update`, …) so workload yaml
261/// uses the same vocabulary the design doc uses.
262#[derive(Debug, Clone, Default, Serialize, Deserialize)]
263pub struct ReadoutsBindings {
264 #[serde(default, skip_serializing_if = "Vec::is_empty")]
265 pub on_session_start: Vec<String>,
266 #[serde(default, skip_serializing_if = "Vec::is_empty")]
267 pub on_session_end: Vec<String>,
268 #[serde(default, skip_serializing_if = "Vec::is_empty")]
269 pub on_phase_start: Vec<String>,
270 #[serde(default, skip_serializing_if = "Vec::is_empty")]
271 pub on_phase_end: Vec<String>,
272 #[serde(default, skip_serializing_if = "Vec::is_empty")]
273 pub on_each_start: Vec<String>,
274 #[serde(default, skip_serializing_if = "Vec::is_empty")]
275 pub on_each_end: Vec<String>,
276 #[serde(default, skip_serializing_if = "Vec::is_empty")]
277 pub on_scope_start: Vec<String>,
278 #[serde(default, skip_serializing_if = "Vec::is_empty")]
279 pub on_scope_end: Vec<String>,
280 #[serde(default, skip_serializing_if = "Vec::is_empty")]
281 pub on_update: Vec<String>,
282}
283
284impl ReadoutsBindings {
285 /// True when no slot has any binding. Workloads in
286 /// this state fall through to the built-in defaults
287 /// activity.rs uses today.
288 pub fn is_empty(&self) -> bool {
289 self.on_session_start.is_empty()
290 && self.on_session_end.is_empty()
291 && self.on_phase_start.is_empty()
292 && self.on_phase_end.is_empty()
293 && self.on_each_start.is_empty()
294 && self.on_each_end.is_empty()
295 && self.on_scope_start.is_empty()
296 && self.on_scope_end.is_empty()
297 && self.on_update.is_empty()
298 }
299
300 /// Look up a slot's body list by its `slot_name` (e.g.
301 /// `"on_update"`). Returns an empty slice when the
302 /// slot has no bindings.
303 pub fn get(&self, slot_name: &str) -> &[String] {
304 match slot_name {
305 "on_session_start" => &self.on_session_start,
306 "on_session_end" => &self.on_session_end,
307 "on_phase_start" => &self.on_phase_start,
308 "on_phase_end" => &self.on_phase_end,
309 "on_each_start" => &self.on_each_start,
310 "on_each_end" => &self.on_each_end,
311 "on_scope_start" => &self.on_scope_start,
312 "on_scope_end" => &self.on_scope_end,
313 "on_update" => &self.on_update,
314 _ => &[],
315 }
316 }
317}
318
319/// Parsed summary report configuration.
320///
321/// Controls which columns, rows, and aggregates appear in the
322/// post-run summary table. Parsed from a semicolon-delimited DSL:
323///
324/// ```text
325/// "recall; mean(recall) over profile~label; details=hide"
326/// ```
327///
328/// Directives:
329/// - Bare words (no `=` or `(`): gauge column filter patterns, comma-separated.
330/// `"all"` shows every discovered gauge.
331/// - `filter=<regex>`: row filter on activity labels.
332/// - `<func>(<col>) over <key>~<pat>`: aggregate expression.
333/// - `details=hide`: suppress individual data rows.
334#[derive(Debug, Clone)]
335pub struct SummaryConfig {
336 /// Gauge column filter patterns (e.g., `["recall", "precision"]`).
337 /// Empty means show all discovered gauges.
338 pub columns: Vec<String>,
339 /// Row filter regex patterns on activity labels.
340 pub row_filters: Vec<String>,
341 /// Aggregate expressions to compute after the data rows.
342 pub aggregates: Vec<AggregateExpr>,
343 /// Whether to show individual data rows (default `true`).
344 pub show_details: bool,
345 /// Raw source string for diagnostics and future Polydat template detection.
346 pub raw: String,
347 /// SRD-46 v2: native MetricsQL columns. When non-empty,
348 /// `summary_command` routes through the metricsql renderer
349 /// instead of the legacy SQL builder. Each entry is
350 /// `(column_name, metricsql_expression)`. Anonymous
351 /// single-column form (`query: <expr>`) lands as
352 /// `("value", expr)`.
353 pub metricsql_columns: Vec<(String, String)>,
354 /// Label key the metricsql results are grouped on (becomes
355 /// the leftmost column of the rendered table). When empty
356 /// AND `metricsql_columns` is non-empty, the renderer falls
357 /// back to a single un-grouped row showing the average value
358 /// across all returned series.
359 ///
360 /// Multi-key form: `group_by: k, r, optimize_for` produces
361 /// one table row per distinct tuple — the same series
362 /// breakdown the matching plot draws.
363 pub group_by: Vec<String>,
364 /// SRD-46 — `state: <expr>`: a per-row COMPLETION test rendered as a word
365 /// rather than a number. A row whose expression yields a value is `complete`;
366 /// a row with none is `active`.
367 ///
368 /// The expression to use is whatever the workload records only at completion,
369 /// so "has this finished" is answered by the presence of that datum rather
370 /// than inferred from a progress percentage — a progress gauge can sit below
371 /// 100 on a row that finished, because the last poll before completion is the
372 /// value that persists.
373 pub state_query: Option<String>,
374 /// Per-column header annotations (`header <col>: <text>`) —
375 /// rendered into the column's header stack under its name, so a
376 /// table can carry each column's DEFINITION (e.g. the SRD-113
377 /// designation `last(result_success)`) where the reader is
378 /// already looking.
379 pub header_notes: Vec<(String, String)>,
380}
381
382/// An aggregate expression: either
383/// `mean(recall) over profile~label` (single-key filter form,
384/// emits one aggregate row) or
385/// `mean(recall) over k,limit,optimize_for` (multi-key grouping
386/// form, emits one aggregate row per distinct value-tuple).
387#[derive(Debug, Clone)]
388pub struct AggregateExpr {
389 /// Aggregation function.
390 pub function: AggFunction,
391 /// Column name pattern — only gauge columns containing this string
392 /// are aggregated; others show `-` in the aggregate row.
393 pub column_pattern: String,
394 /// Label key to filter rows on (e.g., `"profile"`).
395 /// Set in the single-key filter form. Empty when
396 /// `group_by` is non-empty (multi-key grouping form).
397 pub label_key: String,
398 /// Substring pattern matched against the label value (e.g., `"label"`).
399 /// Set in the single-key filter form. Empty when
400 /// `group_by` is non-empty.
401 pub label_pattern: String,
402 /// Multi-key grouping form: when non-empty, rows are
403 /// grouped by every distinct tuple of values across these
404 /// label keys, and the aggregate emits one row per group.
405 /// `label_key` / `label_pattern` are empty when this is set.
406 pub group_by: Vec<String>,
407}
408
409/// Supported aggregation functions for summary report expressions.
410#[derive(Debug, Clone, Copy, PartialEq, Eq)]
411pub enum AggFunction {
412 Mean,
413 Min,
414 Max,
415}
416
417impl std::fmt::Display for AggFunction {
418 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
419 match self {
420 AggFunction::Mean => write!(f, "mean"),
421 AggFunction::Min => write!(f, "min"),
422 AggFunction::Max => write!(f, "max"),
423 }
424 }
425}
426
427impl Serialize for SummaryConfig {
428 fn serialize<S: serde::Serializer>(&self, serializer: S) -> Result<S::Ok, S::Error> {
429 serializer.serialize_str(&self.raw)
430 }
431}
432
433impl<'de> Deserialize<'de> for SummaryConfig {
434 fn deserialize<D: serde::Deserializer<'de>>(deserializer: D) -> Result<Self, D::Error> {
435 let raw = String::deserialize(deserializer)?;
436 Ok(SummaryConfig::parse(&raw))
437 }
438}
439
440impl SummaryConfig {
441 /// Parse a short-form summary DSL string.
442 ///
443 /// Semicolon-separated directives:
444 /// - `"recall,precision"` — column filters
445 /// - `"filter=search_post"` — row filter
446 /// - `"mean(recall) over profile~label"` — aggregate expression
447 /// - `"details=hide"` — hide individual data rows
448 pub fn parse(raw: &str) -> Self {
449 let mut columns = Vec::new();
450 let mut row_filters = Vec::new();
451 let mut aggregates = Vec::new();
452 let mut show_details = true;
453 let mut metricsql_columns: Vec<(String, String)> = Vec::new();
454 let mut group_by: Vec<String> = Vec::new();
455 let mut state_query: Option<String> = None;
456 let mut header_notes: Vec<(String, String)> = Vec::new();
457
458 // Strip `#` line comments before parsing (SRD-46:
459 // report/plot/table bodies all support `#` comments).
460 let cleaned = strip_hash_line_comments(raw);
461
462 // SRD-46 v2 line-pass: native-form directives
463 // (`query: <expr>`, `query <col>: <expr>`, `group_by: <key>`).
464 // Pulled out before the legacy `;`-separator pass so a
465 // metricsql expression containing `;` (rare but legal)
466 // doesn't get sliced apart, and so legacy and native
467 // forms can coexist during migration.
468 let mut residual_lines: Vec<String> = Vec::new();
469 for line in cleaned.lines().map(str::trim).filter(|s| !s.is_empty()) {
470 if let Some(rest) = line
471 .strip_prefix("group_by:")
472 .map(str::trim)
473 .or_else(|| line.strip_prefix("group-by:").map(str::trim))
474 {
475 group_by = rest
476 .split(',')
477 .map(|s| s.trim().to_string())
478 .filter(|s| !s.is_empty())
479 .collect();
480 continue;
481 }
482 // Three surface forms for query columns:
483 // query <col>: <expr> — legacy named (space sep)
484 // query: <col>: <expr> — canonical named (uniform `name: value`)
485 // query: <expr> — single anonymous column
486 // The canonical form factors out as: after `query:`, if the
487 // remainder begins with a bare-identifier followed by `:`,
488 // the leading identifier is the column name; otherwise the
489 // whole remainder is the anonymous expression. Identifiers
490 // here are `[A-Za-z_][A-Za-z0-9_-]*` — anything containing
491 // whitespace, parens, braces, or operators forces the
492 // anonymous interpretation, which is what we want for
493 // metricsql expressions whose label-literal `:` shows up
494 // before a function-call `(`.
495 // `header <col>: <text>` — a column's header annotation.
496 if let Some(rest) = line.strip_prefix("header ") {
497 if let Some((col, note)) = rest.split_once(':') {
498 let (col, note) = (col.trim(), note.trim());
499 if !col.is_empty() && !note.is_empty() {
500 header_notes.push((col.to_string(), note.to_string()));
501 }
502 }
503 continue;
504 }
505 // `state: <expr>` — completion test, rendered as a word.
506 if let Some(rest) = line.strip_prefix("state:") {
507 let expr = rest.trim();
508 if !expr.is_empty() {
509 state_query = Some(expr.to_string());
510 }
511 continue;
512 }
513 if let Some(rest) = line.strip_prefix("query") {
514 let rest = rest.trim_start();
515 if let Some(after_colon) = rest.strip_prefix(':') {
516 let after_colon = after_colon.trim_start();
517 if let Some((col, expr)) = split_named_query(after_colon) {
518 metricsql_columns.push((col, expr));
519 } else {
520 metricsql_columns
521 .push(("value".to_string(), after_colon.trim().to_string()));
522 }
523 continue;
524 }
525 // Legacy `query <col>: <expr>` form — the next
526 // colon terminates the column name.
527 if let Some(colon_idx) = rest.find(':') {
528 let col = rest[..colon_idx].trim().to_string();
529 let expr = rest[colon_idx + 1..].trim().to_string();
530 if !col.is_empty() && !expr.is_empty() {
531 metricsql_columns.push((col, expr));
532 continue;
533 }
534 }
535 }
536 residual_lines.push(line.to_string());
537 }
538 let cleaned: String = residual_lines.join(";");
539
540 for directive in cleaned.split(';').map(str::trim).filter(|s| !s.is_empty()) {
541 if directive == "details=hide" {
542 show_details = false;
543 } else if let Some(filter) = directive.strip_prefix("filter=") {
544 row_filters.push(filter.trim().to_string());
545 } else if let Some(agg) = Self::parse_aggregate(directive) {
546 aggregates.push(agg);
547 } else {
548 // Column filter: comma-separated names. Two
549 // names are recognized as wildcards ("show
550 // every gauge column"): the legacy `all`
551 // keyword and `*` (the bare-`--summary` user
552 // mental model — `nmbrs --summary '*'` means
553 // "default summary of all metrics").
554 for col in directive
555 .split(',')
556 .map(str::trim)
557 .filter(|s| !s.is_empty())
558 {
559 if col != "all" && col != "*" {
560 columns.push(col.to_string());
561 }
562 // `all` / `*` = empty columns vec = show
563 // every gauge with no filtering.
564 }
565 }
566 }
567
568 SummaryConfig {
569 columns,
570 row_filters,
571 aggregates,
572 show_details,
573 raw: raw.to_string(),
574 metricsql_columns,
575 group_by,
576 state_query,
577 header_notes,
578 }
579 }
580
581 /// Try to parse an aggregate directive in either form:
582 /// - `<func>(<col>) over <key>~<pat>` (single-key filter)
583 /// - `<func>(<col>) over <k1>,<k2>,…` (multi-key grouping)
584 fn parse_aggregate(s: &str) -> Option<AggregateExpr> {
585 let paren_open = s.find('(')?;
586 let paren_close = s.find(')')?;
587 if paren_close <= paren_open {
588 return None;
589 }
590
591 let func_name = s[..paren_open].trim();
592 let function = match func_name {
593 "mean" => AggFunction::Mean,
594 "min" => AggFunction::Min,
595 "max" => AggFunction::Max,
596 _ => return None,
597 };
598
599 let column_pattern = s[paren_open + 1..paren_close].trim().to_string();
600
601 let after_paren = s[paren_close + 1..].trim();
602 let over_rest = after_paren.strip_prefix("over")?.trim();
603
604 // Single-key filter form: `<key>~<pat>` (note: `~` may
605 // appear inside multi-key form too, e.g. nobody writes
606 // `k,a~b` — use presence of `~` as the discriminator;
607 // for clean multi-key, no `~` is present).
608 if let Some(tilde) = over_rest.find('~') {
609 let label_key = over_rest[..tilde].trim().to_string();
610 let label_pattern = over_rest[tilde + 1..].trim().to_string();
611 return Some(AggregateExpr {
612 function,
613 column_pattern,
614 label_key,
615 label_pattern,
616 group_by: Vec::new(),
617 });
618 }
619
620 // Multi-key grouping form: comma-separated label keys.
621 let group_by: Vec<String> = over_rest
622 .split(',')
623 .map(str::trim)
624 .filter(|s| !s.is_empty())
625 .map(|s| s.to_string())
626 .collect();
627 if group_by.is_empty() {
628 return None;
629 }
630 Some(AggregateExpr {
631 function,
632 column_pattern,
633 label_key: String::new(),
634 label_pattern: String::new(),
635 group_by,
636 })
637 }
638}
639
640/// Strip `#` line comments from a multi-line spec body. A `#`
641/// Split a `query:` payload into `(column_name, expression)` when
642/// the payload's leading token is a bare identifier followed by
643/// `:`. Returns `None` for the anonymous-column form (the whole
644/// payload is the expression).
645///
646/// An identifier here is `[A-Za-z_][A-Za-z0-9_-]*`. The lookup
647/// fails as soon as any character outside that class appears
648/// before the first `:`, which is what guards a metricsql label
649/// expression like `recall_mean{k="10"}` from being mistaken for
650/// a column name (the `{` ends the candidate identifier before
651/// the eventual `:` inside `k="10"` is reached).
652fn split_named_query(text: &str) -> Option<(String, String)> {
653 let bytes = text.as_bytes();
654 let mut i = 0;
655 while i < bytes.len() {
656 let b = bytes[i];
657 let is_first = i == 0;
658 let ok = if is_first {
659 b.is_ascii_alphabetic() || b == b'_'
660 } else {
661 b.is_ascii_alphanumeric() || b == b'_' || b == b'-'
662 };
663 if !ok {
664 break;
665 }
666 i += 1;
667 }
668 if i == 0 {
669 return None;
670 }
671 // Optional whitespace, then `:` to separate name from value.
672 let mut j = i;
673 while j < bytes.len() && bytes[j].is_ascii_whitespace() {
674 j += 1;
675 }
676 if j >= bytes.len() || bytes[j] != b':' {
677 return None;
678 }
679 let name = text[..i].to_string();
680 let expr = text[j + 1..].trim().to_string();
681 if name.is_empty() || expr.is_empty() {
682 return None;
683 }
684 Some((name, expr))
685}
686
687/// starts a comment only when it's at line-start or preceded by
688/// whitespace — so hex colors (`#117733`) and JSON sub-blocks
689/// (`{"color": "#fff"}`) survive. Quoted strings are honoured.
690fn strip_hash_line_comments(s: &str) -> String {
691 let mut out = String::with_capacity(s.len());
692 for line in s.split_inclusive('\n') {
693 let mut quote: Option<char> = None;
694 let mut prev_ws = true;
695 let mut cut: Option<usize> = None;
696 for (i, ch) in line.char_indices() {
697 match quote {
698 Some(q) if ch == q => {
699 quote = None;
700 prev_ws = false;
701 }
702 Some(_) => {
703 prev_ws = false;
704 }
705 None => match ch {
706 '"' | '\'' => {
707 quote = Some(ch);
708 prev_ws = false;
709 }
710 '#' if prev_ws => {
711 cut = Some(i);
712 break;
713 }
714 c if c.is_whitespace() => {
715 prev_ws = true;
716 }
717 _ => {
718 prev_ws = false;
719 }
720 },
721 }
722 }
723 match cut {
724 Some(idx) => {
725 out.push_str(&line[..idx]);
726 if line.ends_with('\n') {
727 out.push('\n');
728 }
729 }
730 None => out.push_str(line),
731 }
732 }
733 out
734}
735
736#[cfg(test)]
737mod summary_config_tests {
738 use super::*;
739
740 #[test]
741 fn parses_multi_key_grouping() {
742 let cfg = SummaryConfig::parse("recall; mean(recall) over k,limit,optimize_for");
743 assert_eq!(cfg.aggregates.len(), 1, "got: {:?}", cfg.aggregates);
744 let agg = &cfg.aggregates[0];
745 assert_eq!(agg.group_by, vec!["k", "limit", "optimize_for"]);
746 assert!(agg.label_key.is_empty());
747 }
748
749 #[test]
750 fn parses_single_key_filter_form_unchanged() {
751 let cfg = SummaryConfig::parse("mean(recall) over profile~label");
752 assert_eq!(cfg.aggregates.len(), 1);
753 let agg = &cfg.aggregates[0];
754 assert!(agg.group_by.is_empty());
755 assert_eq!(agg.label_key, "profile");
756 assert_eq!(agg.label_pattern, "label");
757 }
758}
759
760/// SRD-83 — a scope-tree *level* a stop condition distributes to. The
761/// `each:` selector names one or more of these; the matter walk binds
762/// the predicate at every node of a named level inside the declaring
763/// subtree (a declared, structural fan-out — never inferred from the
764/// predicate's content). Aligned to the executor's `ScopeKind`.
765#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
766#[serde(rename_all = "snake_case")]
767pub enum ScopeLevel {
768 /// The declaring scope node itself (whatever its kind).
769 #[serde(rename = "self")]
770 SelfScope,
771 /// Every op-template node.
772 Op,
773 /// Every phase node.
774 Phase,
775 /// Every scenario node.
776 Scenario,
777 /// The workload root (the whole-run aggregate).
778 Workload,
779}
780
781fn default_each() -> Vec<ScopeLevel> {
782 // Absent `each:` → the declaring scope only. To fan a workload-level
783 // declaration out per-phase, the author writes `each: phase`.
784 vec![ScopeLevel::SelfScope]
785}
786
787/// Accept either a single level (`each: phase`) or a list
788/// (`each: [phase, scenario]`) — the scalar is sugar for a one-element
789/// set.
790fn de_each<'de, D>(d: D) -> Result<Vec<ScopeLevel>, D::Error>
791where
792 D: serde::Deserializer<'de>,
793{
794 #[derive(Deserialize)]
795 #[serde(untagged)]
796 enum OneOrMany {
797 One(ScopeLevel),
798 Many(Vec<ScopeLevel>),
799 }
800 Ok(match OneOrMany::deserialize(d)? {
801 OneOrMany::One(level) => vec![level],
802 OneOrMany::Many(levels) => levels,
803 })
804}
805
806/// SRD-83 follow-up — the FIRING axis (when a condition is evaluated),
807/// as a tagged-union *value* so the field name can't overclaim
808/// periodicity. `continuous` names the existing inline (per drain-loop
809/// turn) evaluation; `phase_end` names the phase-completion aggregation.
810/// A cadence value (`{every: <duration>}`) driving the metrics
811/// `CadenceReporter` registry is a later step and is intentionally NOT
812/// accepted here yet.
813#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
814#[serde(rename_all = "snake_case")]
815pub enum PulseSpec {
816 /// Evaluate inline, once per drain-loop turn (fine-grained; the only
817 /// scope with attempt-level wires).
818 Continuous,
819 /// Evaluate at phase-completion aggregation.
820 PhaseEnd,
821}
822
823/// SRD-83 — one stop condition declared on a shell. A polydat
824/// `condition:`/`when:` predicate over runtime-state wires (`op_count`,
825/// `error_rate`, `elapsed_ms`, `children_failed`, …); a `per:`/`each:`
826/// **detection** distribution selector (which scope levels it is
827/// evaluated at); a `pulse:` firing axis; an `action:`/`effect:`
828/// (`fail` → Interrupted+Failed, `stop` → Interrupted+Succeeded); and an
829/// `at:` **action** target scope (default = the innermost level of
830/// `per:`). When the predicate trips it stops the `at:` scope with the
831/// effect. Detection scope (`per:`) and action scope (`at:`) are
832/// independent.
833#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
834pub struct StopConditionSpec {
835 /// Polydat predicate over runtime-state wires, evaluating to bool.
836 /// Canonical key `condition:`; `when:` accepted as an alias.
837 #[serde(alias = "condition")]
838 pub when: String,
839 /// Scope levels this predicate is DETECTED/evaluated at (the declared
840 /// fan-out; see [`ScopeLevel`]). Canonical key `per:`; `each:` accepted
841 /// as an alias. Defaults to the declaring scope (`self`).
842 #[serde(default = "default_each", deserialize_with = "de_each", alias = "per")]
843 pub each: Vec<ScopeLevel>,
844 /// Firing trigger (legacy string form). `None` → a sensible default per
845 /// condition kind. Superseded by `pulse:`; retained for compatibility.
846 #[serde(default)]
847 pub trigger: Option<String>,
848 /// Firing pulse — WHEN the predicate is evaluated. `None` → default per
849 /// condition kind (attempt/rate wires → `continuous`; `children_*` →
850 /// `phase_end`).
851 #[serde(default)]
852 pub pulse: Option<PulseSpec>,
853 /// Effect on fire: `"fail"` or `"stop"`. Canonical key `action:`;
854 /// `effect:` accepted as an alias. `None` → `fail`.
855 #[serde(default, alias = "action")]
856 pub effect: Option<String>,
857 /// The ACTION target scope — where the effect lands, independent of
858 /// where it is detected (`per:`). `None` → the innermost (most
859 /// specific) level of `per:`, i.e. act in place (historical
860 /// behaviour). Set e.g. `at: workload` to route a phase-detected stop
861 /// out to the enclosing workload shell.
862 #[serde(default)]
863 pub at: Option<ScopeLevel>,
864}
865
866impl StopConditionSpec {
867 /// SRD-83 Part 5 — the closed `action:`/`effect:` verb vocabulary.
868 /// `stop` = Interrupted+Succeeded (a clean early halt that keeps the
869 /// partial result); `fail` = Interrupted+Failed; `abort` = `fail`
870 /// plus cancelling in-flight ops.
871 pub const EFFECT_VOCABULARY: [&'static str; 3] = ["stop", "fail", "abort"];
872
873 /// Validate the semantic surface serde cannot: the effect verb.
874 /// The runtime's verb→Outcome map resolves any unrecognized string
875 /// to the shell default, so before this check a typo'd
876 /// `effect: sotp` silently became `fail` — an authoring trap.
877 /// Rejected at workload load instead ("never ignore silently").
878 pub fn validate(&self) -> Result<(), String> {
879 match self.effect.as_deref() {
880 Some(e) if !Self::EFFECT_VOCABULARY.contains(&e) => Err(format!(
881 "unknown stop-condition effect '{e}' on `when: {}` — \
882 expected one of stop|fail|abort",
883 self.when
884 )),
885 _ => Ok(()),
886 }
887 }
888}
889
890fn default_continue_if_each() -> Vec<ScopeLevel> {
891 // SRD-101 — a `continue_if` gate defaults to halting the enclosing
892 // SCENARIO sweep (the comprehension loop it rides), NOT the declaring
893 // node only. This intentionally differs from `StopConditionSpec`'s
894 // `default_each` (`self`): the gate's whole purpose is to bound a sweep.
895 vec![ScopeLevel::Scenario]
896}
897
898/// SRD-101 — a `continue_if` pre-entry sweep gate declared on a
899/// comprehension-bearing scenario step or phase. A polydat `when:` predicate
900/// over the iteration's COORDINATE context (`end_of(p)`, `idx_of(p)`,
901/// outer-scope consts like `effective_max_size`) plus an `each:` scope level.
902/// The walker evaluates it per iteration BEFORE entering the body; while it is
903/// true the iteration runs, and the moment it is false the sweep at `each`
904/// ends gracefully (Interrupted+Succeeded — see SRD-101 §4). Aligns with
905/// [`StopConditionSpec`] (shares `when`/`each` and the `ScopedExpr` machinery),
906/// but is a pre-entry gate with continue polarity and a fixed graceful effect.
907#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
908pub struct ContinueIfSpec {
909 /// Polydat predicate over the coordinate context. The sweep continues
910 /// while it is true and halts (gracefully) the moment it is false.
911 pub when: String,
912 /// Scope level whose sweep ends on a false predicate. Defaults to
913 /// `scenario` (the enclosing comprehension); `workload` halts the run.
914 /// Canonical key `per:`; `each:` accepted as an alias. (For a
915 /// `continue_if` gate this level is both detection and action — the
916 /// full `per:`/`at:` split for gates is a later step.)
917 #[serde(
918 default = "default_continue_if_each",
919 deserialize_with = "de_each",
920 alias = "per"
921 )]
922 pub each: Vec<ScopeLevel>,
923}
924
925/// Retry-backoff settings carried by the map form of `tries:`
926/// (`tries: {count: N, backoff: {ratio, min, max}}`). Each field is
927/// optional — a missing key falls back to the op's standalone
928/// `retry_backoff*` param, then the built-in default. Durations are kept
929/// as raw strings (e.g. `"100ms"`, `"10s"`) and parsed at wrap time by the
930/// runtime, so this crate needs no time parser.
931#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Default)]
932pub struct BackoffSpec {
933 /// Geometric growth factor applied per retry: attempt `k`'s wait is
934 /// `min * ratio^(k-1)`, capped at `max`. `2.0` (the default) doubles;
935 /// `1.0` holds the wait constant at `min`. `None` = default.
936 #[serde(default)]
937 pub ratio: Option<f64>,
938 /// Backoff floor — the first retry's wait (a duration string).
939 /// `None` = default (`100ms`). `"0"` disables pacing entirely.
940 #[serde(default)]
941 pub min: Option<String>,
942 /// Backoff ceiling — the wait never exceeds this (a duration string).
943 /// `None` = default (`10s`).
944 #[serde(default)]
945 pub max: Option<String>,
946}
947
948/// `throttle:` — adaptive backpressure governor: boolean sugar
949/// (`throttle: true` = all defaults) or the full spec map.
950#[derive(Debug, Clone, serde::Serialize, serde::Deserialize, PartialEq)]
951#[serde(untagged)]
952pub enum ThrottleField {
953 /// `throttle: true` (defaults) / `throttle: false` (explicit off).
954 Enabled(bool),
955 /// Full spec map.
956 Spec(ThrottleSpec),
957}
958
959impl ThrottleField {
960 /// Normalize: `true` → the default spec, `false` → `None`.
961 pub fn to_spec(&self) -> Option<ThrottleSpec> {
962 match self {
963 ThrottleField::Enabled(true) => Some(ThrottleSpec::default()),
964 ThrottleField::Enabled(false) => None,
965 ThrottleField::Spec(s) => Some(s.clone()),
966 }
967 }
968}
969
970/// Adaptive backpressure governor parameters (SRD-83 §throttle).
971/// The governor keeps the WINDOWED attempt-failure fraction — the
972/// see-through-retries saturation signal — under `high` by walking
973/// the named dynamic control down multiplicatively, and recovers it
974/// toward the authored ceiling while the window stays under `low`.
975#[derive(Debug, Clone, serde::Serialize, serde::Deserialize, PartialEq)]
976#[serde(deny_unknown_fields)]
977pub struct ThrottleSpec {
978 /// Windowed attempt-failure fraction that triggers a back-off.
979 #[serde(default = "throttle_default_high")]
980 pub high: f64,
981 /// Fraction below which the governor recovers toward the
982 /// authored ceiling. Default: `high / 5`.
983 #[serde(default)]
984 pub low: Option<f64>,
985 /// The dynamic control to walk: `concurrency` (default) or `rate`.
986 #[serde(default = "throttle_default_control")]
987 pub control: String,
988 /// Initial offered value (slow-start seed). Default: `floor` —
989 /// the governor assumes the most fragile target and PROVES
990 /// headroom by doubling through clean windows. Declare a higher
991 /// start only for targets known to be robust at phase entry.
992 #[serde(default)]
993 pub start: Option<f64>,
994 /// Never throttle below this value.
995 #[serde(default = "throttle_default_floor")]
996 pub floor: f64,
997 /// Evaluation window (duration string, e.g. "2s").
998 #[serde(default = "throttle_default_window")]
999 pub window: String,
1000}
1001
1002fn throttle_default_high() -> f64 {
1003 0.05
1004}
1005fn throttle_default_control() -> String {
1006 "concurrency".to_string()
1007}
1008fn throttle_default_floor() -> f64 {
1009 1.0
1010}
1011fn throttle_default_window() -> String {
1012 "2s".to_string()
1013}
1014
1015impl Default for ThrottleSpec {
1016 fn default() -> Self {
1017 Self {
1018 high: throttle_default_high(),
1019 low: None,
1020 control: throttle_default_control(),
1021 start: None,
1022 floor: throttle_default_floor(),
1023 window: throttle_default_window(),
1024 }
1025 }
1026}
1027
1028/// SRD-109 — the time-dimension aggregate a key-metric designation
1029/// carries. MANDATORY on every designation: there are no implied
1030/// aggregates, so `rows: result_success` (no qualifier) is a parse
1031/// error naming this vocabulary. Defined over the stored samples of
1032/// one instance within the row scope's activation window — which,
1033/// per the SRD-42 amendment, are last-write-wins point samples with
1034/// PromQL semantics.
1035#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
1036pub enum KeyAgg {
1037 Min,
1038 Max,
1039 Avg,
1040 Last,
1041 First,
1042 Median,
1043 Stddev,
1044 Sum,
1045 Count,
1046 /// Derived: increase over the activation span, per second.
1047 Rate,
1048 /// Derived: the activation's wall clock (family-less — `span()`).
1049 Span,
1050 /// Derived: last − first over the activation.
1051 Delta,
1052}
1053
1054impl KeyAgg {
1055 /// The suggestion list every qualification error carries.
1056 pub const VOCAB: &'static str = "min, max, avg, last, first, median, stddev, sum, count; \
1057 derived: rate(F), span(), delta(F)";
1058
1059 pub fn parse(name: &str) -> Option<Self> {
1060 Some(match name {
1061 "min" => Self::Min,
1062 "max" => Self::Max,
1063 "avg" => Self::Avg,
1064 "last" => Self::Last,
1065 "first" => Self::First,
1066 "median" => Self::Median,
1067 "stddev" => Self::Stddev,
1068 "sum" => Self::Sum,
1069 "count" => Self::Count,
1070 "rate" => Self::Rate,
1071 "span" => Self::Span,
1072 "delta" => Self::Delta,
1073 _ => return None,
1074 })
1075 }
1076}
1077
1078/// SRD-109 — one key-metric designation on an execution node:
1079/// `column: agg(family)`. Designating key metrics both names the
1080/// node's measurables and attaches the node to the table row of its
1081/// nearest enclosing anchor (or the spine). `family` is empty for
1082/// the family-less `span()`.
1083#[derive(Debug, Clone, Serialize, Deserialize)]
1084pub struct KeyMetric {
1085 pub column: String,
1086 pub agg: KeyAgg,
1087 pub family: String,
1088}
1089
1090/// A workload phase: runs as a separate Activity with its own
1091/// cycle count, concurrency, rate limit, and op selection.
1092#[derive(Debug, Clone, Serialize, Deserialize, Default)]
1093pub struct WorkloadPhase {
1094 /// Number of stanzas for this phase. Each stanza executes all
1095 /// ops in sequence once. String type to support Polydat constant
1096 /// references like `"{train_count, dimensions: Default::default(),}"`. Default 1 (one stanza).
1097 #[serde(default)]
1098 pub cycles: Option<String>,
1099 /// Concurrency (async fibers). String type to support Polydat constant
1100 /// or workload param references like `"{concurrency}"`. Default 1.
1101 #[serde(default)]
1102 pub concurrency: Option<String>,
1103 /// Rate limit (ops/sec). A number, or a `{param}` / iter-var
1104 /// reference resolved at the phase gather (the SRD-83
1105 /// `timeout:` discipline: a rate that cannot be resolved
1106 /// fails the phase up front — it never silently becomes
1107 /// "unrated"). Default unlimited.
1108 #[serde(default)]
1109 pub rate: Option<String>,
1110 /// SRD-82 Part 6 — daemon phase. When `true`, this phase runs
1111 /// CONCURRENTLY with its foreground sibling phases (off the
1112 /// scenario's foreground concurrency budget) and is stopped
1113 /// cooperatively when the scope's foreground phases complete (a
1114 /// background "daemon unit" at the phase shell). Its own scope gives
1115 /// it an independent cursor base. Pair with an open-extent cursor
1116 /// (e.g. `until_elapsed`) so it runs for the foreground's duration.
1117 #[serde(default)]
1118 pub daemon: bool,
1119 /// Adapter override for this phase.
1120 #[serde(default)]
1121 pub adapter: Option<String>,
1122 /// Error routing spec override.
1123 #[serde(default)]
1124 pub errors: Option<String>,
1125 /// Total-attempts budget for this phase's ops. `tries:` is the SIGIL for
1126 /// the conditional tries wrapper (SRD-82 Part 3b): when NO budget
1127 /// resolves anywhere in scope (op field, this phase field, the
1128 /// workload-root `tries` param, or an in-scope `tries` wire), the
1129 /// wrapper is not constructed and the op runs single-attempt. `1` = the
1130 /// same single-attempt behaviour, explicitly (shadows an inherited
1131 /// budget); `0` = ops FAIL WITHOUT EXECUTING; `N ≥ 2` = up to N total
1132 /// attempts on adapter-retryable errors (CQL timeouts/overloads).
1133 /// `None` = inherit.
1134 #[serde(default)]
1135 pub tries: Option<u32>,
1136 /// Retry-backoff overrides parsed from the map form of `tries:`
1137 /// (`tries: {count: N, backoff: {ratio, min, max}}`). `None` when the
1138 /// sugared numeric form was used (or `tries` absent) — the wrapper then
1139 /// falls back to the op's standalone `retry_backoff*` params or the
1140 /// built-in defaults (ratio 2.0, min 100ms, max 10s). See
1141 /// [`BackoffSpec`].
1142 #[serde(default)]
1143 pub tries_backoff: Option<BackoffSpec>,
1144 /// SRD-82/92 cross-level wrapper (scoping P0) — phase-execution pacing.
1145 /// `interval:` is the discriminator for a future `WrapperLevel::Phase`
1146 /// interval wrapper: re-run this phase, sleeping `interval` between runs.
1147 /// A raw duration string (e.g. `"5m"`), parsed at wrap time by the
1148 /// runtime. Declarative today — the phase-level cascade that consumes it
1149 /// is not built yet (see `docs/cross-level-wrapper-cascade-scope.md`).
1150 #[serde(default)]
1151 pub interval: Option<String>,
1152 /// Bound for [`interval`](Self::interval) — how many times to run the
1153 /// phase. `None` alongside `interval` = repeat until the session stops.
1154 #[serde(default)]
1155 pub repeat: Option<u64>,
1156 /// OPT-IN error-rate circuit breaker (e.g. `0.1` = fail this phase
1157 /// once >10% of its ops error, after a 50-op floor). Overrides a
1158 /// session-wide `error_rate_max=` param when one was set. There is
1159 /// NO built-in default (SRD-82 §"AggregateGuard retired as a
1160 /// default") — aggregate health belongs to visible `stop_when:`
1161 /// conditions; this field exists as an explicit shorthand only.
1162 #[serde(default)]
1163 pub error_rate_max: Option<f64>,
1164 /// SRD-83 governance timeout (GAP-12). A duration (`"2.5h"`, `"150ms"`,
1165 /// bare fractional seconds) or a `{param}` reference; on expiry the
1166 /// phase ends Interrupted+Failed with reason class `timeout` — the
1167 /// protocol OUT-OF-RANGE disposition: the system is disqualified at
1168 /// this tier, the partial result is not usable. Desugars at the
1169 /// phase gather into a synthesized, logged `elapsed_ms >` stop
1170 /// condition (the `error_rate_max` precedent). Distinct from a
1171 /// BUDGET: a clean time-boxed measurement is a bounded cursor or a
1172 /// `stop_when … effect: stop` (Interrupted+Succeeded), not this.
1173 #[serde(default)]
1174 pub timeout: Option<String>,
1175 /// SRD-83 — stop conditions for this phase shell. Each is a polydat
1176 /// predicate over runtime-state wires, plus a firing trigger and an
1177 /// effect. Evaluated at triggers; a true predicate stops the phase
1178 /// with its effect.
1179 #[serde(default)]
1180 pub stop_when: Vec<StopConditionSpec>,
1181 /// SRD-83 §throttle — adaptive backpressure governor: keep the
1182 /// windowed attempt-failure fraction under a bound by walking a
1183 /// dynamic control (`concurrency`/`rate`) down under overload and
1184 /// back up on recovery. `throttle: true` = defaults.
1185 #[serde(default, skip_serializing_if = "Option::is_none")]
1186 pub throttle: Option<ThrottleField>,
1187 /// Tag filter to select ops from blocks (e.g., `"block:schema"`).
1188 #[serde(default)]
1189 pub tags: Option<String>,
1190 /// Inline ops for this phase (parsed into `ParsedOp` list).
1191 #[serde(default)]
1192 pub ops: Vec<ParsedOp>,
1193 /// Phase template iteration: `"var in expr"`.
1194 /// The phase is instantiated once per element of the Polydat expression
1195 /// result (which must be a comma-separated string). Each instance
1196 /// has `{var}` available as a workload param in its ops and config.
1197 ///
1198 /// Example: `for_each: "profile in matching_profiles('{dataset}', '{prefix}')"`
1199 #[serde(default)]
1200 pub for_each: Option<String>,
1201 /// SRD-101 — optional `continue_if` pre-entry gate bounding this phase's
1202 /// `for_each` sweep (see [`ContinueIfSpec`]). Ignored when `for_each` is
1203 /// absent (no sweep to bound).
1204 #[serde(default, skip_serializing_if = "Option::is_none")]
1205 pub continue_if: Option<ContinueIfSpec>,
1206 /// Loop scope mode for `for_each` phases.
1207 ///
1208 /// Controls how the loop context is seeded from the outer scope:
1209 /// - `clean` (default): snapshot of outer scope at loop entry
1210 /// - `inherit`: outer scope's live state (includes prior phase mutations)
1211 #[serde(default)]
1212 pub loop_scope: Option<String>,
1213 /// Iteration scope mode for `for_each` phases.
1214 ///
1215 /// Controls how each iteration is seeded from the loop scope:
1216 /// - `inherit` (default for for_each): each iteration starts from the loop
1217 /// scope's current state. All loop-level variables are implicitly shared
1218 /// with iterations, so iteration N+1 sees what N wrote.
1219 /// - `clean`: each iteration starts from the loop scope snapshot (isolated)
1220 #[serde(default)]
1221 pub iter_scope: Option<String>,
1222 /// Summary report configuration for this phase.
1223 /// Checkpoint declaration: skip-on-resume eligibility plus
1224 /// optional sub-properties (hashing, verify op). `None` =
1225 /// no declaration = phase always re-runs on resume. See
1226 /// SRD-44 §"Eligibility — `checkpoint:` per-phase declaration".
1227 ///
1228 /// Parsed via [`Checkpoint`]'s custom deserialize from the
1229 /// three YAML forms (short string, disabled string/bool,
1230 /// full mapping).
1231 #[serde(default)]
1232 pub checkpoint: Option<Checkpoint>,
1233 /// Names of metrics to surface on the inline progress line
1234 /// and the per-phase ✓ DONE summary. Empty (default) → no
1235 /// extra metrics shown; the status line carries only the
1236 /// universal counters (pct, throughput, ok-rate, errors,
1237 /// retries, concurrency, duration).
1238 ///
1239 /// Each name is matched against the live relevancy
1240 /// aggregates (`recall_at_10`, `precision_at_10`, …) by exact
1241 /// equality. Workloads that compute custom relevancy metrics
1242 /// list the names they want emphasized; nothing is presumed
1243 /// to be present.
1244 ///
1245 /// Example:
1246 /// ```yaml
1247 /// phases:
1248 /// ann_query:
1249 /// status_metrics: [recall_at_10]
1250 /// ```
1251 #[serde(default)]
1252 pub status_metrics: Vec<String>,
1253 /// Phase-level Polydat `bindings:` block (SRD-13c, SRD-13d).
1254 /// Captured on the phase AST so the scope-tree pre-walk
1255 /// (SRD-13d §3) can classify phase-level Polydat content via
1256 /// [`crate::polydat_matter::HasPolydatMatter`] and so the runtime
1257 /// can compose a phase kernel layered between the
1258 /// workload kernel and any op-template kernels.
1259 ///
1260 /// Today the parser ALSO merges this block into per-op
1261 /// bindings (legacy `parse.rs::parse_phases` behaviour) so
1262 /// the existing runtime keeps working unchanged. Once
1263 /// SRD-13d phases 3–9 land (per-template kernels with
1264 /// proper `bind_outer_scope` chaining through the phase
1265 /// kernel), the per-op merge is removed and ops resolve
1266 /// phase bindings via the Polydat scope chain.
1267 #[serde(default, skip_serializing_if = "BindingsDef::is_empty")]
1268 pub bindings: BindingsDef,
1269 /// Phase-level synthetic-metric declarations. Mirror of
1270 /// [`ParsedOp::metrics`] (same [`MetricSpec`] schema and YAML
1271 /// shapes), but evaluated **once at phase completion** against
1272 /// the phase scope kernel rather than per-cycle. Each entry's
1273 /// `value:` is a Polydat expression over phase-scope wires
1274 /// (bindings, captures, params, iter-vars) plus the
1275 /// executor-injected `phase_start` wire (epoch millis at phase
1276 /// start). The canonical phase-duration metric reads a clock via
1277 /// a `volatile` phase binding and subtracts the injected origin:
1278 /// ```yaml
1279 /// bindings: |
1280 /// volatile now_ms := current_epoch_millis()
1281 /// metrics:
1282 /// time_to_index: { value: now_ms - phase_start }
1283 /// ```
1284 /// yielding the phase's wall-clock duration in millis. Empty when
1285 /// absent. No dedicated clock node is needed: `current_epoch_millis()`
1286 /// is a single read, and `phase_start` arrives as plain data, so the
1287 /// expression re-evaluates correctly at the completion-time pull.
1288 /// Declaring the clock read as its own `volatile` binding (rather
1289 /// than nesting it in the metric value) explicitly acknowledges the
1290 /// non-deterministic node, so the phase kernel stays clean under
1291 /// `--strict`.
1292 ///
1293 /// The synthesiser emits `volatile __metric_<name> := <value>`
1294 /// onto the phase kernel (see
1295 /// `nmbrs_runtime::scope::synthesize_metric_binding_name`); the
1296 /// executor pulls each at completion and records it on the
1297 /// phase component as the declared instrument (gauge by default).
1298 #[serde(default, skip_serializing_if = "HashMap::is_empty")]
1299 pub metrics: HashMap<String, MetricSpec>,
1300 /// Dimensions this phase introduces: label NAME → declaration.
1301 ///
1302 /// Declared at the tier that owns the name, per the component tree's
1303 /// label-ownership rule (a name is set on exactly one tier and
1304 /// inherited downward). Values are not enumerated — they arrive from
1305 /// data via a metric's `cell:`.
1306 ///
1307 /// `BTreeMap` for deterministic synthesis order: a coordinate's
1308 /// rendering is what keys a cell, so an order that varied between runs
1309 /// would key one coordinate two ways.
1310 #[serde(default, skip_serializing_if = "std::collections::BTreeMap::is_empty")]
1311 pub dimensions: std::collections::BTreeMap<String, DimensionSpec>,
1312 /// Phase-level poll spec — when present, the phase's
1313 /// cycle execution runs in a wall-clock loop until a GK
1314 /// predicate over captures returns `true`. SRD-75.
1315 ///
1316 /// The presence of this field carries semantics beyond
1317 /// the data: it forbids `concurrency > 1` (serial-cycle
1318 /// loop is the unit of work), and it triggers
1319 /// scope-synthesis to allocate `shared` cells on the
1320 /// phase scope for capture names referenced by the
1321 /// predicate / `if:` conditions / metric values so
1322 /// cross-op visibility happens through the canonical
1323 /// Polydat chain (no sidecar HashMap; see SRD-75
1324 /// §"Architectural shape").
1325 #[serde(default, skip_serializing_if = "Option::is_none")]
1326 pub poll: Option<PhasePollSpec>,
1327 /// SRD-86 — when present, the executor dispatches the named optimizer over
1328 /// the phase: it writes each axis as an input wire on the phase's binding
1329 /// kernel and reads the `objective` wire back (the objective is *just a
1330 /// wire read*). Workload-local config; `nmbrs-runtime` maps it to its
1331 /// optimizer contract and discovers the optimizer via the link-time
1332 /// registry (`nmbrs describe optimizers`).
1333 ///
1334 /// **Sugar:** a bare **string** value is shorthand for `{ objective: <str> }`
1335 /// with every other field defaulted (`method: sweep`, no `servo:`) — so
1336 /// `optimize: "0 - err_rate"` ≡ `optimize: { objective: "0 - err_rate" }`.
1337 /// See [`de_optimize`].
1338 #[serde(
1339 default,
1340 deserialize_with = "de_optimize",
1341 skip_serializing_if = "Option::is_none"
1342 )]
1343 pub optimize: Option<OptimizeBlock>,
1344 /// SRD-109 — key-metric designations: `key_metrics: {column:
1345 /// "agg(family)", ...}`. Aggregate qualification is mandatory
1346 /// (no implied aggregates); the report synthesizer attaches
1347 /// these columns to the row of the phase's nearest enclosing
1348 /// anchor, or the spine. Empty = spine-only via the SRD-91
1349 /// instrument contract defaults.
1350 #[serde(default)]
1351 pub key_metrics: Vec<KeyMetric>,
1352}
1353
1354/// A phase `optimize:` value: **either** a bare string — sugar for
1355/// `{ objective: <string> }` with every other field defaulted — **or** a full
1356/// [`OptimizeBlock`] map. The untagged enum tries `Inline` first, so a scalar
1357/// value never reaches the map variant. Shared by the serde-derive path
1358/// ([`de_optimize`]) and the hand-rolled phase parser
1359/// (`parse::*` via [`OptimizeBlock::from_yaml_value`]).
1360#[derive(Deserialize)]
1361#[serde(untagged)]
1362enum OptimizeSpec {
1363 Inline(String),
1364 Block(OptimizeBlock),
1365}
1366
1367impl From<OptimizeSpec> for OptimizeBlock {
1368 fn from(spec: OptimizeSpec) -> Self {
1369 match spec {
1370 // SRD-86 string sugar: the whole value IS the objective expression.
1371 OptimizeSpec::Inline(objective) => OptimizeBlock {
1372 method: default_optimize_method(),
1373 objective,
1374 servo: Vec::new(),
1375 max_evals: default_optimize_max_evals(),
1376 seed: default_optimize_seed(),
1377 params: HashMap::new(),
1378 },
1379 OptimizeSpec::Block(b) => b,
1380 }
1381 }
1382}
1383
1384impl OptimizeBlock {
1385 /// Parse a phase `optimize:` value from already-parsed JSON, applying the
1386 /// string sugar (a bare string ≡ `{ objective: <string> }`). Used by the
1387 /// hand-rolled phase parser, which builds [`WorkloadPhase`] field-by-field
1388 /// rather than through the derive (so it doesn't see [`de_optimize`]).
1389 ///
1390 /// Branches explicitly rather than going through the untagged
1391 /// [`OptimizeSpec`] so a malformed *map* keeps its precise serde error
1392 /// (e.g. `missing field 'objective'`) instead of the untagged enum's
1393 /// generic "did not match any variant".
1394 pub fn from_yaml_value(v: &serde_json::Value) -> Result<OptimizeBlock, serde_json::Error> {
1395 if let Some(s) = v.as_str() {
1396 Ok(OptimizeSpec::Inline(s.to_string()).into())
1397 } else {
1398 serde_json::from_value::<OptimizeBlock>(v.clone())
1399 }
1400 }
1401}
1402
1403/// Deserialize a phase `optimize:` value via [`OptimizeSpec`] — a bare string is
1404/// sugar for `{ objective: <string> }`, a map is a full [`OptimizeBlock`]
1405/// (SRD-86):
1406///
1407/// ```yaml
1408/// optimize: |
1409/// 0 - metricsql_scalar("sum(rate(errors_total[3s]))")
1410/// ```
1411fn de_optimize<'de, D>(d: D) -> Result<Option<OptimizeBlock>, D::Error>
1412where
1413 D: serde::Deserializer<'de>,
1414{
1415 Ok(Option::<OptimizeSpec>::deserialize(d)?.map(Into::into))
1416}
1417
1418/// SRD-86 — a phase `optimize:` block. The optimizer **maximizes** the
1419/// `objective` wire by writing the `axes` input wires on the phase kernel.
1420#[derive(Debug, Clone, Serialize, Deserialize)]
1421pub struct OptimizeBlock {
1422 /// Registered optimizer name (`cmaes`, `nelder_mead`, … — see
1423 /// `nmbrs describe optimizers`). Defaults to `sweep` (the identity: evaluate
1424 /// every coordinate and report the best), so a plain "find the best by
1425 /// sweeping" search omits this field; set an adaptive method to search a
1426 /// large or continuous space without enumerating it.
1427 #[serde(default = "default_optimize_method")]
1428 pub method: String,
1429 /// The objective the optimizer **maximizes** (SRD-86 §10). Two forms:
1430 /// a **bare wire reference** — a single identifier naming a phase-kernel
1431 /// output (a `bindings:` entry), read directly; or an **inline polydat
1432 /// expression** (anything with operators / calls, e.g.
1433 /// `objective: "0 - metricsql_scalar(\"sum(rate(errors_total[3s]))\")"`),
1434 /// which is lowered to a synthesized `__objective` binding on the phase
1435 /// kernel (`scope::objective_wire`) so no separate `bindings:` entry is
1436 /// needed. An objective reading a windowed/live metric settles per setting;
1437 /// a deterministic one takes the one-shot read.
1438 pub objective: String,
1439 /// Search axes to actuate as **live controls** — servoed (retargeted without
1440 /// restarting the phase) rather than stepped through by re-running the phase
1441 /// (SRD-86 §4). Every axis is a coordinate (step-through / re-run) by default;
1442 /// naming one here opts it into servoing. Accepts a single name (`servo:
1443 /// concurrency`) or a list (`servo: [concurrency, rate]`). A servoed var
1444 /// resolves to a control either directly (its name IS a control — `servo:
1445 /// concurrency` / `servo: rate`) or indirectly (it feeds one via a `{var}`
1446 /// bind — `concurrency: "{conc}"`, then `servo: conc`). It is validated: it
1447 /// must resolve to a control AND the objective must be a windowed metric the
1448 /// servo can settle — else a clear error, never a silent downgrade.
1449 #[serde(default, deserialize_with = "de_string_or_seq")]
1450 pub servo: Vec<String>,
1451 #[serde(default = "default_optimize_max_evals")]
1452 pub max_evals: usize,
1453 #[serde(default = "default_optimize_seed")]
1454 pub seed: u64,
1455 /// Optimizer-specific knobs (e.g. `{ lambda: 8 }`).
1456 #[serde(default)]
1457 pub params: HashMap<String, f64>,
1458}
1459
1460/// Deserialize a single string OR a sequence of strings into a `Vec<String>` —
1461/// lets `servo: conc` and `servo: [conc, rate]` both parse.
1462fn de_string_or_seq<'de, D>(d: D) -> Result<Vec<String>, D::Error>
1463where
1464 D: serde::Deserializer<'de>,
1465{
1466 #[derive(Deserialize)]
1467 #[serde(untagged)]
1468 enum OneOrMany {
1469 One(String),
1470 Many(Vec<String>),
1471 }
1472 Ok(match OneOrMany::deserialize(d)? {
1473 OneOrMany::One(s) => vec![s],
1474 OneOrMany::Many(v) => v,
1475 })
1476}
1477
1478fn default_optimize_method() -> String {
1479 "sweep".to_string()
1480}
1481fn default_optimize_max_evals() -> usize {
1482 100
1483}
1484fn default_optimize_seed() -> u64 {
1485 1
1486}
1487
1488/// Phase-level `poll:` block (SRD-75). When set on a
1489/// `WorkloadPhase`, the runner wraps the phase's cycle
1490/// execution in a wall-clock loop that re-runs all ops
1491/// per iteration until `until` (a Polydat boolean expression
1492/// over captures) returns `true` or `timeout_ms` elapses.
1493///
1494/// Differs from the per-op `PollingDispenser` (SRD-32):
1495/// per-op poll wraps a SINGLE op with row-count /
1496/// json-path emptiness termination; phase-poll wraps
1497/// MULTIPLE ops with predicate-over-captures termination.
1498/// They coexist; per-op poll is the right tool when a
1499/// single op's response is sufficient to signal
1500/// completion.
1501#[derive(Debug, Clone, Serialize, Deserialize, Default)]
1502pub struct PhasePollSpec {
1503 /// Polydat boolean expression evaluated against the phase
1504 /// scope kernel after each iteration. Compiles into
1505 /// the phase scope as `__poll_until := <until>`;
1506 /// dynamic (re-evaluates per pull) per SRD-11's "two
1507 /// evaluation lifecycles" rule. Required.
1508 pub until: String,
1509 /// Sleep between iterations, milliseconds. A number or a
1510 /// `{param}` / iter-var reference resolved at the phase
1511 /// gather (unresolvable ⇒ the phase fails up front). Default
1512 /// `1000` (one second).
1513 #[serde(default)]
1514 pub interval_ms: Option<String>,
1515 /// Overall wall-clock cap, milliseconds. Same
1516 /// number-or-reference contract as `interval_ms`. The loop
1517 /// returns a `poll_timeout` error if exceeded. Default
1518 /// `300000` (5 minutes).
1519 #[serde(default)]
1520 pub timeout_ms: Option<String>,
1521 /// Consecutive retryable inner-op errors tolerated
1522 /// before propagation. Same number-or-reference contract
1523 /// as `interval_ms`. `0` (default) = strict: any
1524 /// retryable error fails the phase immediately.
1525 /// Mirrors per-op `PollingDispenser` semantics.
1526 #[serde(default)]
1527 pub max_error_retries: Option<String>,
1528 /// Named metric (gauge) written via `ctx.wires.write`
1529 /// when the loop terminates successfully. Value =
1530 /// elapsed wall-clock; unit decoded from the
1531 /// trailing suffix (`_s` / `_ms` / `_ns` / …) per
1532 /// the existing `duration_value_for_metric_name`
1533 /// convention. Same contract as per-op poll's
1534 /// `metric_name`. Default `None` (no metric).
1535 #[serde(default, skip_serializing_if = "Option::is_none")]
1536 pub metric_name: Option<String>,
1537 /// What to do when the wall-clock `timeout_ms`
1538 /// expires without satisfying `until`. SRD-75
1539 /// §"Workload-load validation" — the workload-author
1540 /// declares whether a stuck synchronizer is a
1541 /// recoverable phase-error or a workload-invalidating
1542 /// event.
1543 ///
1544 /// - `error` (default) — phase fails; the outer
1545 /// scenario's error-routing policy decides whether
1546 /// to continue to sibling phases. Suitable when a
1547 /// single cell's failed synchronization doesn't
1548 /// invalidate the rest of the sweep (rare).
1549 /// - `abort` — calls `session_signals::request_stop()`
1550 /// in addition to setting the phase's stop_flag.
1551 /// The scenario walker observes the global stop and
1552 /// terminates the whole run. Use when the
1553 /// predicate's satisfaction is a precondition for
1554 /// any downstream phase being meaningful — e.g.
1555 /// ensure_compacted in the CQL sweep: if the table
1556 /// never reaches `sstables == 1`, every subsequent
1557 /// query phase runs against an un-compacted table
1558 /// and produces meaningless results.
1559 #[serde(default, skip_serializing_if = "Option::is_none")]
1560 pub on_timeout: Option<String>,
1561 /// SRD-75 (C5) — strict-gate selectors. Each entry is a
1562 /// `metric()`-style selector (`"family, key=value, …"`) that MUST
1563 /// resolve to a registered instrument within the gate's grace
1564 /// window (the first poll interval); an unresolved selector is a
1565 /// hard `poll_require` error failing the phase. This is the
1566 /// runtime guard for coordination gates whose `until` reads
1567 /// another phase's live metrics — an unregistered family reads
1568 /// 0.0 silently, so a typo'd selector otherwise passes the gate
1569 /// instantly or hangs it to timeout. Deliberately poll-only:
1570 /// `stop_when`'s lenient reads stay as SRD-83 sanctions them
1571 /// ("family not yet present" is a legitimate not-yet state for a
1572 /// stop predicate; for a gate it is a bug).
1573 #[serde(default, skip_serializing_if = "Vec::is_empty")]
1574 pub require: Vec<String>,
1575}
1576
1577/// Per-phase checkpoint declaration. Three legal forms in YAML:
1578///
1579/// - `checkpoint: idempotent` — short form, equivalent to
1580/// `Checkpoint { idempotent: true, hashed: true, verify: None }`.
1581/// - `checkpoint: none` (or `false`, or `no`) — explicitly not
1582/// skip-eligible. Equivalent to no declaration; the phase
1583/// always re-runs on resume.
1584/// - `checkpoint: { idempotent: true, hashed: true, verify: ... }`
1585/// — full mapping form with sub-properties.
1586///
1587/// See SRD-44 §"Forms" and §"Sub-properties" for the full
1588/// contract. The `Default` is "skip-eligible with hashing on,
1589/// no verify" — what the short form `idempotent` produces.
1590#[derive(Debug, Clone, Serialize, PartialEq)]
1591pub struct Checkpoint {
1592 /// Marks this phase as skip-eligible on resume. `false`
1593 /// here is equivalent to `checkpoint: none` and means the
1594 /// phase always re-runs.
1595 pub idempotent: bool,
1596 /// When `true` (the default for any set checkpoint
1597 /// declaration), the resume planner additionally verifies
1598 /// that the freshly-pre-mapped phase's compiled program
1599 /// hash matches the saved one before honouring the saved
1600 /// status. `false` is the operator opt-out — "trust
1601 /// structural identity (yaml_path + coords) alone".
1602 pub hashed: bool,
1603 /// Optional verify op-template body. When present, the
1604 /// resume planner runs this op against the live system
1605 /// before classifying the phase as Skip; verify failure
1606 /// reclassifies to re-run with wholesale purge.
1607 /// Currently typed as a generic YAML value — the runtime
1608 /// re-parses it through the op-template grammar so the
1609 /// existing pipeline (SRD-32 wrappers, SRD-03 status-
1610 /// determination invariants) governs the verify
1611 /// execution.
1612 pub verify: Option<serde_json::Value>,
1613}
1614
1615impl Default for Checkpoint {
1616 fn default() -> Self {
1617 Self {
1618 idempotent: true,
1619 hashed: true,
1620 verify: None,
1621 }
1622 }
1623}
1624
1625impl<'de> serde::Deserialize<'de> for Checkpoint {
1626 fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
1627 where
1628 D: serde::Deserializer<'de>,
1629 {
1630 // The YAML accepts strings, booleans, and mappings —
1631 // each meaning a different declaration form. Serde's
1632 // visitor pattern lets us handle each input shape
1633 // directly without going through a typed-value
1634 // intermediate, which means this works equally well
1635 // for the YAML parser path and the JSON-staged path
1636 // (parse.rs walks `serde_json::Map` for phases).
1637 struct CheckpointVisitor;
1638 impl<'de> serde::de::Visitor<'de> for CheckpointVisitor {
1639 type Value = Checkpoint;
1640
1641 fn expecting(&self, f: &mut std::fmt::Formatter) -> std::fmt::Result {
1642 f.write_str("checkpoint declaration: short string ('idempotent' / 'none' / etc), bool, or mapping with sub-properties")
1643 }
1644
1645 fn visit_str<E: serde::de::Error>(self, s: &str) -> Result<Checkpoint, E> {
1646 let trimmed = s.trim().to_ascii_lowercase();
1647 match trimmed.as_str() {
1648 "idempotent" => Ok(Checkpoint::default()),
1649 "none" | "no" | "false" | "off" | "" => Ok(Checkpoint {
1650 idempotent: false,
1651 hashed: true,
1652 verify: None,
1653 }),
1654 other => Err(E::custom(format!(
1655 "checkpoint: unknown short form '{other}'; \
1656 expected 'idempotent', 'none', 'no', 'false', or a mapping"
1657 ))),
1658 }
1659 }
1660
1661 fn visit_string<E: serde::de::Error>(self, s: String) -> Result<Checkpoint, E> {
1662 self.visit_str(&s)
1663 }
1664
1665 fn visit_bool<E: serde::de::Error>(self, b: bool) -> Result<Checkpoint, E> {
1666 if b {
1667 Ok(Checkpoint::default())
1668 } else {
1669 Ok(Checkpoint {
1670 idempotent: false,
1671 hashed: true,
1672 verify: None,
1673 })
1674 }
1675 }
1676
1677 fn visit_unit<E: serde::de::Error>(self) -> Result<Checkpoint, E> {
1678 // Bare `null` ≡ `none`.
1679 Ok(Checkpoint {
1680 idempotent: false,
1681 hashed: true,
1682 verify: None,
1683 })
1684 }
1685
1686 fn visit_map<M>(self, mut map: M) -> Result<Checkpoint, M::Error>
1687 where
1688 M: serde::de::MapAccess<'de>,
1689 {
1690 let mut idempotent = true;
1691 let mut hashed = true;
1692 let mut verify: Option<serde_json::Value> = None;
1693 while let Some(key) = map.next_key::<String>()? {
1694 match key.as_str() {
1695 "idempotent" => idempotent = map.next_value::<bool>()?,
1696 "hashed" => hashed = map.next_value::<bool>()?,
1697 "verify" => verify = Some(map.next_value::<serde_json::Value>()?),
1698 other => {
1699 return Err(serde::de::Error::custom(format!(
1700 "checkpoint: unknown key '{other}'; \
1701 expected 'idempotent', 'hashed', or 'verify'"
1702 )));
1703 }
1704 }
1705 }
1706 Ok(Checkpoint {
1707 idempotent,
1708 hashed,
1709 verify,
1710 })
1711 }
1712 }
1713 deserializer.deserialize_any(CheckpointVisitor)
1714 }
1715}
1716
1717/// A node in a scenario execution tree.
1718///
1719/// Scenarios are trees of phases and control flow constructs.
1720/// Nesting is supported to arbitrary depth. All nodes are
1721/// evaluated dynamically at runtime — no pre-flattening.
1722///
1723/// `cycle` is immutable — loop constructs declare their own
1724/// counter variables for iteration indices.
1725///
1726/// All iteration shapes (`for_each` single-clause,
1727/// `for_combinations`, `for_each_union`) collapse into one
1728/// `Comprehension` variant carrying the canonical
1729/// [`polydat::iteration::comprehension::Comprehension`] AST —
1730/// the operator-tree form of the algebra layer. The
1731/// structural variant (`Cartesian` / `Union` / `Clause` /
1732/// `Zip`) is the discriminator; `Filter` and `Order` wrap
1733/// the body when the workload declares `where` / `order`.
1734/// See SRD-18b §"Iteration as a First-Class Concept" and
1735/// `polydat/docs/design/comprehension_cutover_contact_surfaces.md`.
1736#[derive(Debug, Clone, Serialize, Deserialize)]
1737pub enum ScenarioNode {
1738 /// A single phase to execute.
1739 Phase(String),
1740 /// Iteration node — single-clause for_each, multi-clause
1741 /// for_combinations, or for_each_union all map here. The
1742 /// `Comprehension` AST captures the iteration shape; the
1743 /// runtime executes the cross-product of clauses for the
1744 /// Cartesian mode and the concatenation of sub-spaces'
1745 /// products for Union mode.
1746 ///
1747 /// YAML forms (all normalize to this variant):
1748 /// ```yaml
1749 /// # Single clause
1750 /// - for_each: "k in 10,100"
1751 ///
1752 /// # Multi-clause cross product
1753 /// - for_each: "profile in profiles, k in {k_values}"
1754 ///
1755 /// # Multi-clause cross product (map form)
1756 /// - for_combinations:
1757 /// profile: "matching_profiles('{dataset}', '{prefix}')"
1758 /// k: "{k_values}"
1759 ///
1760 /// # Union of sub-spaces
1761 /// - for_each:
1762 /// - "k in 10, limit in 10,20,30"
1763 /// - "k in 100, limit in 100,200,300"
1764 /// ```
1765 Comprehension {
1766 comprehension: polydat::iteration::comprehension::Comprehension,
1767 children: Vec<ScenarioNode>,
1768 /// SRD-101 — optional `continue_if` pre-entry gate bounding this
1769 /// sweep: evaluated per iteration before the body; a false predicate
1770 /// halts the sweep at its `each` scope, gracefully.
1771 #[serde(default, skip_serializing_if = "Option::is_none")]
1772 continue_if: Option<ContinueIfSpec>,
1773 /// SRD-109 — table-row anchor: `anchor: <view>` declares one
1774 /// report-table row per iteration of this sweep, in the view
1775 /// of that name. Views sharing a name must share coordinate
1776 /// label sets. Key metrics designated on phases beneath this
1777 /// node attach to its rows.
1778 #[serde(default, skip_serializing_if = "Option::is_none")]
1779 anchor: Option<String>,
1780 },
1781 /// Execute children while condition is true (test after).
1782 DoWhile {
1783 condition: String,
1784 counter: Option<String>,
1785 children: Vec<ScenarioNode>,
1786 },
1787 /// Execute children until condition becomes true (test after).
1788 DoUntil {
1789 condition: String,
1790 counter: Option<String>,
1791 children: Vec<ScenarioNode>,
1792 },
1793 /// Logical inclusion of another scenario by name.
1794 ///
1795 /// Wherever this node appears (top-level of a scenario, inside
1796 /// a `phases:` list under a `for_each` / `for_combinations` /
1797 /// `for_each_union`, etc.), it expands to the children of the
1798 /// named scenario at execution time. The wrapper is preserved
1799 /// (not flattened) so the scope tree retains the include
1800 /// hierarchy and the renderer can show the operator which
1801 /// scenario each group of phases came from.
1802 ///
1803 /// Resolution happens once after parsing
1804 /// (see `crate::parse::resolve_scenario_includes`); cycles
1805 /// (`A` includes `B` includes `A`) are rejected with a clear
1806 /// error naming the cycle path.
1807 ///
1808 /// YAML form:
1809 /// ```yaml
1810 /// scenarios:
1811 /// smoke:
1812 /// - schema
1813 /// - rampup
1814 /// bench:
1815 /// - scenario: smoke
1816 /// - for_each: "k in 10,100"
1817 /// phases:
1818 /// - scenario: smoke
1819 /// - search
1820 /// ```
1821 IncludedScenario {
1822 name: String,
1823 children: Vec<ScenarioNode>,
1824 },
1825 /// Scenario-tree-level Polydat bindings block — the canonical way
1826 /// to introduce a scope-local layer of bound names anywhere
1827 /// in the scenario tree.
1828 ///
1829 /// `source` is Polydat matter text exactly as a phase-level
1830 /// `bindings:` block would contain. Anything the Polydat grammar
1831 /// accepts is valid: `const NAME := <literal>`, derived
1832 /// bindings (`scaled := mul(workload_limit, 2)`), shared
1833 /// cells, init bindings, etc. Workload-param `{name}` and
1834 /// string-interpolation references resolve through the
1835 /// scope chain at kernel build time — no separate
1836 /// preprocessing pass.
1837 ///
1838 /// `Bindings` is also the canonical lowered form of `set:`.
1839 /// The parser recognizes `set: { name: value, ... }` as
1840 /// syntactic sugar and emits a `Bindings` node whose
1841 /// `source` is `final <name> := <polydat-literal>\n` (one line
1842 /// per pair, declaration order preserved). So
1843 ///
1844 /// ```yaml
1845 /// - set: { mode: verbose }
1846 /// phases:
1847 /// - announce
1848 /// ```
1849 ///
1850 /// is semantically identical to
1851 ///
1852 /// ```yaml
1853 /// - bindings: |
1854 /// const mode := "verbose"
1855 /// phases:
1856 /// - announce
1857 /// ```
1858 ///
1859 /// Both produce one `Bindings` node. Authors keep the
1860 /// short `set:` form for the common override case; the
1861 /// long form unlocks the full Polydat grammar (derived
1862 /// bindings, expressions referencing other in-scope
1863 /// names, etc.) without any new variant.
1864 ///
1865 /// Lexical-shadow semantics are uniform with phase-level
1866 /// `bindings:`: a `const NAME := <value>` shadows any
1867 /// upstream binding for `NAME` over this node's `children`
1868 /// subtree. The shadow is enforced via the local-final
1869 /// transit-suppression rule in `materialize_wiring_from_outer`
1870 /// — the same mechanism every other scope uses.
1871 ///
1872 /// Composition example (two siblings, each defining its
1873 /// own value for the same name; the included subtree is
1874 /// physically cloned per include site so encapsulation is
1875 /// per-instance):
1876 ///
1877 /// ```yaml
1878 /// scenarios:
1879 /// fanout:
1880 /// - set: { mode: verbose }
1881 /// phases:
1882 /// - scenario: load_test
1883 /// - set: { mode: quiet }
1884 /// phases:
1885 /// - scenario: load_test
1886 /// ```
1887 Bindings {
1888 source: String,
1889 children: Vec<ScenarioNode>,
1890 },
1891}
1892
1893/// Legacy alias.
1894pub type ScenarioStep = ScenarioNode;
1895
1896/// How bindings are defined for an op.
1897///
1898/// Two modes:
1899/// - **Map**: Legacy nosqlbench-style `name: "FuncA(); FuncB()"` chains.
1900/// Each binding is independent; inheritance merges at key level.
1901/// - **PolydatSource**: Native Polydat grammar as a multiline string. The entire
1902/// binding block is a single Polydat program with coordinates, named outputs,
1903/// and full DAG wiring. Replaces (not merges with) any inherited bindings.
1904#[derive(Debug, Clone, Serialize, Deserialize)]
1905#[serde(untagged)]
1906pub enum BindingsDef {
1907 /// Legacy nosqlbench-style: name → expression chain.
1908 Map(HashMap<String, String>),
1909 /// Native Polydat grammar source text.
1910 PolydatSource(String),
1911}
1912
1913impl Default for BindingsDef {
1914 fn default() -> Self {
1915 BindingsDef::Map(HashMap::new())
1916 }
1917}
1918
1919impl BindingsDef {
1920 /// Returns true if there are no bindings defined.
1921 pub fn is_empty(&self) -> bool {
1922 match self {
1923 BindingsDef::Map(m) => m.is_empty(),
1924 BindingsDef::PolydatSource(s) => s.trim().is_empty(),
1925 }
1926 }
1927
1928 /// Get the map view (for legacy code). Returns empty map for PolydatSource.
1929 pub fn as_map(&self) -> &HashMap<String, String> {
1930 static EMPTY: std::sync::LazyLock<HashMap<String, String>> =
1931 std::sync::LazyLock::new(HashMap::new);
1932 match self {
1933 BindingsDef::Map(m) => m,
1934 BindingsDef::PolydatSource(_) => &EMPTY,
1935 }
1936 }
1937
1938 /// Insert a key-value pair (legacy map mode). Converts PolydatSource to Map.
1939 pub fn insert(&mut self, key: String, value: String) {
1940 match self {
1941 BindingsDef::Map(m) => {
1942 m.insert(key, value);
1943 }
1944 _ => {
1945 let mut m = HashMap::new();
1946 m.insert(key, value);
1947 *self = BindingsDef::Map(m);
1948 }
1949 }
1950 }
1951}
1952
1953/// SRD-108 Part B — the typed interface an abstract op slot
1954/// declares: wires the blueprint scope guarantees (`needs`),
1955/// wires the bound implementation must deliver via captures
1956/// (`yields`), and wires it must deliver by projecting the
1957/// result body via `result:` bindings (`results`, SRD-109 Part
1958/// 3), each `name -> polydat type name` (`u64`, `f64`, `String`,
1959/// `vec_f32`, `vec_i64`, …). BTreeMaps so the SRD-107 config
1960/// digest serializes stably.
1961#[derive(Debug, Clone, Serialize, Deserialize, Default, PartialEq, Eq)]
1962pub struct OpInterface {
1963 #[serde(default, skip_serializing_if = "std::collections::BTreeMap::is_empty")]
1964 pub needs: std::collections::BTreeMap<String, String>,
1965 #[serde(default, skip_serializing_if = "std::collections::BTreeMap::is_empty")]
1966 pub yields: std::collections::BTreeMap<String, String>,
1967 #[serde(default, skip_serializing_if = "std::collections::BTreeMap::is_empty")]
1968 pub results: std::collections::BTreeMap<String, String>,
1969}
1970
1971/// A normalized op template — the canonical form.
1972#[derive(Debug, Clone, Serialize, Deserialize)]
1973pub struct ParsedOp {
1974 pub name: String,
1975 #[serde(default, skip_serializing_if = "Option::is_none")]
1976 pub description: Option<String>,
1977 /// The operation payload: field name → value.
1978 /// The statement is always under `"stmt"` after normalization.
1979 /// The original field name that carried the statement (e.g., `"raw"`,
1980 /// `"simple"`, `"prepared"`, `"stmt"`) is preserved in `stmt_type`
1981 /// for adapters that dispatch on execution mode.
1982 pub op: HashMap<String, serde_json::Value>,
1983 /// Binding definitions: either a name→expression map (legacy) or
1984 /// a Polydat grammar source string (native).
1985 #[serde(default, skip_serializing_if = "BindingsDef::is_empty")]
1986 pub bindings: BindingsDef,
1987 /// Configuration parameters.
1988 #[serde(default, skip_serializing_if = "HashMap::is_empty")]
1989 pub params: HashMap<String, serde_json::Value>,
1990 /// Tags for filtering and metadata.
1991 #[serde(default)]
1992 pub tags: HashMap<String, String>,
1993 /// Optional condition expression (from YAML `if:` field).
1994 /// Evaluated per cycle before the op executes. If the result
1995 /// is falsy (false, 0, empty string, None), the op is skipped.
1996 #[serde(default, skip_serializing_if = "Option::is_none")]
1997 pub condition: Option<String>,
1998 /// Optional delay specification (from YAML `delay:` field).
1999 /// Two surface forms:
2000 /// - Bare string: `delay: <name>` — a Polydat binding name
2001 /// producing the pre-op delay value (u64 = ns, f64 = ms).
2002 /// - Map: `delay: { before: <name>, after: <name> }` —
2003 /// independent pre-op and post-op delays; both subkeys
2004 /// are optional but at least one must be set.
2005 #[serde(default, skip_serializing_if = "Option::is_none")]
2006 pub delay: Option<DelaySpec>,
2007 /// SRD-40b synthetic-metric declarations. Each entry
2008 /// publishes one metric family per cycle, valued by a GK
2009 /// expression evaluated in the op's bound scope. Empty
2010 /// when absent. Map key is the metric name and the
2011 /// **default family name**; `MetricSpec::family` overrides
2012 /// it when set. See SRD-40b §1 for the schema, §2 for
2013 /// sugared forms (bare-string / list with wire-expression
2014 /// entries).
2015 #[serde(default, skip_serializing_if = "HashMap::is_empty")]
2016 pub metrics: HashMap<String, MetricSpec>,
2017 /// SRD-66 result-bindings. Vari-structured: string is
2018 /// Polydat source, list is a sequence of fragments, map is
2019 /// named-key short-forms with a composite-map output.
2020 /// `None` ⇒ no result wires; the result wrapper is a
2021 /// no-op for this op.
2022 #[serde(default, skip_serializing_if = "Option::is_none")]
2023 pub result: Option<ResultSpec>,
2024 /// Optional `traverse:` block — CUSTOMISES the always-installed result
2025 /// traversal layer; it does not select it. Like `result:`, absence means
2026 /// "default behaviour", not "no traversal".
2027 #[serde(default, skip_serializing_if = "Option::is_none")]
2028 pub traverse: Option<TraverseSpec>,
2029 /// SRD-32a Push 3 — per-op wrapper-composition override.
2030 /// When present, this op uses the named order instead of
2031 /// the workload-root or runtime-default tiebreaker order.
2032 /// Shadows the workload-root `wrappers:` block entirely
2033 /// (no merge).
2034 #[serde(default, skip_serializing_if = "Option::is_none")]
2035 pub wrappers: Option<WrappersConfig>,
2036 /// Capture-point specs extracted at parse time from any
2037 /// string-valued entry in `op`. Each spec names a column the
2038 /// result body carries and the wire it should be written to
2039 /// via `ctx.wires.write` at cycle time. The `slurp` flag
2040 /// selects between single-row (`[name]`) and all-rows
2041 /// (`[@name]`) extraction.
2042 ///
2043 /// The parser strips the bracket syntax from the source op
2044 /// fields after harvesting the spec, so adapters consume
2045 /// clean text (e.g. `SELECT [key] FROM ...` becomes
2046 /// `SELECT key FROM ...`). Downstream wrappers read this
2047 /// list directly — no re-parsing of the op's text fields.
2048 #[serde(default, skip_serializing_if = "Vec::is_empty")]
2049 pub captures: Vec<crate::bindpoints::CapturePoint>,
2050 /// SRD-108 Part B — present when this op was declared as an
2051 /// ABSTRACT slot (`abstract:` body): the interface a bound
2052 /// implementation must satisfy. Retained on the bound op so
2053 /// pre-map synthesis can verify `yields` against the compiled
2054 /// op-template program.
2055 #[serde(default, skip_serializing_if = "Option::is_none")]
2056 pub abstract_interface: Option<OpInterface>,
2057 /// SRD-108 Part B — `true` once an implementation has been
2058 /// bound into this slot. An op with an interface but
2059 /// `interface_bound == false` at run initiation is a load
2060 /// error ("abstract op unbound").
2061 #[serde(default, skip_serializing_if = "std::ops::Not::not")]
2062 pub interface_bound: bool,
2063 /// Daemon-fiber declaration. When set to a non-`Disabled`
2064 /// value, dispatches of this op from the cycle-pool spawn
2065 /// onto a daemon fiber instead of running inline. Daemons
2066 /// stay in scope of the phase so their failures bubble up;
2067 /// the cycle-pool fiber continues immediately to the next
2068 /// op without awaiting.
2069 ///
2070 /// YAML forms accepted:
2071 /// - `daemon: false` / `daemon: 0` / `daemon: "off"`
2072 /// → not a daemon (cycle-pool op).
2073 /// - `daemon: true` / `daemon: 1` / `daemon: "on"`
2074 /// → daemon, max 1 concurrent fiber per phase activation.
2075 /// - `daemon: <N>` (N ≥ 1) → daemon, max N concurrent
2076 /// fibers per phase activation.
2077 ///
2078 /// The cap is enforced at spawn time: when the cycle-pool
2079 /// dispatches the op and the live-fiber count is already at
2080 /// N, the spawn errors and the phase fails with a clear
2081 /// stop_reason. There's no queuing — exceeding the cap is a
2082 /// workload-design error, not backpressure to absorb.
2083 ///
2084 /// Per-op-name dedup: daemons are tracked by op-template
2085 /// name. Subsequent dispatches of the same name silently
2086 /// succeed up to the cap; over the cap they fail loud.
2087 /// Natural daemon exit (Completed / Cancelled / Errored /
2088 /// TimedOut / Panicked) decrements the count, freeing a
2089 /// slot for the next dispatch.
2090 ///
2091 /// Use case: a long-running server call (e.g.
2092 /// `forceKeyspaceCompaction`) that the workload wants to
2093 /// fire while a sibling op concurrently observes progress.
2094 /// The daemon op stays in scope so its failures bubble up;
2095 /// the cycle-pool op runs alongside without serialisation.
2096 ///
2097 /// Parser invariants when `daemon` is `MaxFibers(_)`:
2098 /// - `cycles:` and `ratio:` on this op are rejected — the
2099 /// daemon's dispatch cadence is governed by the cycle-pool
2100 /// walks, not by cycles-per-second.
2101 /// - `if:` / `while:` apply as usual; a falsy guard skips
2102 /// the dispatch with no spawn.
2103 #[serde(default, skip_serializing_if = "DaemonSpec::is_disabled")]
2104 pub daemon: DaemonSpec,
2105 /// How long the phase waits for this daemon's in-flight
2106 /// future to drop after sending the stop signal. Past the
2107 /// window, the phase records a daemon-shutdown failure and
2108 /// fails. Per-op override of the activity-level default
2109 /// (5000 ms).
2110 ///
2111 /// Only meaningful when `daemon` is `MaxFibers(_)`.
2112 #[serde(default, skip_serializing_if = "Option::is_none")]
2113 pub daemon_cancel_grace_ms: Option<u64>,
2114 /// Polydat boolean expression evaluated on the daemon's fiber
2115 /// inside the wrapper stack. When set, the daemon body runs
2116 /// a loop: while the condition is truthy, dispatch the
2117 /// inner op; when falsy or stop-signalled, exit.
2118 ///
2119 /// `while:` composes inside `if:` (which gates whether the
2120 /// loop starts at all) and outside the per-op-rate
2121 /// throttling (which paces iterations).
2122 ///
2123 /// Parser invariants:
2124 /// - `while:` on a non-daemon op is currently allowed but
2125 /// blocks the cycle-pool fiber until the loop exits. The
2126 /// common use case is daemon + while.
2127 #[serde(default, skip_serializing_if = "Option::is_none")]
2128 pub while_cond: Option<String>,
2129 /// Per-op rate spec governing the iteration cadence of the
2130 /// while-loop (or per-cycle dispatch for non-loop ops).
2131 /// Format: `"<N>"` (N/s, bare integer), `"<N>/s"`,
2132 /// `"<N>/m"`, `"<N>/h"`. Each op with `rate:` gets its own
2133 /// `RateLimiter` at phase init — INDEPENDENT of the
2134 /// activity-level `rate:` and of other ops' rate limiters.
2135 ///
2136 /// `rate: 0` is "unlimited" (no throttling).
2137 #[serde(default, skip_serializing_if = "Option::is_none")]
2138 pub rate: Option<String>,
2139}
2140
2141/// Daemon-fiber capacity declaration. The on-disk surface
2142/// accepts multiple YAML scalar shapes (bool / int / string)
2143/// that all map into this two-variant enum.
2144///
2145/// `Disabled` is the default — the op runs on the cycle-pool
2146/// fiber inline like everything else. `MaxFibers(N)` opts in
2147/// to daemon-fiber dispatch with a per-op-name cap of N.
2148#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
2149pub enum DaemonSpec {
2150 #[default]
2151 Disabled,
2152 MaxFibers(u32),
2153}
2154
2155impl DaemonSpec {
2156 pub fn is_disabled(&self) -> bool {
2157 matches!(self, DaemonSpec::Disabled)
2158 }
2159 pub fn max_fibers(&self) -> Option<u32> {
2160 match self {
2161 DaemonSpec::Disabled => None,
2162 DaemonSpec::MaxFibers(n) => Some(*n),
2163 }
2164 }
2165}
2166
2167impl serde::Serialize for DaemonSpec {
2168 fn serialize<S: serde::Serializer>(&self, s: S) -> Result<S::Ok, S::Error> {
2169 match self {
2170 DaemonSpec::Disabled => s.serialize_bool(false),
2171 DaemonSpec::MaxFibers(1) => s.serialize_bool(true),
2172 DaemonSpec::MaxFibers(n) => s.serialize_u32(*n),
2173 }
2174 }
2175}
2176
2177impl<'de> serde::Deserialize<'de> for DaemonSpec {
2178 fn deserialize<D: serde::Deserializer<'de>>(d: D) -> Result<Self, D::Error> {
2179 let v = serde_json::Value::deserialize(d)?;
2180 parse_daemon_spec_value(&v).map_err(serde::de::Error::custom)
2181 }
2182}
2183
2184/// Parse a YAML/JSON scalar into a `DaemonSpec`.
2185///
2186/// Accepted forms:
2187/// - `true` / `1` / `"true"` / `"on"` → `MaxFibers(1)`
2188/// - `false` / `0` / `"false"` / `"off"` → `Disabled`
2189/// - non-negative integer N ≥ 1 → `MaxFibers(N)`
2190/// - non-negative integer 0 → `Disabled`
2191///
2192/// Rejected forms (returns descriptive error):
2193/// - negative integers
2194/// - non-integer numbers (1.5, NaN, ...)
2195/// - strings other than the accepted set
2196/// - null, arrays, objects
2197///
2198/// Public so the unit + proptest layers can exercise it
2199/// directly without a full YAML round-trip.
2200pub fn parse_daemon_spec_value(v: &serde_json::Value) -> Result<DaemonSpec, String> {
2201 match v {
2202 serde_json::Value::Bool(true) => Ok(DaemonSpec::MaxFibers(1)),
2203 serde_json::Value::Bool(false) => Ok(DaemonSpec::Disabled),
2204 serde_json::Value::Number(n) => {
2205 if let Some(u) = n.as_u64() {
2206 if u == 0 {
2207 Ok(DaemonSpec::Disabled)
2208 } else if u <= u32::MAX as u64 {
2209 Ok(DaemonSpec::MaxFibers(u as u32))
2210 } else {
2211 Err(format!(
2212 "daemon: {u} exceeds u32::MAX — caps above {} \
2213 aren't supported (and don't make practical sense)",
2214 u32::MAX,
2215 ))
2216 }
2217 } else if n.as_i64().is_some_and(|i| i < 0) {
2218 Err(format!(
2219 "daemon: {n} — negative integers are invalid. \
2220 Use 0 / false / \"off\" to disable, or a positive \
2221 integer for the max-fibers cap.",
2222 ))
2223 } else {
2224 Err(format!(
2225 "daemon: {n} — non-integer numbers are invalid. \
2226 Use a boolean or a non-negative integer.",
2227 ))
2228 }
2229 }
2230 serde_json::Value::String(s) => match s.trim().to_ascii_lowercase().as_str() {
2231 "true" | "on" => Ok(DaemonSpec::MaxFibers(1)),
2232 "false" | "off" => Ok(DaemonSpec::Disabled),
2233 other => Err(format!(
2234 "daemon: \"{other}\" — unknown string form. \
2235 Accepted: \"on\" / \"off\" / \"true\" / \"false\", \
2236 or use a boolean / non-negative integer directly.",
2237 )),
2238 },
2239 serde_json::Value::Null => {
2240 Err("daemon: null is not a valid value. Use false to disable.".into())
2241 }
2242 other => Err(format!(
2243 "daemon: {other:?} — only boolean, integer, or string forms \
2244 are accepted.",
2245 )),
2246 }
2247}
2248
2249/// Per-op delay specification. Two surface forms on YAML:
2250/// - Bare string `delay: <name>` → `Before(<name>)`: a GK
2251/// binding name producing the pre-op delay value (u64 ns,
2252/// f64 ms).
2253/// - Map `delay: { before: <name>, after: <name> }` →
2254/// `BeforeAfter`: independent pre-op and post-op delays.
2255/// Both subkeys are optional; an empty map is a parse error.
2256#[derive(Debug, Clone, PartialEq, Eq)]
2257pub enum DelaySpec {
2258 /// Single pre-op delay. Equivalent to BeforeAfter { Some, None }
2259 /// at runtime; the discriminant exists so YAML round-trip
2260 /// preserves the author's chosen surface form.
2261 Before(String),
2262 BeforeAfter {
2263 before: Option<String>,
2264 after: Option<String>,
2265 },
2266}
2267
2268impl DelaySpec {
2269 /// Pre-op delay binding name, if any.
2270 pub fn before(&self) -> Option<&str> {
2271 match self {
2272 DelaySpec::Before(n) => Some(n.as_str()),
2273 DelaySpec::BeforeAfter { before, .. } => before.as_deref(),
2274 }
2275 }
2276 /// Post-op delay binding name, if any.
2277 pub fn after(&self) -> Option<&str> {
2278 match self {
2279 DelaySpec::Before(_) => None,
2280 DelaySpec::BeforeAfter { after, .. } => after.as_deref(),
2281 }
2282 }
2283 /// Every binding name this spec references. Used by scope
2284 /// synthesis to ensure the names land on the per-op kernel.
2285 pub fn names(&self) -> Vec<&str> {
2286 match self {
2287 DelaySpec::Before(n) => vec![n.as_str()],
2288 DelaySpec::BeforeAfter { before, after } => {
2289 let mut out = Vec::with_capacity(2);
2290 if let Some(b) = before.as_deref() {
2291 out.push(b);
2292 }
2293 if let Some(a) = after.as_deref() {
2294 out.push(a);
2295 }
2296 out
2297 }
2298 }
2299 }
2300}
2301
2302impl serde::Serialize for DelaySpec {
2303 fn serialize<S: serde::Serializer>(&self, s: S) -> Result<S::Ok, S::Error> {
2304 use serde::ser::SerializeMap;
2305 match self {
2306 DelaySpec::Before(name) => s.serialize_str(name),
2307 DelaySpec::BeforeAfter { before, after } => {
2308 let mut m = s.serialize_map(None)?;
2309 if let Some(b) = before {
2310 m.serialize_entry("before", b)?;
2311 }
2312 if let Some(a) = after {
2313 m.serialize_entry("after", a)?;
2314 }
2315 m.end()
2316 }
2317 }
2318 }
2319}
2320
2321impl<'de> serde::Deserialize<'de> for DelaySpec {
2322 fn deserialize<D: serde::Deserializer<'de>>(d: D) -> Result<Self, D::Error> {
2323 let v = serde_json::Value::deserialize(d)?;
2324 parse_delay_spec_value(&v).map_err(serde::de::Error::custom)
2325 }
2326}
2327
2328/// Parse a YAML/JSON value into a `DelaySpec`.
2329///
2330/// Accepted:
2331/// - non-empty string → `Before(<string>)`
2332/// - object with `before` and/or `after` keys → `BeforeAfter`
2333///
2334/// Rejected (descriptive error):
2335/// - empty string
2336/// - empty object
2337/// - object with unknown keys
2338/// - object whose `before` / `after` values aren't strings
2339/// - null, arrays, numbers, booleans
2340///
2341/// Public for unit tests + smoke validation.
2342pub fn parse_delay_spec_value(v: &serde_json::Value) -> Result<DelaySpec, String> {
2343 match v {
2344 serde_json::Value::String(s) => {
2345 let trimmed = s.trim();
2346 if trimmed.is_empty() {
2347 Err("delay: empty string is not a valid binding name".into())
2348 } else {
2349 Ok(DelaySpec::Before(trimmed.to_string()))
2350 }
2351 }
2352 serde_json::Value::Object(map) => {
2353 let mut before: Option<String> = None;
2354 let mut after: Option<String> = None;
2355 for (k, v) in map {
2356 match k.as_str() {
2357 "before" => {
2358 before = match v {
2359 serde_json::Value::String(s) => {
2360 let t = s.trim();
2361 if t.is_empty() {
2362 return Err(
2363 "delay.before: empty string is not a valid binding name"
2364 .into(),
2365 );
2366 }
2367 Some(t.to_string())
2368 }
2369 other => {
2370 return Err(format!(
2371 "delay.before: expected string, got {other:?}",
2372 ));
2373 }
2374 };
2375 }
2376 "after" => {
2377 after = match v {
2378 serde_json::Value::String(s) => {
2379 let t = s.trim();
2380 if t.is_empty() {
2381 return Err(
2382 "delay.after: empty string is not a valid binding name"
2383 .into(),
2384 );
2385 }
2386 Some(t.to_string())
2387 }
2388 other => {
2389 return Err(
2390 format!("delay.after: expected string, got {other:?}",),
2391 );
2392 }
2393 };
2394 }
2395 other => {
2396 return Err(format!(
2397 "delay: unknown key `{other}` — accepted: `before`, `after`",
2398 ));
2399 }
2400 }
2401 }
2402 if before.is_none() && after.is_none() {
2403 Err("delay: map form must set at least one of `before` / `after`".into())
2404 } else {
2405 Ok(DelaySpec::BeforeAfter { before, after })
2406 }
2407 }
2408 serde_json::Value::Null => {
2409 Err("delay: null is not a valid value. Omit the field instead.".into())
2410 }
2411 other => Err(format!(
2412 "delay: {other:?} — accepted forms are a binding-name string or \
2413 a map `{{ before: <name>, after: <name> }}`",
2414 )),
2415 }
2416}
2417
2418#[cfg(test)]
2419mod delay_spec_tests {
2420 use super::*;
2421 use serde_json::json;
2422
2423 #[test]
2424 fn parse_bare_string_is_before() {
2425 let spec = parse_delay_spec_value(&json!("ticks")).unwrap();
2426 assert_eq!(spec, DelaySpec::Before("ticks".into()));
2427 }
2428
2429 #[test]
2430 fn parse_trimmed_string() {
2431 let spec = parse_delay_spec_value(&json!(" ticks ")).unwrap();
2432 assert_eq!(spec, DelaySpec::Before("ticks".into()));
2433 }
2434
2435 #[test]
2436 fn parse_map_with_both() {
2437 let spec = parse_delay_spec_value(&json!({
2438 "before": "pre", "after": "post"
2439 }))
2440 .unwrap();
2441 assert_eq!(
2442 spec,
2443 DelaySpec::BeforeAfter {
2444 before: Some("pre".into()),
2445 after: Some("post".into()),
2446 }
2447 );
2448 }
2449
2450 #[test]
2451 fn parse_map_before_only() {
2452 let spec = parse_delay_spec_value(&json!({ "before": "pre" })).unwrap();
2453 assert_eq!(
2454 spec,
2455 DelaySpec::BeforeAfter {
2456 before: Some("pre".into()),
2457 after: None,
2458 }
2459 );
2460 }
2461
2462 #[test]
2463 fn parse_map_after_only() {
2464 let spec = parse_delay_spec_value(&json!({ "after": "post" })).unwrap();
2465 assert_eq!(
2466 spec,
2467 DelaySpec::BeforeAfter {
2468 before: None,
2469 after: Some("post".into()),
2470 }
2471 );
2472 }
2473
2474 #[test]
2475 fn parse_rejects_empty_string() {
2476 assert!(parse_delay_spec_value(&json!("")).is_err());
2477 assert!(parse_delay_spec_value(&json!(" ")).is_err());
2478 }
2479
2480 #[test]
2481 fn parse_rejects_empty_map() {
2482 assert!(parse_delay_spec_value(&json!({})).is_err());
2483 }
2484
2485 #[test]
2486 fn parse_rejects_unknown_key() {
2487 let e = parse_delay_spec_value(&json!({
2488 "before": "pre", "during": "mid"
2489 }))
2490 .unwrap_err();
2491 assert!(e.contains("during"));
2492 }
2493
2494 #[test]
2495 fn parse_rejects_non_string_value() {
2496 assert!(parse_delay_spec_value(&json!({ "before": 5 })).is_err());
2497 assert!(parse_delay_spec_value(&json!({ "after": true })).is_err());
2498 assert!(parse_delay_spec_value(&json!({ "before": null })).is_err());
2499 }
2500
2501 #[test]
2502 fn parse_rejects_null_top_level() {
2503 assert!(parse_delay_spec_value(&json!(null)).is_err());
2504 }
2505
2506 #[test]
2507 fn parse_rejects_array() {
2508 assert!(parse_delay_spec_value(&json!(["pre"])).is_err());
2509 }
2510
2511 #[test]
2512 fn parse_rejects_number() {
2513 assert!(parse_delay_spec_value(&json!(1.5)).is_err());
2514 assert!(parse_delay_spec_value(&json!(100)).is_err());
2515 }
2516
2517 #[test]
2518 fn parse_rejects_empty_after_value() {
2519 assert!(parse_delay_spec_value(&json!({ "after": "" })).is_err());
2520 }
2521
2522 #[test]
2523 fn round_trip_before_serializes_as_string() {
2524 let spec = DelaySpec::Before("ticks".into());
2525 let v = serde_json::to_value(&spec).unwrap();
2526 assert_eq!(v, json!("ticks"));
2527 let parsed: DelaySpec = serde_json::from_value(v).unwrap();
2528 assert_eq!(parsed, spec);
2529 }
2530
2531 #[test]
2532 fn round_trip_map_serializes_as_object() {
2533 let spec = DelaySpec::BeforeAfter {
2534 before: Some("pre".into()),
2535 after: Some("post".into()),
2536 };
2537 let v = serde_json::to_value(&spec).unwrap();
2538 assert_eq!(v.get("before"), Some(&json!("pre")));
2539 assert_eq!(v.get("after"), Some(&json!("post")));
2540 let parsed: DelaySpec = serde_json::from_value(v).unwrap();
2541 assert_eq!(parsed, spec);
2542 }
2543
2544 #[test]
2545 fn names_returns_all_referenced() {
2546 assert_eq!(DelaySpec::Before("x".into()).names(), vec!["x"]);
2547 let spec = DelaySpec::BeforeAfter {
2548 before: Some("a".into()),
2549 after: Some("b".into()),
2550 };
2551 assert_eq!(spec.names(), vec!["a", "b"]);
2552 let only_before = DelaySpec::BeforeAfter {
2553 before: Some("a".into()),
2554 after: None,
2555 };
2556 assert_eq!(only_before.names(), vec!["a"]);
2557 }
2558
2559 #[test]
2560 fn accessors_return_correct_names() {
2561 let s = DelaySpec::Before("x".into());
2562 assert_eq!(s.before(), Some("x"));
2563 assert_eq!(s.after(), None);
2564 let s = DelaySpec::BeforeAfter {
2565 before: Some("a".into()),
2566 after: Some("b".into()),
2567 };
2568 assert_eq!(s.before(), Some("a"));
2569 assert_eq!(s.after(), Some("b"));
2570 }
2571}
2572
2573#[cfg(test)]
2574mod daemon_spec_tests {
2575 use super::*;
2576 use serde_json::json;
2577
2578 #[test]
2579 fn bool_true_is_max_1() {
2580 assert_eq!(
2581 parse_daemon_spec_value(&json!(true)).unwrap(),
2582 DaemonSpec::MaxFibers(1)
2583 );
2584 }
2585 #[test]
2586 fn bool_false_is_disabled() {
2587 assert_eq!(
2588 parse_daemon_spec_value(&json!(false)).unwrap(),
2589 DaemonSpec::Disabled
2590 );
2591 }
2592 #[test]
2593 fn int_0_is_disabled() {
2594 assert_eq!(
2595 parse_daemon_spec_value(&json!(0)).unwrap(),
2596 DaemonSpec::Disabled
2597 );
2598 }
2599 #[test]
2600 fn int_1_is_max_1() {
2601 assert_eq!(
2602 parse_daemon_spec_value(&json!(1)).unwrap(),
2603 DaemonSpec::MaxFibers(1)
2604 );
2605 }
2606 #[test]
2607 fn int_n_is_max_n() {
2608 assert_eq!(
2609 parse_daemon_spec_value(&json!(10)).unwrap(),
2610 DaemonSpec::MaxFibers(10)
2611 );
2612 }
2613 #[test]
2614 fn str_on_is_max_1() {
2615 assert_eq!(
2616 parse_daemon_spec_value(&json!("on")).unwrap(),
2617 DaemonSpec::MaxFibers(1)
2618 );
2619 assert_eq!(
2620 parse_daemon_spec_value(&json!("true")).unwrap(),
2621 DaemonSpec::MaxFibers(1)
2622 );
2623 assert_eq!(
2624 parse_daemon_spec_value(&json!("ON")).unwrap(),
2625 DaemonSpec::MaxFibers(1)
2626 );
2627 }
2628 #[test]
2629 fn str_off_is_disabled() {
2630 assert_eq!(
2631 parse_daemon_spec_value(&json!("off")).unwrap(),
2632 DaemonSpec::Disabled
2633 );
2634 assert_eq!(
2635 parse_daemon_spec_value(&json!("false")).unwrap(),
2636 DaemonSpec::Disabled
2637 );
2638 assert_eq!(
2639 parse_daemon_spec_value(&json!("OFF")).unwrap(),
2640 DaemonSpec::Disabled
2641 );
2642 }
2643 #[test]
2644 fn negative_int_rejected() {
2645 assert!(parse_daemon_spec_value(&json!(-1)).is_err());
2646 assert!(parse_daemon_spec_value(&json!(-100)).is_err());
2647 }
2648 #[test]
2649 fn float_rejected() {
2650 assert!(parse_daemon_spec_value(&json!(1.5)).is_err());
2651 }
2652 #[test]
2653 fn unknown_string_rejected() {
2654 assert!(parse_daemon_spec_value(&json!("garbage")).is_err());
2655 assert!(parse_daemon_spec_value(&json!("yes")).is_err());
2656 }
2657 #[test]
2658 fn null_rejected() {
2659 assert!(parse_daemon_spec_value(&json!(null)).is_err());
2660 }
2661 #[test]
2662 fn array_object_rejected() {
2663 assert!(parse_daemon_spec_value(&json!([1, 2])).is_err());
2664 assert!(parse_daemon_spec_value(&json!({"max": 5})).is_err());
2665 }
2666 #[test]
2667 fn round_trip_disabled() {
2668 let s = serde_json::to_value(DaemonSpec::Disabled).unwrap();
2669 assert_eq!(parse_daemon_spec_value(&s).unwrap(), DaemonSpec::Disabled);
2670 }
2671 #[test]
2672 fn round_trip_max_1_serialises_as_bool() {
2673 let s = serde_json::to_value(DaemonSpec::MaxFibers(1)).unwrap();
2674 assert_eq!(s, json!(true));
2675 assert_eq!(
2676 parse_daemon_spec_value(&s).unwrap(),
2677 DaemonSpec::MaxFibers(1)
2678 );
2679 }
2680 #[test]
2681 fn round_trip_max_n_serialises_as_int() {
2682 let s = serde_json::to_value(DaemonSpec::MaxFibers(5)).unwrap();
2683 assert_eq!(s, json!(5));
2684 assert_eq!(
2685 parse_daemon_spec_value(&s).unwrap(),
2686 DaemonSpec::MaxFibers(5)
2687 );
2688 }
2689}
2690
2691/// `traverse:` — knobs on the result-traversal layer.
2692///
2693/// The layer itself is always installed (`result:` and `metrics:` declare it
2694/// as `requires_inner`, and the composition resolver enforces that at init),
2695/// so this block only tunes behaviour that is otherwise fixed.
2696#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
2697pub struct TraverseSpec {
2698 /// Base JSON Pointer (RFC 6901) applied to the result body BEFORE any
2699 /// capture is resolved, re-rooting the document.
2700 ///
2701 /// Every capture otherwise repeats the same prefix — `/value/0/x`,
2702 /// `/value/0/y` — which is the shape envelope responses force
2703 /// (`{value, status, …}`). `poll.json_path` already exists for exactly
2704 /// this reason on exactly this data; this is the same idea for captures.
2705 #[serde(default, skip_serializing_if = "Option::is_none")]
2706 pub path: Option<String>,
2707 /// What to do when a declared capture resolves to nothing.
2708 #[serde(default)]
2709 pub on_missing: OnMissing,
2710}
2711
2712/// Policy for a capture that resolved to nothing.
2713#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, Serialize, Deserialize)]
2714#[serde(rename_all = "lowercase")]
2715pub enum OnMissing {
2716 /// Bind `None` and continue — the historical behaviour, and the default
2717 /// so existing workloads are unaffected.
2718 #[default]
2719 Ignore,
2720 /// Log it. For a capture that is legitimately sometimes absent but whose
2721 /// absence you still want to see.
2722 Warn,
2723 /// Fail the op. For a capture whose absence means the query or the schema
2724 /// changed under you — which otherwise reads identically to "measured,
2725 /// and it was absent".
2726 Error,
2727}
2728
2729/// SRD-40b §1 schema for one synthetic-metric declaration on
2730/// an op template.
2731#[derive(Debug, Clone, Serialize, Deserialize)]
2732pub struct MetricSpec {
2733 /// Required. A Polydat expression evaluated in the op's bound
2734 /// scope. A bare binding name is the canonical form when
2735 /// the formula belongs in a `bindings:` block; any GK
2736 /// expression that produces a numeric result is also
2737 /// valid. See SRD-40b §4.
2738 pub value: String,
2739 /// Optional override of the family name. Defaults to the
2740 /// map key on `ParsedOp.metrics`. SRD-40b §1.
2741 #[serde(default, skip_serializing_if = "Option::is_none")]
2742 pub family: Option<String>,
2743 /// Optional metric type. Defaults to `Gauge` per SRD-40b §1.
2744 #[serde(default, skip_serializing_if = "Option::is_none")]
2745 pub kind: Option<MetricKind>,
2746 /// Optional OpenMetrics unit suffix (`ms`, `bytes`, …).
2747 /// When set, lands in BOTH the family-name suffix and the
2748 /// `metric_family.unit` column per SRD-40a §4.3.
2749 #[serde(default, skip_serializing_if = "Option::is_none")]
2750 pub unit: Option<String>,
2751 /// Optional generation-time numeric sanitiser using
2752 /// Excel-style hash patterns (`#.##`, `0.000`, etc.).
2753 /// Translated at registration time into a round op that
2754 /// runs before the value is recorded on the instrument;
2755 /// storage holds the sanitised number. SRD-40b §1.
2756 #[serde(default, skip_serializing_if = "Option::is_none")]
2757 pub format: Option<String>,
2758 /// Optional dimensional placement: dimension name → a Polydat
2759 /// expression producing that dimension's value for this sample.
2760 ///
2761 /// A metric identity is its label set with the family name promoted
2762 /// into it, a closed 1:1 association. `cell:` therefore does not
2763 /// attach a label to a sample — it selects the dimensional CELL the
2764 /// sample belongs to, REFINING the identity its registration site
2765 /// already composes. One value per cell, one family per cell, and the
2766 /// existing duplicate-family check keeps working unchanged.
2767 ///
2768 /// `BTreeMap` so synthesis order is deterministic: the coordinate's
2769 /// rendering is what keys a cell, and a map that iterated differently
2770 /// between runs would key the same coordinate two ways.
2771 #[serde(default, skip_serializing_if = "std::collections::BTreeMap::is_empty")]
2772 pub cell: std::collections::BTreeMap<String, String>,
2773}
2774
2775/// A dimension a scope introduces: the label NAME whose values arrive from
2776/// data, declared once at the tier that owns it.
2777///
2778/// The name is a structural declaration; only the value varies per cell.
2779/// `Component::attach` enforces the same rule on the runtime tree — a label
2780/// name is owned by exactly one tier and inherited downward — so declaring
2781/// the name here is what lets that be checked against the program before a
2782/// cycle runs, rather than surfacing as an attach-time panic.
2783#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
2784pub struct DimensionSpec {
2785 /// Value type. `str` today — label values are strings, and a
2786 /// dimension whose values came from a float would key cells on
2787 /// formatting rather than on identity.
2788 #[serde(default)]
2789 pub value_type: DimensionType,
2790}
2791
2792/// Declared value type of a [`DimensionSpec`].
2793#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, Serialize, Deserialize)]
2794#[serde(rename_all = "lowercase")]
2795pub enum DimensionType {
2796 #[default]
2797 Str,
2798}
2799
2800/// Metric type discriminator. SRD-40b §1.
2801#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
2802#[serde(rename_all = "lowercase")]
2803pub enum MetricKind {
2804 /// Current-state observation; per-cycle `set(value)`.
2805 /// Default per SRD-40b §1.
2806 Gauge,
2807 /// Distribution sample; per-cycle `record(value)`.
2808 Histogram,
2809 /// Monotonic running total; per-cycle `inc_by(value)`.
2810 Counter,
2811}
2812
2813impl Default for MetricKind {
2814 /// SRD-40b §1: gauge is the default — synthetic values are
2815 /// most often current-state observations.
2816 fn default() -> Self {
2817 MetricKind::Gauge
2818 }
2819}
2820
2821/// SRD-66 result-bindings declaration. Vari-structured to
2822/// match the three YAML shapes the user can write:
2823///
2824/// - **String**: a multi-line Polydat source block. Each
2825/// `<name> := <expr>` assignment declares one result wire.
2826/// - **List**: a sequence of nested `ResultSpec`s; each
2827/// element processes in order and contributes its
2828/// declarations.
2829/// - **Map**: named-key short-forms (`count`, `ok`,
2830/// path-expr, or any other string treated as a GK
2831/// expression). Map shape additionally produces a
2832/// composite-map wire (deferred — see Push 2 follow-ups).
2833///
2834/// SRD-40b §5.1's mapping form is preserved as the map
2835/// shape with two refinements: any non-built-in non-path
2836/// string is a Polydat expression (no `(`-detector magic), and
2837/// the composite-map output is added.
2838#[derive(Debug, Clone, Serialize, Deserialize)]
2839#[serde(untagged)]
2840pub enum ResultSpec {
2841 /// Polydat source block — one or more `<name> := <expr>`
2842 /// assignments separated by newlines. The pre-bound
2843 /// wires (`body`, `count`, `ok`, captures) are
2844 /// available; references resolve via the standard
2845 /// closure-binding rule (polydat module matter detects
2846 /// linkages).
2847 String(String),
2848 /// Sequence of fragments. Each element is itself a
2849 /// `ResultSpec` (string or map; nested lists are parsed
2850 /// but unconventional). Fragments concatenate into one
2851 /// result-bindings scope; key collisions across map-
2852 /// shape fragments are a hard error.
2853 List(Vec<ResultSpec>),
2854 /// Named-key short-forms. Each value is one of:
2855 /// `"count"`, `"ok"`, a path expression (no parens), or
2856 /// any other string treated as a Polydat expression. Map
2857 /// shape also produces a composite-map wire keyed by
2858 /// the YAML keys.
2859 Map(std::collections::BTreeMap<String, String>),
2860}
2861
2862/// Legacy alias for backwards compatibility during Push 2.
2863/// Drops once every consumer migrates to `ResultSpec`.
2864pub type ResultWireSpec = LegacyResultWireSpec;
2865
2866#[derive(Debug, Clone, Serialize, Deserialize)]
2867#[serde(untagged)]
2868pub enum LegacyResultWireSpec {
2869 String(String),
2870 Object {
2871 source: String,
2872 #[serde(default, skip_serializing_if = "Option::is_none")]
2873 default: Option<String>,
2874 },
2875}
2876
2877impl LegacyResultWireSpec {
2878 pub fn source(&self) -> &str {
2879 match self {
2880 LegacyResultWireSpec::String(s) => s,
2881 LegacyResultWireSpec::Object { source, .. } => source,
2882 }
2883 }
2884}
2885
2886impl ResultSpec {
2887 /// Walk the spec tree (handling list-shape recursion)
2888 /// and yield every (wire-name, source-expr) pair the
2889 /// spec ultimately declares.
2890 ///
2891 /// For map-shape entries, the source is whatever the
2892 /// user wrote (`count` / `ok` / path-expr / Polydat expr).
2893 /// For string-shape entries, the source is the entire
2894 /// Polydat block — the caller compiles it as a unit and
2895 /// extracts wire names from the LHS of each `:=`
2896 /// assignment.
2897 pub fn walk_fragments<F: FnMut(ResultFragment<'_>)>(&self, mut on: F) {
2898 self.walk_fragments_inner(&mut on);
2899 }
2900
2901 fn walk_fragments_inner<F: FnMut(ResultFragment<'_>)>(&self, on: &mut F) {
2902 match self {
2903 ResultSpec::String(s) => on(ResultFragment::Source(s)),
2904 ResultSpec::List(items) => {
2905 for item in items {
2906 item.walk_fragments_inner(on);
2907 }
2908 }
2909 ResultSpec::Map(entries) => {
2910 for (name, source) in entries {
2911 on(ResultFragment::Named { name, source });
2912 }
2913 }
2914 }
2915 }
2916
2917 /// True when the spec declares no wires. Used by the
2918 /// wrapper-trigger to skip wrapping when `result:` was
2919 /// explicitly empty.
2920 pub fn is_empty(&self) -> bool {
2921 match self {
2922 ResultSpec::String(s) => s.trim().is_empty(),
2923 ResultSpec::List(items) => items.iter().all(|i| i.is_empty()),
2924 ResultSpec::Map(entries) => entries.is_empty(),
2925 }
2926 }
2927}
2928
2929/// One step of `ResultSpec::walk_fragments`. Either a
2930/// string-shape source block (compile as a Polydat module) or a
2931/// map-shape `(name, source)` pair (compile as a single
2932/// `name := source` binding).
2933pub enum ResultFragment<'a> {
2934 Source(&'a str),
2935 Named { name: &'a str, source: &'a str },
2936}
2937
2938impl ParsedOp {
2939 /// Create a minimal ParsedOp with just a name and stmt.
2940 pub fn simple(name: &str, stmt: &str) -> Self {
2941 let mut op = HashMap::new();
2942 op.insert(
2943 "stmt".to_string(),
2944 serde_json::Value::String(stmt.to_string()),
2945 );
2946 Self {
2947 traverse: None,
2948 name: name.to_string(),
2949 description: None,
2950 op,
2951 bindings: BindingsDef::default(),
2952 params: HashMap::new(),
2953 tags: HashMap::new(),
2954 condition: None,
2955 delay: None,
2956 metrics: HashMap::new(),
2957 result: None,
2958 wrappers: None,
2959 captures: Vec::new(),
2960 abstract_interface: None,
2961 interface_bound: false,
2962 daemon: DaemonSpec::Disabled,
2963 daemon_cancel_grace_ms: None,
2964 while_cond: None,
2965 rate: None,
2966 }
2967 }
2968}