nmbrs_runtime/synthesis.rs
1// Copyright 2024-2026 Jonathan Shook
2// SPDX-License-Identifier: Apache-2.0
3
4//! Per-fiber Polydat Kernel construction.
5//!
6//! `OpBuilder` owns the activity's source kernel and seeds each
7//! per-fiber [`FiberBuilder`] with a kernel bound under it, on the
8//! fiber engine where it has an image ([`crate::fiber_engine`]), plus
9//! the named scope overrides. The
10//! adapter-facing cycle-time bind-point resolution path
11//! historically lived here too, but SRD-68 Push 5 retired it in
12//! favour of the generic [`crate::wires::WireSource`] surface;
13//! this module now scopes to fiber construction and bind-point
14//! validation.
15//!
16//! See `docs/SRD/68_dispenser_owned_polydat_context.md` for the
17//! resolution model.
18
19use std::sync::{Arc, OnceLock};
20
21use crate::fiber_engine::OpTemplateModule;
22use crate::scope_kernel::ScopeKernel;
23use nmbrs_workload::bindpoints::{self, BindPoint, BindQualifier};
24use nmbrs_workload::model::ParsedOp;
25use polydat::Kernel;
26use polydat::ast::Value;
27use polydat::kernel::{KernelProgram, PolydatProgram};
28
29/// Cached `NMBRS_DIRTY_DEBUG` flag — per-cycle `std::env::var`
30/// reads measured at ~30% of single-fiber CPU; the OnceLock
31/// makes the gate a single atomic load on the hot path. See
32/// the matching helpers in `wires.rs` / polydat's `engines.rs`.
33fn nmbrs_dirty_debug_enabled() -> bool {
34 static FLAG: OnceLock<bool> = OnceLock::new();
35 *FLAG.get_or_init(|| std::env::var("NMBRS_DIRTY_DEBUG").is_ok())
36}
37
38/// Shared op builder that distributes per-fiber builders.
39///
40/// Holds the activity's source kernel (immutable, shared). Each
41/// executor fiber calls `create_fiber_builder()` to get its own
42/// `FiberBuilder` with private kernels — no locks, no contention on
43/// the hot path.
44pub struct OpBuilder {
45 /// Name-keyed scope values (per SRD-13c), written into every new
46 /// fiber's kernels. Stored by name rather than `(input_idx, value)`
47 /// because each kernel — the fiber main kernel, every per-op-template
48 /// kernel — owns its own input layout, and an index captured against
49 /// the source kernel doesn't translate. The previous
50 /// `Vec<(usize, Value)>` shape silently mis-routed writes across
51 /// kernels (e.g. `table` value landing in the `profile` slot of an
52 /// op-template kernel whose extern declaration order differed from
53 /// the phase scope).
54 scope_values: Vec<(String, Value)>,
55 /// The source kernel — the activity's own kernel that each
56 /// per-fiber `FiberBuilder` binds its main kernel under. Owning the
57 /// kernel (not just its program) carries the activity's full cell
58 /// state — own input-slot cells plus transit cells inherited from
59 /// ancestors — to every fiber's main kernel.
60 source_kernel: Arc<ScopeKernel>,
61 /// SRD-13d Phase 9 — per-op-template kernel programs keyed
62 /// by op name. Populated by [`Self::with_op_template_programs`]
63 /// when the runner has materialised op-template kernels in
64 /// the scope tree. Wrappers (e.g. `MetricsDispenser`) look up
65 /// the program for their template via [`Self::program_for_op`]
66 /// and build their `ScopeFixture` against it; flattened
67 /// op-templates fall through to the activity-wide `program`.
68 op_template_programs: std::collections::HashMap<String, Arc<PolydatProgram>>,
69 /// What each canonical kernel this builder hands out stands for,
70 /// keyed by its `program_id`: a dispenser's canonical may be on any
71 /// engine, and each fiber finds here the interpreter program that
72 /// resolves its indices and the module it instantiates the per-op
73 /// kernel from.
74 canonicals: Canonicals,
75}
76
77/// What a canonical kernel's program is to a fiber: the interpreter
78/// program that resolves its indices (the analysis program; a native
79/// image of it shares its indices, `fiber_engine::agrees`), and the
80/// op-template module a per-op kernel is instantiated from on the fiber
81/// engine, when it has one, or else the image a per-op kernel is bound
82/// from — the scope's own engine program, the interpreter program when
83/// nothing better stands in for it.
84#[derive(Clone)]
85struct CanonicalSource {
86 program: Arc<PolydatProgram>,
87 image: Arc<dyn KernelProgram>,
88 module: Option<Arc<OpTemplateModule>>,
89}
90
91impl CanonicalSource {
92 /// A program with no other image: per-op kernels bind it on the
93 /// interpreter.
94 fn interpreted(program: Arc<PolydatProgram>) -> Self {
95 Self {
96 image: program.clone(),
97 program,
98 module: None,
99 }
100 }
101}
102
103/// Canonical sources by the `program_id` of every kernel that stands for
104/// them.
105type Canonicals = Arc<std::collections::HashMap<polydat::kernel::ProgramId, CanonicalSource>>;
106
107/// The identity an interpreter program's kernels report.
108fn program_id_of(program: &Arc<PolydatProgram>) -> polydat::kernel::ProgramId {
109 KernelProgram::program_id(program.as_ref())
110}
111
112impl OpBuilder {
113 /// Create an OpBuilder from a kernel.
114 ///
115 /// If the kernel has scope values (set via `materialize_wiring_from_outer`
116 /// or directly via `kernel.state().set_input`), they are
117 /// captured and propagated into every fiber's kernels.
118 pub fn new(kernel: impl crate::scope_kernel::IntoSharedScope) -> Self {
119 let kernel: Arc<ScopeKernel> = kernel.into_shared_scope();
120 // Scope values seed PLAIN slots only. A CELL-BOUND slot's value
121 // is the live shared cell — snapshotting it here freezes the
122 // cell's activity-start value, and every downstream
123 // re-application (fiber seeding, the stanza-boundary
124 // `reset_captures`) would `set_input` that stale snapshot
125 // THROUGH the cell, clobbering later writes for every kernel
126 // sharing it. Same exclusion `reset_inputs` documents:
127 // cells are cross-kernel shared state with their own
128 // lifecycle.
129 let scope_values: Vec<(String, Value)> = kernel
130 .scope_values()
131 .into_iter()
132 .filter(|(name, _)| {
133 kernel
134 .program()
135 .find_input(name)
136 .map(|idx| !kernel.input_is_cell_bound(idx))
137 .unwrap_or(true)
138 })
139 .collect();
140 // The source kernel is the canonical of every flattened op.
141 let canonicals = std::iter::once((
142 kernel.program_id(),
143 CanonicalSource {
144 program: kernel.program().clone(),
145 image: kernel.image().clone(),
146 module: None,
147 },
148 ))
149 .collect();
150 Self {
151 scope_values,
152 source_kernel: kernel,
153 op_template_programs: std::collections::HashMap::new(),
154 canonicals: Arc::new(canonicals),
155 }
156 }
157
158 /// Install per-op-template kernel programs (SRD-13d Phase 9).
159 /// The runner builds these from the scope tree's
160 /// `cached_kernel` slots for materialised op-template scopes
161 /// and threads them here so wrappers can look up the right
162 /// program when constructing their fixtures.
163 pub fn with_op_template_programs(
164 mut self,
165 programs: std::collections::HashMap<String, Arc<PolydatProgram>>,
166 ) -> Self {
167 let canonicals = Arc::make_mut(&mut self.canonicals);
168 for program in programs.values() {
169 canonicals
170 .entry(program_id_of(program))
171 .or_insert_with(|| CanonicalSource::interpreted(program.clone()));
172 }
173 self.op_template_programs = programs;
174 self
175 }
176
177 /// Instantiate per-op kernels from `modules`, the op-template scope
178 /// modules of this phase. A module whose fiber-engine image
179 /// disagrees with its program is left out, and its per-op kernels
180 /// stay on the interpreter.
181 pub fn with_op_template_modules(
182 mut self,
183 modules: impl IntoIterator<Item = (String, Arc<OpTemplateModule>)>,
184 ) -> Self {
185 let canonicals = Arc::make_mut(&mut self.canonicals);
186 for (name, module) in modules {
187 let Some(image) =
188 crate::fiber_engine::module_image(&module, &format!("op template '{name}'"))
189 else {
190 continue;
191 };
192 let source = CanonicalSource {
193 program: module.program().clone(),
194 image: image.clone(),
195 module: Some(module.clone()),
196 };
197 // Both the interpreter program and its native image stand
198 // for the module: a canonical kernel reports either.
199 canonicals.insert(program_id_of(module.program()), source.clone());
200 canonicals.insert(image.program_id(), source);
201 }
202 self
203 }
204
205 /// Look up the kernel program for op `name`. Returns the
206 /// per-op-template program if Phase 9 produced one for this
207 /// op (i.e. `materialised` and bindings non-empty); otherwise
208 /// returns the activity-wide program (the flatten path).
209 pub fn program_for_op(&self, name: &str) -> Arc<PolydatProgram> {
210 self.op_template_programs
211 .get(name)
212 .cloned()
213 .unwrap_or_else(|| self.source_kernel.program().clone())
214 }
215
216 /// The activity-wide kernel program. Used by callers that
217 /// need the source program shape (output names, manifest)
218 /// without rebuilding a fresh kernel.
219 pub fn program(&self) -> Arc<PolydatProgram> {
220 self.source_kernel.program().clone()
221 }
222
223 /// The activity-wide source kernel — the Polydat context every
224 /// op-template subscope is built upon. Adapters' `map_op`
225 /// implementations receive a kernel of this scope (see
226 /// [`Self::canonical_kernel_for_op`]) as the `parent` argument, to
227 /// bind their own canonical op-template kernel under (SRD-68
228 /// invariant I-3) or to retain when their op has no matter to add.
229 pub fn source_kernel(&self) -> &Arc<ScopeKernel> {
230 &self.source_kernel
231 }
232
233 /// Build the canonical op-template kernel for `op_name` —
234 /// the Polydat context the dispenser owns and that per-fiber
235 /// instances are materialised from (SRD-68 invariants I-3,
236 /// I-4). Built once at dispenser construction time.
237 ///
238 /// When `op_name` has a registered op-template program (phase
239 /// `bindings:`, op-level `bindings:`, `result:` block — the
240 /// matter assembled by the synthesis pipeline before the
241 /// activity runs), the canonical is that program bound under
242 /// `source_kernel`: instantiated from its module on the fiber
243 /// engine when it has one, on the interpreter otherwise.
244 /// Otherwise the canonical is a fork of the source kernel (its
245 /// program, state and cells), which covers the flattened-op-template
246 /// path (no per-op matter). Either way the adapter holds a kernel of any
247 /// engine, whose `program_id` this builder recognizes.
248 pub fn canonical_kernel_for_op(&self, op_name: &str) -> Arc<dyn Kernel> {
249 let Some(program) = self.op_template_programs.get(op_name) else {
250 return Arc::from(self.source_kernel.fork().into_kernel());
251 };
252 let module = self
253 .canonicals
254 .get(&program_id_of(program))
255 .and_then(|source| source.module.clone());
256 match module {
257 Some(module) => Arc::from(
258 module
259 .instantiate_under(
260 self.source_kernel.kernel(),
261 crate::fiber_engine::fiber_engine(),
262 &[],
263 )
264 .unwrap_or_else(|e| {
265 panic!("op '{op_name}': canonical kernel failed to instantiate: {e}")
266 }),
267 ),
268 None => Arc::from(
269 polydat::kernel::bind_under(
270 self.source_kernel.kernel(),
271 program.clone() as Arc<dyn KernelProgram>,
272 &[],
273 )
274 .unwrap_or_else(|e| panic!("op '{op_name}': canonical kernel failed to bind: {e}")),
275 ),
276 }
277 }
278
279 /// Create a per-fiber builder. No locks, no sharing — the fiber
280 /// owns its kernels exclusively. Scope values (per-iteration
281 /// inputs from `for_each` / `for_combinations` / outer scope
282 /// constants) are written into the main kernel's inputs and
283 /// remembered on the builder so `reset_captures` (called at
284 /// stanza boundaries) can re-apply them — otherwise the
285 /// blanket "reset all non-coord inputs" pass would clobber
286 /// the iteration's bound values.
287 pub fn create_fiber_builder(&self) -> FiberBuilder {
288 // The fiber's main kernel is bound under the activity's source
289 // kernel: cells, transit cells, and value-copy bindings flow in
290 // automatically, and the scope-init constants materialize, so
291 // the fiber observes the same cell handles as the workload-root
292 // through the chain.
293 //
294 // Scope values are bound with it, as iteration bindings, so the
295 // main kernel's consts initialize from them (a const is fixed at
296 // init). The indices cached here feed the per-cycle
297 // `reset_captures` re-application.
298 let mut fb = FiberBuilder::with_scope(&self.source_kernel, None, self.scope_values.clone());
299 fb.canonicals = self.canonicals.clone();
300 // SRD-68: per-fiber op-template kernels are populated by
301 // `attach_dispenser_kernels`, which runs right after this
302 // function returns (see executor cycle dispatch), each from
303 // the firing dispenser's `OpDispenser::canonical_kernel()`.
304 fb
305 }
306}
307
308/// The input a scope value is written to on `program`, or `None` when
309/// the program has no such input or it is a coordinate. A coordinate
310/// (the `cycle` a scope carries among its values) advances with every
311/// cycle through `set_inputs` and is never written by name; the stanza
312/// reset leaves it untouched, so there is nothing to re-apply.
313fn scope_value_index(program: &PolydatProgram, name: &str) -> Option<usize> {
314 let idx = program.find_input(name)?;
315 (program.input_kind(idx) != Some(polydat::kernel::InputKind::Coordinate)).then_some(idx)
316}
317
318/// The scope values a kernel of `program` takes, as the `iter_bindings`
319/// its binder writes before the kernel's consts are initialized: each
320/// value converted to its slot's declared type, skipping the values
321/// `program` has no (non-coordinate) input for.
322///
323/// # Panics
324/// On a scope value its slot's type cannot take, converted or not: a
325/// fail-loud condition, since it would otherwise corrupt every read.
326fn scope_bindings(
327 program: &PolydatProgram,
328 scope_values: &[(String, Value)],
329) -> Vec<(String, Value)> {
330 scope_values
331 .iter()
332 .filter(|(name, _)| scope_value_index(program, name).is_some())
333 .map(|(name, value)| {
334 let value = match program.input_port_type(name) {
335 Some(port) => polydat::convert::to_port(value.clone(), port).unwrap_or_else(|e| {
336 panic!("scope value '{name}' failed typed write at scope-init: {e}")
337 }),
338 None => value.clone(),
339 };
340 (name.clone(), value)
341 })
342 .collect()
343}
344
345/// Per-fiber op builder. Owns its own kernels.
346/// No locks, no synchronization, no contention.
347///
348/// Created via `OpBuilder::create_fiber_builder()` at fiber startup.
349///
350/// Every kernel here may be on any engine: the main kernel and the
351/// per-op kernels run on the fiber engine where an image is available
352/// ([`crate::fiber_engine`]) and on the interpreter otherwise. Each is
353/// paired with the interpreter program it runs or was imaged from,
354/// which shares its input and output indices: names resolve on the
355/// program once, and the kernel is driven by index.
356pub struct FiberBuilder {
357 /// The fiber's main kernel — typically the activity-wide
358 /// (workload / phase) program.
359 main_kernel: Box<dyn Kernel>,
360 /// The interpreter program [`Self::main_kernel`] runs or was
361 /// imaged from.
362 main_program: Arc<PolydatProgram>,
363 /// Scope-bound input values (per-iteration extern bindings)
364 /// that should persist across stanza-level `reset_inputs`
365 /// resets. Empty for a builder created via plain
366 /// [`FiberBuilder::new`]; populated by
367 /// [`OpBuilder::create_fiber_builder`].
368 scope_values: Vec<(String, Value)>,
369 /// SRD-68 invariant I-4 — per-fiber kernel instances, indexed
370 /// parallel to the activity's dispenser registry. Each entry
371 /// is the corresponding dispenser's canonical program bound under
372 /// [`Self::main_kernel`]; `None` for dispensers that don't expose
373 /// a canonical kernel (adapters with no Polydat needs, or wrappers
374 /// that delegate). Populated by
375 /// [`Self::attach_dispenser_kernels`] right after fiber spawn,
376 /// before any cycles run; read at cycle dispatch to populate
377 /// `ExecCtx::wires` for the firing dispenser.
378 per_op_kernels: Vec<Option<Box<dyn Kernel>>>,
379 /// The interpreter program each per-op kernel runs or was imaged
380 /// from, parallel to [`Self::per_op_kernels`].
381 per_op_programs: Vec<Option<Arc<PolydatProgram>>>,
382 /// Per-op-kernel side-effecting output indices (parallels
383 /// [`Self::per_op_kernels`]). The subset of each op-template
384 /// kernel's outputs whose cone contains a `Purity::SideChannel`
385 /// node — the only outputs the per-cycle "fire side effects" pass
386 /// pulls. Computed once at [`Self::attach_dispenser_kernels`] so
387 /// volatile metric-reader outputs (the objective bindings) are NOT
388 /// re-evaluated every cycle just to fire a non-existent effect.
389 per_op_side_effecting: Vec<Vec<usize>>,
390 /// Pre-resolved input indices for each entry in
391 /// [`Self::scope_values`] against [`Self::main_kernel`].
392 /// `None` slots are scope values the main program doesn't
393 /// declare, or declares as a coordinate (silently skipped at
394 /// write time — matches the historical name-based-skip semantics).
395 scope_value_main_idx: Vec<Option<usize>>,
396 /// Per-op-kernel mirror of [`Self::scope_value_main_idx`].
397 /// Outer Vec parallels [`Self::per_op_kernels`]; inner Vec
398 /// parallels [`Self::scope_values`]. `None` outer slots
399 /// match per_op_kernels' `None` entries.
400 scope_value_per_op_idx: Vec<Option<Vec<Option<usize>>>>,
401 /// Whether some per-op kernel runs a program other than the main
402 /// kernel's, and so reads the main kernel's outputs through its
403 /// broadcast cells (see [`Self::set_source_item`]).
404 needs_broadcast: bool,
405 /// What each dispenser's canonical kernel stands for — its
406 /// interpreter program and op-template module — by the kernel's
407 /// `program_id` (see [`OpBuilder::canonical_kernel_for_op`]).
408 canonicals: Canonicals,
409}
410
411/// Validate that all bind points in op templates can be resolved.
412///
413/// Called at init time. Warns for each unresolvable `{name}` reference.
414/// A bind point is resolvable if it matches a Polydat output, input name,
415/// or a known capture declaration from another op. Workload params are
416/// injected into the Polydat source as constant bindings before compilation,
417/// so they resolve as Polydat outputs.
418/// Validate that all bind points in op templates can be resolved.
419///
420/// Returns `Err` with a descriptive message if any bind point is
421/// unresolvable. Callers should treat this as a fatal error —
422/// unresolved bind points produce broken ops at runtime.
423pub fn validate_bind_points(
424 templates: &[ParsedOp],
425 program_for_op: &dyn Fn(&str) -> Arc<PolydatProgram>,
426) -> Result<(), String> {
427 // Collect all capture declarations across templates. Captures
428 // are extracted at workload-parse time and live on
429 // `ParsedOp.captures`; the op-text fields have brackets
430 // stripped by then, so re-parsing the text wouldn't surface
431 // them.
432 let mut capture_names: std::collections::HashSet<String> = std::collections::HashSet::new();
433 for template in templates {
434 for cap in &template.captures {
435 capture_names.insert(cap.as_name.clone());
436 }
437 }
438
439 let mut errors: Vec<String> = Vec::new();
440
441 for template in templates {
442 // THIS op's program, not the activity-wide one. An op that owns a
443 // kernel (its own `bindings:`, or an adapter that materialises one)
444 // has its bindings in that program and nowhere else, so validating
445 // every template against the activity kernel reported an op's own
446 // binding as unresolvable — `{x}` in `stmt:`/`raw:` failed at RUNTIME
447 // while the same name still rendered fine in `memo:`/`gutter:`, which
448 // read the live wires instead of this check.
449 let program = program_for_op(&template.name);
450 for (field_name, value) in &template.op {
451 if let serde_json::Value::String(s) = value {
452 let bps = bindpoints::extract_bind_points(s);
453 for bp in &bps {
454 if let BindPoint::Reference {
455 name, qualifier, ..
456 } = bp
457 {
458 let resolvable = match qualifier {
459 BindQualifier::Bind => program.resolve_output(name).is_some(),
460 BindQualifier::Capture => capture_names.contains(name),
461 BindQualifier::Input => {
462 program.input_names().contains(&name.to_string())
463 }
464 BindQualifier::None => {
465 program.resolve_output(name).is_some()
466 || capture_names.contains(name)
467 || program.input_names().contains(&name.to_string())
468 }
469 };
470 if !resolvable {
471 errors.push(format!(
472 "unresolved bind point '{{{name}}}' in op '{}' field '{field_name}'. \
473 Not found in Polydat bindings, captures, or inputs.",
474 template.name
475 ));
476 }
477 }
478 }
479 }
480 }
481 }
482
483 if errors.is_empty() {
484 Ok(())
485 } else {
486 for e in &errors {
487 crate::observer::log(crate::observer::LogLevel::Error, &format!("error: {e}"));
488 }
489 Err(format!("{} unresolved bind point(s)", errors.len()))
490 }
491}
492
493impl FiberBuilder {
494 /// Create a new fiber builder whose main kernel is the parent's own
495 /// image bound under the parent, on the parent's engine.
496 pub fn new(parent: &ScopeKernel) -> Self {
497 Self::with_image(parent, None)
498 }
499
500 /// Create a new fiber builder whose main kernel runs `image` — the
501 /// fiber engine's image of `parent`'s program — bound under
502 /// `parent`; the parent's own image when `image` is `None`.
503 /// Per-fiber state is fresh; cell handles are Arc-shared with the
504 /// parent so writes propagate to the workload-root through the
505 /// cascade.
506 pub fn with_image(parent: &ScopeKernel, image: Option<Arc<dyn KernelProgram>>) -> Self {
507 Self::with_scope(parent, image, Vec::new())
508 }
509
510 /// [`Self::with_image`], binding `scope_values` into the main kernel
511 /// as it is built, before its consts initialize, and remembering
512 /// them for the stanza-boundary [`Self::reset_captures`].
513 pub fn with_scope(
514 parent: &ScopeKernel,
515 image: Option<Arc<dyn KernelProgram>>,
516 scope_values: Vec<(String, Value)>,
517 ) -> Self {
518 let main_program = parent.program().clone();
519 let image: Arc<dyn KernelProgram> = image.unwrap_or_else(|| parent.image().clone());
520 let bindings = scope_bindings(&main_program, &scope_values);
521 let main_kernel = polydat::kernel::bind_under(parent.kernel(), image, &bindings)
522 .unwrap_or_else(|e| panic!("a fiber's main kernel failed to bind: {e}"));
523 let scope_value_main_idx = scope_values
524 .iter()
525 .map(|(name, _)| scope_value_index(&main_program, name))
526 .collect();
527 // Standing alone, the fiber knows one canonical: its parent's own
528 // program, the flattened op's. `OpBuilder::create_fiber_builder`
529 // replaces this with the activity's full set.
530 let canonicals = std::iter::once((
531 parent.program_id(),
532 CanonicalSource {
533 program: main_program.clone(),
534 image: parent.image().clone(),
535 module: None,
536 },
537 ))
538 .collect();
539 Self {
540 main_kernel,
541 main_program,
542 scope_values,
543 per_op_kernels: Vec::new(),
544 per_op_programs: Vec::new(),
545 per_op_side_effecting: Vec::new(),
546 scope_value_main_idx,
547 scope_value_per_op_idx: Vec::new(),
548 needs_broadcast: false,
549 canonicals: Arc::new(canonicals),
550 }
551 }
552
553 /// SRD-68 Push 3 — populate this fiber's per-op kernel slots
554 /// from the activity's dispenser registry. Walks each
555 /// dispenser, calls `dispenser.canonical_kernel()` to get the
556 /// dispenser-owned canonical kernel (when present), and binds a
557 /// per-fiber kernel of its program under this fiber's main kernel.
558 /// Slot positions match the dispenser registry's order so
559 /// cycle-time dispatch can index by `template_idx`. Dispensers
560 /// that return `None` (no GK needs) get a `None` slot —
561 /// `ExecCtx::wires` falls back to the `NullWireSource` baseline
562 /// for those cycles.
563 ///
564 /// A per-op kernel is instantiated from its op-template module on
565 /// the fiber engine when the fiber has one for the canonical's
566 /// program, carrying the module's `result:` write-throughs; any
567 /// other canonical program runs on the interpreter.
568 ///
569 /// Called once per fiber, right after spawn, before any cycles
570 /// run. Idempotent: re-attaching with the same registry is
571 /// a no-op since canonical kernels are stable across phase
572 /// activation.
573 pub fn attach_dispenser_kernels(
574 &mut self,
575 dispensers: &[std::sync::Arc<dyn crate::adapter::OpDispenser>],
576 ) {
577 let scope_values = self.scope_values.clone();
578 // SRD-13f Stage 1: per-op kernels descend from
579 // `fiber.main_kernel` (this fiber's per-fiber scope kernel for
580 // the current phase), NOT from the dispenser's shared
581 // `canonical_kernel`. The dispenser's canonical_kernel becomes
582 // a *program source*, so there is one consistent per-fiber
583 // chain:
584 // fiber.main_kernel → per_op_kernel
585 // Computed outputs on main_kernel are reachable from
586 // per_op_kernel via the standard scope-chain mechanism;
587 // per-fiber state (cycle, scope values) propagates
588 // correctly without external refresh.
589 // A canonical kernel may be on any engine: what it stands for —
590 // the interpreter program resolving its indices, and the module
591 // a per-op kernel is instantiated from — is looked up by its
592 // program's identity. A canonical this activity's builder did
593 // not hand out — an adapter's own kernel — stands for its
594 // interpreter program when it has one; a compiled one stands for
595 // nothing it can bind, and its dispenser runs against the fiber's
596 // main kernel.
597 let sources: Vec<Option<CanonicalSource>> = dispensers
598 .iter()
599 .map(|d| {
600 let kernel = d.canonical_kernel()?;
601 if let Some(source) = self.canonicals.get(&kernel.program_id()) {
602 return Some(source.clone());
603 }
604 let program = kernel.fork().into_program().as_interpreter();
605 if program.is_none() {
606 crate::diag!(
607 crate::observer::LogLevel::Warn,
608 "a dispenser's canonical kernel ({}) was not built by this \
609 activity; its op reads the fiber's main kernel",
610 kernel.engine()
611 );
612 }
613 program.map(CanonicalSource::interpreted)
614 })
615 .collect();
616 let dispenser_programs: Vec<Option<Arc<PolydatProgram>>> = sources
617 .iter()
618 .map(|s| s.as_ref().map(|s| s.program.clone()))
619 .collect();
620 let mut per_op_kernels: Vec<Option<Box<dyn Kernel>>> =
621 Vec::with_capacity(dispenser_programs.len());
622 let mut per_op_idx: Vec<Option<Vec<Option<usize>>>> =
623 Vec::with_capacity(dispenser_programs.len());
624 let mut per_op_side_effecting: Vec<Vec<usize>> =
625 Vec::with_capacity(dispenser_programs.len());
626 for maybe_source in &sources {
627 let Some(CanonicalSource {
628 program,
629 image,
630 module,
631 }) = maybe_source
632 else {
633 per_op_kernels.push(None);
634 per_op_idx.push(None);
635 per_op_side_effecting.push(Vec::new());
636 continue;
637 };
638 // Scope values are bound with the kernel, before its consts
639 // initialize; the indices are cached for `reset_captures`.
640 let bindings = scope_bindings(program, &scope_values);
641 let mut op_kernel = match module {
642 Some(module) => module
643 .instantiate_under(
644 self.main_kernel.as_ref(),
645 crate::fiber_engine::fiber_engine(),
646 &bindings,
647 )
648 .unwrap_or_else(|e| panic!("per-op kernel failed to instantiate: {e}")),
649 None => {
650 polydat::kernel::bind_under(self.main_kernel.as_ref(), image.clone(), &bindings)
651 .unwrap_or_else(|e| panic!("per-op kernel failed to bind: {e}"))
652 }
653 };
654 let idx_vec: Vec<Option<usize>> = scope_values
655 .iter()
656 .map(|(name, _)| scope_value_index(program, name))
657 .collect();
658 for init_name in program.const_outputs() {
659 let Some(idx) = program.output_index(init_name) else {
660 continue;
661 };
662 // Const warmup is best-effort: a const whose freeze
663 // fails here stays dirty and re-evals (or fails
664 // visibly) at its first per-cycle use. But the failure
665 // is never discarded silently — the enriched payload
666 // goes to the session log so a broken const is
667 // diagnosable before the per-cycle path trips over it.
668 if let Err(payload) = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
669 op_kernel.pull_at(idx);
670 })) {
671 let msg = payload
672 .downcast_ref::<&'static str>()
673 .map(|s| (*s).to_string())
674 .or_else(|| payload.downcast_ref::<String>().cloned())
675 .unwrap_or_else(|| "<non-string panic payload>".into());
676 crate::diag!(
677 crate::observer::LogLevel::Warn,
678 "const warmup pull '{init_name}' panicked during \
679 op-template kernel init (deferring to first \
680 per-cycle use): {msg}"
681 );
682 }
683 }
684 let side_effecting = program
685 .outputs_with_side_effects()
686 .iter()
687 .filter_map(|name| program.output_index(name))
688 .collect();
689 per_op_kernels.push(Some(op_kernel));
690 per_op_idx.push(Some(idx_vec));
691 per_op_side_effecting.push(side_effecting);
692 }
693 // SRD-13f Push D: the main kernel publishes its outputs through
694 // broadcast cells only when a per-op kernel has a *different*
695 // program from it. When the per-op kernel reuses the main
696 // program (the flattened-op-template path — no per-op matter),
697 // evaluating outputs on main is both redundant (the descendant
698 // evaluates the same wires locally) and harmful (side-effecting
699 // nodes like `testkit_throw_at` fire outside the per-op cascade
700 // surface, losing the panic-to-error pipeline). A per-op kernel
701 // with its own program carries `extern <name>` slots for
702 // cross-scope wires it doesn't replicate locally; those slots
703 // need main to compute and broadcast the values through the
704 // cell.
705 self.needs_broadcast = dispenser_programs
706 .iter()
707 .flatten()
708 .any(|p| !Arc::ptr_eq(p, &self.main_program));
709 self.per_op_kernels = per_op_kernels;
710 self.per_op_programs = dispenser_programs;
711 self.per_op_side_effecting = per_op_side_effecting;
712 self.scope_value_per_op_idx = per_op_idx;
713 }
714
715 /// The cycle-time wire surface for the firing dispenser at
716 /// `template_idx`: its per-fiber kernel, or this fiber's main
717 /// kernel when the dispenser exposes no canonical kernel (the
718 /// flattened path).
719 pub fn cycle_wires(&mut self, template_idx: usize) -> crate::wires::CycleWires<'_> {
720 let program = self
721 .per_op_programs
722 .get(template_idx)
723 .and_then(|p| p.clone());
724 match (
725 self.per_op_kernels
726 .get_mut(template_idx)
727 .and_then(|s| s.as_mut()),
728 program,
729 ) {
730 (Some(kernel), Some(program)) => {
731 crate::wires::CycleWires::over(kernel.as_mut(), program)
732 }
733 _ => {
734 crate::wires::CycleWires::over(self.main_kernel.as_mut(), self.main_program.clone())
735 }
736 }
737 }
738
739 /// The cycle-time wire surface over this fiber's main kernel.
740 pub fn main_wires(&mut self) -> crate::wires::CycleWires<'_> {
741 crate::wires::CycleWires::over(self.main_kernel.as_mut(), self.main_program.clone())
742 }
743
744 /// Get the per-fiber kernel for the firing dispenser at
745 /// `template_idx`. Returns `None` when the dispenser exposes
746 /// no canonical kernel (adapters with no Polydat needs); callers
747 /// fall back to the `NullWireSource` baseline.
748 pub fn per_op_kernel(&self, template_idx: usize) -> Option<&dyn Kernel> {
749 self.per_op_kernels
750 .get(template_idx)
751 .and_then(|s| s.as_deref())
752 }
753
754 /// This fiber's main kernel.
755 pub fn main_kernel(&self) -> &dyn Kernel {
756 self.main_kernel.as_ref()
757 }
758
759 /// The interpreter program this fiber's main kernel runs or was
760 /// imaged from.
761 pub fn program(&self) -> &Arc<PolydatProgram> {
762 &self.main_program
763 }
764
765 /// Set coordinates and begin a new evaluation scope.
766 ///
767 /// Bounded by each kernel's coordinate count: the slice is
768 /// truncated to the program's declared coordinate count
769 /// before being written. A phase kernel with no
770 /// coordinates (e.g. all bindings are invariant within the
771 /// stanza, only externs declared) has `coord_count = 0`,
772 /// so this becomes a no-op rather than clobbering the
773 /// extern slots that follow.
774 ///
775 /// SRD-13d Phase 9: the same coordinates are also written to
776 /// every per-op-template kernel that declares them as
777 /// coords. Each kernel binds its own input slot for `cycle`
778 /// (cascaded from parent) so per-cycle propagation is a
779 /// per-kernel `set_inputs`, not a chain walk.
780 pub fn set_inputs(&mut self, coords: &[u64]) {
781 let main_n = coords.len().min(self.main_kernel.coord_count());
782 if main_n > 0 {
783 self.main_kernel.set_inputs(&coords[..main_n]);
784 }
785 for kernel in self.per_op_kernels.iter_mut().flatten() {
786 let n = coords.len().min(kernel.coord_count());
787 if n > 0 {
788 kernel.set_inputs(&coords[..n]);
789 }
790 }
791 }
792
793 /// Feed a source item into the fiber's kernels.
794 ///
795 /// Sets the ordinal as the coordinate input and injects field
796 /// projections into the appropriate input slots (e.g.
797 /// `base__ordinal`, `base__vector`). The ordinal write is
798 /// skipped by kernels whose programs declare no coordinates
799 /// (only externs and stanza-invariant bindings) rather than
800 /// clobbering an extern slot. Field projections always write by
801 /// name, so they're safe regardless of coordinate count.
802 ///
803 /// SRD-13d Phase 9: ordinal + fields propagate to every
804 /// op-template kernel that declares matching slots.
805 pub fn set_source_item(&mut self, item: &polydat::iteration::source::SourceItem) {
806 if self.main_kernel.coord_count() > 0 {
807 self.main_kernel.set_inputs(&[item.ordinal]);
808 }
809 // Source items carry typed values from their upstream
810 // DataSource; one the slot's type cannot take, converted or
811 // not, means the source produced an incompatible value, which
812 // is a fail-loud condition.
813 for (name, value) in &item.fields {
814 if let Some(idx) = self.main_program.find_input(name) {
815 crate::wires::write_input(self.main_kernel.as_mut(), idx, name, value.clone())
816 .unwrap_or_else(|e| {
817 panic!("source item field '{name}' failed typed write: {e}")
818 });
819 }
820 }
821 // Cell-bound cross-fiber visibility is substrate-owned via the
822 // per-cell revision counter + per-scope intent-dirty vector
823 // (polydat/docs/design/cross_fiber_invalidation.md): ancestor
824 // writes through a `SharedCell` bump the cell's revision and
825 // set its intent bit, and the evaluator re-checks both before
826 // serving a memoized result. No host-side refresh call needed.
827 for (kernel, program) in self
828 .per_op_kernels
829 .iter_mut()
830 .zip(self.per_op_programs.iter())
831 {
832 let (Some(kernel), Some(program)) = (kernel, program) else {
833 continue;
834 };
835 if kernel.coord_count() > 0 {
836 kernel.set_inputs(&[item.ordinal]);
837 }
838 for (name, value) in &item.fields {
839 if let Some(idx) = program.find_input(name) {
840 crate::wires::write_input(kernel.as_mut(), idx, name, value.clone())
841 .unwrap_or_else(|e| {
842 panic!("source item field '{name}' failed typed write: {e}")
843 });
844 }
845 }
846 }
847 if self.needs_broadcast {
848 self.main_kernel.publish_broadcasts();
849 }
850 }
851
852 /// Reset capture inputs to defaults. Called at stanza
853 /// boundaries to prevent capture leakage across stanzas.
854 /// Coordinates and cell-bound slots are not reset. Scope-bound
855 /// iter-var inputs (set by [`OpBuilder::create_fiber_builder`])
856 /// are re-applied after the reset so the iteration's bound
857 /// values survive the boundary.
858 pub fn reset_captures(&mut self) {
859 self.main_kernel.reset_inputs();
860 for ((name, value), idx_opt) in self
861 .scope_values
862 .iter()
863 .zip(self.scope_value_main_idx.iter())
864 {
865 if let Some(idx) = idx_opt {
866 crate::wires::write_input(self.main_kernel.as_mut(), *idx, name, value.clone())
867 .unwrap_or_else(|e| panic!("scope value '{name}' failed typed write: {e}"));
868 }
869 }
870 for (slot, idx_slot) in self
871 .per_op_kernels
872 .iter_mut()
873 .zip(self.scope_value_per_op_idx.iter())
874 {
875 if let (Some(kernel), Some(idx_vec)) = (slot, idx_slot) {
876 kernel.reset_inputs();
877 for ((name, value), idx_opt) in self.scope_values.iter().zip(idx_vec.iter()) {
878 if let Some(idx) = idx_opt {
879 crate::wires::write_input(kernel.as_mut(), *idx, name, value.clone())
880 .unwrap_or_else(|e| {
881 panic!("scope value '{name}' failed typed write: {e}")
882 });
883 }
884 }
885 }
886 }
887 }
888
889 /// Invalidate all state: every step, a side channel included,
890 /// runs again when next pulled. Provides "clean slate" semantics.
891 pub fn invalidate_all(&mut self) {
892 self.main_kernel.invalidate_all();
893 for kernel in self.per_op_kernels.iter_mut().flatten() {
894 kernel.invalidate_all();
895 }
896 }
897
898 /// Store a captured value into the main kernel's input slot `name`.
899 /// Returns `true` when the slot exists and took the value, `false`
900 /// when the program has no such input or the value could not be
901 /// converted to its type (value dropped).
902 pub fn capture(&mut self, name: &str, value: Value) -> bool {
903 match self.main_program.find_input(name) {
904 Some(idx) => {
905 crate::wires::write_input(self.main_kernel.as_mut(), idx, name, value).is_ok()
906 }
907 None => false,
908 }
909 }
910
911 /// SRD-68 Push 5d: position-indexed write into the per-fiber
912 /// op-template kernel slot. Used by the cycle dispatch's
913 /// post-execute capture flow to feed result-binding inputs
914 /// (`body` / `count` / `ok` and any captures) into the kernel
915 /// before [`Self::commit_op_template_write_throughs_for_idx`]
916 /// fans the computed values up through parent `shared` cells.
917 ///
918 /// No-op returning `false` when (a) the op didn't materialise a
919 /// kernel (flattened op-template), or (b) the kernel doesn't
920 /// declare an input slot for `name` (the closure-binding economy
921 /// dropped it because the source doesn't reference it), or (c)
922 /// the value could not be converted to the slot's type.
923 pub fn write_op_template_input_for_idx(
924 &mut self,
925 template_idx: usize,
926 name: &str,
927 value: Value,
928 ) -> bool {
929 let (Some(Some(kernel)), Some(Some(program))) = (
930 self.per_op_kernels.get_mut(template_idx),
931 self.per_op_programs.get(template_idx),
932 ) else {
933 if nmbrs_dirty_debug_enabled() && name == "body" {
934 eprintln!("DIRTY: write body template={template_idx} NO_KERNEL");
935 }
936 return false;
937 };
938 let Some(idx) = program.find_input(name) else {
939 if nmbrs_dirty_debug_enabled() && name == "body" {
940 let inputs = program.input_names();
941 eprintln!(
942 "DIRTY: write body template={template_idx} NO_SLOT in_count={} names={:?}",
943 inputs.len(),
944 inputs
945 );
946 }
947 return false;
948 };
949 if nmbrs_dirty_debug_enabled() && name == "body" {
950 let display = value.to_display_string();
951 let head: String = display.chars().take(48).collect();
952 eprintln!(
953 "DIRTY: write body template={template_idx} idx={idx} in_count={} \
954 head=\"{head}\"",
955 program.input_names().len()
956 );
957 }
958 crate::wires::write_input(kernel.as_mut(), idx, name, value).is_ok()
959 }
960
961 /// SRD-68 Push 5d: position-indexed Rule 2 write-through commit.
962 /// Pulls every `__write_<X>` and stores its value through the
963 /// cell-bound input slot for `<X>`, propagating each result-
964 /// binding LHS value to the parent's `SharedCell` (and from
965 /// there to any sibling phase that imports the same name).
966 /// No-op when the kernel carries no write-throughs (typical
967 /// for ops without `result:`).
968 ///
969 /// Errors when a write-through violates cell type stability
970 /// (scope_model.md §"Type stability") — surfaced by the fiber
971 /// loop as a phase-stopping workload bug (it is deterministic:
972 /// every subsequent cycle would repeat it).
973 pub fn commit_op_template_write_throughs_for_idx(
974 &mut self,
975 template_idx: usize,
976 ) -> Result<(), String> {
977 let debug = polydat::library::debug_nodes_enabled();
978 let Some(kernel) = self
979 .per_op_kernels
980 .get_mut(template_idx)
981 .and_then(|s| s.as_mut())
982 else {
983 if debug {
984 crate::observer::log(
985 crate::observer::LogLevel::Debug,
986 &format!(
987 "commit_op_template_write_throughs_for_idx: template_idx {template_idx} \
988 has no per-fiber kernel slot"
989 ),
990 );
991 }
992 return Ok(());
993 };
994 if debug {
995 crate::observer::log(
996 crate::observer::LogLevel::Debug,
997 &format!(
998 "commit_op_template_write_throughs_for_idx: idx {template_idx} kernel found"
999 ),
1000 );
1001 }
1002 kernel.commit_write_throughs()
1003 }
1004
1005 /// Per-cycle: pull the op-template kernel's **side-effecting**
1006 /// outputs at `template_idx` so side-effecting nodes (`log_info`
1007 /// and friends) actually evaluate. Without this, captured wires
1008 /// whose only consumer is a write-through are pulled by
1009 /// `commit_write_throughs`, but captured wires that aren't shared
1010 /// with a parent never get pulled — their compute chain (including
1011 /// any side-effecting nodes inside it) stays dormant, and the
1012 /// diagnostic the workload asked for never fires.
1013 ///
1014 /// Only outputs whose cone contains a `Purity::SideChannel` node are
1015 /// pulled (the set is precomputed at attach time,
1016 /// [`PolydatProgram::outputs_with_side_effects`]). A side-effect-free
1017 /// output is **not** pulled here: re-evaluating it would do nothing,
1018 /// and for a volatile metric reader (a `metricsql_*` / `metric`
1019 /// objective binding) it would issue a live metrics query *every
1020 /// cycle*. Such values are evaluated only when actually consumed.
1021 /// No-op when the kernel has no side-effecting outputs.
1022 pub fn pull_all_op_template_outputs_for_idx(&mut self, template_idx: usize) {
1023 let (Some(indices), Some(Some(kernel))) = (
1024 self.per_op_side_effecting.get(template_idx),
1025 self.per_op_kernels.get_mut(template_idx),
1026 ) else {
1027 return;
1028 };
1029 for &idx in indices {
1030 let _ = kernel.pull_at(idx);
1031 }
1032 }
1033
1034 /// Materialize a [`PullPlan`] against this fiber's main kernel.
1035 /// O(plan_len) on the hot path, no name hashing — the plan
1036 /// holds pre-resolved indices.
1037 ///
1038 /// This is the cycle-time read path used by every wrapper that
1039 /// holds [`PullHandle`]s registered into the corresponding
1040 /// [`ScopeFixture`] at init (SRD 31 §"Pull plan vs bind plan",
1041 /// SRD 32 §"Init-Time Fixture and Consumer Self-Registration").
1042 ///
1043 /// [`PullPlan`]: crate::fixture::PullPlan
1044 /// [`PullHandle`]: crate::fixture::PullHandle
1045 /// [`ScopeFixture`]: crate::fixture::ScopeFixture
1046 pub fn resolve_pulls(
1047 &mut self,
1048 plan: &crate::fixture::PullPlan,
1049 ) -> crate::fixture::ResolvedPulls {
1050 plan.resolve(self.main_kernel.as_mut())
1051 }
1052
1053 /// SRD-68 Push 5d resolve path — picks the right kernel for the
1054 /// dispenser at `template_idx` and resolves the plan against it.
1055 /// When a per-fiber op-template kernel was instanced for that
1056 /// position (every adapter exposes `canonical_kernel()` so this is
1057 /// the typical case), it is used; otherwise the plan resolves
1058 /// against the fiber's main kernel (the flattened op-template path).
1059 /// Either way the plan is checked against the interpreter program
1060 /// the kernel was imaged from, whose indices it holds.
1061 pub fn resolve_pulls_for_idx(
1062 &mut self,
1063 template_idx: usize,
1064 plan: &crate::fixture::PullPlan,
1065 ) -> crate::fixture::ResolvedPulls {
1066 let program = self
1067 .per_op_programs
1068 .get(template_idx)
1069 .and_then(|p| p.clone());
1070 match (
1071 self.per_op_kernels
1072 .get_mut(template_idx)
1073 .and_then(|s| s.as_mut()),
1074 program,
1075 ) {
1076 (Some(kernel), Some(program)) => {
1077 plan.check_program_match(&program, template_idx);
1078 plan.resolve(kernel.as_mut())
1079 }
1080 _ => {
1081 plan.check_program_match(&self.main_program, template_idx);
1082 plan.resolve(self.main_kernel.as_mut())
1083 }
1084 }
1085 }
1086}
1087
1088#[cfg(test)]
1089mod tests {
1090 use super::*;
1091 use polydat::compile::assembly::{PolydatAssembler, WireRef};
1092 use polydat::library::arithmetic::Mod;
1093 use polydat::library::hash::Hash;
1094
1095 fn make_kernel() -> ScopeKernel {
1096 let mut asm = PolydatAssembler::new(vec!["cycle".into()]);
1097 asm.add_node(
1098 "hashed",
1099 Box::new(Hash::new()),
1100 vec![WireRef::input("cycle")],
1101 );
1102 asm.add_node(
1103 "user_id",
1104 Box::new(Mod::new(1_000_000)),
1105 vec![WireRef::node("hashed")],
1106 );
1107 asm.add_output("user_id", WireRef::node("user_id"));
1108 asm.add_output("hashed", WireRef::node("hashed"));
1109 asm.compile().unwrap().into()
1110 }
1111
1112 /// SRD-13f Stage 1 gate — verify that the per-fiber
1113 /// kernel chain is a single linear descent. The
1114 /// op-template (per-op) kernel must be built as a
1115 /// subscope of the fiber's main kernel, NOT of a
1116 /// separate shared canonical. Inner reads of cross-scope
1117 /// wires depend on this being a per-fiber-consistent
1118 /// chain so cell-attached values reach inner via outer's
1119 /// per-fiber pull writes.
1120 #[test]
1121 fn per_fiber_chain_is_linear_and_unshared() {
1122 use crate::adapter::OpDispenser;
1123 // The structural property under test: when a fiber
1124 // attaches per-op kernels from a dispenser-shaped
1125 // canonical, each per-op kernel must be a NEW
1126 // per-fiber instance (subscope of the fiber's main
1127 // kernel), not the shared canonical kernel itself.
1128 //
1129 // Two fibers attaching against the same canonical
1130 // must produce different per-op kernel instances.
1131 let workload_kernel = make_kernel();
1132 let builder = OpBuilder::new(workload_kernel);
1133
1134 // Stand up a shared canonical via the public API.
1135 let workload_src = "input cycle: u64\nfolded := 42\n";
1136 let canonical_program = polydat::dsl::compile::compile_polydat_interpreter(workload_src)
1137 .expect("compile probe canonical")
1138 .program()
1139 .clone();
1140 let canonical_kernel: std::sync::Arc<dyn polydat::Kernel> =
1141 builder.canonical_kernel_for_op("nonexistent");
1142 // For this probe we only need the canonical to expose
1143 // a program; reuse builder's source_kernel program.
1144 let _ = canonical_program;
1145
1146 struct ProbeDispenser(std::sync::Arc<dyn polydat::Kernel>);
1147 impl OpDispenser for ProbeDispenser {
1148 fn canonical_kernel(&self) -> Option<&std::sync::Arc<dyn polydat::Kernel>> {
1149 Some(&self.0)
1150 }
1151 fn execute<'a>(
1152 &'a self,
1153 _cycle: u64,
1154 _ctx: &'a crate::fixture::ExecCtx<'a>,
1155 ) -> std::pin::Pin<
1156 Box<
1157 dyn std::future::Future<
1158 Output = Result<
1159 crate::adapter::OpResult,
1160 crate::adapter::ExecutionError,
1161 >,
1162 > + Send
1163 + 'a,
1164 >,
1165 > {
1166 Box::pin(async move { Ok(crate::adapter::OpResult::default()) })
1167 }
1168 }
1169 let dispensers: Vec<std::sync::Arc<dyn OpDispenser>> = vec![std::sync::Arc::new(
1170 ProbeDispenser(canonical_kernel.clone()),
1171 )];
1172
1173 let mut fiber_a = builder.create_fiber_builder();
1174 fiber_a.attach_dispenser_kernels(&dispensers);
1175 let mut fiber_b = builder.create_fiber_builder();
1176 fiber_b.attach_dispenser_kernels(&dispensers);
1177
1178 let per_op_a = fiber_a.per_op_kernel(0).expect("per-op A attached");
1179 let per_op_b = fiber_b.per_op_kernel(0).expect("per-op B attached");
1180
1181 // SRD-13f invariant: each fiber's per-op kernel is a
1182 // distinct per-fiber instance, neither of them
1183 // pointing to the shared canonical kernel.
1184 assert!(
1185 !std::ptr::eq(
1186 per_op_a as *const dyn Kernel as *const (),
1187 canonical_kernel.as_ref() as *const dyn polydat::Kernel as *const ()
1188 ),
1189 "per_op_a must be a distinct per-fiber instance, \
1190 not the shared canonical",
1191 );
1192 assert!(
1193 !std::ptr::eq(
1194 per_op_b as *const dyn Kernel as *const (),
1195 canonical_kernel.as_ref() as *const dyn polydat::Kernel as *const ()
1196 ),
1197 "per_op_b must be a distinct per-fiber instance, \
1198 not the shared canonical",
1199 );
1200 assert!(
1201 !std::ptr::eq(
1202 per_op_a as *const dyn Kernel as *const (),
1203 per_op_b as *const dyn Kernel as *const ()
1204 ),
1205 "fiber A and fiber B must each have their own \
1206 per-op kernel instance",
1207 );
1208 }
1209
1210 /// SRD 11 §"Init Binding Contract" Plan B verification:
1211 /// after the activation kernel pulls an init binding, the
1212 /// pulled value must propagate to every fiber via
1213 /// `init_overrides`, and per-fiber pulls must read the seeded
1214 /// buffer rather than re-firing the eval.
1215 #[test]
1216 fn init_binding_fires_once_across_many_fibers() {
1217 use std::sync::Arc as StdArc;
1218 use std::sync::atomic::{AtomicU64, Ordering};
1219
1220 // Counting custom node: bumps a shared counter on every
1221 // eval call, returns U64(42). Tracks how many times its
1222 // eval body actually runs across the test.
1223 struct CountingNode {
1224 meta: polydat::ast::NodeMeta,
1225 calls: StdArc<AtomicU64>,
1226 }
1227 impl polydat::ast::PolydatNode for CountingNode {
1228 fn meta(&self) -> &polydat::ast::NodeMeta {
1229 &self.meta
1230 }
1231 fn eval(&self, _inputs: &[Value], outputs: &mut [Value]) {
1232 self.calls.fetch_add(1, Ordering::Relaxed);
1233 outputs[0] = Value::U64(42);
1234 }
1235 }
1236
1237 let calls = StdArc::new(AtomicU64::new(0));
1238 let mut asm = PolydatAssembler::new(vec!["cycle".into()]);
1239 // Compile-const seed expression — wires empty.
1240 asm.add_node(
1241 "ticks",
1242 Box::new(CountingNode {
1243 meta: polydat::ast::NodeMeta {
1244 name: "ticks".into(),
1245 outs: vec![polydat::ast::Port::new(
1246 "output",
1247 polydat::ast::PortType::U64,
1248 )],
1249 ins: vec![],
1250 },
1251 calls: calls.clone(),
1252 }),
1253 vec![],
1254 );
1255 asm.add_output("ticks", WireRef::node("ticks"));
1256 asm.mark_const_output("ticks");
1257
1258 let mut kernel = asm.compile().expect("compile");
1259 // Plan B normally runs in the executor; for this unit test
1260 // we simulate it by pulling the init binding once on the
1261 // activation kernel.
1262 let v = kernel.pull_ref("ticks").clone();
1263 assert_eq!(v, Value::U64(42));
1264 let after_pull = calls.load(Ordering::Relaxed);
1265 // The fold pass evaluates the node once, then ConstU64
1266 // replaces it (init binding is pure compile-const here),
1267 // so subsequent state pulls return the leaf const without
1268 // re-eval. Assert the call count never grows from here.
1269 let builder = OpBuilder::new(kernel);
1270
1271 // Spawn many fibers, each pulls the init binding. None
1272 // should trigger an eval — the post-fold leaf-const path
1273 // returns the constant directly.
1274 for _ in 0..32 {
1275 let mut fiber = builder.create_fiber_builder();
1276 fiber.set_inputs(&[0]);
1277 let pulled = {
1278 use crate::wires::WireSource as _;
1279 fiber.main_wires().get("ticks").expect("ticks resolves")
1280 };
1281 assert_eq!(pulled, Value::U64(42));
1282 }
1283 let after_fibers = calls.load(Ordering::Relaxed);
1284 assert_eq!(
1285 after_pull, after_fibers,
1286 "init binding 'ticks' eval must not re-fire across fibers \
1287 (eval calls before fibers: {after_pull}, after 32 fibers: {after_fibers})"
1288 );
1289 // Independent: confirm the eval ran at most once during
1290 // compile-time fold + the activation pull.
1291 assert!(
1292 after_fibers <= 1,
1293 "expected at most one eval across compile fold + activation pull, got {after_fibers}"
1294 );
1295 }
1296
1297 /// The fiber engine carries the scope tree: the phase scope, a
1298 /// fiber's main kernel bound under it, and a per-op kernel
1299 /// instantiated from its op-template module all run off the
1300 /// interpreter — and every value they compute is the interpreter's.
1301 #[test]
1302 fn fiber_kernels_run_on_the_fiber_engine_with_interpreter_values() {
1303 use crate::adapter::OpDispenser;
1304 use crate::wires::WireSource as _;
1305 use polydat::kernel::subcontext::{BodyFragment, SubcontextBuilder};
1306
1307 let phase_src = "input cycle: u64
1308h := mod(hash(cycle), 1000)
1309";
1310 let phase = Arc::new(ScopeKernel::compile(phase_src).expect("phase"));
1311 let reference_phase = Arc::new(ScopeKernel::from(
1312 polydat::dsl::compile::compile_polydat_interpreter(phase_src).expect("reference phase"),
1313 ));
1314
1315 let mut op = SubcontextBuilder::under(phase.kernel());
1316 op.body(BodyFragment::PolydatSource(
1317 "extern h: u64
1318scaled := h * 3 + 1
1319"
1320 .into(),
1321 ));
1322 let module = Arc::new(op.finalize().expect("op-template module"));
1323 let canonical: Arc<dyn polydat::Kernel> = Arc::from(
1324 polydat::kernel::bind_under(phase.kernel(), module.program().clone(), &[])
1325 .expect("canonical per-op kernel"),
1326 );
1327
1328 struct Probe(Arc<dyn polydat::Kernel>);
1329 impl OpDispenser for Probe {
1330 fn canonical_kernel(&self) -> Option<&Arc<dyn polydat::Kernel>> {
1331 Some(&self.0)
1332 }
1333 fn execute<'a>(
1334 &'a self,
1335 _cycle: u64,
1336 _ctx: &'a crate::fixture::ExecCtx<'a>,
1337 ) -> std::pin::Pin<
1338 Box<
1339 dyn std::future::Future<
1340 Output = Result<
1341 crate::adapter::OpResult,
1342 crate::adapter::ExecutionError,
1343 >,
1344 > + Send
1345 + 'a,
1346 >,
1347 > {
1348 Box::pin(async move { Ok(crate::adapter::OpResult::default()) })
1349 }
1350 }
1351 let dispensers: Vec<Arc<dyn OpDispenser>> = vec![Arc::new(Probe(canonical))];
1352
1353 let native =
1354 OpBuilder::new(phase.clone()).with_op_template_modules([("op".to_string(), module)]);
1355 let interpreted = OpBuilder::new(reference_phase);
1356 let mut fiber = native.create_fiber_builder();
1357 fiber.attach_dispenser_kernels(&dispensers);
1358 let mut reference = interpreted.create_fiber_builder();
1359 reference.attach_dispenser_kernels(&dispensers);
1360
1361 let on_interpreter = |k: &dyn Kernel| matches!(k.engine(), polydat::Engine::Interpreter(_));
1362 assert!(
1363 !on_interpreter(phase.kernel()),
1364 "phase scope on {}",
1365 phase.engine()
1366 );
1367 assert!(
1368 !on_interpreter(fiber.main_kernel()),
1369 "main kernel on {}",
1370 fiber.main_kernel().engine()
1371 );
1372 let per_op = fiber.per_op_kernel(0).expect("per-op kernel attached");
1373 assert!(
1374 !on_interpreter(per_op),
1375 "per-op kernel on {}",
1376 per_op.engine()
1377 );
1378 assert!(on_interpreter(reference.main_kernel()));
1379 assert!(on_interpreter(
1380 reference.per_op_kernel(0).expect("reference per-op")
1381 ));
1382
1383 for cycle in [0, 1, 7, 1_000_003] {
1384 fiber.set_inputs(&[cycle]);
1385 reference.set_inputs(&[cycle]);
1386 fiber.set_source_item(&polydat::iteration::source::SourceItem {
1387 ordinal: cycle,
1388 fields: Vec::new(),
1389 });
1390 reference.set_source_item(&polydat::iteration::source::SourceItem {
1391 ordinal: cycle,
1392 fields: Vec::new(),
1393 });
1394 for name in ["h", "scaled"] {
1395 assert_eq!(
1396 fiber.cycle_wires(0).get(name),
1397 reference.cycle_wires(0).get(name),
1398 "'{name}' at cycle {cycle}"
1399 );
1400 }
1401 }
1402 }
1403}