Skip to main content

sva_engine/render/stream/
mod.rs

1// Concern: renders a target block after block, each node carrying its state across blocks | Non-concern: one node's rows or machine, a whole-horizon render | IO: (&Graph, target, bindings) -> blocks
2
3mod live;
4mod node;
5
6use std::collections::BTreeMap;
7
8use sva_ast::Graph;
9use sva_samples::Tape;
10use sva_samples::physics::chaigne_askenfelt::landing_step;
11
12use super::silent::{Kept, Silent, bound_from, not_silent_by};
13use super::{Render, RenderConfig, prepared};
14use crate::error::{Diagnostic, EngineError, Located};
15use crate::flops::Work;
16use crate::instantiate::RELEASE;
17use crate::schedule;
18use node::Streamed;
19
20pub const STREAMED: &str = "streamed";
21
22#[derive(Clone, Copy, Debug, PartialEq)]
23pub struct StreamConfig {
24    pub rate: u32,
25    pub block: usize,
26    /// Ends the stream where the tail bound proves silence; `None` streams on forever.
27    pub silent: Option<Silent>,
28}
29
30/// A target rendered from the grid's first sample on, one block at a time: every block is the
31/// samples a whole render over the same rows writes there, bit for bit.
32pub struct Stream {
33    graph: Graph,
34    target: String,
35    bindings: Vec<(String, f64)>,
36    config: StreamConfig,
37    shell: Render,
38    nodes: Vec<Streamed>,
39    root: usize,
40    at: usize,
41    /// Silence to prove by `limit`, and the last bound a proof found.
42    silent: Option<(Silent, usize)>,
43    bound: Option<f64>,
44    kept: Kept,
45    end: Option<usize>,
46    work: Work,
47}
48
49/// The samples one `next` wrote, `[start, start + len)` of the grid.
50pub struct Block<'a> {
51    tape: &'a Tape,
52    start: usize,
53}
54
55impl Block<'_> {
56    pub fn start(&self) -> usize {
57        self.start
58    }
59
60    pub fn len(&self) -> usize {
61        self.tape.end() - self.start
62    }
63
64    pub fn is_empty(&self) -> bool {
65        self.len() == 0
66    }
67
68    pub fn width(&self) -> usize {
69        self.tape.width()
70    }
71
72    pub fn plane(&self, c: usize) -> &[f64] {
73        self.tape.since(c, self.start)
74    }
75}
76
77impl Stream {
78    /// `bindings` are the named arguments `@target(t, name=value, ...)` is written with.
79    /// Proves no silence until a block has run.
80    pub fn open(
81        graph: &Graph,
82        target: &str,
83        bindings: &[(String, f64)],
84        config: StreamConfig,
85    ) -> Result<Stream, EngineError> {
86        Stream::opened_at(graph, target, bindings, config, 0)
87    }
88
89    /// Silence is proven by `max_secs` past `from`.
90    fn opened_at(
91        graph: &Graph,
92        target: &str,
93        bindings: &[(String, f64)],
94        config: StreamConfig,
95        from: usize,
96    ) -> Result<Stream, EngineError> {
97        if config.block == 0 {
98            return Err(refusal(target, "a block of no samples".to_string()));
99        }
100        let wrapped = bound(graph, target, bindings)?;
101        let held = prepared(&wrapped, STREAMED)?;
102        let rate = f64::from(config.rate);
103        let silent = config
104            .silent
105            .map(|silent| (silent, from + (silent.max_secs * rate).ceil() as usize));
106        // No row a stream takes reads a horizon.
107        let render_config = RenderConfig::seconds(config.rate, 1.0);
108        let schedule = schedule::plan(&held.tys, &held.order, held.root, &[]);
109        let shell = Render {
110            root: held.root,
111            tys: held.tys,
112            buffers: BTreeMap::new(),
113            frames: BTreeMap::new(),
114            symbolic: BTreeMap::new(),
115            labels: BTreeMap::new(),
116            traces: Vec::new(),
117            config: render_config,
118            schedule,
119            bindings: BTreeMap::new(),
120            cache_stats: None,
121            proofs: 0,
122        };
123        let nodes = node::built(&shell, config.block)?;
124        let kept = Kept::new(shell.tys.clone(), shell.config.clone());
125        let root = nodes
126            .iter()
127            .position(|n| n.id == shell.root)
128            .ok_or_else(|| EngineError::UnknownNode(shell.tys.name(shell.root).to_string()))?;
129        Ok(Stream {
130            graph: graph.clone(),
131            target: target.to_string(),
132            bindings: bindings.to_vec(),
133            config,
134            shell,
135            nodes,
136            root,
137            at: 0,
138            silent,
139            bound: None,
140            kept,
141            end: None,
142            work: Work {
143                waves: Some(0),
144                ..Work::default()
145            },
146        })
147    }
148
149    /// Ends the stream here where its states now prove every later sample under the threshold.
150    fn settle(&mut self) -> Result<(), EngineError> {
151        let (Some((silent, _)), None) = (self.silent, self.end) else {
152            return Ok(());
153        };
154        let live = live::View::of(&self.nodes, self.at);
155        self.work.proofs += 1;
156        let bound = bound_from(&self.kept, self.shell.root, silent, &live, self.at)?;
157        self.bound = Some(bound);
158        if bound < silent.threshold() {
159            self.end = Some(self.at.max(1));
160        }
161        Ok(())
162    }
163
164    /// No block past the limit is taken while silence is unproven.
165    fn in_time(&self) -> Result<(), EngineError> {
166        match (self.silent, self.end) {
167            (Some((silent, limit)), None) if self.at >= limit => {
168                let max_secs = limit as f64 / f64::from(self.config.rate);
169                let (tys, root) = (&self.shell.tys, self.shell.root);
170                Err(not_silent_by(
171                    tys,
172                    root,
173                    self.bound,
174                    Silent { max_secs, ..silent },
175                ))
176            }
177            _ => Ok(()),
178        }
179    }
180
181    /// The next block, cut where silence was proven; `None` from there on.
182    pub fn next_block(&mut self) -> Result<Option<Block<'_>>, EngineError> {
183        self.in_time()?;
184        let from = self.at;
185        let to = match self.end {
186            Some(end) if from >= end => return Ok(None),
187            Some(end) => end.min(from + self.config.block),
188            None => from + self.config.block,
189        };
190        for at in 0..self.nodes.len() {
191            let (done, rest) = self.nodes.split_at_mut(at);
192            rest[0].run(&self.shell, done, from, to)?;
193            let (priced, waves) = rest[0].work(from, to);
194            self.work.priced_flops += priced;
195            self.work.waves = self.work.waves.zip(waves).map(|(held, more)| held + more);
196        }
197        self.at = to;
198        self.work.samples += (to - from) as u64;
199        self.settle()?;
200        Ok(Some(Block {
201            tape: &self.nodes[self.root].tape,
202            start: from,
203        }))
204    }
205
206    pub fn position(&self) -> usize {
207        self.at
208    }
209
210    /// Since this stream opened, or since the checkpoint it resumed from.
211    pub fn work(&self) -> Work {
212        self.work
213    }
214
215    /// Where silence ends the stream, once proven.
216    pub fn end(&self) -> Option<usize> {
217        self.end
218    }
219
220    pub fn width(&self) -> usize {
221        self.nodes[self.root].width
222    }
223
224    pub fn config(&self) -> StreamConfig {
225        self.config
226    }
227
228    pub fn checkpoint(&self) -> Checkpoint {
229        Checkpoint {
230            at: self.at,
231            target: self.target.clone(),
232            bindings: self.bindings.clone(),
233            rate: self.config.rate,
234            block: self.config.block,
235            nodes: self.nodes.iter().map(Streamed::held).collect(),
236        }
237    }
238
239    /// This stream's target from `checkpoint` on, `bindings` in force, ending at `silent`
240    /// proven by `max_secs` past the checkpoint. A binding may move only where no sample
241    /// before the checkpoint can hear it: `release`, at or past it.
242    pub fn resume(
243        &self,
244        checkpoint: &Checkpoint,
245        bindings: &[(String, f64)],
246        silent: Option<Silent>,
247    ) -> Result<Stream, EngineError> {
248        let (rate, block) = (self.config.rate, self.config.block);
249        if checkpoint.target != self.target || (checkpoint.rate, checkpoint.block) != (rate, block)
250        {
251            return Err(mismatch(&self.target, "another stream"));
252        }
253        causal(checkpoint, bindings, rate, &self.target)?;
254        let config = StreamConfig {
255            rate,
256            block,
257            silent,
258        };
259        let mut resumed =
260            Stream::opened_at(&self.graph, &self.target, bindings, config, checkpoint.at)?;
261        if resumed.nodes.len() != checkpoint.nodes.len() {
262            return Err(mismatch(&self.target, "a graph of another shape"));
263        }
264        for (node, held) in resumed.nodes.iter_mut().zip(&checkpoint.nodes) {
265            node.resume(&resumed.shell, held, checkpoint.at)?;
266        }
267        resumed.at = checkpoint.at;
268        resumed.settle()?;
269        Ok(resumed)
270    }
271}
272
273/// Every node's state at a block's end, and what its stream was opened with.
274#[derive(Clone)]
275pub struct Checkpoint {
276    at: usize,
277    target: String,
278    bindings: Vec<(String, f64)>,
279    rate: u32,
280    block: usize,
281    nodes: Vec<Option<node::NodeState>>,
282}
283
284impl Checkpoint {
285    pub fn position(&self) -> usize {
286        self.at
287    }
288}
289
290/// A moved `release` changes nothing before the step a felt released then lands on.
291fn causal(
292    checkpoint: &Checkpoint,
293    bindings: &[(String, f64)],
294    rate: u32,
295    target: &str,
296) -> Result<(), EngineError> {
297    let value =
298        |set: &[(String, f64)], name: &str| set.iter().find(|(n, _)| n == name).map(|(_, v)| *v);
299    let names = checkpoint.bindings.iter().chain(bindings).map(|(n, _)| n);
300    for name in names {
301        let (was, now) = (value(&checkpoint.bindings, name), value(bindings, name));
302        if was == now {
303            continue;
304        }
305        let lands = |v: Option<f64>| {
306            landing_step(v.unwrap_or(f64::INFINITY), f64::from(rate))
307                .is_none_or(|at| at >= checkpoint.at as u64)
308        };
309        if name != RELEASE || !lands(was) || !lands(now) {
310            return Err(EngineError::refused(Diagnostic {
311                code: "engine.binding_not_causal".to_string(),
312                message: format!(
313                    "`{name}` moves from {} to {} at sample {}, and a sample before it can \
314                     hear that",
315                    shown(was),
316                    shown(now),
317                    checkpoint.at
318                ),
319                location: Located::at(target, None),
320                help: "move only release, to a key-up at or past the checkpoint".to_string(),
321            }));
322        }
323    }
324    Ok(())
325}
326
327fn shown(value: Option<f64>) -> String {
328    value.map_or("unbound".to_string(), |v| v.to_string())
329}
330
331fn mismatch(target: &str, what: &str) -> EngineError {
332    EngineError::refused(Diagnostic {
333        code: "engine.checkpoint_mismatch".to_string(),
334        message: format!("this checkpoint was taken of {what}, not of this `{target}` stream"),
335        location: Located::at(target, None),
336        help: "resume a checkpoint on the stream it was taken of".to_string(),
337    })
338}
339
340/// The graph with `streamed` defined as the target called with `bindings`.
341fn bound(graph: &Graph, target: &str, bindings: &[(String, f64)]) -> Result<Graph, EngineError> {
342    if !graph.defines(target) {
343        return Err(EngineError::UnknownNode(target.to_string()));
344    }
345    let mut named = String::new();
346    for (name, value) in bindings {
347        let word = name.starts_with(|c: char| c.is_ascii_alphabetic() || c == '_')
348            && name.chars().all(|c| c.is_ascii_alphanumeric() || c == '_');
349        if !word || !value.is_finite() {
350            return Err(refusal(target, format!("`{name}` bound to {value}")));
351        }
352        named.push_str(&format!(", {name}={value}"));
353    }
354    let call = sva_ast::parse_expr(&format!("@{target}(t{named})"))
355        .map_err(|d| refusal(target, d.message))?;
356    let mut wrapped = graph.clone();
357    if !wrapped.define(STREAMED, call) {
358        return Err(refusal(
359            target,
360            format!("this composition already has a node named `{STREAMED}`"),
361        ));
362    }
363    Ok(wrapped)
364}
365
366fn refusal(target: &str, what: String) -> EngineError {
367    EngineError::refused(Diagnostic {
368        code: "engine.no_stream".to_string(),
369        message: format!("`{target}` opens no stream: {what}"),
370        location: Located::at(target, None),
371        help: "name a node the composition defines, and bind each name to a finite number"
372            .to_string(),
373    })
374}