Skip to main content

nmbrs_runtime/
fixture.rs

1// Copyright 2024-2026 Jonathan Shook
2// SPDX-License-Identifier: Apache-2.0
3
4//! Init-time scope fixture and cycle-time pull plan.
5//!
6//! Implements SRD 32 §"Init-Time Fixture and Consumer Self-
7//! Registration". Each op-template consumer (validation,
8//! conditional, throttle, …) registers the Polydat names it will read
9//! at cycle time into a shared [`ScopeFixture`]. The fixture is
10//! sealed once construction completes and yields a [`PullPlan`]
11//! whose entries are materialized per cycle into a [`ResolvedPulls`]
12//! buffer indexed by the [`PullHandle`]s the consumer received at
13//! registration.
14//!
15//! This is the **only** path by which wrappers gain access to GK
16//! values at cycle time (SRD 31 §"Pull plan vs bind plan"). The
17//! prior side channel — wrapper names threaded through
18//! `synthesis::resolve_with_extras` into `ResolvedFields` — has
19//! been removed; the resolver now exposes `resolve_with_field_pulls`,
20//! which carries op-field bind points only (the adapter-facing
21//! read).
22
23use std::collections::HashMap;
24use std::sync::Arc;
25
26use crate::scope_kernel::ScopeKernel;
27use polydat::ast::Value;
28use polydat::kernel::PolydatProgram;
29
30/// What the kernel reports a registered name resolves to.
31///
32/// Captures live in `input_defs` (the input array) at indices
33/// `>= coord_count`, so they share the `Input` variant — no
34/// separate capture variant is needed.
35#[derive(Debug, Clone)]
36enum PlanEntry {
37    Output { name: String, output_idx: usize },
38    Input { name: String, input_idx: usize },
39}
40
41impl PlanEntry {
42    fn name(&self) -> &str {
43        match self {
44            PlanEntry::Output { name, .. } => name,
45            PlanEntry::Input { name, .. } => name,
46        }
47    }
48}
49
50/// Init-time accumulator for consumer-declared pulls.
51///
52/// Holds an `Arc<PolydatProgram>` clone of the per-template canonical
53/// kernel's program (SRD 16 §"Per-Scope Canonical Kernel Cache").
54/// The program carries every fact this fixture needs at
55/// registration time — the output map, input definitions, and
56/// auto-extern-provisioned externs — so no per-fiber state is
57/// touched until the sealed [`PullPlan`] is resolved at cycle
58/// time. Each consumer registers names via [`register_pull`];
59/// the fixture deduplicates by name and yields a stable
60/// [`PullHandle`] for each unique registration.
61///
62/// [`register_pull`]: ScopeFixture::register_pull
63pub struct ScopeFixture {
64    program: Arc<PolydatProgram>,
65    handles: HashMap<String, PullHandle>,
66    plan: Vec<PlanEntry>,
67}
68
69impl ScopeFixture {
70    /// Open a new fixture against a per-template canonical
71    /// program. The Arc clone is cheap; the activity construction
72    /// loop typically already holds the program for adapter
73    /// dispatch and bind-plan synthesis.
74    pub fn new(program: Arc<PolydatProgram>) -> Self {
75        Self {
76            program,
77            handles: HashMap::new(),
78            plan: Vec::new(),
79        }
80    }
81
82    /// The program this fixture is scoped against. Useful for
83    /// consumers that need to inspect the manifest (e.g. type-
84    /// aware strict parsing).
85    pub fn program(&self) -> &Arc<PolydatProgram> {
86        &self.program
87    }
88
89    /// Register a name. Resolves it against the program's output
90    /// map (folded constants) first, then against `find_input`
91    /// (extern slots, capture inputs, coordinates).
92    ///
93    /// Returns a memoized handle. Re-registering the same name
94    /// returns the existing handle — registrations are idempotent.
95    ///
96    /// **Errors** when the program does not know the name. The
97    /// Polydat compiler is responsible for provisioning every name
98    /// referenced anywhere in the op template (op fields and
99    /// params; SRD 16 §"Auto-Extern Generation"). An unknown
100    /// name here therefore signals a workload bug — typically a
101    /// reference to a binding that was never declared, or a
102    /// typo that the compiler should have already caught.
103    pub fn register_pull(&mut self, name: &str) -> Result<PullHandle, String> {
104        if let Some(&h) = self.handles.get(name) {
105            return Ok(h);
106        }
107        let entry = if let Some((node_idx, _port_idx)) = self.program.resolve_output(name) {
108            // Output: prefer the output index over (node_idx, port_idx)
109            // because state.pull_by_index resolves the pair internally
110            // and matches OutputAccessor's existing pattern.
111            let output_idx = self.program.output_index(name).ok_or_else(|| {
112                format!(
113                    "fixture: name '{name}' resolved to output node {node_idx} \
114                     but had no entry in the output_list — this is a kernel \
115                     consistency bug, please report",
116                )
117            })?;
118            PlanEntry::Output {
119                name: name.to_string(),
120                output_idx,
121            }
122        } else if let Some(input_idx) = self.program.find_input(name) {
123            PlanEntry::Input {
124                name: name.to_string(),
125                input_idx,
126            }
127        } else {
128            return Err(format!(
129                "fixture: name '{name}' is not known to the program — neither \
130                 a declared output nor an input slot. The Polydat compiler should \
131                 have provisioned it from a bind-point reference somewhere in \
132                 the op template; if it didn't, the workload is referencing a \
133                 binding that doesn't exist. Available outputs: [{outs}]; \
134                 inputs: [{ins}].",
135                outs = self.program.output_names().join(", "),
136                ins = self.program.input_names().join(", "),
137            ));
138        };
139        let handle = PullHandle(self.plan.len());
140        self.plan.push(entry);
141        self.handles.insert(name.to_string(), handle);
142        Ok(handle)
143    }
144
145    /// Seal and yield the immutable plan. The plan owns its own
146    /// `Arc<PolydatProgram>` clone, so cycle-time `resolve` only needs
147    /// the per-fiber state.
148    pub fn seal(self) -> PullPlan {
149        PullPlan {
150            program: self.program,
151            entries: self.plan,
152        }
153    }
154}
155
156/// Opaque, copy-able handle into a [`PullPlan`]. The only way to
157/// turn a handle into a value is [`ResolvedPulls::get`].
158#[derive(Copy, Clone, Debug, Eq, PartialEq, Hash)]
159pub struct PullHandle(usize);
160
161impl PullHandle {
162    /// Internal: index into the plan / resolved buffer. Public
163    /// only to the crate so test scaffolding can build synthetic
164    /// pulls; outside this crate, handles are opaque.
165    pub(crate) fn index(self) -> usize {
166        self.0
167    }
168}
169
170/// Sealed init-time plan. One entry per unique name registered.
171/// Owns an `Arc<PolydatProgram>` so cycle-time resolve only needs the
172/// per-fiber state.
173pub struct PullPlan {
174    program: Arc<PolydatProgram>,
175    entries: Vec<PlanEntry>,
176}
177
178impl PullPlan {
179    /// Assert this plan is being resolved against the program it was BUILT
180    /// against.
181    ///
182    /// A plan holds resolved indices, not names — `Output { output_idx: 7 }`.
183    /// Resolving it against a different program silently reindexes every read:
184    /// where the index happens to be in range it returns a NEIGHBOURING wire's
185    /// value (a different metric, reported under this one's name), and where it
186    /// is out of range it panics deep in the engine as
187    /// "index out of bounds", naming nothing that points back here.
188    ///
189    /// Both happened: a CQL dispenser that returned `None` from
190    /// `canonical_kernel()` left its per-op slot empty, so plans built against
191    /// the op-template program fell back to the phase kernel. One metric read
192    /// its phase's `__metric_time_to_index`; the next two panicked.
193    ///
194    /// The check is an `Arc` pointer comparison per op per cycle — cheaper than
195    /// the first index it protects.
196    pub(crate) fn check_program_match(&self, program: &Arc<PolydatProgram>, template_idx: usize) {
197        if Arc::ptr_eq(&self.program, program) {
198            return;
199        }
200        panic!(
201            "pull plan for op #{template_idx} was built against a different \
202             program than the kernel it is being resolved against — every \
203             index in it addresses the wrong wire.\n  \
204             plan program outputs:   {:?}\n  \
205             kernel program outputs: {:?}\n\
206             This is a wiring bug, not a workload error. The usual cause is a \
207             dispenser whose `canonical_kernel()` returns `None` (so no per-op \
208             kernel is materialised and the plan falls back to the phase \
209             kernel) — every leaf dispenser must return the op-template kernel \
210             it was mapped against.",
211            self.program.output_names(),
212            program.output_names(),
213        );
214    }
215
216    /// Number of distinct names in this plan.
217    pub fn len(&self) -> usize {
218        self.entries.len()
219    }
220
221    /// Whether this plan has no entries.
222    pub fn is_empty(&self) -> bool {
223        self.entries.is_empty()
224    }
225
226    /// The names in the plan, in registration order. Diagnostic only.
227    pub fn names(&self) -> Vec<&str> {
228        self.entries.iter().map(|e| e.name()).collect()
229    }
230
231    /// The program this plan was sealed against.
232    pub fn program(&self) -> &Arc<PolydatProgram> {
233        &self.program
234    }
235
236    /// Materialize every entry against `kernel`, a kernel of this
237    /// plan's program on any engine (its indices are the program's).
238    /// Output entries go through `pull_at` (eval cone if not current);
239    /// input entries through `input_value_at` (cell-aware for shared
240    /// slots).
241    ///
242    /// O(plan_len) on the hot path — no name hashing.
243    pub fn resolve(&self, kernel: &mut dyn polydat::Kernel) -> ResolvedPulls {
244        let mut values = Vec::with_capacity(self.entries.len());
245        for entry in &self.entries {
246            let v = match entry {
247                PlanEntry::Output { output_idx, .. } => kernel.pull_at(*output_idx),
248                PlanEntry::Input { input_idx, .. } => {
249                    kernel.input_value_at(*input_idx).unwrap_or(Value::None)
250                }
251            };
252            values.push(v);
253        }
254        ResolvedPulls { values }
255    }
256
257    /// Convenience: resolve against a scope kernel. Provided so
258    /// test scaffolding and other ergonomic call sites don't have to
259    /// name the trait object.
260    pub fn resolve_with(&self, kernel: &mut ScopeKernel) -> ResolvedPulls {
261        self.resolve(kernel.kernel_mut())
262    }
263}
264
265/// Cycle-time materialization of a [`PullPlan`].
266///
267/// Read-only, indexed by [`PullHandle`]. Created once per cycle
268/// during the resolve phase (SRD 31 §"Cycle-Time Pipeline"); the
269/// values inside reflect the PolydatState snapshot at the moment of
270/// resolution and are not invalidated by subsequent state changes.
271pub struct ResolvedPulls {
272    values: Vec<Value>,
273}
274
275impl ResolvedPulls {
276    /// Empty pulls — no consumer registered anything. Used as the
277    /// default in the wrapper construction path before α.4 deletes
278    /// the legacy `extras` flow.
279    pub fn empty() -> Self {
280        Self { values: Vec::new() }
281    }
282
283    /// Resolve a handle to a borrowed value. Panics on a handle
284    /// from a different plan (the index is out of range). This is
285    /// a programming error — a handle's plan provenance is
286    /// statically associated with the wrapper that owns it, and
287    /// `ExecCtx` carries the matching `ResolvedPulls`.
288    pub fn get(&self, h: PullHandle) -> &Value {
289        &self.values[h.index()]
290    }
291
292    /// Number of resolved values.
293    pub fn len(&self) -> usize {
294        self.values.len()
295    }
296
297    /// Whether this resolution is empty.
298    pub fn is_empty(&self) -> bool {
299        self.values.is_empty()
300    }
301}
302
303/// Trait every cross-cutting wrapper that reads Polydat values must
304/// implement. The activity construction loop calls `fixture` once
305/// per template per consumer; failures (closed-vocab violation,
306/// missing required field, unresolvable name) are returned as
307/// `Err` and abort construction loudly.
308pub trait OpConsumer: Sized {
309    /// Inspect `template`, register the names this consumer will
310    /// read at cycle time, and return a fully-configured instance
311    /// holding handles plus parsed (strict) config.
312    fn fixture(
313        template: &nmbrs_workload::model::ParsedOp,
314        fx: &mut ScopeFixture,
315    ) -> Result<Self, String>;
316}
317
318/// Cycle-time bundle handed to every dispenser via
319/// `OpDispenser::execute`. Adapters use `fields` exclusively;
320/// wrappers read `pulls` via stored handles.
321///
322/// Bundling rather than passing two parameters keeps the trait
323/// surface forward-compatible for per-cycle component context,
324/// streaming-capture hooks, and diagnostic taps (SRD 32
325/// §"`ExecCtx` — cycle-time bundle").
326pub struct ExecCtx<'a> {
327    pub fields: &'a crate::adapter::ResolvedFields,
328    pub pulls: &'a ResolvedPulls,
329    /// Narrow read surface for op-template name resolution against
330    /// the dispenser's bound Polydat context (SRD-68 invariants I-1 + I-2).
331    /// During the SRD-68 migration this defaults to a no-op
332    /// `NullWireSource` for legacy call sites; adapters that own a
333    /// kernel construct via [`Self::with_wires`].
334    pub wires: &'a dyn crate::wires::WireSource,
335    /// Number of consecutive wire ordinals this invocation should
336    /// cover — the ACTUAL length of the cursor sub-run the executor
337    /// reserved for this op (`base .. actual_end`). Equals the op's
338    /// [`crate::adapter::OpDispenser::rows_per_op`] except at the
339    /// cursor tail, where the final reservation is short: a batch op
340    /// reads exactly `[cycle, cycle + run_len)` so the partial tail
341    /// of `M % N` rows is inserted too — never over-read, never
342    /// dropped. Ordinary (non-batch) ops ignore it. Defaults to `1`
343    /// for every legacy / test call site (single-row semantics).
344    pub run_len: usize,
345}
346
347impl<'a> ExecCtx<'a> {
348    /// Legacy constructor — defaults `wires` to a no-op
349    /// [`crate::wires::NullWireSource`]. Used by call sites that
350    /// haven't migrated to SRD-68's dispenser-owned-kernel model
351    /// yet. Adapters that own a kernel handle should call
352    /// [`Self::with_wires`] instead.
353    pub fn new(fields: &'a crate::adapter::ResolvedFields, pulls: &'a ResolvedPulls) -> Self {
354        Self {
355            fields,
356            pulls,
357            wires: &crate::wires::NULL_WIRES,
358            run_len: 1,
359        }
360    }
361
362    /// Construct an `ExecCtx` with an explicit `WireSource` — the
363    /// SRD-68 path. The `wires` value should be the per-fiber
364    /// kernel slot for the firing dispenser, narrowed to the
365    /// `WireSource` trait so adapter code never sees `ScopeKernel`
366    /// internals.
367    pub fn with_wires(
368        fields: &'a crate::adapter::ResolvedFields,
369        pulls: &'a ResolvedPulls,
370        wires: &'a dyn crate::wires::WireSource,
371    ) -> Self {
372        Self {
373            fields,
374            pulls,
375            wires,
376            run_len: 1,
377        }
378    }
379}
380
381#[cfg(test)]
382mod tests {
383    use super::*;
384
385    fn k() -> ScopeKernel {
386        crate::bindings::compile_scope_kernel(
387            "input cycle: u64\n\
388             folded := 42\n\
389             cyc_dep := hash(cycle)\n",
390            &Default::default(),
391        )
392        .expect("compile_scope_kernel")
393    }
394
395    #[test]
396    fn register_resolves_folded_output() {
397        let kernel = k();
398        let mut fx = ScopeFixture::new(kernel.program().clone());
399        let h = fx.register_pull("folded").expect("folded should resolve");
400        let plan = fx.seal();
401        assert_eq!(plan.len(), 1);
402        assert_eq!(plan.names(), vec!["folded"]);
403        let _ = h; // handle is opaque from outside
404    }
405
406    #[test]
407    fn register_resolves_cycle_dependent_output() {
408        let kernel = k();
409        let mut fx = ScopeFixture::new(kernel.program().clone());
410        fx.register_pull("cyc_dep").expect("cyc_dep should resolve");
411        let plan = fx.seal();
412        assert_eq!(plan.len(), 1);
413    }
414
415    #[test]
416    fn register_resolves_input() {
417        let kernel = k();
418        let mut fx = ScopeFixture::new(kernel.program().clone());
419        // 'cycle' is the coordinate input, not an output.
420        fx.register_pull("cycle")
421            .expect("cycle input should resolve");
422        let plan = fx.seal();
423        assert_eq!(plan.names(), vec!["cycle"]);
424    }
425
426    #[test]
427    fn register_unknown_name_errors() {
428        let kernel = k();
429        let mut fx = ScopeFixture::new(kernel.program().clone());
430        let err = fx.register_pull("nonexistent").unwrap_err();
431        assert!(
432            err.contains("nonexistent"),
433            "error should name the missing binding: {err}"
434        );
435        assert!(
436            err.contains("Available outputs"),
437            "error should list available outputs: {err}"
438        );
439    }
440
441    #[test]
442    fn register_is_idempotent_per_name() {
443        let kernel = k();
444        let mut fx = ScopeFixture::new(kernel.program().clone());
445        let h1 = fx.register_pull("folded").unwrap();
446        let h2 = fx.register_pull("folded").unwrap();
447        assert_eq!(h1, h2, "same name should yield same handle");
448        let plan = fx.seal();
449        assert_eq!(
450            plan.len(),
451            1,
452            "duplicate registrations should not grow the plan"
453        );
454    }
455
456    #[test]
457    fn register_assigns_distinct_handles_for_distinct_names() {
458        let kernel = k();
459        let mut fx = ScopeFixture::new(kernel.program().clone());
460        let h_folded = fx.register_pull("folded").unwrap();
461        let h_cyc = fx.register_pull("cyc_dep").unwrap();
462        assert_ne!(h_folded, h_cyc);
463        let plan = fx.seal();
464        assert_eq!(plan.len(), 2);
465    }
466
467    #[test]
468    fn resolve_pulls_folded_output_value() {
469        let mut kernel = k();
470        let mut fx = ScopeFixture::new(kernel.program().clone());
471        let h = fx.register_pull("folded").unwrap();
472        let plan = fx.seal();
473
474        kernel.set_inputs(&[0]);
475        let pulls = plan.resolve_with(&mut kernel);
476        let v = pulls.get(h);
477        assert_eq!(v.as_u64(), 42);
478    }
479
480    #[test]
481    fn resolve_pulls_cycle_dependent_value_per_cycle() {
482        let mut kernel = k();
483        let mut fx = ScopeFixture::new(kernel.program().clone());
484        let h = fx.register_pull("cyc_dep").unwrap();
485        let plan = fx.seal();
486
487        kernel.set_inputs(&[0]);
488        let v0 = plan.resolve_with(&mut kernel).get(h).as_u64();
489        kernel.set_inputs(&[1]);
490        let v1 = plan.resolve_with(&mut kernel).get(h).as_u64();
491        assert_ne!(v0, v1, "cycle-dependent output should change per cycle");
492    }
493
494    #[test]
495    fn resolve_pulls_input_slot_value() {
496        let mut kernel = k();
497        let mut fx = ScopeFixture::new(kernel.program().clone());
498        let h = fx.register_pull("cycle").unwrap();
499        let plan = fx.seal();
500
501        kernel.set_inputs(&[7]);
502        let pulls = plan.resolve_with(&mut kernel);
503        assert_eq!(pulls.get(h).as_u64(), 7);
504    }
505
506    #[test]
507    fn empty_plan_resolves_to_empty_pulls() {
508        let mut kernel = k();
509        let fx = ScopeFixture::new(kernel.program().clone());
510        let plan = fx.seal();
511        kernel.set_inputs(&[0]);
512        let pulls = plan.resolve_with(&mut kernel);
513        assert!(pulls.is_empty());
514    }
515}