Skip to main content

polydat_core/kernel/
activation.rs

1// Copyright 2024-2026 Jonathan Shook
2// SPDX-License-Identifier: Apache-2.0
3
4//! The activation runtime for `for` traversals (for_traversal.md §3.4,
5//! §3.6, §5).
6//!
7//! A [`TraversalStream`] dispenses one [`Activation`] per tuple of a
8//! compiled traversal. An activation is a fresh state over the body's
9//! program, compiled once at parent compile time, with the tuple's
10//! elements and the parent's cascaded wires bound and every `over`
11//! cursor narrowed. Activation never compiles: the cost of the second
12//! activation is the cost of the first minus nothing, because the first
13//! did not compile either.
14//!
15//! Cycles follow the rule in §3.4: a body with a cursor iterates its
16//! narrowest cursor extent, one cycle per ordinal, with the cursor's
17//! ordinal written before each pull; a body without a cursor has one
18//! cycle per activation.
19
20use std::sync::Arc;
21
22use crate::ast::Value;
23use crate::dsl::traversal::Traversal;
24use crate::iteration::comprehension::runtime::{IndexedTuples, evaluate_indexed};
25use crate::iteration::cursor_partition::{cursor_extent_on, cursor_over_partitions_on};
26use crate::kernel::Kernel;
27use crate::kernel::interp::Layered;
28
29use super::{PolydatKernel, PolydatProgram};
30
31/// The interval of ordinals an activation's cursor iterates.
32#[derive(Debug, Clone, PartialEq, Eq)]
33pub struct CursorSlice {
34    /// The cursor's name.
35    pub cursor: String,
36    /// The first ordinal of the slice.
37    pub start: u64,
38    /// One past the last ordinal.
39    pub end: u64,
40}
41
42impl CursorSlice {
43    /// Ordinals in the slice.
44    pub fn len(&self) -> u64 {
45        self.end.saturating_sub(self.start)
46    }
47
48    /// Whether the slice has no ordinal.
49    pub fn is_empty(&self) -> bool {
50        self.len() == 0
51    }
52}
53
54/// One child scope of a traversal: the tuple it was activated for, a
55/// fresh kernel over the body's program, and its cursor slice if the
56/// body declares a cursor.
57///
58/// The kernel is on the stream's engine — the engine of the kernel
59/// that opened the traversal — and is driven through the [`Kernel`]
60/// trait. [`TraversalStream::activation_on`] names another engine for
61/// a caller that wants one (for_traversal.md §5.2).
62pub struct Activation<K = PolydatKernel> {
63    /// Position of this activation's tuple in the traversal's dispense
64    /// order.
65    pub index: u64,
66    /// The tuple, in element order.
67    pub coords: Vec<(String, Value)>,
68    /// A fresh kernel over the body's shared program.
69    pub kernel: K,
70    /// The narrowest cursor slice, when the body declares a cursor.
71    pub cursor: Option<CursorSlice>,
72}
73
74impl std::fmt::Debug for Activation {
75    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
76        f.debug_struct("Activation")
77            .field("index", &self.index)
78            .field("coords", &self.coords)
79            .field("cursor", &self.cursor)
80            .field("program_nodes", &self.kernel.program().node_count())
81            .finish()
82    }
83}
84
85impl std::fmt::Debug for Activation<Box<dyn Kernel>> {
86    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
87        f.debug_struct("Activation")
88            .field("index", &self.index)
89            .field("coords", &self.coords)
90            .field("cursor", &self.cursor)
91            .field("engine", &self.kernel.engine())
92            .finish()
93    }
94}
95
96impl<K> Activation<K> {
97    /// Cycles this activation runs under §3.4: the cursor slice length,
98    /// or one when the body has no cursor.
99    pub fn cycle_count(&self) -> u64 {
100        match &self.cursor {
101            Some(slice) => slice.len(),
102            None => 1,
103        }
104    }
105
106    /// The value of one coordinate.
107    pub fn coord(&self, name: &str) -> Option<&Value> {
108        self.coords.iter().find(|(n, _)| n == name).map(|(_, v)| v)
109    }
110}
111
112impl Activation<Box<dyn Kernel>> {
113    /// Position the kernel at cycle `i` and return it ready to pull, as
114    /// the interpreter activation's `cycle` does.
115    pub fn cycle(&mut self, i: u64) -> &mut dyn Kernel {
116        self.kernel.set_inputs(&[i]);
117        if let Some(slice) = &self.cursor {
118            let ordinal = slice.start.saturating_add(i);
119            let slot = format!("{}__ordinal", slice.cursor);
120            // A body without the projection has no slot to write.
121            let _ = self.kernel.set_input(&slot, Value::U64(ordinal));
122        }
123        self.kernel.as_mut()
124    }
125
126    /// Run `f` once per cycle, in order.
127    pub fn for_each_cycle(&mut self, mut f: impl FnMut(u64, &mut dyn Kernel)) {
128        for i in 0..self.cycle_count() {
129            let kernel = self.cycle(i);
130            f(i, kernel);
131        }
132    }
133}
134
135impl Activation {
136    /// Position the kernel at cycle `i` and return it ready to pull.
137    /// The body's `cycle` coordinate is the local index; the cursor's
138    /// ordinal slot receives the absolute ordinal.
139    pub fn cycle(&mut self, i: u64) -> &mut PolydatKernel {
140        self.kernel.set_inputs(&[i]);
141        if let Some(slice) = &self.cursor {
142            let ordinal = slice.start.saturating_add(i);
143            let slot = format!("{}__ordinal", slice.cursor);
144            if let Some(idx) = self.kernel.program().find_input(&slot) {
145                self.kernel.state().set_input(idx, Value::U64(ordinal));
146            }
147        }
148        &mut self.kernel
149    }
150
151    /// Run `f` once per cycle, in order.
152    pub fn for_each_cycle(&mut self, mut f: impl FnMut(u64, &mut PolydatKernel)) {
153        for i in 0..self.cycle_count() {
154            let kernel = self.cycle(i);
155            f(i, kernel);
156        }
157    }
158}
159
160/// Dispenses activations for one traversal of a kernel.
161///
162/// The stream holds the comprehension's tuples addressed by position
163/// ([`IndexedTuples`]): opening evaluated every source, and each
164/// activation computes its own tuple, so `len`, `seek`, and
165/// `activation(index)` cost the same at any index, and an order's
166/// selection over a large product holds only the selected positions.
167pub struct TraversalStream {
168    traversal: Traversal,
169    tuples: IndexedTuples,
170    cascade: Vec<(String, Value)>,
171    next: usize,
172    /// The engine of the kernel that opened this traversal. Its
173    /// activations run there: a body belongs to the kernel that opened
174    /// it, and runs where that kernel runs.
175    engine: crate::Engine,
176}
177
178impl TraversalStream {
179    /// Number of activations the traversal dispenses.
180    pub fn len(&self) -> usize {
181        usize::try_from(self.tuples.len()).unwrap_or(usize::MAX)
182    }
183
184    /// Whether the traversal dispenses no activation.
185    pub fn is_empty(&self) -> bool {
186        self.tuples.is_empty()
187    }
188
189    /// The traversal this stream dispenses.
190    pub fn traversal(&self) -> &Traversal {
191        &self.traversal
192    }
193
194    /// Move the dispense position. Every strategy is a decidable
195    /// permutation, so seeking costs nothing beyond the index.
196    pub fn seek(&mut self, index: usize) {
197        self.next = index.min(self.len());
198    }
199
200    /// Current dispense position.
201    pub fn position(&self) -> usize {
202        self.next
203    }
204
205    /// The engine this stream's activations run on: the engine of the
206    /// kernel that opened the traversal.
207    pub fn engine(&self) -> crate::Engine {
208        self.engine
209    }
210
211    /// The body's program on this stream's engine — the one every
212    /// activation is a kernel over, compiled once and shared, which is
213    /// what makes activation cost no compile.
214    pub fn body_program(&self) -> Result<std::sync::Arc<dyn crate::kernel::KernelProgram>, String> {
215        self.traversal
216            .program_on(self.engine)
217            .map_err(|e| e.to_string())
218    }
219
220    /// The next activation, or `None` when exhausted.
221    pub fn advance(&mut self) -> Result<Option<Activation<Box<dyn Kernel>>>, String> {
222        if self.next >= self.len() {
223            return Ok(None);
224        }
225        let i = self.next;
226        self.next += 1;
227        self.activation(i).map(Some)
228    }
229
230    /// Build the activation at `index` without moving the dispense
231    /// position. Fibers partition a traversal by calling this over
232    /// disjoint index ranges.
233    ///
234    /// The kernel is on the stream's engine — the one the kernel that
235    /// opened the traversal runs on.
236    pub fn activation(&self, index: usize) -> Result<Activation<Box<dyn Kernel>>, String> {
237        self.activation_on(index, self.engine)
238    }
239
240    /// [`Self::activation`] on `engine` (engines.md §3.6): a fresh
241    /// kernel over the body's program for that engine, compiled once
242    /// per engine and shared by every activation after, driven through
243    /// the [`Kernel`] trait with the same elements, cascade, and cursor
244    /// narrowing. An engine that cannot run the body says so by name.
245    pub fn activation_on(
246        &self,
247        index: usize,
248        engine: crate::Engine,
249    ) -> Result<Activation<Box<dyn Kernel>>, String> {
250        let tuple = self.tuples.get(index as u64).ok_or_else(|| {
251            format!(
252                "activation index {index} is out of range; traversal has {} tuples",
253                self.tuples.len()
254            )
255        })?;
256        let program = self
257            .traversal
258            .program_on(engine)
259            .map_err(|e| e.to_string())?;
260        let mut kernel = program.create_uninitialized();
261        bind_by_name_on(kernel.as_mut(), &tuple)?;
262        bind_by_name_on(kernel.as_mut(), &self.cascade)?;
263        let cursor = narrow_cursors_on(kernel.as_mut())?;
264        // The activation's consts are evaluated once its tuple, cascade,
265        // and cursor slice are bound.
266        kernel.init().map_err(|e| e.to_string())?;
267        Ok(Activation {
268            index: index as u64,
269            coords: tuple,
270            kernel,
271            cursor,
272        })
273    }
274}
275
276/// Bind the inputs the body declares among `values`, through the trait.
277fn bind_by_name_on(kernel: &mut dyn Kernel, values: &[(String, Value)]) -> Result<(), String> {
278    let declared: std::collections::HashSet<String> = kernel.input_names().into_iter().collect();
279    for (name, value) in values {
280        if declared.contains(name) {
281            kernel
282                .set_input(name, value.clone())
283                .map_err(|e| e.to_string())?;
284        }
285    }
286    Ok(())
287}
288
289/// Resolve every `over` clause in the body and narrow its cursor.
290/// Returns the narrowest slice, or the full extent of the first cursor
291/// when none has an `over` clause. One routine for every engine,
292/// through the trait.
293fn narrow_cursors_on(kernel: &mut dyn Kernel) -> Result<Option<CursorSlice>, String> {
294    let schemas: Vec<crate::iteration::source::SourceSchema> = kernel.cursor_schemas().to_vec();
295    let mut narrowest: Option<CursorSlice> = None;
296    for schema in &schemas {
297        let slice = if schema.partition_output.is_some() {
298            let parts = cursor_over_partitions_on(kernel, schema)?;
299            let partition = match parts.len() {
300                1 => parts[0],
301                0 => {
302                    return Err(format!(
303                        "cursor '{}': its `over` value resolved to no partitions",
304                        schema.name
305                    ));
306                }
307                n => {
308                    return Err(format!(
309                        "cursor '{}': its `over` value resolved to {n} partitions; inside a traversal, bind the list \
310                     with an enclosing `for p in ...` and declare the cursor `over p`",
311                        schema.name
312                    ));
313                }
314            };
315            kernel
316                .set_cursor(&schema.name, &partition)
317                .map_err(|e| e.to_string())?;
318            CursorSlice {
319                cursor: schema.name.clone(),
320                start: partition.start_ord,
321                end: partition.end_ord,
322            }
323        } else {
324            let extent = cursor_extent_on(kernel, schema);
325            CursorSlice {
326                cursor: schema.name.clone(),
327                start: 0,
328                end: extent,
329            }
330        };
331        narrowest = Some(match narrowest {
332            Some(prev) if prev.len() <= slice.len() => prev,
333            _ => slice,
334        });
335    }
336    Ok(narrowest)
337}
338
339impl PolydatKernel {
340    /// A fresh kernel over a shared, already compiled program: the host
341    /// side of one program, many states.
342    pub fn over(program: Arc<PolydatProgram>) -> Self {
343        PolydatKernel::from_program(program)
344    }
345
346    /// Open the traversal at `index` among this program's top-level
347    /// `for` statements, evaluated against this kernel's current values.
348    ///
349    /// Comprehension sources that reference this kernel's wires see the
350    /// values currently set on it. Cascade externs are snapshotted from
351    /// this kernel now and bound into every activation.
352    pub fn traverse(&mut self, index: usize) -> Result<TraversalStream, String> {
353        let program = self.program().clone();
354        let traversal = program.traversals().get(index).cloned().ok_or_else(|| {
355            format!(
356                "no traversal at index {index}; the program declares {}",
357                program.traversals().len()
358            )
359        })?;
360        open_traversal(self, traversal)
361    }
362}
363
364/// Identity of the program an activation runs over, for callers that
365/// want to assert the one-program-per-position property.
366pub fn program_identity(kernel: &PolydatKernel) -> *const PolydatProgram {
367    Arc::as_ptr(kernel.program())
368}
369
370/// Open `traversal` against `parent`'s current values, on any engine
371/// (engines.md §3.6): the cascaded wires and the wires the
372/// sources reference are snapshotted through
373/// the [`Kernel`] trait, and the comprehension is evaluated in the
374/// body's scope, the body's program with those wires bound, where a
375/// source or predicate resolves every name it can reference and a
376/// tuple's own elements are layered in front as it is built. Nothing
377/// here needs the opening kernel beyond the snapshot.
378pub fn open_traversal(
379    parent: &mut dyn Kernel,
380    traversal: Traversal,
381) -> Result<TraversalStream, String> {
382    let mut cascade = Vec::with_capacity(traversal.cascade.len());
383    for (name, _) in &traversal.cascade {
384        let value = if parent.output_type(name).is_some() {
385            parent.pull(name)
386        } else {
387            parent.input_value(name).unwrap_or(Value::None)
388        };
389        cascade.push((name.clone(), value));
390    }
391    // What the comprehension's sources resolve against: the cascaded
392    // wires, over the body program's ledger, which is what a source
393    // that has to compile is charged to. Opening allocates no state
394    // over the body's program, because a source reads the cascade and
395    // the enclosing scope's wires and never the body's own constants.
396    let base = crate::kernel::interp::NoScope::charged_to(traversal.program.ledger().clone());
397    let cascaded = Layered {
398        prefix: &cascade,
399        inner: &base,
400    };
401    // The names the sources and predicates read and the comprehension
402    // does not bind are captured from the parent here, whatever their
403    // provenance (a coordinate input as much as an extern) and even
404    // where the body declares the same name, as every body declares
405    // `cycle`: a source or predicate belongs to the enclosing scope, and
406    // a traversal reads that frame once, when it opens
407    // (for_traversal.md §3.1). A name the parent has no wire for is one
408    // nothing binds, which the compile refused under `pragma strict` and
409    // otherwise warned about: it is not captured, and reads None
410    // (comprehension_forms.md §5 V3).
411    let mut captured: Vec<(String, Value)> = Vec::new();
412    let mut outer: Vec<String> =
413        crate::iteration::comprehension::validate::outer_reads(&traversal.comprehension)
414            .into_iter()
415            .filter(|r| !r.bare)
416            .map(|r| r.name)
417            .collect();
418    outer.sort();
419    outer.dedup();
420    for name in outer {
421        let value = if parent.output_type(&name).is_some() {
422            parent.pull(&name)
423        } else if let Some(value) = parent.input_value(&name) {
424            value
425        } else {
426            continue;
427        };
428        captured.push((name, value));
429    }
430    let scope = Layered {
431        prefix: &captured,
432        inner: &cascaded,
433    };
434    let tuples = evaluate_indexed(&traversal.comprehension, &scope).map_err(|e| {
435        format!(
436            "`for {}` at line {}, col {}: {e}",
437            traversal.source_text, traversal.span.line, traversal.span.col
438        )
439    })?;
440    Ok(TraversalStream {
441        traversal,
442        tuples,
443        cascade,
444        next: 0,
445        engine: parent.engine(),
446    })
447}