Skip to main content

nmbrs_runtime/wrappers/
traverse.rs

1// Copyright 2024-2026 Jonathan Shook
2// SPDX-License-Identifier: Apache-2.0
3
4//! Default result-traversal wrapper. Always wraps the inner
5//! adapter dispenser: counts result elements + bytes, walks
6//! declared capture points, and writes extracted values onto
7//! the per-fiber op-template kernel via `ctx.wires.write`.
8
9use std::collections::HashMap;
10use std::sync::Arc;
11
12use crate::adapter::WrappingDispenser;
13use crate::adapter::{ExecutionError, OpDispenser, OpResult};
14use crate::wrapper_registry::{WrapperName, WrapperRegistration, WrapperSubject};
15use nmbrs_workload::bindpoints;
16
17/// SRD-32a wrapper name. Innermost layer; always present.
18pub const NAME: WrapperName = WrapperName::new("traverse");
19
20/// Trigger: always — every op gets a traversal layer so result
21/// bodies are consumed (and per-row metrics record element /
22/// byte counts) even if the workload didn't ask for anything.
23fn triggers(s: WrapperSubject) -> bool {
24    s.op().is_some()
25}
26
27/// No per-op assignment text — the wrapper has no operator-
28/// configurable knobs.
29fn describe_assignment(_: WrapperSubject) -> Option<String> {
30    None
31}
32
33inventory::submit! {
34    WrapperRegistration {
35        name: NAME,
36        owned_fields: &[],
37        triggers,
38        requires_inner: &[],
39        forbids_outer: &[],
40        mutually_exclusive_with: &[],
41        describe_assignment,
42        levels: &[crate::wrapper_registry::WrapperLevel::Op],
43    }
44}
45
46/// Result traversal statistics, backed by activity metrics counters.
47pub struct TraversalStats {
48    pub metrics: Arc<crate::activity::ActivityMetrics>,
49}
50
51/// Wraps an inner OpDispenser with result traversal and optional
52/// capture extraction.
53///
54/// This is the default wrapper, always applied unless disabled.
55/// It ensures that:
56/// 1. The result body is fully consumed (element/byte counting)
57/// 2. Captures are extracted from the result (if declared)
58/// 3. Traversal metrics are recorded
59pub struct TraversingDispenser {
60    inner: Arc<dyn OpDispenser>,
61    stats: Arc<TraversalStats>,
62    /// Capture points parsed from the template at init time.
63    /// Empty if no captures are declared.
64    captures: Vec<bindpoints::CapturePoint>,
65    /// The op's `traverse:` block, if it declared one. Absent means defaults —
66    /// the layer is always installed, so this only tunes it.
67    spec: Option<nmbrs_workload::model::TraverseSpec>,
68}
69
70impl TraversingDispenser {
71    /// Wrap an inner dispenser with traversal.
72    ///
73    /// Reads `template.captures` (the parse-time-extracted capture
74    /// specs) directly. The op-template parser has already stripped
75    /// `[name]` / `[@name]` brackets from the op text fields, so
76    /// adapters see clean SQL/URL/body strings.
77    pub fn wrap(
78        inner: Arc<dyn OpDispenser>,
79        template: &nmbrs_workload::model::ParsedOp,
80        stats: Arc<TraversalStats>,
81    ) -> Arc<dyn OpDispenser> {
82        Arc::new(Self {
83            inner,
84            stats,
85            captures: template.captures.clone(),
86            spec: template.traverse.clone(),
87        })
88    }
89}
90
91/// Extract captures from a result body's JSON.
92///
93/// Walks each declared capture spec against the body. Two modes:
94///
95/// - **Single** (`[name]`): take the first matching value. For an
96///   array-of-rows body shape (CQL's standard JSON form), reads
97///   row[0].name. For an object body, reads top-level `name`. For
98///   wildcard `*`, captures every top-level field.
99/// - **Slurp** (`[@name]`): walks every row of an array-of-rows
100///   body and collects each row's column into a single
101///   `Value::Json(array)`. Object bodies produce a single-element
102///   list. This is the convenient shape for downstream consumers
103///   that need all per-row values as a list (e.g. recall
104///   evaluator's `actual:` reads).
105///
106/// The body's `.to_json()` form is the source of truth — adapters
107/// that produce typed-row data render to a JSON array of row
108/// objects.
109#[cfg(test)]
110fn extract_captures_from_json(
111    body: &dyn crate::adapter::ResultBody,
112    specs: &[bindpoints::CapturePoint],
113) -> HashMap<String, polydat::ast::Value> {
114    extract_captures_rooted(body, specs, None)
115}
116
117/// As [`extract_captures_from_json`], with an optional `traverse.path:` base
118/// pointer applied to the body FIRST.
119///
120/// Re-roots the document rather than prefixing each capture: every capture
121/// then addresses relative to the base, which is the whole point of the knob
122/// (`{value, status, …}` envelopes force the same prefix onto every capture).
123/// A base that does not resolve yields no captures at all — the caller's
124/// `on_missing` policy decides whether that is worth saying out loud.
125fn extract_captures_rooted(
126    body: &dyn crate::adapter::ResultBody,
127    specs: &[bindpoints::CapturePoint],
128    base: Option<&str>,
129) -> HashMap<String, polydat::ast::Value> {
130    if specs.is_empty() {
131        return HashMap::new();
132    }
133    let full = body.to_json();
134    let json = match base.filter(|b| !b.is_empty()) {
135        None => full,
136        Some(b) => match full.pointer(b) {
137            Some(sub) => sub.clone(),
138            None => return HashMap::new(),
139        },
140    };
141    let mut captures = HashMap::new();
142    for spec in specs {
143        // Declarative `capture:` block form: JSON-Pointer path
144        // takes precedence over the bracket-source-name path.
145        // Lets a workload address Jolokia bulk-POST responses
146        // (`[{value:N}, {value:[...]}, {value:K}]`) by index +
147        // nested field without re-shaping the response.
148        if let Some(path) = spec.path.as_deref() {
149            let sub = json.pointer(path);
150            let value = if spec.count {
151                polydat::ast::Value::U64(count_of_subtree(sub))
152            } else if let Some(agg) = &spec.agg {
153                aggregate_rows(sub, agg, spec.row_filter.as_ref())
154            } else {
155                match sub {
156                    Some(v) => json_subtree_to_value(v),
157                    None => polydat::ast::Value::None,
158                }
159            };
160            captures.insert(spec.as_name.clone(), value);
161            continue;
162        }
163        if spec.slurp {
164            // Slurp form: collect across all rows.
165            let collected = slurp_column(&json, &spec.source_name);
166            captures.insert(
167                spec.as_name.clone(),
168                polydat::ast::Value::Json(std::sync::Arc::new(serde_json::Value::Array(collected))),
169            );
170            continue;
171        }
172        // Single form.
173        if spec.source_name == "*" {
174            // Wildcard: capture every top-level field. Falls
175            // through to scalar-form per field.
176            let target = match &json {
177                serde_json::Value::Array(rows) => {
178                    rows.first().cloned().unwrap_or(serde_json::Value::Null)
179                }
180                other => other.clone(),
181            };
182            if let serde_json::Value::Object(map) = target {
183                for (k, v) in map {
184                    captures.insert(k, json_to_value(&v));
185                }
186            }
187            continue;
188        }
189        if let Some(val) = first_row_field(&json, &spec.source_name) {
190            captures.insert(spec.as_name.clone(), json_to_value(&val));
191        }
192    }
193    captures
194}
195
196/// Reduce an addressed JSON sub-tree to a u64 count, mirroring
197/// the polling wrapper's `count_from_json_pointer` semantics:
198/// array → length, object → key count, scalar → 1 (non-empty)
199/// or 0 (false / empty string / zero number), null / missing →
200/// 0. Use cases: `capture: { active_count: "/value:count" }`
201/// reduces a list-of-running-jobs to a numeric gauge in one
202/// step.
203fn count_of_subtree(v: Option<&serde_json::Value>) -> u64 {
204    let Some(v) = v else { return 0 };
205    match v {
206        serde_json::Value::Array(a) => a.len() as u64,
207        serde_json::Value::Object(m) => m.len() as u64,
208        serde_json::Value::Number(n) => n
209            .as_u64()
210            .or_else(|| n.as_i64().map(|i| i.max(0) as u64))
211            .or_else(|| n.as_f64().map(|f| f.max(0.0) as u64))
212            .unwrap_or(0),
213        serde_json::Value::Bool(b) => {
214            if *b {
215                1
216            } else {
217                0
218            }
219        }
220        serde_json::Value::String(s) if s.is_empty() => 0,
221        serde_json::Value::String(_) => 1,
222        serde_json::Value::Null => 0,
223    }
224}
225
226/// Project a JSON sub-tree into a Polydat [`Value`]. Scalars go
227/// through [`json_to_value`]'s typed coercion; structural
228/// shapes (array / object) are kept as `Value::Json` so the
229/// kernel can carry the original shape without lossy
230/// stringification.
231/// Folds one numeric field across the rows of an addressed array
232/// (`:min(f)` / `:max(f)` / `:sum(f)` capture aggregation). Rows
233/// missing the field, and non-numeric values, are skipped; an
234/// empty fold yields `Value::None` (renders like a missing
235/// capture). Sum stays integral (`U64`) while every contributing
236/// value is a non-negative integer, else widens to `F64`; min/max
237/// return the winning row's value with its original JSON type.
238/// The capture values an op should publish when it SUCCEEDED but returned no
239/// body at all — an empty CQL result set, an HTTP 204, an accepted timeout.
240///
241/// Skipping the write (the original behaviour) leaves every capture holding
242/// its PREVIOUS value, so a wire that is being watched for "went to zero" can
243/// never read zero once it has been nonzero. That pinned a compaction drain's
244/// task count above zero permanently and hung the poll: cqlsh showed an empty
245/// table while the workload's wire still read the last nonzero count.
246///
247/// An empty measurement is a measurement, and every fold is a commutative
248/// monoid — the empty fold publishes the monoid's identity, TOTALLY:
249///
250/// - `count` / `sum` — the identity is mathematically forced: **0**.
251/// - `min` / `max` — over an unbounded domain the identity is an extreme
252///   (⊤ / ⊥) that would poison displays and metric series; the AUTHOR's
253///   identity is the wire's own declared initial value. These return
254///   `Value::None` here as a *reset marker*, and [`publish_capture`]
255///   translates it into `wires.reset(name)` — the wire returns to its
256///   declared unit instead of holding a raw `None` that every downstream
257///   consumer would have to be None-aware of.
258///
259/// Folding any later real observation into a reset wire behaves as folding
260/// into the unit — the monoid law the drain polls and the SRD-93 pressure
261/// servo both lean on.
262fn captures_for_empty_result(
263    specs: &[bindpoints::CapturePoint],
264) -> std::collections::HashMap<String, polydat::ast::Value> {
265    use bindpoints::CaptureAgg;
266    let mut out = std::collections::HashMap::new();
267    for spec in specs {
268        let value = if spec.count {
269            polydat::ast::Value::U64(0)
270        } else {
271            match &spec.agg {
272                Some(CaptureAgg::Sum(_)) => polydat::ast::Value::U64(0),
273                Some(CaptureAgg::Min(_)) | Some(CaptureAgg::Max(_)) => polydat::ast::Value::None,
274                None => polydat::ast::Value::None,
275            }
276        };
277        out.insert(spec.as_name.clone(), value);
278    }
279    out
280}
281
282/// The single capture-publish chokepoint: a real value writes; the
283/// `Value::None` reset marker (an empty min/max fold, or a plain capture
284/// that resolved to nothing) RESTORES the wire to its declared initial
285/// value. No path parks a raw `None` on a typed wire.
286fn publish_capture(ctx: &crate::adapter::ExecCtx<'_>, name: &str, value: polydat::ast::Value) {
287    if matches!(value, polydat::ast::Value::None) {
288        let _ = ctx.wires.reset(name);
289    } else {
290        let _ = ctx.wires.write(name, value);
291    }
292}
293
294fn aggregate_rows(
295    sub: Option<&serde_json::Value>,
296    agg: &bindpoints::CaptureAgg,
297    row_filter: Option<&(String, String)>,
298) -> polydat::ast::Value {
299    use bindpoints::CaptureAgg;
300    let rows = match sub {
301        Some(serde_json::Value::Array(rows)) => rows,
302        _ => return polydat::ast::Value::None,
303    };
304    let field = match agg {
305        CaptureAgg::Min(f) | CaptureAgg::Max(f) | CaptureAgg::Sum(f) => f.as_str(),
306    };
307    // `where <field>='<value>'` — drop rows that are not commensurable before
308    // folding. A result set can mix them: `system_views.sstable_tasks` lists a
309    // data compaction in `unit=bytes` beside an index build in
310    // `unit=token range parts`, and summing across both adds unlike
311    // quantities. Compared as strings so it works on any scalar column
312    // without the capture layer needing the column's type.
313    let keep = |row: &serde_json::Value| -> bool {
314        let Some((k, want)) = row_filter else {
315            return true;
316        };
317        match row.get(k.as_str()) {
318            Some(serde_json::Value::String(s)) => s == want,
319            Some(other) => other.to_string().trim_matches('"') == want,
320            None => false,
321        }
322    };
323    let nums: Vec<&serde_json::Number> = rows
324        .iter()
325        .filter(|row| keep(row))
326        .filter_map(|row| row.get(field))
327        .filter_map(|v| v.as_number())
328        .collect();
329    if nums.is_empty() {
330        // Same totality contract as `captures_for_empty_result`, applied to
331        // the rows-present-but-none-match case — the two empty paths used to
332        // disagree (`:sum` gave 0 with no body but None past a `where` that
333        // matched nothing), and a `None` parked on a typed wire forced every
334        // downstream consumer to be None-aware. `sum` publishes its forced
335        // identity 0; `min`/`max` return the `None` RESET MARKER, which
336        // `publish_capture` turns into a reset to the wire's declared
337        // initial value (the author's identity element).
338        return match agg {
339            CaptureAgg::Sum(_) => polydat::ast::Value::U64(0),
340            CaptureAgg::Min(_) | CaptureAgg::Max(_) => polydat::ast::Value::None,
341        };
342    }
343    match agg {
344        CaptureAgg::Sum(_) => {
345            if nums.iter().all(|n| n.as_u64().is_some()) {
346                polydat::ast::Value::U64(nums.iter().map(|n| n.as_u64().unwrap()).sum())
347            } else {
348                polydat::ast::Value::F64(nums.iter().filter_map(|n| n.as_f64()).sum())
349            }
350        }
351        CaptureAgg::Min(_) | CaptureAgg::Max(_) => {
352            let want_min = matches!(agg, CaptureAgg::Min(_));
353            let mut best = nums[0];
354            for n in &nums[1..] {
355                let (a, b) = (
356                    n.as_f64().unwrap_or(f64::NAN),
357                    best.as_f64().unwrap_or(f64::NAN),
358                );
359                if (want_min && a < b) || (!want_min && a > b) {
360                    best = n;
361                }
362            }
363            json_to_value(&serde_json::Value::Number((*best).clone()))
364        }
365    }
366}
367
368fn json_subtree_to_value(v: &serde_json::Value) -> polydat::ast::Value {
369    match v {
370        serde_json::Value::Array(_) | serde_json::Value::Object(_) => {
371            polydat::ast::Value::Json(std::sync::Arc::new(v.clone()))
372        }
373        scalar => json_to_value(scalar),
374    }
375}
376
377/// First-row lookup: for an array body, read `rows[0].name`; for
378/// an object body, read `obj.name`. Returns `None` when the field
379/// isn't present.
380fn first_row_field(json: &serde_json::Value, name: &str) -> Option<serde_json::Value> {
381    match json {
382        serde_json::Value::Array(rows) => rows.first().and_then(|row| row.get(name)).cloned(),
383        serde_json::Value::Object(_) => json.get(name).cloned(),
384        _ => None,
385    }
386}
387
388/// Slurp helper: walk an array body and collect each row's `name`
389/// field. Object bodies produce a single-element list. Non-object,
390/// non-array bodies produce an empty list.
391fn slurp_column(json: &serde_json::Value, name: &str) -> Vec<serde_json::Value> {
392    match json {
393        serde_json::Value::Array(rows) => rows
394            .iter()
395            .filter_map(|row| row.get(name).cloned())
396            .collect(),
397        serde_json::Value::Object(_) => json.get(name).map(|v| vec![v.clone()]).unwrap_or_default(),
398        _ => Vec::new(),
399    }
400}
401
402/// Convert a serde_json::Value to a Polydat Value. JSON `null`
403/// maps to [`Value::None`] (per SRD-74 — None propagates
404/// rather than coercing to the string `"null"`); arrays and
405/// objects stringify only when reached via the row-shape
406/// extraction path (the JSON-Pointer extraction route uses
407/// [`json_subtree_to_value`] which preserves them as
408/// `Value::Json`).
409pub(crate) fn json_to_value(v: &serde_json::Value) -> polydat::ast::Value {
410    match v {
411        serde_json::Value::Null => polydat::ast::Value::None,
412        serde_json::Value::Number(n) => {
413            if let Some(i) = n.as_u64() {
414                polydat::ast::Value::U64(i)
415            } else if let Some(f) = n.as_f64() {
416                polydat::ast::Value::F64(f)
417            } else {
418                polydat::ast::Value::Str(n.to_string().into())
419            }
420        }
421        serde_json::Value::Bool(b) => polydat::ast::Value::Bool(*b),
422        serde_json::Value::String(s) => polydat::ast::Value::Str(s.as_str().into()),
423        other => polydat::ast::Value::Str(other.to_string().into()),
424    }
425}
426
427impl WrappingDispenser for TraversingDispenser {}
428
429impl OpDispenser for TraversingDispenser {
430    fn execute<'a>(
431        &'a self,
432        cycle: u64,
433        ctx: &'a crate::fixture::ExecCtx<'a>,
434    ) -> std::pin::Pin<
435        Box<dyn std::future::Future<Output = Result<OpResult, ExecutionError>> + Send + 'a>,
436    > {
437        Box::pin(async move {
438            // Execute the inner dispenser
439            let result = self.inner.execute(cycle, ctx).await?;
440
441            // Traverse: count elements and bytes
442            if let Some(body) = &result.body {
443                self.stats
444                    .metrics
445                    .result_elements
446                    .inc_by(body.element_count());
447                if let Some(bytes) = body.byte_count() {
448                    self.stats.metrics.result_bytes.inc_by(bytes);
449                }
450            }
451
452            // Extract captures from result if declared. Values land
453            // on the per-fiber kernel's input slot via ctx.wires.write;
454            // wrappers above this layer (e.g. MetricsDispenser) see
455            // them through wires.get on the same cycle.
456            // A SUCCEEDING op that returned no body still measured something:
457            // zero rows. Publish the identity values rather than skipping the
458            // write, or every capture silently keeps its previous reading and a
459            // "wait for it to reach zero" poll can never succeed.
460            if !self.captures.is_empty() && result.body.is_none() {
461                for (name, value) in captures_for_empty_result(&self.captures) {
462                    publish_capture(&ctx, &name, value);
463                }
464            }
465            if !self.captures.is_empty()
466                && let Some(body) = &result.body
467            {
468                let base = self.spec.as_ref().and_then(|s| s.path.as_deref());
469                let extracted = extract_captures_rooted(body.as_ref(), &self.captures, base);
470                // `on_missing:` — a capture that resolved to nothing reads
471                // exactly like one that measured an absence, which is how a
472                // renamed column or a changed schema hides. Default stays
473                // `ignore` so existing workloads are unaffected.
474                let policy = self.spec.as_ref().map(|s| s.on_missing).unwrap_or_default();
475                if !matches!(policy, nmbrs_workload::model::OnMissing::Ignore) {
476                    for cp in &self.captures {
477                        let missing = match extracted.get(&cp.as_name) {
478                            None => true,
479                            Some(polydat::ast::Value::None) => true,
480                            Some(_) => false,
481                        };
482                        if !missing {
483                            continue;
484                        }
485                        let detail = format!(
486                            "op capture '{}' resolved to nothing{}",
487                            cp.as_name,
488                            match base {
489                                Some(b) => format!(" (traverse path '{b}')"),
490                                None => String::new(),
491                            }
492                        );
493                        match policy {
494                            nmbrs_workload::model::OnMissing::Warn => {
495                                crate::diag!(crate::observer::LogLevel::Warn, "{detail}")
496                            }
497                            nmbrs_workload::model::OnMissing::Error => {
498                                return Err(ExecutionError::Op(crate::adapter::AdapterError {
499                                    error_name: "capture_missing".into(),
500                                    message: detail,
501                                    retryable: false,
502                                }));
503                            }
504                            nmbrs_workload::model::OnMissing::Ignore => {}
505                        }
506                    }
507                }
508                for (name, value) in extracted {
509                    publish_capture(&ctx, &name, value);
510                }
511            }
512
513            Ok(result)
514        })
515    }
516    fn inner_dispenser(&self) -> Option<&dyn OpDispenser> {
517        Some(self.inner.as_ref())
518    }
519}
520
521#[cfg(test)]
522mod tests {
523    use super::*;
524    use crate::adapter::ResultBody;
525
526    fn agg_cap(
527        alias: &str,
528        agg: Option<bindpoints::CaptureAgg>,
529        count: bool,
530    ) -> bindpoints::CapturePoint {
531        bindpoints::CapturePoint {
532            row_filter: None,
533            source_name: alias.into(),
534            as_name: alias.into(),
535            cast_type: None,
536            slurp: false,
537            path: Some(String::new()),
538            count,
539            agg,
540        }
541    }
542
543    /// An op that SUCCEEDS with no body measured zero rows, and must publish
544    /// that. Skipping the write leaves every capture holding its previous
545    /// reading, so a wire watched for "went to zero" can never reach zero once
546    /// it has been nonzero — which pinned a compaction drain's task count above
547    /// zero permanently and hung its poll for 275s+ against a 48h timeout,
548    /// while cqlsh showed the table empty.
549    #[test]
550    fn empty_result_publishes_identity_for_count_and_sum() {
551        let specs = vec![
552            agg_cap("n", None, true),
553            agg_cap("s", Some(bindpoints::CaptureAgg::Sum("x".into())), false),
554        ];
555        let out = super::captures_for_empty_result(&specs);
556        assert!(
557            matches!(out.get("n"), Some(polydat::ast::Value::U64(0))),
558            "counting nothing is 0, not 'leave the old count': {:?}",
559            out.get("n")
560        );
561        assert!(
562            matches!(out.get("s"), Some(polydat::ast::Value::U64(0))),
563            "summing nothing is 0: {:?}",
564            out.get("s")
565        );
566    }
567
568    /// min/max of nothing has NO honest value, so they publish None — which
569    /// clears the wire rather than leaving a stale reading that is
570    /// indistinguishable from a real observation. Writing 0 would be a lie:
571    /// zero is a plausible minimum.
572    #[test]
573    fn empty_result_clears_min_and_max_rather_than_inventing_a_value() {
574        let specs = vec![
575            agg_cap("lo", Some(bindpoints::CaptureAgg::Min("x".into())), false),
576            agg_cap("hi", Some(bindpoints::CaptureAgg::Max("x".into())), false),
577        ];
578        let out = super::captures_for_empty_result(&specs);
579        assert!(
580            matches!(out.get("lo"), Some(polydat::ast::Value::None)),
581            "min of nothing must not invent 0: {:?}",
582            out.get("lo")
583        );
584        assert!(
585            matches!(out.get("hi"), Some(polydat::ast::Value::None)),
586            "max of nothing must not invent 0: {:?}",
587            out.get("hi")
588        );
589    }
590
591    /// A plain (non-aggregate) capture over an empty result is simply absent.
592    #[test]
593    fn empty_result_leaves_plain_captures_absent() {
594        let specs = vec![agg_cap("v", None, false)];
595        let out = super::captures_for_empty_result(&specs);
596        assert!(matches!(out.get("v"), Some(polydat::ast::Value::None)));
597    }
598
599    /// Rows that EXIST but match nothing stay `None`, deliberately — a
600    /// different statement from "there was no result at all", and the
601    /// pre-existing rule (`aggregate_filter_matching_no_rows_yields_none`).
602    /// That case was never the staleness bug: this value IS written to the
603    /// wire, so it clears. Only the no-body path skipped the write.
604    #[test]
605    fn empty_fold_totality_sum_is_zero_min_max_are_reset_markers() {
606        // The monoid contract (SRD-93-era fold totality): sum-of-nothing
607        // publishes its forced identity 0; min/max-of-nothing return the
608        // `None` RESET MARKER that `publish_capture` translates into a
609        // reset to the wire's declared initial value — no path parks a
610        // raw None on a typed wire, and no path leaves a stale reading.
611        let rows = serde_json::json!([{"other": 1}, {"other": 2}]);
612        let v = super::aggregate_rows(
613            Some(&rows),
614            &bindpoints::CaptureAgg::Sum("missing".into()),
615            None,
616        );
617        assert!(
618            matches!(v, polydat::ast::Value::U64(0)),
619            "sum-of-nothing is its identity 0: {v:?}"
620        );
621        for agg in [
622            bindpoints::CaptureAgg::Min("missing".into()),
623            bindpoints::CaptureAgg::Max("missing".into()),
624        ] {
625            let v = super::aggregate_rows(Some(&rows), &agg, None);
626            assert!(
627                matches!(v, polydat::ast::Value::None),
628                "min/max-of-nothing is the reset marker: {v:?}"
629            );
630        }
631    }
632
633    fn cap(source: &str, alias: &str, slurp: bool) -> bindpoints::CapturePoint {
634        bindpoints::CapturePoint {
635            row_filter: None,
636            source_name: source.into(),
637            as_name: alias.into(),
638            cast_type: None,
639            slurp,
640            path: None,
641            count: false,
642            agg: None,
643        }
644    }
645
646    #[test]
647    fn parse_captures_from_template() {
648        let parsed =
649            bindpoints::parse_capture_points("SELECT [username], [age as user_age] FROM users");
650        assert_eq!(parsed.captures.len(), 2);
651        assert_eq!(parsed.captures[0].source_name, "username");
652        assert_eq!(parsed.captures[0].as_name, "username");
653        assert!(!parsed.captures[0].slurp);
654        assert_eq!(parsed.captures[1].source_name, "age");
655        assert_eq!(parsed.captures[1].as_name, "user_age");
656        assert_eq!(parsed.raw_template, "SELECT username, age FROM users");
657    }
658
659    #[test]
660    fn parse_slurp_capture() {
661        let parsed = bindpoints::parse_capture_points("SELECT [@keys] FROM t");
662        assert_eq!(parsed.captures.len(), 1);
663        assert_eq!(parsed.captures[0].source_name, "keys");
664        assert!(parsed.captures[0].slurp);
665        assert_eq!(parsed.raw_template, "SELECT keys FROM t");
666    }
667
668    #[derive(Debug)]
669    struct JsonBody(serde_json::Value);
670    impl ResultBody for JsonBody {
671        fn to_json(&self) -> serde_json::Value {
672            self.0.clone()
673        }
674        fn as_any(&self) -> &dyn std::any::Any {
675            self
676        }
677    }
678
679    #[test]
680    fn extract_from_json_top_level() {
681        let body = JsonBody(serde_json::json!({
682            "user_id": 42,
683            "name": "alice",
684            "balance": 99.5
685        }));
686        let specs = vec![cap("user_id", "uid", false), cap("name", "name", false)];
687        let captures = extract_captures_from_json(&body, &specs);
688        assert_eq!(captures.len(), 2);
689        assert_eq!(captures["uid"].as_u64(), 42);
690        match &captures["name"] {
691            polydat::ast::Value::Str(s) => assert_eq!(&**s, "alice"),
692            other => panic!("expected Str, got {other:?}"),
693        }
694    }
695
696    #[test]
697    fn extract_wildcard() {
698        let body = JsonBody(serde_json::json!({"a": 1, "b": 2}));
699        let specs = vec![cap("*", "*", false)];
700        let captures = extract_captures_from_json(&body, &specs);
701        assert_eq!(captures.len(), 2);
702    }
703
704    #[test]
705    fn extract_slurp_array_of_rows() {
706        let body = JsonBody(serde_json::json!([
707            {"key": 4, "value": 0.5},
708            {"key": 17, "value": 0.4},
709            {"key": 42, "value": 0.3},
710        ]));
711        let specs = vec![cap("key", "key", true)];
712        let captures = extract_captures_from_json(&body, &specs);
713        assert_eq!(captures.len(), 1);
714        match &captures["key"] {
715            polydat::ast::Value::Json(arc) => {
716                let serde_json::Value::Array(items) = arc.as_ref() else {
717                    panic!("expected Value::Json(array), got {arc:?}");
718                };
719                assert_eq!(items.len(), 3);
720                assert_eq!(items[0], serde_json::json!(4));
721                assert_eq!(items[1], serde_json::json!(17));
722                assert_eq!(items[2], serde_json::json!(42));
723            }
724            other => panic!("expected Value::Json(array), got {other:?}"),
725        }
726    }
727
728    #[test]
729    fn extract_single_first_row_of_array() {
730        let body = JsonBody(serde_json::json!([
731            {"key": 4}, {"key": 17}, {"key": 42},
732        ]));
733        let specs = vec![cap("key", "first_key", false)];
734        let captures = extract_captures_from_json(&body, &specs);
735        assert_eq!(captures.len(), 1);
736        assert_eq!(captures["first_key"].as_u64(), 4);
737    }
738
739    fn cap_path(name: &str, path: &str, count: bool) -> bindpoints::CapturePoint {
740        bindpoints::CapturePoint {
741            source_name: name.into(),
742            as_name: name.into(),
743            cast_type: None,
744            slurp: false,
745            path: Some(path.into()),
746            count,
747            agg: None,
748            row_filter: None,
749        }
750    }
751
752    /// `traverse.path:` re-roots the body, so a capture addresses relative to
753    /// it instead of repeating the envelope prefix on every line.
754    #[test]
755    fn traverse_path_re_roots_the_document() {
756        let body = JsonBody(serde_json::json!({
757            "status": 200,
758            "value": [{"progress": 7u64}],
759        }));
760        let spec = |name: &str, path: &str| bindpoints::CapturePoint {
761            source_name: name.into(),
762            as_name: name.into(),
763            cast_type: None,
764            slurp: false,
765            path: Some(path.to_string()),
766            count: false,
767            agg: None,
768            row_filter: None,
769        };
770        // Without the base, the capture must carry the whole prefix.
771        let long = extract_captures_rooted(&body, &[spec("p", "/value/0/progress")], None);
772        assert_eq!(long["p"].as_u64(), 7);
773
774        // With it, the same value is addressed relative to `/value`.
775        let short = extract_captures_rooted(&body, &[spec("p", "/0/progress")], Some("/value"));
776        assert_eq!(short["p"].as_u64(), 7);
777    }
778
779    /// A base path that does not resolve yields no captures — the caller's
780    /// `on_missing` decides whether that is worth saying out loud.
781    #[test]
782    fn an_unresolvable_traverse_path_yields_nothing() {
783        let body = JsonBody(serde_json::json!({"value": 1}));
784        let spec = bindpoints::CapturePoint {
785            source_name: "x".into(),
786            as_name: "x".into(),
787            cast_type: None,
788            slurp: false,
789            path: Some("/x".into()),
790            count: false,
791            agg: None,
792            row_filter: None,
793        };
794        let caps = extract_captures_rooted(&body, &[spec], Some("/nonesuch"));
795        assert!(caps.is_empty(), "unresolvable base must not invent values");
796    }
797
798    /// The same shape, folded PER KIND. Summing across the two rows adds
799    /// bytes to token range parts — dimensionally meaningless, and its
800    /// magnitude depends on how many tasks were registered at sample time.
801    /// `where kind='…'` keeps each fold to one unit.
802    #[test]
803    fn aggregate_captures_filter_rows_by_kind() {
804        let body = JsonBody(serde_json::json!([
805            {"kind": "compaction",            "progress": 47048035855u64},
806            {"kind": "secondary index build", "progress": 4855601u64},
807        ]));
808        fn cap_filtered(
809            name: &str,
810            field: &str,
811            filter: Option<(&str, &str)>,
812        ) -> bindpoints::CapturePoint {
813            bindpoints::CapturePoint {
814                source_name: name.into(),
815                as_name: name.into(),
816                cast_type: None,
817                slurp: false,
818                path: Some(String::new()),
819                count: false,
820                agg: Some(bindpoints::CaptureAgg::Sum(field.into())),
821                row_filter: filter.map(|(k, v)| (k.to_string(), v.to_string())),
822            }
823        }
824        let specs = vec![
825            cap_filtered("all", "progress", None),
826            cap_filtered("index", "progress", Some(("kind", "secondary index build"))),
827            cap_filtered("data", "progress", Some(("kind", "compaction"))),
828        ];
829        let caps = extract_captures_from_json(&body, &specs);
830        let num = |k: &str| caps[k].as_u64();
831        assert_eq!(num("index"), 4_855_601, "index build only");
832        assert_eq!(num("data"), 47_048_035_855, "data compaction only");
833        // The unfiltered fold is the sum of unlike quantities — kept as the
834        // record of what the old capture actually produced.
835        assert_eq!(num("all"), num("index") + num("data"));
836    }
837
838    /// A predicate that matches nothing folds to absent, not to zero — zero
839    /// would read as "measured, and it was none".
840    #[test]
841    fn aggregate_filter_matching_no_rows_folds_to_sum_identity() {
842        let body = JsonBody(serde_json::json!([
843            {"kind": "compaction", "progress": 5u64},
844        ]));
845        let specs = vec![bindpoints::CapturePoint {
846            source_name: "x".into(),
847            as_name: "x".into(),
848            cast_type: None,
849            slurp: false,
850            path: Some(String::new()),
851            count: false,
852            agg: Some(bindpoints::CaptureAgg::Sum("progress".into())),
853            row_filter: Some(("kind".into(), "nonesuch".into())),
854        }];
855        let caps = extract_captures_from_json(&body, &specs);
856        // Fold totality: a `where` that matches nothing sums to the
857        // identity 0 — the same statement the no-body path makes, where
858        // the two paths used to disagree (0 vs a stale-prone None).
859        assert!(
860            matches!(caps["x"], polydat::ast::Value::U64(0)),
861            "a filtered-empty sum folds to its identity 0; got {:?}",
862            caps["x"]
863        );
864    }
865
866    #[test]
867    fn aggregate_captures_fold_rows_with_mixed_units() {
868        // The sstable_tasks shape that motivated aggregation: a
869        // byte-denominated task saturated at 1.0 alongside an
870        // ordinal-denominated task mid-flight. min(ratio) must
871        // surface the merge's 0.4, not row 0's 1.0.
872        let body = JsonBody(serde_json::json!([
873            {"completion_ratio": 1.0, "progress": 31357891323u64, "total": 31357891323u64},
874            {"completion_ratio": 0.4, "progress": 8000000u64,     "total": 20000000u64},
875        ]));
876        fn cap_agg(name: &str, agg: bindpoints::CaptureAgg) -> bindpoints::CapturePoint {
877            bindpoints::CapturePoint {
878                source_name: name.into(),
879                as_name: name.into(),
880                cast_type: None,
881                slurp: false,
882                path: Some(String::new()),
883                count: false,
884                agg: Some(agg),
885                row_filter: None,
886            }
887        }
888        let specs = vec![
889            cap_agg(
890                "completion_ratio",
891                bindpoints::CaptureAgg::Min("completion_ratio".into()),
892            ),
893            cap_agg(
894                "max_ratio",
895                bindpoints::CaptureAgg::Max("completion_ratio".into()),
896            ),
897            cap_agg("progress", bindpoints::CaptureAgg::Sum("progress".into())),
898        ];
899        let captures = extract_captures_from_json(&body, &specs);
900        assert_eq!(captures["completion_ratio"].as_f64(), 0.4);
901        assert_eq!(captures["max_ratio"].as_f64(), 1.0);
902        assert_eq!(captures["progress"].as_u64(), 31357891323 + 8000000);
903        // Empty result set folds to None (missing-capture semantics).
904        let empty = JsonBody(serde_json::json!([]));
905        let c2 = extract_captures_from_json(&empty, &specs[..1]);
906        assert!(matches!(c2["completion_ratio"], polydat::ast::Value::None));
907    }
908
909    #[test]
910    fn extract_json_pointer_scalar_from_bulk_response() {
911        let body = JsonBody(serde_json::json!([
912            {"value": 7,       "status": 200},
913            {"value": [],      "status": 200},
914            {"value": 0,       "status": 200},
915        ]));
916        let specs = vec![
917            cap_path("sstables", "/0/value", false),
918            cap_path("pending_for_cf", "/2/value", false),
919        ];
920        let captures = extract_captures_from_json(&body, &specs);
921        assert_eq!(captures["sstables"].as_u64(), 7);
922        assert_eq!(captures["pending_for_cf"].as_u64(), 0);
923    }
924
925    #[test]
926    fn extract_json_pointer_resolved_null_yields_none_not_string() {
927        let body = JsonBody(serde_json::json!([
928            {"value": null, "status": 200},
929        ]));
930        let specs = vec![cap_path("pending_for_cf", "/0/value", false)];
931        let captures = extract_captures_from_json(&body, &specs);
932        assert!(
933            matches!(captures["pending_for_cf"], polydat::ast::Value::None),
934            "JSON null at path should yield Value::None, got {:?}",
935            captures["pending_for_cf"],
936        );
937    }
938
939    #[test]
940    fn extract_json_pointer_count_on_resolved_null_returns_zero() {
941        let body = JsonBody(serde_json::json!([
942            {"value": null, "status": 200},
943        ]));
944        let specs = vec![cap_path("pending_for_cf", "/0/value", true)];
945        let captures = extract_captures_from_json(&body, &specs);
946        assert_eq!(
947            captures["pending_for_cf"].as_u64(),
948            0,
949            "`:count` on resolved-null should return 0, got {:?}",
950            captures["pending_for_cf"],
951        );
952    }
953
954    #[test]
955    fn extract_json_pointer_count_collapses_array_to_length() {
956        let body = JsonBody(serde_json::json!([
957            {"value": 7},
958            {"value": [
959                {"compactionId":"a", "keyspace":"ks", "columnfamily":"cf"},
960                {"compactionId":"b", "keyspace":"ks", "columnfamily":"cf"},
961            ]},
962        ]));
963        let specs = vec![cap_path("active_count", "/1/value", true)];
964        let captures = extract_captures_from_json(&body, &specs);
965        assert_eq!(captures["active_count"].as_u64(), 2);
966    }
967
968    #[test]
969    fn extract_json_pointer_missing_path_yields_none() {
970        let body = JsonBody(serde_json::json!({"value": 7}));
971        let specs = vec![cap_path("not_there", "/missing/path", false)];
972        let captures = extract_captures_from_json(&body, &specs);
973        assert!(
974            matches!(captures["not_there"], polydat::ast::Value::None),
975            "expected Value::None for unresolvable JSON-Pointer, got {:?}",
976            captures["not_there"],
977        );
978    }
979
980    #[test]
981    fn extract_json_pointer_count_of_missing_path_is_zero() {
982        let body = JsonBody(serde_json::json!({"value": 7}));
983        let specs = vec![cap_path("active_count", "/missing/value", true)];
984        let captures = extract_captures_from_json(&body, &specs);
985        assert_eq!(captures["active_count"].as_u64(), 0);
986    }
987
988    #[test]
989    fn extract_json_pointer_structural_sub_tree_captured_as_json() {
990        let body = JsonBody(serde_json::json!([
991            {"value": {"keyspace": "ks", "table": "cf", "ssTables": 3}},
992        ]));
993        let specs = vec![cap_path("state", "/0/value", false)];
994        let captures = extract_captures_from_json(&body, &specs);
995        match &captures["state"] {
996            polydat::ast::Value::Json(arc) => {
997                assert_eq!(arc.get("keyspace").and_then(|v| v.as_str()), Some("ks"));
998                assert_eq!(arc.get("ssTables").and_then(|v| v.as_u64()), Some(3));
999            }
1000            other => panic!("expected Value::Json, got {other:?}"),
1001        }
1002        // Silence dead-code warnings if helpers are unused in some subset
1003        let _ = count_of_subtree;
1004        let _ = json_subtree_to_value;
1005    }
1006}