Skip to main content

polydat_core/kernel/
engines.rs

1// Copyright 2024-2026 Jonathan Shook
2// SPDX-License-Identifier: Apache-2.0
3
4//! Polydat evaluation engines: EngineCore (shared eval loop) and the two
5//! P1 engine types — PolydatState (dependent-list) and ProvScanState
6//! (provenance-scan).
7
8use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
9use std::sync::{Arc, Mutex, OnceLock};
10
11use super::WireSource;
12use super::program::PolydatProgram;
13use crate::ast::Value;
14
15/// Cached lookup of the `NBRS_DIRTY_DEBUG` env var. Called from
16/// the per-cycle hot path (`PolydatState::set_input`); reading the
17/// real `std::env::var` on every cycle costs ~30% of CPU on
18/// single-fiber dryrun benches (it walks the libc env table and
19/// formats a fresh CString each call). The OnceLock evaluates
20/// once on first touch and every subsequent call is one atomic
21/// load.
22fn nbrs_dirty_debug_enabled() -> bool {
23    static FLAG: OnceLock<bool> = OnceLock::new();
24    *FLAG.get_or_init(|| std::env::var("NBRS_DIRTY_DEBUG").is_ok())
25}
26
27/// A cross-kernel mutable cell for a `shared`-modifier wire.
28///
29/// When a `shared` output in an outer scope is bound into an
30/// inner kernel via `materialize_wiring_from_outer`, both kernels' input
31/// slots reference the same `SharedCell`. Writes from inner via
32/// `set_input` flow through to the cell; reads on either side
33/// pick up the latest value.
34///
35/// Concurrent writers serialize at the Mutex, and the observable
36/// value is **last-write-wins** in lock-acquisition order
37/// (scope_model.md §6.2).
38///
39/// ## Cross-fiber validity tracking
40///
41/// Each cell carries its own validity-tracking handles per
42/// `polydat/docs/design/cross_fiber_invalidation.md`:
43///
44/// - `revision: AtomicU64` — monotonic counter, bumped on every
45///   write. Consumer fibers cache the last revision they
46///   observed in their per-fiber `last_seen` map; a mismatch
47///   tells the cone walker to re-evaluate.
48/// - `scope_intent_dirty: Arc<AtomicU64>` — one intent word, shared
49///   with every other cell allocated from the same word. The
50///   cell's `bit` position is set on every write, allowing
51///   consumers to do an O(1) bulk check ("any cell in this
52///   scope dirty?") before drilling down to the per-cell
53///   revision compare.
54/// - `bit: u8` — this cell's position within its word. The scope
55///   keeps one `Arc<AtomicU64>` per 64 cells, grown on demand by
56///   the defining scope's `EngineCore::allocate_cell_bit`.
57///
58/// The reader contract (S5 §1.1) is preserved: a producer's
59/// `publish` writes value + revision + intent bit in three
60/// Release stores; a consumer's `check_clean` walk on its next
61/// read observes the change without any host-side ceremony.
62pub struct SharedCellInner {
63    /// Cell value. The mutex serialises concurrent writers and
64    /// gives readers single-value atomicity.
65    pub value: Mutex<Value>,
66    /// Monotonic revision counter. Bumped on every write
67    /// (Release); compared by consumers (Acquire) against
68    /// per-fiber `last_seen`.
69    pub revision: AtomicU64,
70    /// Defining scope's intent-dirty bit-vector. Shared by Arc
71    /// across every cell allocated by the same scope. On every
72    /// write the producer ORs `1 << self.bit` into this
73    /// (Release) so consumers' bulk-mask check sees the scope
74    /// as dirty.
75    pub scope_intent_dirty: Arc<AtomicU64>,
76    /// This cell's bit position in `scope_intent_dirty`. Stable
77    /// for the cell's lifetime.
78    pub bit: u8,
79}
80
81impl std::fmt::Debug for SharedCellInner {
82    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
83        f.debug_struct("SharedCellInner")
84            .field("revision", &self.revision.load(Ordering::Relaxed))
85            .field("bit", &self.bit)
86            .finish_non_exhaustive()
87    }
88}
89
90impl SharedCellInner {
91    /// Construct a new cell with the given initial value, bound
92    /// to the defining scope's intent-dirty word at the given
93    /// bit position within that word. Callers must allocate
94    /// `(word, bit)` via `EngineCore::allocate_cell_bit` —
95    /// the bit is not reusable for the cell's lifetime.
96    pub fn new(initial: Value, scope_intent_dirty: Arc<AtomicU64>, bit: u8) -> Self {
97        debug_assert!(
98            bit < 64,
99            "bit-within-word {bit} must be < 64; the allocator splits >64-bit \
100             scope vectors across multiple words"
101        );
102        Self {
103            value: Mutex::new(initial),
104            revision: AtomicU64::new(0),
105            scope_intent_dirty,
106            bit,
107        }
108    }
109
110    /// Producer-side write: replace the cell's value, bump the
111    /// revision, set the intent bit. Three Release stores
112    /// publish the write across any consumer fiber per
113    /// `cross_fiber_invalidation.md` §6. The mutex critical
114    /// section is held only for the value swap; the atomics
115    /// run outside it.
116    pub fn publish(&self, value: Value) {
117        {
118            let mut guard = self.value.lock().unwrap();
119            *guard = value;
120        }
121        self.revision.fetch_add(1, Ordering::Release);
122        self.scope_intent_dirty
123            .fetch_or(1u64 << self.bit, Ordering::Release);
124    }
125
126    /// Consumer-side read: snapshot the cell's value and the
127    /// revision it was published at. Returns a pair so the
128    /// caller can update its `last_seen[cell] = revision`
129    /// alongside taking the value, without a second cell access.
130    pub fn snapshot(&self) -> (Value, u64) {
131        // Acquire-load the revision first so the value read
132        // synchronises-with the producer's value publication.
133        // The mutex itself provides the memory barrier for the
134        // value, but the revision is read with explicit Acquire
135        // for the cross-fiber happens-before relation.
136        let value = self.value.lock().unwrap().clone();
137        let revision = self.revision.load(Ordering::Acquire);
138        (value, revision)
139    }
140}
141
142/// Externally-held handle to a shared cell. `Arc<SharedCellInner>`
143/// so a single cell can be referenced from many kernels at
144/// once. The handle is cheap to clone (Arc bump).
145pub type SharedCell = Arc<SharedCellInner>;
146
147/// Per-node cone metadata for cell-bound input dependencies.
148///
149/// Built lazily on first `check_cell_clean` per node and
150/// cached in [`EngineCore::cell_cones`]; invalidated by
151/// clearing the cache whenever cells are attached or detached.
152///
153/// The structure groups a node's cell-bound input dependencies
154/// by the defining scope's `intent_dirty` Arc (compared by
155/// `Arc::ptr_eq`). Each group carries the bulk-check
156/// `interest_mask` for the scope plus per-cell drill-down
157/// entries — implementing the bulk-mask + per-cell-revision
158/// protocol from `cross_fiber_invalidation.md` §5.
159#[derive(Debug, Default, Clone)]
160pub(crate) struct CellCone {
161    /// Cells grouped by defining-scope's `intent_dirty`. Empty
162    /// = no cell-bound deps; check returns trivially clean.
163    pub(crate) groups: Vec<CellConeGroup>,
164}
165
166#[derive(Debug, Clone)]
167pub(crate) struct CellConeGroup {
168    /// Defining scope's `intent_dirty` vector (Arc cloned from
169    /// the cells). Bulk-mask check: AND this against
170    /// `interest_mask`; if zero, every cell in this group is
171    /// clean for this consumer (modulo last_seen) — skip the
172    /// drill-down.
173    pub(crate) intent_dirty: Arc<AtomicU64>,
174    /// OR of `1 << cell.bit` for every cell in this group.
175    pub(crate) interest_mask: u64,
176    /// Per-cell drill-down entries. Each gives the bit position
177    /// in `intent_dirty` plus the input slot index where the
178    /// cell is attached, for revision compare against
179    /// `last_seen`.
180    pub(crate) cells: Vec<CellConeEntry>,
181}
182
183#[derive(Debug, Clone, Copy)]
184pub(crate) struct CellConeEntry {
185    pub(crate) bit: u8,
186    pub(crate) input_slot: usize,
187}
188
189/// One named shared cell propagated through the parent → child
190/// scope chain. Carried on `PolydatKernel` (and surfaced through
191/// `ScopeKernel::shared_cells_in_scope`) so a descendant whose
192/// program declares a matching input slot can attach the cell —
193/// even when intermediate scopes' bodies never name it and so
194/// have no input slot for it themselves.
195///
196/// Without this carrier, an ancestral `shared X := …` cell
197/// becomes invisible past the first intermediate scope under
198/// the closure-binding economy. With it, every spawn step
199/// computes "every cell visible at this scope" and threads the
200/// full set forward — the cascade is transitive by
201/// construction.
202#[derive(Clone, Debug)]
203pub struct SharedCellEntry {
204    /// The binding's name.
205    pub name: String,
206    /// The cell's declared type.
207    pub port_type: crate::ast::PortType,
208    /// The cell.
209    pub cell: SharedCell,
210}
211
212/// Set by a host runtime that catches worker panics and renders the full
213/// enriched diagnostic itself (the `errors:` block). When set,
214/// the re-raise hook below prints a single first-line notice
215/// instead of the full body; bare polydat consumers never set it
216/// and keep the full print.
217static PANIC_REPORTING_DOWNSTREAM: std::sync::atomic::AtomicBool =
218    std::sync::atomic::AtomicBool::new(false);
219
220/// Declare that a downstream reporter will render eval-panic
221/// diagnostics in full (see `PANIC_REPORTING_DOWNSTREAM`).
222pub fn set_panic_reporting_downstream(on: bool) {
223    PANIC_REPORTING_DOWNSTREAM.store(on, std::sync::atomic::Ordering::Relaxed);
224}
225
226thread_local! {
227    /// True while a node eval runs inside the enrichment
228    /// catch_unwind in `eval_node`. The suppression hook checks
229    /// this to swallow the raw std panic-hook print (bare payload
230    /// + backtrace pointer at the original panic site) — that
231    /// same panic is about to be caught, enriched with
232    /// node/output/input context, and re-raised via `panic_any`,
233    /// which fires the hook again with this flag clear. Net
234    /// effect: exactly ONE hook print, and it's the enriched one.
235    static EVAL_PANIC_CAPTURE: std::cell::Cell<bool> =
236        const { std::cell::Cell::new(false) };
237    /// Original panic location captured by the suppression hook
238    /// while the flag above is set. Folded into the enriched
239    /// message so the true `file:line` survives the re-raise
240    /// (the re-raised panic's own location points at the
241    /// re-raise site, which is useless).
242    static EVAL_PANIC_LOCATION: std::cell::RefCell<Option<String>> =
243        const { std::cell::RefCell::new(None) };
244    /// One-shot marker armed just before the enriched re-raise
245    /// when a downstream reporter exists: the hook prints a short
246    /// first-line notice for that panic instead of the full body.
247    static RERAISE_SHORT: std::cell::Cell<bool> =
248        const { std::cell::Cell::new(false) };
249}
250
251/// Install (once, process-wide) a panic hook that chains to the
252/// previously installed hook unless the current thread is inside
253/// the wrapped node eval, in which case it records the panic
254/// location and stays quiet.
255fn install_eval_panic_hook() {
256    static HOOK: std::sync::Once = std::sync::Once::new();
257    HOOK.call_once(|| {
258        let prev = std::panic::take_hook();
259        std::panic::set_hook(Box::new(move |info| {
260            if EVAL_PANIC_CAPTURE.with(|c| c.get()) {
261                let loc = info.location().map(|l| l.to_string());
262                EVAL_PANIC_LOCATION.with(|slot| *slot.borrow_mut() = loc);
263            } else if RERAISE_SHORT.with(|c| c.replace(false)) {
264                // The runtime will render the full enriched
265                // diagnostic in the phase error list; one short
266                // line keeps the terminal signal without the
267                // four-fold repeat.
268                let first = info
269                    .payload()
270                    .downcast_ref::<String>()
271                    .map(String::as_str)
272                    .and_then(|m| m.lines().next())
273                    .unwrap_or("<non-string panic payload>");
274                eprintln!("op eval panic (detail in phase errors): {first}");
275            } else {
276                prev(info);
277            }
278        }));
279    });
280}
281
282/// RAII guard arming the suppression hook for one wrapped eval.
283/// Saves and restores the previous flag value: nodes that drive
284/// sub-kernels (comprehensions, gk-call) nest evals, and each
285/// level's catch_unwind must see its own panics suppressed.
286pub(crate) struct EvalPanicCaptureGuard {
287    prev: bool,
288}
289
290impl EvalPanicCaptureGuard {
291    pub(crate) fn arm() -> Self {
292        install_eval_panic_hook();
293        let prev = EVAL_PANIC_CAPTURE.with(|c| c.replace(true));
294        EVAL_PANIC_LOCATION.with(|slot| slot.borrow_mut().take());
295        Self { prev }
296    }
297}
298
299impl Drop for EvalPanicCaptureGuard {
300    fn drop(&mut self) {
301        EVAL_PANIC_CAPTURE.with(|c| c.set(self.prev));
302    }
303}
304
305/// The text of a panic payload: a `String` or a `&str`, else a marker.
306pub(crate) fn panic_payload_text(payload: &(dyn std::any::Any + Send)) -> String {
307    payload
308        .downcast_ref::<&'static str>()
309        .map(|s| (*s).to_string())
310        .or_else(|| payload.downcast_ref::<String>().cloned())
311        .unwrap_or_else(|| "<non-string panic payload>".into())
312}
313
314/// Build the rich diagnostic message for a node-level eval panic, on
315/// every engine (engines.md §3.4): the original payload, the
316/// panic location the capture guard recorded, the node's function
317/// name, every output it feeds, the program's diagnostic context
318/// (typically the source path / scope label), and the input values,
319/// already formatted (the interpreter's `Value`s through
320/// [`format_value_for_diag`], a compiled kernel's slots through its
321/// decoder). This is what the user sees instead of the bare panic
322/// payload, and it reads the same whichever engine raised it.
323pub(crate) fn enrich_panic(
324    payload: Box<dyn std::any::Any + Send>,
325    node_name: &str,
326    output_names: &[&str],
327    context: &str,
328    inputs: &[String],
329) -> String {
330    let original = panic_payload_text(payload.as_ref());
331    // A payload that already carries node context came from a
332    // nested wrapped eval's re-raise; its captured "location" is
333    // the re-raise site, not the original panic — skip it.
334    let location_line = if original.contains("↳ in node") {
335        String::new()
336    } else {
337        EVAL_PANIC_LOCATION
338            .with(|slot| slot.borrow_mut().take())
339            .map(|loc| format!("\n  ↳ panicked at {loc}"))
340            .unwrap_or_default()
341    };
342    let outputs_label = if output_names.is_empty() {
343        "no declared output".to_string()
344    } else {
345        format!(
346            "output{} {}",
347            if output_names.len() == 1 { "" } else { "s" },
348            output_names.join(", ")
349        )
350    };
351    let mut input_label = String::new();
352    for (i, v) in inputs.iter().enumerate() {
353        if i > 0 {
354            input_label.push_str(", ");
355        }
356        input_label.push_str(&format!("[{i}]={v}"));
357    }
358    format!(
359        "{original}{location_line}\n  ↳ in node `{node_name}` ({outputs_label}) \
360         while evaluating {context}\n  \
361         ↳ inputs: [{input_label}]"
362    )
363}
364
365/// Drop the `↳ panicked at <file>:<line>` line from an enriched
366/// message, keeping the node, outputs and inputs.
367///
368/// For a failure at *build* — the compile-constant fold — the reason is
369/// a property of the program, and polydat's own source location reads
370/// as an internal defect rather than the diagnosis it is: the user sees
371/// "bad spec `s97`" and a line number in a file they do not have. The
372/// three compiled engines never carried it there, so stripping it on
373/// the interpreter is what makes the fold error the same sentence on
374/// every engine ([`crate::KernelError::ConstantFold`]).
375///
376/// Evaluation failures keep the location: a panic at run time is as
377/// likely to be a node's bug as a program's, and then it is the first
378/// thing worth knowing.
379pub(crate) fn without_panic_location(message: String) -> String {
380    message
381        .lines()
382        .filter(|l| !l.trim_start().starts_with("↳ panicked at "))
383        .collect::<Vec<_>>()
384        .join("\n")
385}
386
387/// Re-raise an enriched message as the interpreter does: through
388/// `panic_any`, so the hook prints it once, or prints the short notice
389/// when a downstream reporter renders the full body.
390pub(crate) fn reraise_enriched(enriched: String) -> ! {
391    if PANIC_REPORTING_DOWNSTREAM.load(std::sync::atomic::Ordering::Relaxed) {
392        RERAISE_SHORT.with(|c| c.set(true));
393    }
394    std::panic::panic_any(enriched)
395}
396
397/// The interpreter's enrichment: the node's name and outputs from the
398/// program, the inputs as the `Value`s it was called with.
399fn enrich_eval_panic(
400    payload: Box<dyn std::any::Any + Send>,
401    program: &PolydatProgram,
402    node_idx: usize,
403    inputs: &[Value],
404) -> String {
405    let node_name = program
406        .nodes
407        .get(node_idx)
408        .map(|n| n.meta().name.to_string())
409        .unwrap_or_else(|| format!("<unknown node #{node_idx}>"));
410    let mut output_names: Vec<&str> = program
411        .output_map_iter()
412        .filter_map(|(name, (n_idx, _))| {
413            if *n_idx == node_idx {
414                Some(name.as_str())
415            } else {
416                None
417            }
418        })
419        .collect();
420    output_names.sort();
421    let inputs: Vec<String> = inputs.iter().map(format_value_for_diag).collect();
422    enrich_panic(
423        payload,
424        &node_name,
425        &output_names,
426        program.context(),
427        &inputs,
428    )
429}
430
431/// Format a `Value` into a short diagnostic string. Strings are
432/// quoted + truncated; vectors print their length not contents.
433pub(crate) fn format_value_for_diag(v: &Value) -> String {
434    match v {
435        Value::U64(n) => format!("U64({n})"),
436        Value::F64(n) => format!("F64({n})"),
437        Value::Bool(b) => format!("Bool({b})"),
438        Value::Str(s) => {
439            let trimmed: String = s.chars().take(40).collect();
440            if s.chars().count() > 40 {
441                format!("Str({trimmed:?}…)")
442            } else {
443                format!("Str({trimmed:?})")
444            }
445        }
446        Value::None => "None".to_string(),
447        other => format!("{:?}", other.port_type()),
448    }
449}
450
451/// Shared evaluation state for all Polydat engines. Contains the node
452/// output buffers, input values, and the eval loop.
453/// Engine types wrap this and provide their own invalidation strategy.
454pub struct EngineCore {
455    /// Per-node output value buffers, reused across evaluations.
456    pub(crate) buffers: Vec<Vec<Value>>,
457    /// Per-node: true = cached output is valid, false = needs eval.
458    pub(crate) node_clean: Vec<bool>,
459    /// Current input values (coordinates + captures, all unified).
460    /// For a cell-bound slot this entry is unused: the cell is the
461    /// slot's only register (`read_input` reads it, `set_input`
462    /// publishes to it).
463    pub(crate) inputs: Vec<Value>,
464    /// Default values for each input (used by reset_inputs).
465    pub(crate) input_defaults: Vec<Value>,
466    /// Optional cross-kernel shared cell per input slot. `None`
467    /// = local-only input (the common case). `Some(cell)` =
468    /// the slot is bound to a shared cell; writes propagate
469    /// through the cell to whatever other kernels share it.
470    pub(crate) shared_cells: Vec<Option<SharedCell>>,
471    /// Per-output broadcast cell (cross_fiber_invalidation.md §3.1). Indexed
472    /// by output position in `program.output_list`. `Some(cell)`
473    /// = the output broadcasts its value to descendants via
474    /// the cell whenever the owner pulls the output; `None` =
475    /// no broadcast subscribers were set up (no descendant
476    /// scope binds against this output's name).
477    ///
478    /// `materialize_wiring_from_outer` plumbs the same `Arc<SharedCell>`
479    /// onto the matching input slot on the inner kernel — at
480    /// that point both ends share the storage. Inner reads
481    /// transparently through the cell on every `read_input`;
482    /// outer's `pull` writes the freshly computed value into
483    /// the cell so subsequent inner reads return the current
484    /// value with no traversal.
485    pub(crate) output_cells: Vec<Option<SharedCell>>,
486    /// Whether a descendant has taken one of `output_cells`: the one
487    /// check a pull makes before publishing, false for every kernel
488    /// with no subscope under it.
489    pub(crate) broadcasting: AtomicBool,
490    /// Pre-allocated scratch buffer for node input gathering.
491    pub(crate) input_scratch: Vec<Value>,
492    /// Per node, the scratch entries the node declared through
493    /// `scratch_layout` (a native cone's own slot buffer): storage
494    /// belongs to the state, never to the node, which is shared by
495    /// every state of the program (axiom S3).
496    pub(crate) node_scratch: Vec<Vec<crate::ast::ScratchBuf>>,
497    /// This scope's intent-dirty bit-vector. One `AtomicU64`
498    /// word per 64 cells allocated by this scope; new words
499    /// are appended on demand by [`Self::allocate_cell_bit`].
500    /// Each cell carries a clone of the specific `Arc<AtomicU64>`
501    /// for its word (and its bit-within-word). Consumer fibers'
502    /// bulk-mask check (per `cross_fiber_invalidation.md` §5)
503    /// groups cells by `Arc::ptr_eq` of their word and ANDs
504    /// the loaded word against the cone's interest mask for
505    /// that word.
506    ///
507    /// The `Vec<Arc<...>>` shape — rather than a single
508    /// `Arc<Vec<AtomicU64>>` — lets cells take a stable
509    /// per-word handle that the scope can grow without
510    /// invalidating any existing cell's reference.
511    pub(crate) scope_intent_words: Vec<Arc<AtomicU64>>,
512    /// Next bit position to allocate from
513    /// [`Self::scope_intent_words`]. Word index is
514    /// `next_cell_bit / 64`; bit within word is
515    /// `next_cell_bit % 64`. Monotonic; bits are never reused
516    /// within a scope's lifetime.
517    pub(crate) next_cell_bit: u32,
518    /// Per-fiber cache of the last revision this engine observed
519    /// for each cell it has read. Keyed by `Arc::as_ptr` of the
520    /// `SharedCellInner`. Sparse; entries are inserted lazily
521    /// on first observation via `check_cell_clean`.
522    ///
523    /// Per-fiber state — no contention. Pointer keys are stable
524    /// for the cell's lifetime; orphaned entries for dropped
525    /// cells are harmless (the handle is never observed again).
526    pub(crate) last_seen: std::collections::HashMap<*const SharedCellInner, u64>,
527    /// Per-node cone metadata for cell-bound input deps. Lazy:
528    /// `None` until first `check_cell_clean` for that node;
529    /// then built once and reused. Cleared in bulk on any
530    /// attach/detach of shared cells.
531    pub(crate) cell_cones: Vec<Option<CellCone>>,
532}
533
534// SAFETY: the only fields Rust will not mark Send/Sync itself are
535// `last_seen`'s `*const SharedCellInner` keys, which are compared by
536// identity and never dereferenced. Sync rests on one invariant: **no
537// `&self` method mutates the core.** Every mutation, `last_seen` and
538// `cell_cones` included, goes through `&mut self`, so any number of
539// threads may read one core at once. Hosts rely on that: a scope
540// parent is an `Arc` shared by every fiber bound under it, and
541// `Kernel: Sync` promises it on every engine (native_scope_trees.md
542// §4). A cache or counter reached through `&self` would break it
543// silently; put it behind a lock or an atomic, or take `&mut self`.
544unsafe impl Send for EngineCore {}
545unsafe impl Sync for EngineCore {}
546
547impl EngineCore {
548    /// Allocate the next bit position from this scope's
549    /// intent-dirty vector for a newly-created cell. Returns
550    /// the specific word's `Arc<AtomicU64>` plus the bit
551    /// position within that word. Grows
552    /// [`Self::scope_intent_words`] on demand — each new word
553    /// is a freshly-allocated `Arc<AtomicU64>` so existing
554    /// cells' references stay stable.
555    pub(crate) fn allocate_cell_bit(&mut self) -> (Arc<AtomicU64>, u8) {
556        let bit = self.next_cell_bit;
557        let word_idx = (bit / 64) as usize;
558        let bit_in_word = (bit % 64) as u8;
559        while self.scope_intent_words.len() <= word_idx {
560            self.scope_intent_words.push(Arc::new(AtomicU64::new(0)));
561        }
562        let word = self.scope_intent_words[word_idx].clone();
563        self.next_cell_bit += 1;
564        (word, bit_in_word)
565    }
566
567    /// Construct a new `SharedCell` bound to this scope's
568    /// intent-dirty vector. Convenience wrapper that allocates
569    /// a fresh bit and builds the cell — every cell creation
570    /// site goes through here so the scope's bit allocator
571    /// stays the single source of truth.
572    pub(crate) fn make_shared_cell(&mut self, initial: Value) -> SharedCell {
573        let (word, bit) = self.allocate_cell_bit();
574        Arc::new(SharedCellInner::new(initial, word, bit))
575    }
576}
577
578impl EngineCore {
579    /// Read an input slot's current value, transparent to whether
580    /// it's a plain slot or backed by a `SharedCell`. The
581    /// canonical read path used by both `eval_node` and
582    /// `PolydatState::get_input` — there's no separate "refresh" step
583    /// the caller must remember; the cell is queried on every
584    /// read.
585    ///
586    /// Cost: one Mutex lock per read on shared slots; a clone of
587    /// `inputs[idx]` on plain slots (Value's clone is cheap —
588    /// Arc-based for vectors, primitive copy otherwise).
589    #[inline]
590    pub(crate) fn read_input(&self, idx: usize) -> Value {
591        if let Some(cell) = self.shared_cells.get(idx).and_then(|c| c.as_ref()) {
592            return cell.value.lock().unwrap().clone();
593        }
594        self.inputs[idx].clone()
595    }
596
597    /// Build the cone metadata for `node_idx` — the per-scope
598    /// groups of cell-bound input dependencies, derived from
599    /// `program.input_provenance[node_idx]` and the cells
600    /// currently attached on this engine.
601    ///
602    /// Returns an empty `CellCone { groups: [] }` for nodes
603    /// with no cell-bound deps (the common case).
604    fn build_cell_cone(&self, program: &PolydatProgram, node_idx: usize) -> CellCone {
605        let empty = crate::kernel::ProvMask::empty();
606        let prov = program.input_provenance.get(node_idx).unwrap_or(&empty);
607        let mut groups: Vec<CellConeGroup> = Vec::new();
608        // Iterate set bits of `prov` directly: each bit is an
609        // input slot that flows into this node transitively.
610        for input_idx in prov.iter_ones() {
611            let Some(Some(cell)) = self.shared_cells.get(input_idx) else {
612                continue;
613            };
614            // Group by Arc-pointer identity of scope_intent_dirty.
615            let group_idx = groups
616                .iter()
617                .position(|g| Arc::ptr_eq(&g.intent_dirty, &cell.scope_intent_dirty));
618            let i = match group_idx {
619                Some(i) => i,
620                None => {
621                    groups.push(CellConeGroup {
622                        intent_dirty: cell.scope_intent_dirty.clone(),
623                        interest_mask: 0,
624                        cells: Vec::new(),
625                    });
626                    groups.len() - 1
627                }
628            };
629            groups[i].interest_mask |= 1u64 << cell.bit;
630            groups[i].cells.push(CellConeEntry {
631                bit: cell.bit,
632                input_slot: input_idx,
633            });
634        }
635        CellCone { groups }
636    }
637
638    /// Cross-fiber check: return `true` if this fiber's
639    /// `last_seen` is up-to-date for every cell in `node_idx`'s
640    /// cone (no cross-fiber writes since last observation).
641    /// Returns `false` if any cell's revision has advanced,
642    /// updating `last_seen` to reflect the new revisions in
643    /// preparation for the caller's re-evaluation.
644    ///
645    /// Per cross_fiber_invalidation.md §5: bulk-mask check
646    /// (one Acquire load + AND per scope group) early-outs
647    /// when nothing in the scope is dirty; per-cell drill-down
648    /// runs only on set bits.
649    fn check_cell_clean(&mut self, program: &PolydatProgram, node_idx: usize) -> bool {
650        // Lazy build the cone metadata.
651        if self.cell_cones.len() <= node_idx {
652            self.cell_cones.resize_with(node_idx + 1, || None);
653        }
654        if self.cell_cones[node_idx].is_none() {
655            let cone = self.build_cell_cone(program, node_idx);
656            self.cell_cones[node_idx] = Some(cone);
657        }
658
659        // First pass: walk the cone, collect mismatches. The
660        // immutable borrow of `self.cell_cones`,
661        // `self.shared_cells`, and `self.last_seen` coexist
662        // because they're disjoint fields of `self`.
663        let mut dirty: Vec<(*const SharedCellInner, u64, usize)> = Vec::new();
664        {
665            let cone = self.cell_cones[node_idx].as_ref().unwrap();
666            for group in &cone.groups {
667                let intent = group.intent_dirty.load(Ordering::Acquire);
668                let masked = intent & group.interest_mask;
669                if masked == 0 {
670                    continue;
671                }
672                for entry in &group.cells {
673                    if masked & (1u64 << entry.bit) == 0 {
674                        continue;
675                    }
676                    let Some(Some(cell)) = self.shared_cells.get(entry.input_slot) else {
677                        continue;
678                    };
679                    let r = cell.revision.load(Ordering::Acquire);
680                    let ptr = Arc::as_ptr(cell);
681                    let prev = self.last_seen.get(&ptr).copied().unwrap_or(0);
682                    if r != prev {
683                        dirty.push((ptr, r, entry.input_slot));
684                    }
685                }
686            }
687        }
688        let clean = dirty.is_empty();
689        // Second pass: update last_seen for every cell whose
690        // revision we observed has advanced. Done in a
691        // separate pass to release the cone borrow above.
692        //
693        // Updating `last_seen` CONSUMES the dirty signal for this
694        // fiber, so the re-evaluation it triggers must reach every
695        // memoized node between the dirty slot and any consumer —
696        // not just the node that happened to check first. The
697        // caller only re-evaluates the CHECKED node; its recursive
698        // upstream walk re-checks each parent's own cone, which
699        // reads the just-updated `last_seen` and comes back clean,
700        // so without this pass the intermediate buffers stay stale
701        // and the checked node recomputes from stale parents.
702        // Mirror `set_input`'s write-side rule on the read side: a
703        // detected cross-fiber write invalidates every node whose
704        // transitive input provenance covers the dirty slot.
705        if !clean {
706            // Exact multi-word mask, so slots >= 64 invalidate too.
707            let mut dirty_mask = crate::kernel::ProvMask::empty();
708            for (ptr, r, slot) in dirty {
709                self.last_seen.insert(ptr, r);
710                dirty_mask.set(slot);
711            }
712            for node_idx in 0..program.nodes.len() {
713                if program
714                    .input_provenance
715                    .get(node_idx)
716                    .is_some_and(|prov| prov.intersects(&dirty_mask))
717                {
718                    self.node_clean[node_idx] = false;
719                }
720            }
721        }
722        clean
723    }
724
725    /// Mark `cell_cones` as stale. Called after any change to
726    /// `shared_cells` that could affect the per-node cone
727    /// metadata (attach, detach). Next `check_cell_clean` on
728    /// any node will rebuild on demand.
729    pub(crate) fn invalidate_cell_cones(&mut self) {
730        for cone in self.cell_cones.iter_mut() {
731            *cone = None;
732        }
733    }
734
735    /// Evaluate a node by index. Shared by all engines.
736    /// Checks the clean flag, recursively evaluates upstream, gathers
737    /// inputs, calls node.eval(), marks clean.
738    pub fn eval_node(&mut self, program: &PolydatProgram, node_idx: usize) {
739        if self.node_clean[node_idx] {
740            // Memoization hit candidate — confirm cell-bound
741            // inputs in this node's cone are still at the
742            // revisions this fiber last observed. If any
743            // producer fiber has bumped a cell's revision since
744            // then, force a re-eval (the cache is stale even
745            // though `node_clean` is true) per
746            // cross_fiber_invalidation.md §5.
747            if self.check_cell_clean(program, node_idx) {
748                return;
749            }
750            self.node_clean[node_idx] = false;
751        }
752
753        let wiring = &program.wiring[node_idx];
754        for source in wiring.iter() {
755            if let WireSource::NodeOutput(upstream_idx, _) = source {
756                self.eval_node(program, *upstream_idx);
757            }
758        }
759
760        for (i, source) in wiring.iter().enumerate() {
761            self.input_scratch[i] = match source {
762                // `read_input` transparently reads the cell for
763                // `shared`-bound slots, so per-cycle eval picks
764                // up cross-kernel writes without any explicit
765                // refresh.
766                WireSource::Input(idx) => self.read_input(*idx),
767                WireSource::NodeOutput(upstream_idx, port_idx) => {
768                    self.buffers[*upstream_idx][*port_idx].clone()
769                }
770            };
771        }
772
773        let input_count = wiring.len();
774
775        // none_semantics.md Rule 1 — None propagation lifted to the kernel
776        // level. Any node whose inputs include `Value::None`
777        // emits `Value::None` on every output without invoking
778        // the node's `eval`. This holds the SQL-NULL / Rust
779        // `Option::?` propagation rule uniformly for ALL GK
780        // nodes, avoiding the dozens of duplicate per-node
781        // `if matches!(input, Value::None)` checks. Individual
782        // nodes (e.g. `Printf`) keep their checks redundant but
783        // harmless — the kernel guard fires first.
784        //
785        // Opt-out: nodes whose semantics explicitly consume
786        // `Value::None` (coalesce-style `default_or`, explicit
787        // optionality handlers per none_semantics.md Rule 2) override
788        // `PolydatNode::accepts_none_inputs` to skip this guard. Such
789        // nodes handle `None` in their own `eval`.
790        let node_ref = &*program.nodes[node_idx];
791        if !node_ref.accepts_none_inputs()
792            && self.input_scratch[..input_count]
793                .iter()
794                .any(|v| matches!(v, Value::None))
795        {
796            for slot in &mut self.buffers[node_idx] {
797                *slot = Value::None;
798            }
799            self.node_clean[node_idx] = true;
800            return;
801        }
802
803        // Wrap the node's eval in catch_unwind so a node-level
804        // panic (e.g. `Value::as_u64` on a Str) can be re-raised
805        // with the diagnostic context the user actually needs:
806        // which node panicked, which output(s) it feeds, what
807        // the input values were, and where in the source the
808        // node came from. Without this, the fiber-level catcher
809        // sees only the bare message — "expected U64, got Str"
810        // — and the user has no way to find the offending
811        // binding short of bisecting the workload.
812        //
813        // Cost: one catch_unwind frame per slow-path node eval.
814        // The JIT path doesn't go through here. On the success
815        // path the frame is a few stack words.
816        //
817        // The capture guard suppresses the std panic hook for
818        // the duration: without it, the hook prints the BARE
819        // payload ("expected U64, got F64" + backtrace) at the
820        // original panic site, before enrichment exists, and
821        // that raw print is the loudest thing the user sees.
822        // Re-raising with `panic_any` (not `resume_unwind`)
823        // fires the hook again — now unsuppressed — so the one
824        // message that prints is the enriched one.
825        let guard = EvalPanicCaptureGuard::arm();
826        let payload = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
827            program.nodes[node_idx].eval_in(
828                &mut self.node_scratch[node_idx],
829                &self.input_scratch[..input_count],
830                &mut self.buffers[node_idx],
831            );
832        }));
833        drop(guard);
834        if let Err(e) = payload {
835            // A native cone re-raises its member's failure already
836            // enriched with the member's name, its inputs, and this
837            // program's context (engines.md §3.4); the cone itself is not a frame,
838            // so the report reads as it does on every other engine.
839            if program.nodes[node_idx].fusion_subgraph().is_some()
840                && e.downcast_ref::<String>()
841                    .is_some_and(|s| s.contains("↳ in node"))
842            {
843                let enriched = *e.downcast::<String>().expect("checked above");
844                reraise_enriched(enriched);
845            }
846            let enriched =
847                enrich_eval_panic(e, program, node_idx, &self.input_scratch[..input_count]);
848            reraise_enriched(enriched);
849        }
850        self.node_clean[node_idx] = true;
851    }
852
853    /// Pull a named output.
854    pub fn pull(&mut self, program: &PolydatProgram, output_name: &str) -> &Value {
855        let (node_idx, port_idx) = *program
856            .output_map
857            .get(output_name)
858            .unwrap_or_else(|| panic!("unknown output variate: {output_name}"));
859        self.eval_node(program, node_idx);
860        if let Some(output_idx) = program.output_index(output_name) {
861            self.publish_output(output_idx, node_idx, port_idx);
862        }
863        &self.buffers[node_idx][port_idx]
864    }
865
866    /// Broadcast an output's freshly computed value
867    /// through its cell, so a descendant that bound its matching input
868    /// to the cell reads the current value next. Every pull does this,
869    /// by name or by index (cross_fiber_invalidation.md §3.1).
870    ///
871    /// Only once a descendant has taken a cell, and only to a cell a
872    /// descendant still holds: publishing clones the value and takes a
873    /// lock, which a pull with no reader should not pay. A descendant
874    /// bound later gets the current value when it asks for the cell
875    /// (`output_cell`).
876    #[inline]
877    pub(crate) fn publish_output(&self, output_idx: usize, node_idx: usize, port_idx: usize) {
878        if self.broadcasting.load(Ordering::Acquire) {
879            self.publish_output_cold(output_idx, node_idx, port_idx);
880        }
881    }
882
883    /// `publish_output` past its flag, out of line so the pull path
884    /// keeps only the check.
885    #[cold]
886    #[inline(never)]
887    fn publish_output_cold(&self, output_idx: usize, node_idx: usize, port_idx: usize) {
888        if let Some(Some(cell)) = self.output_cells.get(output_idx)
889            && Arc::strong_count(cell) > 1
890        {
891            // `publish` does the mutex write + revision bump +
892            // intent-bit set in three Release stores so the
893            // descendant's cone walker observes the change on
894            // its next read (cross_fiber_invalidation.md §5).
895            cell.publish(self.buffers[node_idx][port_idx].clone());
896        }
897    }
898
899    /// Allocate broadcast cells for every output in `program`
900    /// (cross_fiber_invalidation.md §3.1). Idempotent: if cells are already
901    /// allocated (size matches the program's output count),
902    /// the call is a no-op. Initial cell value is taken from
903    /// the current buffer (typically `Value::None` at
904    /// construction, before any pull has fired).
905    ///
906    /// Called from kernel constructors and from
907    /// `materialize_wiring_from_outer`-style operations that materialize
908    /// new descendants — the inner side needs the cell to
909    /// exist before it can attach to its input slot.
910    pub(crate) fn seed_output_cells(&mut self, program: &PolydatProgram) {
911        let n = program.output_names().len();
912        if self.output_cells.len() == n {
913            return;
914        }
915        // Two-pass to avoid borrowing `self` immutably (for
916        // buffer lookups) while also borrowing it mutably (for
917        // `make_shared_cell`). First collect initial values,
918        // then construct the cells.
919        let initials: Vec<Value> = (0..n)
920            .map(|i| {
921                let name = &program.output_list()[i].0;
922                let (node_idx, port_idx) = program.output_map[name];
923                // Defensive bounds-check: an output whose node has
924                // no buffer is seeded with `Value::None` rather than
925                // panicking.
926                self.buffers
927                    .get(node_idx)
928                    .and_then(|b| b.get(port_idx))
929                    .cloned()
930                    .unwrap_or(Value::None)
931            })
932            .collect();
933        self.output_cells = initials
934            .into_iter()
935            .map(|init| Some(self.make_shared_cell(init)))
936            .collect();
937    }
938
939    /// Output broadcast cell for the named output, if seeded, holding
940    /// the output's current value: pulls publish only while a
941    /// descendant holds the cell (`publish_output`), so a descendant
942    /// asking for it now is handed it up to date.
943    pub(crate) fn output_cell(&self, program: &PolydatProgram, name: &str) -> Option<SharedCell> {
944        let idx = program.output_index(name)?;
945        let cell = self.output_cells.get(idx)?.clone()?;
946        self.broadcasting.store(true, Ordering::Release);
947        let (node_idx, port_idx) = program.output_map[name];
948        if self.node_clean.get(node_idx).copied().unwrap_or(false) {
949            let current = &self.buffers[node_idx][port_idx];
950            if *cell.value.lock().unwrap() != *current {
951                cell.publish(current.clone());
952            }
953        }
954        Some(cell)
955    }
956}
957
958// =================================================================
959// PolydatState: dependent-list engine (default, O(affected) invalidation)
960// =================================================================
961
962/// Polydat evaluation engine using precomputed per-input dependent lists.
963///
964/// On `set_input()`, only nodes that depend on the changed input
965/// are dirtied. O(affected_nodes) per input change.
966/// This is the interpreter's state, the default of its three; the
967/// default engine for engine-less entry points is the compiled P3 engine.
968pub struct PolydatState {
969    /// Shared evaluation core (buffers, clean flags, inputs).
970    pub core: EngineCore,
971    /// Per-input dependent node lists for O(affected) invalidation.
972    input_dependents: Vec<Vec<usize>>,
973    /// Indices of non-deterministic nodes (zero-provenance, no declared inputs).
974    ///
975    /// These nodes produce a different value on every evaluation (e.g.,
976    /// `counter()`, `current_epoch_millis()`). They are unconditionally
977    /// marked dirty on every `set_input()` call so they are never cached.
978    nondeterministic_nodes: Vec<usize>,
979}
980
981impl PolydatState {
982    /// Construct a PolydatState from its component parts.
983    pub(crate) fn from_parts(
984        core: EngineCore,
985        input_dependents: Vec<Vec<usize>>,
986        nondeterministic_nodes: Vec<usize>,
987    ) -> Self {
988        Self {
989            core,
990            input_dependents,
991            nondeterministic_nodes,
992        }
993    }
994
995    /// Set all coordinate inputs at once. Wraps each u64 as
996    /// `Value::U64` and sets them at indices 0..N with per-input
997    /// change detection.
998    pub fn set_inputs(&mut self, coords: &[u64]) {
999        self.write_coordinates(coords);
1000    }
1001
1002    /// Write the coordinates: what construction does to seed a state's
1003    /// folded constants. A host's write is [`Self::set_inputs`].
1004    pub(crate) fn seed_inputs(&mut self, coords: &[u64]) {
1005        self.write_coordinates(coords);
1006    }
1007
1008    fn write_coordinates(&mut self, coords: &[u64]) {
1009        for (i, &c) in coords.iter().enumerate().take(self.core.inputs.len()) {
1010            self.core.inputs[i] = Value::U64(c);
1011            // Unconditional invalidation: the write itself is the
1012            // signal — see `set_input` for the rationale.
1013            if i < self.input_dependents.len() {
1014                for &node_idx in &self.input_dependents[i] {
1015                    self.core.node_clean[node_idx] = false;
1016                }
1017            }
1018        }
1019    }
1020
1021    /// Set a single input by index, dirtying only dependent nodes.
1022    ///
1023    /// Single-register semantics: a cell-bound slot's only
1024    /// register IS the cell — `set_input` writes through the
1025    /// cell. A non-cell slot's register is the local
1026    /// `inputs[idx]` array. There's no second snapshot kept in
1027    /// lockstep with the cell; reads always go to whichever is
1028    /// the slot's register.
1029    ///
1030    /// Dependents-marking is the dependent-list invalidation
1031    /// strategy carried by `PolydatState`; it's the write-side
1032    /// half of the engine's dirty-tracking. Other engines
1033    /// (`ProvScanState`) implement different
1034    /// strategies — see their own `set_inputs` impls.
1035    pub fn set_input(&mut self, idx: usize, value: Value) {
1036        if let Some(cell) = self.core.shared_cells.get(idx).and_then(|c| c.as_ref()) {
1037            // Cell-bound slot: the cell is the register. We do
1038            // NOT mirror the value into `inputs[idx]`; that
1039            // array slot is unused for cell-bound inputs.
1040            //
1041            // `publish` does the mutex write + revision bump +
1042            // intent-bit set in three Release stores so the
1043            // any other fiber's cone walker observes the
1044            // change on its next read
1045            // (cross_fiber_invalidation.md §5).
1046            cell.publish(value);
1047        } else {
1048            self.core.inputs[idx] = value;
1049        }
1050        // Mark every transitive dependent dirty unconditionally.
1051        // The act of writing an input IS the invalidation
1052        // signal — we don't gate on value equality because (a)
1053        // structural equality on rich Value variants
1054        // (Json/Bytes/VecF32) is expensive enough to defeat
1055        // the purpose of the optimisation, and (b) a same-
1056        // value rewrite is still a legitimate "the upstream
1057        // owner asked for a re-evaluation" signal that
1058        // downstream side-effecting nodes (`log_*`, audit
1059        // emitters, time-stamped observers) MUST honour.
1060        let dirty_debug = nbrs_dirty_debug_enabled();
1061        if idx < self.input_dependents.len() {
1062            if dirty_debug {
1063                eprintln!(
1064                    "DIRTY: set_input idx={idx} input_count={} dependents_for_idx={} \
1065                     total_input_dependents_len={}",
1066                    self.core.inputs.len(),
1067                    self.input_dependents[idx].len(),
1068                    self.input_dependents.len()
1069                );
1070            }
1071            for &node_idx in &self.input_dependents[idx] {
1072                self.core.node_clean[node_idx] = false;
1073            }
1074        } else if dirty_debug {
1075            eprintln!(
1076                "DIRTY: set_input idx={idx} OUT_OF_RANGE input_dependents_len={}",
1077                self.input_dependents.len()
1078            );
1079        }
1080    }
1081
1082    /// Begin a read: every volatile step is not current again, so the
1083    /// read evaluates each one the pulled cone reaches, once, and the
1084    /// steps downstream of it (runtime_model.md R1.v). Steps upstream of
1085    /// a volatile step keep their currency. A write does not re-arm a
1086    /// volatile step; only a read does.
1087    #[inline]
1088    pub(crate) fn rearm_volatile(&mut self) {
1089        for &idx in &self.nondeterministic_nodes {
1090            self.core.node_clean[idx] = false;
1091        }
1092    }
1093
1094    /// A pull within a read already begun with [`Self::rearm_volatile`]:
1095    /// several outputs read together see one evaluation of each
1096    /// volatile step.
1097    pub(crate) fn pull_in_read(&mut self, program: &PolydatProgram, output_name: &str) -> &Value {
1098        self.core.pull(program, output_name)
1099    }
1100
1101    /// Read the value of an input by index.
1102    ///
1103    /// Single-register read: cell-bound slots return the cell's
1104    /// current value; non-cell slots return the local register.
1105    /// One canonical value per slot, no stale snapshot.
1106    pub fn get_input(&self, idx: usize) -> Value {
1107        self.core.read_input(idx)
1108    }
1109
1110    /// Alias for [`Self::get_input`] under a more explicit name.
1111    /// Both read the cell when one is attached.
1112    pub fn read_input_value(&self, idx: usize) -> Value {
1113        self.core.read_input(idx)
1114    }
1115
1116    /// Attach a `SharedCell` to an input slot.
1117    ///
1118    /// After this call the cell becomes the slot's sole
1119    /// register: reads via `read_input` go through the cell,
1120    /// `set_input` writes through the cell. The local
1121    /// `inputs[idx]` array entry for this slot is unused for
1122    /// cell-bound slots — there is no second register kept in
1123    /// lockstep.
1124    ///
1125    /// Dependents are dirtied because the slot's effective
1126    /// value just changed from the local default to whatever
1127    /// the cell currently holds.
1128    pub fn attach_shared_cell(&mut self, idx: usize, cell: SharedCell) {
1129        if idx >= self.core.shared_cells.len() {
1130            self.core.shared_cells.resize(idx + 1, None);
1131        }
1132        self.core.shared_cells[idx] = Some(cell);
1133        if idx < self.input_dependents.len() {
1134            for &node_idx in &self.input_dependents[idx] {
1135                self.core.node_clean[node_idx] = false;
1136            }
1137        }
1138        // Cone metadata depends on which slots have cells; the
1139        // new attachment invalidates any cached cone groups.
1140        // Next `check_cell_clean` per node rebuilds on demand
1141        // per cross_fiber_invalidation.md §3.1.
1142        self.core.invalidate_cell_cones();
1143    }
1144
1145    /// Returns the `SharedCell` attached to an input slot, if any.
1146    /// Used by `materialize_wiring_from_outer` to share an existing cell with
1147    /// inner kernels.
1148    pub fn shared_cell(&self, idx: usize) -> Option<SharedCell> {
1149        self.core.shared_cells.get(idx).and_then(|c| c.clone())
1150    }
1151
1152    /// Reset a range of inputs to their defaults. Used at stanza
1153    /// boundaries to prevent capture leakage across stanzas.
1154    /// `from_idx` is typically `coord_count` (skip coordinates,
1155    /// reset only capture inputs).
1156    ///
1157    /// Cell-bound slots are skipped: the cell is cross-kernel
1158    /// shared state with its own lifecycle (managed by the
1159    /// owning ancestor scope), and a stanza-local reset must
1160    /// not clobber other kernels' views.
1161    pub fn reset_inputs_from(&mut self, from_idx: usize) {
1162        for i in from_idx..self.core.inputs.len() {
1163            // Cell-bound slots: the cell is the register, owned
1164            // by the ancestor that declared `shared X := init`.
1165            // Don't touch.
1166            if self.core.shared_cells.get(i).is_some_and(|c| c.is_some()) {
1167                continue;
1168            }
1169            if self.core.inputs[i] != self.core.input_defaults[i] {
1170                self.core.inputs[i] = self.core.input_defaults[i].clone();
1171                if i < self.input_dependents.len() {
1172                    for &node_idx in &self.input_dependents[i] {
1173                        self.core.node_clean[node_idx] = false;
1174                    }
1175                }
1176            }
1177        }
1178    }
1179
1180    /// Mark every node dirty and leave the inputs as they are: every
1181    /// node reruns at the next pull, as if the cycle had moved. What
1182    /// `Kernel::invalidate_all` means on every engine; a host that
1183    /// wants the inputs back at their defaults calls
1184    /// [`Self::reset_inputs_from`] as well.
1185    pub fn invalidate_all(&mut self) {
1186        self.core.node_clean.fill(false);
1187    }
1188
1189    /// Pull a named output variate from the program: one read.
1190    pub fn pull(&mut self, program: &PolydatProgram, output_name: &str) -> &Value {
1191        self.rearm_volatile();
1192        self.core.pull(program, output_name)
1193    }
1194
1195    /// Pre-populate a node's output buffer slot and mark it clean,
1196    /// suppressing on-demand evaluation. A caller seeds a state
1197    /// with a value another state over the same program already
1198    /// evaluated, so the node does not run again at first pull.
1199    pub fn seed_node_buffer(&mut self, node_idx: usize, port_idx: usize, value: Value) {
1200        if node_idx >= self.core.buffers.len() {
1201            return;
1202        }
1203        if port_idx >= self.core.buffers[node_idx].len() {
1204            return;
1205        }
1206        self.core.buffers[node_idx][port_idx] = value;
1207        self.core.node_clean[node_idx] = true;
1208    }
1209
1210    /// Read a node's output buffer slot, the counterpart of
1211    /// [`Self::seed_node_buffer`] for carrying an evaluated value
1212    /// from one state into another.
1213    pub fn node_buffer(&self, node_idx: usize, port_idx: usize) -> Option<&Value> {
1214        self.core
1215            .buffers
1216            .get(node_idx)
1217            .and_then(|ports| ports.get(port_idx))
1218    }
1219
1220    /// Pull an output by index (declaration order). Only evaluates
1221    /// the computation cone for this specific output.
1222    pub fn pull_by_index(&mut self, program: &PolydatProgram, output_idx: usize) -> &Value {
1223        self.rearm_volatile();
1224        let (node_idx, port_idx) = program.resolve_output_by_index(output_idx);
1225        self.core.eval_node(program, node_idx);
1226        // A pull by index publishes as a pull by name does.
1227        self.core.publish_output(output_idx, node_idx, port_idx);
1228        &self.core.buffers[node_idx][port_idx]
1229    }
1230
1231    /// Pull all outputs in declaration order, as one read.
1232    pub fn pull_all<'a>(&'a mut self, program: &PolydatProgram) -> Vec<&'a Value> {
1233        self.rearm_volatile();
1234        for i in 0..program.output_count() {
1235            let (node_idx, _) = program.resolve_output_by_index(i);
1236            self.core.eval_node(program, node_idx);
1237        }
1238        (0..program.output_count())
1239            .map(|i| {
1240                let (ni, pi) = program.resolve_output_by_index(i);
1241                &self.core.buffers[ni][pi]
1242            })
1243            .collect()
1244    }
1245
1246    /// Create a memoized accessor for a named subset of outputs.
1247    /// Resolves names to indices once; subsequent access uses indices only.
1248    pub fn accessor(program: &PolydatProgram, names: &[&str]) -> OutputAccessor {
1249        let indices: Vec<usize> = names
1250            .iter()
1251            .filter_map(|n| program.output_index(n))
1252            .collect();
1253        OutputAccessor { indices }
1254    }
1255
1256    /// Evaluate a node by index (exposed for constant folding in PolydatProgram).
1257    pub(crate) fn eval_node_public(&mut self, program: &PolydatProgram, node_idx: usize) {
1258        self.core.eval_node(program, node_idx);
1259    }
1260}
1261
1262/// Memoized output accessor for a named subset of outputs.
1263///
1264/// Created once from output names via `PolydatState::accessor()`.
1265/// Subsequent pulls use pre-resolved indices — no name lookups.
1266pub struct OutputAccessor {
1267    indices: Vec<usize>,
1268}
1269
1270impl OutputAccessor {
1271    /// Pull all outputs in this accessor from the given state.
1272    pub fn pull_all<'a>(
1273        &self,
1274        state: &'a mut PolydatState,
1275        program: &PolydatProgram,
1276    ) -> Vec<&'a Value> {
1277        for &idx in &self.indices {
1278            let (node_idx, _) = program.resolve_output_by_index(idx);
1279            state.core.eval_node(program, node_idx);
1280        }
1281        self.indices
1282            .iter()
1283            .map(|&idx| {
1284                let (ni, pi) = program.resolve_output_by_index(idx);
1285                &state.core.buffers[ni][pi]
1286            })
1287            .collect()
1288    }
1289
1290    /// Number of outputs in this accessor.
1291    pub fn len(&self) -> usize {
1292        self.indices.len()
1293    }
1294
1295    /// Whether this accessor has no outputs.
1296    pub fn is_empty(&self) -> bool {
1297        self.indices.is_empty()
1298    }
1299}
1300
1301// =================================================================
1302// ProvScanState: provenance-scan engine (O(all) invalidation)
1303// =================================================================
1304
1305/// Polydat evaluation engine using provenance bitmask scanning.
1306///
1307/// On `set_inputs()`, scans ALL nodes and checks each node's
1308/// provenance bitmask against the changed-inputs mask.
1309/// O(all_nodes) per input change regardless of how many changed.
1310pub struct ProvScanState {
1311    /// Shared evaluation core.
1312    pub core: EngineCore,
1313    input_provenance: Vec<crate::kernel::ProvMask>,
1314    /// Indices of non-deterministic nodes.
1315    nondeterministic_nodes: Vec<usize>,
1316}
1317
1318impl ProvScanState {
1319    /// Construct a ProvScanState from its component parts.
1320    pub(crate) fn from_parts(
1321        core: EngineCore,
1322        input_provenance: Vec<crate::kernel::ProvMask>,
1323        nondeterministic_nodes: Vec<usize>,
1324    ) -> Self {
1325        Self {
1326            core,
1327            input_provenance,
1328            nondeterministic_nodes,
1329        }
1330    }
1331
1332    /// Set new input values and invalidate affected nodes. Volatile
1333    /// nodes are re-armed by the read, not here.
1334    pub fn set_inputs(&mut self, coords: &[u64]) {
1335        let mut mask = crate::kernel::ProvMask::empty();
1336        for (i, &c) in coords.iter().enumerate().take(self.core.inputs.len()) {
1337            self.core.inputs[i] = Value::U64(c);
1338            // Unconditional: writing the input IS the
1339            // invalidation signal regardless of value equality.
1340            mask.set(i);
1341        }
1342        if !mask.is_zero() {
1343            for (i, clean) in self.core.node_clean.iter_mut().enumerate() {
1344                if *clean && self.input_provenance[i].intersects(&mask) {
1345                    *clean = false;
1346                }
1347            }
1348        }
1349    }
1350
1351    /// Pull a named output variate from the program: one read, which
1352    /// re-arms every volatile node first.
1353    pub fn pull(&mut self, program: &PolydatProgram, output_name: &str) -> &Value {
1354        for &idx in &self.nondeterministic_nodes {
1355            self.core.node_clean[idx] = false;
1356        }
1357        self.core.pull(program, output_name)
1358    }
1359}