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