Skip to main content

nmbrs_runtime/wrappers/
result.rs

1// Copyright 2024-2026 Jonathan Shook
2// SPDX-License-Identifier: Apache-2.0
3
4//! Result-as-GK adapter (SRD-40b §5). After the inner adapter
5//! returns its `OpResult`, this wrapper exposes declared
6//! op-result fields plus the magic externs (`body`, `count`,
7//! `ok`) as Polydat named wires on the per-fiber op-template kernel
8//! via `ctx.wires.write`. Sits between the inner adapter and
9//! the metrics layer in the wrapper stack.
10
11use std::sync::Arc;
12
13use super::traverse::json_to_value;
14use crate::adapter::WrappingDispenser;
15use crate::adapter::{ExecutionError, OpDispenser, OpResult};
16use crate::wrapper_registry::{WrapperName, WrapperRegistration, WrapperSubject};
17
18/// SRD-32a wrapper name. Always present — no-op when the op
19/// declares no `result:` wires.
20pub const NAME: WrapperName = WrapperName::new("result");
21
22fn triggers(s: WrapperSubject) -> bool {
23    s.op().is_some()
24}
25
26fn describe_assignment(s: WrapperSubject) -> Option<String> {
27    let template = s.op()?;
28    let spec = template.result.as_ref()?;
29    if spec.is_empty() {
30        return None;
31    }
32    let mut names: Vec<String> = Vec::new();
33    spec.walk_fragments(|frag| match frag {
34        nmbrs_workload::model::ResultFragment::Named { name, .. } => {
35            names.push(name.to_string());
36        }
37        nmbrs_workload::model::ResultFragment::Source(source) => {
38            for line in source.lines() {
39                let line = line.trim();
40                if let Some((name, _)) = line.split_once(":=") {
41                    names.push(name.trim().to_string());
42                }
43            }
44        }
45    });
46    if names.is_empty() {
47        return None;
48    }
49    names.sort();
50    names.dedup();
51    Some(format!("result: captures {}", names.join(", ")))
52}
53
54inventory::submit! {
55    WrapperRegistration {
56        name: NAME,
57        // `result:` is parsed into ParsedOp.result, not into
58        // params, so this wrapper has no owned `params`-keys
59        // to declare — the trigger always fires (the cascade
60        // wraps unconditionally, no-op when result map is
61        // empty).
62        owned_fields: &[],
63        triggers,
64        requires_inner: &[super::traverse::NAME],
65        forbids_outer: &[],
66        mutually_exclusive_with: &[],
67        describe_assignment,
68        levels: &[crate::wrapper_registry::WrapperLevel::Op],
69    }
70}
71
72/// Wraps an inner OpDispenser to expose declared op-result fields
73/// as Polydat named wires (SRD-40b §5).
74///
75/// Per cycle, after the inner adapter returns its `OpResult`, this
76/// wrapper walks the op template's `result: HashMap<String,
77/// ResultWireSpec>` declarations, computes each value from the
78/// result body, and writes it through `ctx.wires.write(name, value)`
79/// onto the per-fiber op-template kernel's input slot. Wrappers
80/// later in the stack (e.g. `MetricsDispenser`) read freshly through
81/// `ctx.wires.get(name)`, so SRD-40b §5.2 metric evaluation sees
82/// the values landed this cycle — no HashMap intermediary.
83///
84/// Insertion order in the wrapper stack (SRD-40b §5.2): inner
85/// adapter → ResultDispenser → MetricsDispenser (Phase E).
86///
87/// Source grammars (SRD-40b §5.1):
88/// - `count` — built-in; `OpResult::body.element_count()`.
89/// - `ok` — built-in; `true` iff the inner adapter returned `Ok(_)`
90///   (errors short-circuit before this wrapper's body runs, so
91///   reaching this code already implies success — `ok` is `true`
92///   unless the result was a skip).
93/// - `<path-expr>` — JSON-pointer-style lookup into the result
94///   body. Supports bare names (`field`), dotted paths
95///   (`rows.0.field`), and bracketed indices (`rows[0].field`).
96/// - `<polydat-call>` — DEFERRED. Recognized as anything containing a
97///   `(` token; currently logged once and skipped. Phase E or a
98///   follow-up adds the GK-eval-against-result-context path.
99pub struct ResultDispenser {
100    inner: Arc<dyn OpDispenser>,
101    /// Map-shape `count` / `ok` / path-expr declarations that
102    /// stay on the dispenser's evaluator (SRD-40b §5.1
103    /// backwards compat). SRD-66 string-shape and polydat-call
104    /// entries are NOT here — they're compiled into the
105    /// op-template kernel's body via SRD-67 Phase 5
106    /// `add_result_bindings` and evaluated by GK; the dispenser
107    /// just feeds them inputs through the
108    /// `populate_kernel_inputs` flag.
109    specs: Vec<ResultSlot>,
110    /// SRD-67 Phase 5 — when the op's `result:` source contains
111    /// any string-shape or polydat-call entries, the dispenser writes
112    /// the magic pre-bound inputs (`body` / `count` / `ok`)
113    /// through `ctx.wires.write` onto the op-template kernel's
114    /// input slots before result-binding expressions evaluate.
115    /// The kernel's closure-binding economy returns `NoSlot` for
116    /// slots it doesn't reference, so unconditional writes are
117    /// safe. Under `KernelOptLevel::Diagnostic` every slot is
118    /// allocated regardless.
119    populate_kernel_inputs: bool,
120}
121
122/// Parsed form of one `result:` declaration.
123struct ResultSlot {
124    /// Wire name (the map key in `ParsedOp.result`). Drives the
125    /// `ctx.wires.write(name, …)` call inside `execute` — the
126    /// kernel's matching input slot receives the value.
127    wire: String,
128    /// Decoded source grammar.
129    source: ResultSource,
130    /// Optional default rendered as a Polydat Value (string fallback)
131    /// when the source resolves to nothing.
132    default: Option<polydat::ast::Value>,
133    /// SRD-109 Part 3 — the declared interface type when this wire
134    /// fills an `abstract: results:` contract. Drives wildcard
135    /// projection coercion (the projection lands as exactly this
136    /// type, empty form included) instead of the inference ladder.
137    target: Option<polydat::ast::PortType>,
138}
139
140/// Decoded SRD-40b §5.1 source grammar.
141enum ResultSource {
142    /// `count` — element count of the result body.
143    Count,
144    /// `ok` — success boolean.
145    Ok,
146    /// `<path-expr>` — JSON path into the result body, pre-parsed
147    /// into segments.
148    Path(Vec<PathSeg>),
149    /// `<polydat-call>` — deferred. Carries the raw source for the
150    /// follow-up implementation.
151    #[allow(dead_code)]
152    PolydatCall(String),
153}
154
155/// One segment of a parsed path expression.
156#[derive(Debug, Clone)]
157pub(crate) enum PathSeg {
158    Field(String),
159    Index(usize),
160    /// SRD-70 `[*]` — column projection across every element of
161    /// the array at this position. At most one per path.
162    Wildcard,
163}
164
165/// Parse `rows[0].field` / `rows.0.field` / `field` into segments.
166/// Returns `Err` for empty paths only — anything else parses as a
167/// best-effort sequence of identifiers + indices.
168pub(crate) fn parse_path_expr(src: &str) -> Result<Vec<PathSeg>, String> {
169    let trimmed = src.trim();
170    if trimmed.is_empty() {
171        return Err("empty path".into());
172    }
173    let mut segs = Vec::new();
174    let mut cur = String::new();
175    let mut iter = trimmed.chars().peekable();
176    let push_field = |segs: &mut Vec<PathSeg>, cur: &mut String| {
177        if !cur.is_empty() {
178            // Numeric bareword (after a dot) becomes an index; a
179            // bare `*` bareword is the SRD-70 wildcard.
180            if cur == "*" {
181                segs.push(PathSeg::Wildcard);
182            } else if let Ok(n) = cur.parse::<usize>() {
183                segs.push(PathSeg::Index(n));
184            } else {
185                segs.push(PathSeg::Field(std::mem::take(cur)));
186            }
187            cur.clear();
188        }
189    };
190    while let Some(&c) = iter.peek() {
191        match c {
192            '.' => {
193                push_field(&mut segs, &mut cur);
194                iter.next();
195            }
196            '[' => {
197                push_field(&mut segs, &mut cur);
198                iter.next();
199                let mut idx = String::new();
200                for c2 in iter.by_ref() {
201                    if c2 == ']' {
202                        break;
203                    }
204                    idx.push(c2);
205                }
206                if idx.trim() == "*" {
207                    segs.push(PathSeg::Wildcard);
208                } else {
209                    let n: usize = idx
210                        .trim()
211                        .parse()
212                        .map_err(|_| format!("path '{src}': invalid index '[{idx}]'"))?;
213                    segs.push(PathSeg::Index(n));
214                }
215            }
216            _ => {
217                cur.push(c);
218                iter.next();
219            }
220        }
221    }
222    push_field(&mut segs, &mut cur);
223    if segs.is_empty() {
224        return Err(format!("path '{src}': no segments"));
225    }
226    // SRD-70 first wave: one wildcard per path. Two-level nested
227    // projection is parked until a use case drives it.
228    let wildcards = segs
229        .iter()
230        .filter(|s| matches!(s, PathSeg::Wildcard))
231        .count();
232    if wildcards > 1 {
233        return Err(format!(
234            "path '{src}': at most one [*] wildcard per path \
235             (nested projection is not supported)"
236        ));
237    }
238    Ok(segs)
239}
240
241/// Split a parsed path at its wildcard, if any: `(prefix,
242/// suffix)` — segments before and after the `[*]`.
243fn wildcard_split(segs: &[PathSeg]) -> Option<(&[PathSeg], &[PathSeg])> {
244    let pos = segs.iter().position(|s| matches!(s, PathSeg::Wildcard))?;
245    Some((&segs[..pos], &segs[pos + 1..]))
246}
247
248/// SRD-70 column projection: walk `prefix` to the array, then
249/// apply `suffix` to each element, collecting every match.
250/// Elements whose suffix misses are dropped from the projection
251/// (SRD-70: "skip / drop the row"). Returns `None` when the
252/// prefix itself misses or lands on a non-array — the projection
253/// target does not exist in this body.
254fn resolve_path_projection<'a>(
255    json: &'a serde_json::Value,
256    prefix: &[PathSeg],
257    suffix: &[PathSeg],
258) -> Option<Vec<&'a serde_json::Value>> {
259    let at = resolve_path(json, prefix)?;
260    let arr = at.as_array()?;
261    Some(
262        arr.iter()
263            .filter_map(|el| resolve_path(el, suffix))
264            .collect(),
265    )
266}
267
268/// Coerce one JSON leaf to i64: native integer, or numeric string.
269fn leaf_as_i64(v: &serde_json::Value) -> Option<i64> {
270    v.as_i64().or_else(|| v.as_str()?.trim().parse().ok())
271}
272
273/// Coerce one JSON leaf to f64: native number, or numeric string.
274fn leaf_as_f64(v: &serde_json::Value) -> Option<f64> {
275    v.as_f64().or_else(|| v.as_str()?.trim().parse().ok())
276}
277
278/// Evaluate a parsed path against a result-body JSON: wildcard
279/// paths run the SRD-70 column projection (always landing a
280/// value — an empty body / short array projects to the empty
281/// form of the wire's type, never a stale previous-cycle read);
282/// scalar paths walk to a single value (`None` on miss). Shared
283/// with the validation wrapper, which evaluates `actual:`
284/// projections directly from the op result — it sits INSIDE the
285/// `result` wrapper in the cascade, before the wire write lands.
286pub(crate) fn evaluate_path_value(
287    json: &serde_json::Value,
288    segs: &[PathSeg],
289    target: Option<polydat::ast::PortType>,
290) -> Option<polydat::ast::Value> {
291    if let Some((prefix, suffix)) = wildcard_split(segs) {
292        let collected = resolve_path_projection(json, prefix, suffix).unwrap_or_default();
293        return Some(coerce_projection(&collected, target));
294    }
295    resolve_path(json, segs).map(json_to_value)
296}
297
298/// Type a collected projection as a Polydat value. With a declared
299/// `target` (an SRD-109 `results:` interface type) the column
300/// coerces to exactly that type — empty projection included — and
301/// non-coercible elements are dropped per SRD-70. Without a
302/// target, the inference ladder applies: all-integer → `VecI64`,
303/// all-float → `VecF64`, anything else → the collected array as
304/// `Json`.
305fn coerce_projection(
306    collected: &[&serde_json::Value],
307    target: Option<polydat::ast::PortType>,
308) -> polydat::ast::Value {
309    use polydat::ast::{PortType, SliceArc, Value};
310    let json_array = || {
311        Value::Json(std::sync::Arc::new(serde_json::Value::Array(
312            collected.iter().map(|v| (*v).clone()).collect(),
313        )))
314    };
315    match target {
316        Some(PortType::VecI64) => Value::VecI64(SliceArc::from_vec(
317            collected.iter().filter_map(|v| leaf_as_i64(v)).collect(),
318        )),
319        Some(PortType::VecI32) => Value::VecI32(SliceArc::from_vec(
320            collected
321                .iter()
322                .filter_map(|v| leaf_as_i64(v).map(|n| n as i32))
323                .collect(),
324        )),
325        Some(PortType::VecF64) => Value::VecF64(SliceArc::from_vec(
326            collected.iter().filter_map(|v| leaf_as_f64(v)).collect(),
327        )),
328        Some(PortType::VecF32) => Value::VecF32(SliceArc::from_vec(
329            collected
330                .iter()
331                .filter_map(|v| leaf_as_f64(v).map(|n| n as f32))
332                .collect(),
333        )),
334        Some(PortType::Json) => json_array(),
335        // A declared scalar/other type on a wildcard projection is
336        // rejected at load by the binder; reaching here means an
337        // untyped (non-interface) wire — fall through to inference.
338        _ => {
339            let ints: Vec<i64> = collected.iter().filter_map(|v| leaf_as_i64(v)).collect();
340            if !collected.is_empty() && ints.len() == collected.len() {
341                return Value::VecI64(SliceArc::from_vec(ints));
342            }
343            let floats: Vec<f64> = collected.iter().filter_map(|v| leaf_as_f64(v)).collect();
344            if !collected.is_empty() && floats.len() == collected.len() {
345                return Value::VecF64(SliceArc::from_vec(floats));
346            }
347            json_array()
348        }
349    }
350}
351
352/// Walk a parsed path against a JSON value. Returns `None` when
353/// any segment misses (object lacks key, array shorter than index,
354/// scalar where a container was expected).
355fn resolve_path<'a>(
356    json: &'a serde_json::Value,
357    segs: &[PathSeg],
358) -> Option<&'a serde_json::Value> {
359    let mut cur = json;
360    for seg in segs {
361        cur = match (cur, seg) {
362            (serde_json::Value::Object(m), PathSeg::Field(k)) => m.get(k)?,
363            (serde_json::Value::Array(a), PathSeg::Index(i)) => a.get(*i)?,
364            // Bareword field on an array, or numeric on an object —
365            // we don't try to coerce; the path doesn't match.
366            _ => return None,
367        };
368    }
369    Some(cur)
370}
371
372/// Decode one `(name, source-string)` pair from a SRD-66
373/// map-shape `result:` fragment into a `ResultSlot`. Unknown
374/// / unparseable sources land as `None` (caller logs and
375/// drops them — SRD-40b §5.1 calls for "log a warning and
376/// skip" over a hard failure here, since the value mechanism
377/// is supposed to be best-effort per cycle).
378///
379/// String-shape and list-shape `result:` fragments don't go
380/// through this function — they compile into the auxiliary
381/// kernel under SRD-66 §"Compilation lifecycle" (TBD when
382/// the structural-body Value variant lands).
383fn decode_slot(
384    name: &str,
385    raw_source: &str,
386    target: Option<polydat::ast::PortType>,
387) -> Option<ResultSlot> {
388    let raw = raw_source.trim();
389    let source = if raw == "count" {
390        ResultSource::Count
391    } else if raw == "ok" {
392        ResultSource::Ok
393    } else if raw.contains('(') {
394        // SRD-66 Surface 1 — GK-expression form. The full
395        // kernel-driven path (compile auxiliary kernel,
396        // wire body/count/ok externs via the closure-binding
397        // rule, evaluate per-cycle) is staged behind this
398        // diagnostic until the structural-body Value variant
399        // and op-template kernel extension land. Today the
400        // form is recognised but evaluates to its default.
401        crate::diag!(
402            crate::observer::LogLevel::Warn,
403            "result wire '{name}': GK-expression source '{raw}' is not yet \
404             evaluated end-to-end — slot will resolve to its default. \
405             SRD-66 Push 2 follow-up wires the kernel-driven path.",
406        );
407        ResultSource::PolydatCall(raw.to_string())
408    } else {
409        // Path expression. Parse failures degrade to skip.
410        match parse_path_expr(raw) {
411            Ok(segs) => ResultSource::Path(segs),
412            Err(e) => {
413                crate::diag!(
414                    crate::observer::LogLevel::Warn,
415                    "result wire '{name}': source '{raw}' is not parseable as a \
416                     path expression ({e}) — slot will be skipped.",
417                );
418                return None;
419            }
420        }
421    };
422    Some(ResultSlot {
423        wire: name.to_string(),
424        source,
425        default: None,
426        target,
427    })
428}
429
430impl ResultDispenser {
431    /// Wrap an inner dispenser with result-as-GK exposure
432    /// (SRD-40b §5.2's result-as-GK adapter layer).
433    ///
434    /// **Always wraps.** Per the SRD, this layer is part of the
435    /// canonical per-cycle pipeline: it writes the magic externs
436    /// (`body` / `count` / `ok`) into the op-template kernel's
437    /// input slots after the inner adapter returns, so any
438    /// downstream wrapper (metrics, validation, conditional next
439    /// op) reading those names through `ctx.wires.get` sees fresh
440    /// values. Writes to slots the kernel didn't allocate (the
441    /// closure-binding economy's DCE) silently no-op via
442    /// `WriteOutcome::NoSlot` — no overhead for ops whose
443    /// op-template kernel doesn't reference any magic extern.
444    ///
445    /// The optional `result_spec` adds *additional* dispenser-side
446    /// dispatch slots (legacy SRD-40b §5.1 path-expr / `count` /
447    /// `ok` map-shape forms). Kernel-driven entries (string-shape
448    /// source blocks, polydat-call entries) need no per-cycle code
449    /// here — `add_result_bindings` compiled them into the
450    /// op-template kernel; the magic-extern population this
451    /// wrapper always performs is what makes them resolve.
452    pub fn wrap(
453        inner: Arc<dyn OpDispenser>,
454        result_spec: Option<&nmbrs_workload::model::ResultSpec>,
455        interface_results: Option<&std::collections::BTreeMap<String, String>>,
456    ) -> Arc<dyn OpDispenser> {
457        let mut specs: Vec<ResultSlot> = Vec::new();
458
459        // SRD-109 Part 3 — a wire filling an `abstract: results:`
460        // contract projects at exactly its declared type.
461        let target_for = |name: &str| -> Option<polydat::ast::PortType> {
462            interface_results?
463                .get(name)
464                .and_then(|kw| polydat::ast::PortType::from_keyword(kw))
465        };
466
467        if let Some(spec) = result_spec {
468            spec.walk_fragments(|frag| match frag {
469                nmbrs_workload::model::ResultFragment::Named { name, source } => {
470                    let raw = source.trim();
471                    if raw == "count" || raw == "ok" {
472                        if let Some(slot) = decode_slot(name, source, target_for(name)) {
473                            specs.push(slot);
474                        }
475                    } else if !raw.contains('(') {
476                        // Path expression — keep the legacy
477                        // JSON-path path for SRD-40b §5.1 back-compat.
478                        if let Some(slot) = decode_slot(name, source, target_for(name)) {
479                            specs.push(slot);
480                        }
481                    }
482                    // polydat-call entries (raw.contains('(')) are
483                    // kernel-driven; nothing per-cycle to do here.
484                }
485                nmbrs_workload::model::ResultFragment::Source(_source) => {
486                    // String-shape — fully kernel-driven. No
487                    // per-cycle code here.
488                }
489            });
490        }
491
492        // Stable order so wire-resolution warnings (and the
493        // per-cycle insertion order) are reproducible.
494        specs.sort_by(|a, b| a.wire.cmp(&b.wire));
495        Arc::new(Self {
496            inner,
497            specs,
498            // Magic-extern population always fires. The
499            // populate_kernel_inputs field is retained for the
500            // diagnostic-trace conditional below but its value
501            // is now always-true.
502            populate_kernel_inputs: true,
503        })
504    }
505
506    /// Compute the Polydat value for one slot from the cycle's result.
507    /// Returns `None` when the slot resolves to nothing and has no
508    /// default — caller logs at debug and moves on.
509    fn evaluate(slot: &ResultSlot, result: &OpResult) -> Option<polydat::ast::Value> {
510        match &slot.source {
511            ResultSource::Count => {
512                let n = result.body.as_ref().map(|b| b.element_count()).unwrap_or(0);
513                Some(polydat::ast::Value::U64(n))
514            }
515            ResultSource::Ok => {
516                // Reached only on Ok(_) from the inner adapter; a
517                // skipped op also counts as "not a failure" — we
518                // treat skip as ok=true, matching the SRD-40b §5
519                // intent that this is a binary success signal.
520                Some(polydat::ast::Value::Bool(true))
521            }
522            ResultSource::Path(segs) => {
523                let body = result.body.as_ref()?;
524                let json = body.to_json();
525                evaluate_path_value(&json, segs, slot.target).or_else(|| slot.default.clone())
526            }
527            ResultSource::PolydatCall(_) => slot.default.clone(),
528        }
529    }
530}
531
532impl WrappingDispenser for ResultDispenser {}
533
534impl OpDispenser for ResultDispenser {
535    fn execute<'a>(
536        &'a self,
537        cycle: u64,
538        ctx: &'a crate::fixture::ExecCtx<'a>,
539    ) -> std::pin::Pin<
540        Box<dyn std::future::Future<Output = Result<OpResult, ExecutionError>> + Send + 'a>,
541    > {
542        Box::pin(async move {
543            let result = self.inner.execute(cycle, ctx).await?;
544            // Skipped ops carry no body; per SRD-40b §5.2 the
545            // metric pipeline doesn't fire on skips either. The
546            // Phase E wrappers observe `result.skipped` and bail
547            // before evaluating.
548            if result.skipped {
549                return Ok(result);
550            }
551            for slot in &self.specs {
552                if let Some(v) = Self::evaluate(slot, &result) {
553                    // Canonical write — lands directly on the op-template
554                    // kernel's input slot via ctx.wires. Subsequent
555                    // wrapper reads (e.g. MetricsDispenser) see the
556                    // fresh value through wires.get on the same cycle.
557                    let _ = ctx.wires.write(&slot.wire, v);
558                }
559                // Per-cycle missing-wire is silent. If a downstream
560                // consumer (e.g. MetricsDispenser) references a wire
561                // that didn't land on the kernel, that consumer
562                // surfaces the failure as a hard ExecutionError —
563                // logging it here would just add per-cycle session.log
564                // spam without telling the user anything actionable.
565            }
566
567            // SRD-67 Phase 5 — magic-extern population. When the
568            // op declares any kernel-driven result-bindings
569            // (string-shape OR map-shape polydat-call), inject the
570            // standard `body` / `count` / `ok` inputs through
571            // ctx.wires so the op-template kernel's input slots
572            // are populated before any wrapper above this one in
573            // the stack reads them. The closure-binding economy
574            // drops slots the kernel doesn't reference, so this
575            // is safe even if the user's source only references
576            // a subset — NoSlot writes are silently ignored.
577            // Under KernelOptLevel::Diagnostic every magic extern
578            // gets a slot, so all three always land.
579            if self.populate_kernel_inputs {
580                let count = result.body.as_ref().map(|b| b.element_count()).unwrap_or(0);
581                // SRD-66 §"Surface 4 §Open: body type" resolved
582                // to `Value::Json` — body rides the kernel as a
583                // structural value so `exactly_one_value(body)`
584                // can walk row × column shape (per
585                // `polydat::library::exactly_one`). For ops
586                // whose body has no structural projection the
587                // adapter's `to_json()` returns a JSON String,
588                // which `exactly_one_value` collapses to
589                // `Value::Str`.
590                let body_json = result
591                    .body
592                    .as_ref()
593                    .map(|b| b.to_json())
594                    .unwrap_or(serde_json::Value::Null);
595                let _ = ctx.wires.write(
596                    "body",
597                    polydat::ast::Value::Json(std::sync::Arc::new(body_json)),
598                );
599                let _ = ctx.wires.write("count", polydat::ast::Value::U64(count));
600                let _ = ctx.wires.write("ok", polydat::ast::Value::Bool(true));
601            }
602
603            let _ = cycle;
604            Ok(result)
605        })
606    }
607    fn inner_dispenser(&self) -> Option<&dyn OpDispenser> {
608        Some(self.inner.as_ref())
609    }
610}
611
612#[cfg(test)]
613mod tests {
614    use super::*;
615    use crate::adapter::{AdapterError, ExecutionError, OpResult, ResultBody};
616    use crate::fixture::{ExecCtx, ResolvedPulls};
617    use std::sync::Arc;
618
619    #[derive(Debug)]
620    struct ResultDispBody {
621        value: serde_json::Value,
622        count: u64,
623    }
624    impl ResultBody for ResultDispBody {
625        fn to_json(&self) -> serde_json::Value {
626            self.value.clone()
627        }
628        fn as_any(&self) -> &dyn std::any::Any {
629            self
630        }
631        fn element_count(&self) -> u64 {
632            self.count
633        }
634    }
635
636    struct FakeInner {
637        body: Option<ResultDispBody>,
638        error: Option<&'static str>,
639    }
640    impl OpDispenser for FakeInner {
641        fn execute<'a>(
642            &'a self,
643            _cycle: u64,
644            _ctx: &'a ExecCtx<'a>,
645        ) -> std::pin::Pin<
646            Box<dyn std::future::Future<Output = Result<OpResult, ExecutionError>> + Send + 'a>,
647        > {
648            Box::pin(async move {
649                if let Some(msg) = self.error {
650                    return Err(ExecutionError::Op(AdapterError {
651                        error_name: "test".into(),
652                        message: msg.into(),
653                        retryable: false,
654                    }));
655                }
656                Ok(OpResult {
657                    body: self.body.as_ref().map(|b| {
658                        Box::new(ResultDispBody {
659                            value: b.value.clone(),
660                            count: b.count,
661                        }) as Box<dyn ResultBody>
662                    }),
663                    skipped: false,
664                })
665            })
666        }
667    }
668
669    fn empty_ctx() -> (crate::adapter::ResolvedFields, ResolvedPulls) {
670        let fields = crate::adapter::ResolvedFields::new(vec![], vec![]);
671        let pulls = ResolvedPulls::empty();
672        (fields, pulls)
673    }
674
675    fn kernel_with_extern_inputs(names: &[(&str, &str)]) -> crate::scope_kernel::ScopeKernel {
676        let mut src = String::from("input cycle: u64\n");
677        for (n, ty) in names {
678            src.push_str(&format!("extern {n}: {ty}\n"));
679        }
680        let mut k = crate::bindings::compile_scope_kernel(&src, &Default::default())
681            .expect("kernel_with_extern_inputs compile");
682        k.set_inputs(&[0]);
683        k
684    }
685
686    fn run_with_wires(
687        dispenser: Arc<dyn OpDispenser>,
688        kernel: &mut crate::scope_kernel::ScopeKernel,
689    ) -> Result<OpResult, ExecutionError> {
690        let fields = crate::adapter::ResolvedFields::new(vec![], vec![]);
691        let pulls = ResolvedPulls::empty();
692        let cw = crate::wires::CycleWires::new(kernel);
693        let ctx = ExecCtx::with_wires(&fields, &pulls, &cw);
694        let rt = tokio::runtime::Builder::new_current_thread()
695            .build()
696            .unwrap();
697        rt.block_on(dispenser.execute(0, &ctx))
698    }
699
700    #[test]
701    fn parse_path_dotted_and_bracketed_equivalent() {
702        let a = parse_path_expr("rows[0].value").unwrap();
703        let b = parse_path_expr("rows.0.value").unwrap();
704        assert_eq!(a.len(), 3);
705        assert_eq!(b.len(), 3);
706        match (&a[0], &b[0]) {
707            (PathSeg::Field(f1), PathSeg::Field(f2)) => assert_eq!(f1, f2),
708            _ => panic!("expected leading field"),
709        }
710        match (&a[1], &b[1]) {
711            (PathSeg::Index(0), PathSeg::Index(0)) => {}
712            _ => panic!("expected index 0"),
713        }
714    }
715
716    fn map_spec(entries: &[(&str, &str)]) -> nmbrs_workload::model::ResultSpec {
717        let mut m = std::collections::BTreeMap::new();
718        for (k, v) in entries {
719            m.insert((*k).to_string(), (*v).to_string());
720        }
721        nmbrs_workload::model::ResultSpec::Map(m)
722    }
723
724    #[test]
725    fn result_dispenser_count_and_path() {
726        let inner = Arc::new(FakeInner {
727            body: Some(ResultDispBody {
728                value: serde_json::json!({"rows": [{"value": 42}]}),
729                count: 1,
730            }),
731            error: None,
732        });
733        let decl = map_spec(&[("row_count", "count"), ("first_value", "rows[0].value")]);
734
735        let wrapped = ResultDispenser::wrap(inner, Some(&decl), None);
736        let mut kernel = kernel_with_extern_inputs(&[("row_count", "u64"), ("first_value", "u64")]);
737        let _ = run_with_wires(wrapped, &mut kernel).unwrap();
738
739        let cw = crate::wires::CycleWires::new(&mut kernel);
740        let w: &dyn crate::wires::WireSource = &cw;
741        assert_eq!(w.get("row_count").map(|v| v.as_u64()), Some(1));
742        assert_eq!(w.get("first_value").map(|v| v.as_u64()), Some(42));
743    }
744
745    #[test]
746    fn result_dispenser_ok_builtin_on_success() {
747        let inner = Arc::new(FakeInner {
748            body: Some(ResultDispBody {
749                value: serde_json::json!({}),
750                count: 0,
751            }),
752            error: None,
753        });
754        let decl = map_spec(&[("succeeded", "ok")]);
755
756        let wrapped = ResultDispenser::wrap(inner, Some(&decl), None);
757        let mut kernel = kernel_with_extern_inputs(&[("succeeded", "bool")]);
758        let _ = run_with_wires(wrapped, &mut kernel).unwrap();
759
760        let cw = crate::wires::CycleWires::new(&mut kernel);
761        let w: &dyn crate::wires::WireSource = &cw;
762        match w.get("succeeded") {
763            Some(polydat::ast::Value::Bool(b)) => assert!(b),
764            other => panic!("expected Bool(true), got {other:?}"),
765        }
766    }
767
768    #[test]
769    fn result_dispenser_error_propagates_no_capture_write() {
770        let inner = Arc::new(FakeInner {
771            body: None,
772            error: Some("boom"),
773        });
774        let decl = map_spec(&[("succeeded", "ok")]);
775
776        let wrapped = ResultDispenser::wrap(inner, Some(&decl), None);
777        let (fields, pulls) = empty_ctx();
778        let ctx = ExecCtx::new(&fields, &pulls);
779        let rt = tokio::runtime::Builder::new_current_thread()
780            .build()
781            .unwrap();
782        let err = rt.block_on(wrapped.execute(0, &ctx)).unwrap_err();
783        assert!(format!("{err}").contains("boom"));
784    }
785
786    #[test]
787    fn result_dispenser_unresolved_path_skips_silently() {
788        let inner = Arc::new(FakeInner {
789            body: Some(ResultDispBody {
790                value: serde_json::json!({"rows": []}),
791                count: 0,
792            }),
793            error: None,
794        });
795        let decl = map_spec(&[("missing", "rows[0].value")]);
796
797        let wrapped = ResultDispenser::wrap(inner, Some(&decl), None);
798        let mut kernel = kernel_with_extern_inputs(&[("missing", "u64")]);
799        let _ = run_with_wires(wrapped, &mut kernel).unwrap();
800
801        let cw = crate::wires::CycleWires::new(&mut kernel);
802        let w: &dyn crate::wires::WireSource = &cw;
803        assert!(matches!(
804            w.get("missing"),
805            Some(polydat::ast::Value::None) | None
806        ));
807    }
808
809    #[test]
810    fn result_dispenser_always_wraps_per_srd_40b() {
811        let inner: Arc<dyn OpDispenser> = Arc::new(FakeInner {
812            body: Some(ResultDispBody {
813                value: serde_json::json!({}),
814                count: 0,
815            }),
816            error: None,
817        });
818        let inner_ptr = Arc::as_ptr(&inner);
819        let wrapped = ResultDispenser::wrap(inner.clone(), None, None);
820        assert_ne!(
821            Arc::as_ptr(&wrapped),
822            inner_ptr,
823            "ResultDispenser must always wrap so magic-extern population fires"
824        );
825    }
826
827    #[test]
828    fn result_dispenser_skipped_op_writes_no_captures() {
829        struct SkipInner;
830        impl OpDispenser for SkipInner {
831            fn execute<'a>(
832                &'a self,
833                _cycle: u64,
834                _ctx: &'a ExecCtx<'a>,
835            ) -> std::pin::Pin<
836                Box<dyn std::future::Future<Output = Result<OpResult, ExecutionError>> + Send + 'a>,
837            > {
838                Box::pin(async move { Ok(OpResult::skipped()) })
839            }
840        }
841        let decl = map_spec(&[("c", "count")]);
842        let wrapped = ResultDispenser::wrap(Arc::new(SkipInner), Some(&decl), None);
843        let mut kernel = kernel_with_extern_inputs(&[("c", "u64")]);
844        let result = run_with_wires(wrapped, &mut kernel).unwrap();
845        assert!(result.skipped);
846        let cw = crate::wires::CycleWires::new(&mut kernel);
847        let w: &dyn crate::wires::WireSource = &cw;
848        assert!(matches!(w.get("c"), Some(polydat::ast::Value::None) | None));
849    }
850
851    /// Silence dead-code warnings for `decode_slot` when only
852    /// the path-expr branch is exercised by the public tests.
853    #[test]
854    fn decode_slot_reachable() {
855        let slot = decode_slot("c", "count", None).expect("count slot decodes");
856        assert!(matches!(slot.source, ResultSource::Count));
857    }
858
859    // ── SRD-70 wildcard projection ──
860
861    #[test]
862    fn parse_wildcard_forms_and_reject_nested() {
863        for src in [
864            "rows[*].key",
865            "rows.*.key",
866            "[*].key",
867            "[*]",
868            "result[*].id",
869        ] {
870            let segs = parse_path_expr(src).unwrap_or_else(|e| panic!("'{src}' should parse: {e}"));
871            assert_eq!(
872                segs.iter()
873                    .filter(|s| matches!(s, PathSeg::Wildcard))
874                    .count(),
875                1,
876                "'{src}' carries one wildcard"
877            );
878        }
879        let err = parse_path_expr("rows[*].items[*].id").unwrap_err();
880        assert!(err.contains("at most one"), "err: {err}");
881    }
882
883    #[test]
884    fn projection_collects_column_as_target_type() {
885        let inner = Arc::new(FakeInner {
886            body: Some(ResultDispBody {
887                // Zero-padded numeric strings — the vector-suite
888                // key shape.
889                value: serde_json::json!([
890                    {"key": "000000000007"}, {"key": "000000000003"}]),
891                count: 2,
892            }),
893            error: None,
894        });
895        let decl = map_spec(&[("keys", "[*].key")]);
896        let mut iface = std::collections::BTreeMap::new();
897        iface.insert("keys".to_string(), "vec_i64".to_string());
898
899        let wrapped = ResultDispenser::wrap(inner, Some(&decl), Some(&iface));
900        let mut kernel = kernel_with_extern_inputs(&[("keys", "vec_i64")]);
901        let _ = run_with_wires(wrapped, &mut kernel).unwrap();
902
903        let cw = crate::wires::CycleWires::new(&mut kernel);
904        let w: &dyn crate::wires::WireSource = &cw;
905        match w.get("keys") {
906            Some(polydat::ast::Value::VecI64(slice)) => assert_eq!(slice.as_slice(), &[7, 3]),
907            other => panic!("expected VecI64([7,3]), got {other:?}"),
908        }
909    }
910
911    #[test]
912    fn projection_nested_prefix_and_inference_ladder() {
913        let inner = Arc::new(FakeInner {
914            body: Some(ResultDispBody {
915                value: serde_json::json!({"result": [
916                    {"id": 12, "score": 0.9},
917                    {"id": 5,  "score": 0.7}]}),
918                count: 2,
919            }),
920            error: None,
921        });
922        // No interface types — the inference ladder applies.
923        let decl = map_spec(&[("ids", "result[*].id"), ("scores", "result[*].score")]);
924        let wrapped = ResultDispenser::wrap(inner, Some(&decl), None);
925        let mut kernel = kernel_with_extern_inputs(&[("ids", "vec_i64"), ("scores", "vec_f64")]);
926        let _ = run_with_wires(wrapped, &mut kernel).unwrap();
927
928        let cw = crate::wires::CycleWires::new(&mut kernel);
929        let w: &dyn crate::wires::WireSource = &cw;
930        match w.get("ids") {
931            Some(polydat::ast::Value::VecI64(slice)) => assert_eq!(slice.as_slice(), &[12, 5]),
932            other => panic!("expected VecI64, got {other:?}"),
933        }
934        match w.get("scores") {
935            Some(polydat::ast::Value::VecF64(slice)) => assert_eq!(slice.as_slice(), &[0.9, 0.7]),
936            other => panic!("expected VecF64, got {other:?}"),
937        }
938    }
939
940    #[test]
941    fn projection_on_missing_prefix_lands_empty_not_stale() {
942        let inner = Arc::new(FakeInner {
943            body: Some(ResultDispBody {
944                value: serde_json::json!({"error": "no result key"}),
945                count: 0,
946            }),
947            error: None,
948        });
949        let decl = map_spec(&[("keys", "result[*].id")]);
950        let mut iface = std::collections::BTreeMap::new();
951        iface.insert("keys".to_string(), "vec_i64".to_string());
952        let wrapped = ResultDispenser::wrap(inner, Some(&decl), Some(&iface));
953        let mut kernel = kernel_with_extern_inputs(&[("keys", "vec_i64")]);
954        let _ = run_with_wires(wrapped, &mut kernel).unwrap();
955
956        let cw = crate::wires::CycleWires::new(&mut kernel);
957        let w: &dyn crate::wires::WireSource = &cw;
958        match w.get("keys") {
959            Some(polydat::ast::Value::VecI64(slice)) => assert!(
960                slice.as_slice().is_empty(),
961                "missing projection target must land the EMPTY typed form"
962            ),
963            other => panic!("expected empty VecI64, got {other:?}"),
964        }
965    }
966}