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::cell::RefCell;
7use std::collections::BTreeMap;
8
9use sva_ast::Graph;
10use sva_samples::Tape;
11use sva_samples::physics::chaigne_askenfelt::landing_step;
12
13use super::silent::{Silent, bound_from, not_silent_by};
14use super::{Render, RenderConfig, prepared};
15use crate::error::{Diagnostic, EngineError, Located};
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 at the threshold, to be proven by `limit`, and the last bound found.
42    silent: Option<(Silent, usize)>,
43    bound: f64,
44    heard: RefCell<BTreeMap<(sva_formula::NodeId, u64), sva_samples::Buffer>>,
45    end: Option<usize>,
46}
47
48/// The samples one `next` wrote, `[start, start + len)` of the grid.
49pub struct Block<'a> {
50    tape: &'a Tape,
51    start: usize,
52}
53
54impl Block<'_> {
55    pub fn start(&self) -> usize {
56        self.start
57    }
58
59    pub fn len(&self) -> usize {
60        self.tape.end() - self.start
61    }
62
63    pub fn is_empty(&self) -> bool {
64        self.len() == 0
65    }
66
67    pub fn width(&self) -> usize {
68        self.tape.width()
69    }
70
71    pub fn plane(&self, c: usize) -> &[f64] {
72        self.tape.since(c, self.start)
73    }
74}
75
76impl Stream {
77    /// `bindings` are the named arguments `@target(t, name=value, ...)` is written with.
78    pub fn open(
79        graph: &Graph,
80        target: &str,
81        bindings: &[(String, f64)],
82        config: StreamConfig,
83    ) -> Result<Stream, EngineError> {
84        let mut stream = Stream::opened_at(graph, target, bindings, config, 0)?;
85        stream.settle()?;
86        Ok(stream)
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        };
122        let nodes = node::built(&shell, config.block)?;
123        let root = nodes
124            .iter()
125            .position(|n| n.id == shell.root)
126            .ok_or_else(|| EngineError::UnknownNode(shell.tys.name(shell.root).to_string()))?;
127        Ok(Stream {
128            graph: graph.clone(),
129            target: target.to_string(),
130            bindings: bindings.to_vec(),
131            config,
132            shell,
133            nodes,
134            root,
135            at: 0,
136            silent,
137            bound: f64::INFINITY,
138            heard: RefCell::new(BTreeMap::new()),
139            end: None,
140        })
141    }
142
143    /// Ends the stream here where every later sample is proven under the threshold from the
144    /// states it holds now.
145    fn settle(&mut self) -> Result<(), EngineError> {
146        let (Some((silent, _)), None) = (self.silent, self.end) else {
147            return Ok(());
148        };
149        let live = live::View::of(&self.nodes, self.at);
150        let (tys, root) = (&self.shell.tys, self.shell.root);
151        let config = &self.shell.config;
152        self.bound = bound_from(tys, root, config, silent, &live, self.at, &self.heard)?;
153        if self.bound < silent.threshold() {
154            self.end = Some(self.at.max(1));
155        }
156        Ok(())
157    }
158
159    /// No block past the limit is taken while silence is unproven.
160    fn in_time(&self) -> Result<(), EngineError> {
161        match (self.silent, self.end) {
162            (Some((silent, limit)), None) if self.at >= limit => {
163                let max_secs = limit as f64 / f64::from(self.config.rate);
164                let (tys, root) = (&self.shell.tys, self.shell.root);
165                Err(not_silent_by(
166                    tys,
167                    root,
168                    self.bound,
169                    Silent { max_secs, ..silent },
170                ))
171            }
172            _ => Ok(()),
173        }
174    }
175
176    /// The next block, cut where silence was proven; `None` from there on.
177    pub fn next_block(&mut self) -> Result<Option<Block<'_>>, EngineError> {
178        self.in_time()?;
179        let from = self.at;
180        let to = match self.end {
181            Some(end) if from >= end => return Ok(None),
182            Some(end) => end.min(from + self.config.block),
183            None => from + self.config.block,
184        };
185        for at in 0..self.nodes.len() {
186            let (done, rest) = self.nodes.split_at_mut(at);
187            rest[0].run(&self.shell, done, from, to)?;
188        }
189        self.at = to;
190        self.settle()?;
191        Ok(Some(Block {
192            tape: &self.nodes[self.root].tape,
193            start: from,
194        }))
195    }
196
197    pub fn position(&self) -> usize {
198        self.at
199    }
200
201    /// Where silence ends the stream, once a block has proven it.
202    pub fn end(&self) -> Option<usize> {
203        self.end
204    }
205
206    pub fn width(&self) -> usize {
207        self.nodes[self.root].width
208    }
209
210    pub fn config(&self) -> StreamConfig {
211        self.config
212    }
213
214    pub fn checkpoint(&self) -> Checkpoint {
215        Checkpoint {
216            at: self.at,
217            target: self.target.clone(),
218            bindings: self.bindings.clone(),
219            rate: self.config.rate,
220            block: self.config.block,
221            nodes: self.nodes.iter().map(Streamed::held).collect(),
222        }
223    }
224
225    /// This stream's target from `checkpoint` on, `bindings` in force, ending at `silent`
226    /// proven by `max_secs` past the checkpoint. A binding may move only where no sample
227    /// before the checkpoint can hear it: `release`, at or past it.
228    pub fn resume(
229        &self,
230        checkpoint: &Checkpoint,
231        bindings: &[(String, f64)],
232        silent: Option<Silent>,
233    ) -> Result<Stream, EngineError> {
234        let (rate, block) = (self.config.rate, self.config.block);
235        if checkpoint.target != self.target || (checkpoint.rate, checkpoint.block) != (rate, block)
236        {
237            return Err(mismatch(&self.target, "another stream"));
238        }
239        causal(checkpoint, bindings, rate, &self.target)?;
240        let config = StreamConfig {
241            rate,
242            block,
243            silent,
244        };
245        let mut resumed =
246            Stream::opened_at(&self.graph, &self.target, bindings, config, checkpoint.at)?;
247        if resumed.nodes.len() != checkpoint.nodes.len() {
248            return Err(mismatch(&self.target, "a graph of another shape"));
249        }
250        for (node, held) in resumed.nodes.iter_mut().zip(&checkpoint.nodes) {
251            node.resume(&resumed.shell, held, checkpoint.at)?;
252        }
253        resumed.at = checkpoint.at;
254        resumed.settle()?;
255        Ok(resumed)
256    }
257}
258
259/// Every node's state at a block's end, and what its stream was opened with.
260#[derive(Clone)]
261pub struct Checkpoint {
262    at: usize,
263    target: String,
264    bindings: Vec<(String, f64)>,
265    rate: u32,
266    block: usize,
267    nodes: Vec<Option<node::NodeState>>,
268}
269
270impl Checkpoint {
271    pub fn position(&self) -> usize {
272        self.at
273    }
274}
275
276/// A moved `release` changes nothing before the step a felt released then lands on.
277fn causal(
278    checkpoint: &Checkpoint,
279    bindings: &[(String, f64)],
280    rate: u32,
281    target: &str,
282) -> Result<(), EngineError> {
283    let value =
284        |set: &[(String, f64)], name: &str| set.iter().find(|(n, _)| n == name).map(|(_, v)| *v);
285    let names = checkpoint.bindings.iter().chain(bindings).map(|(n, _)| n);
286    for name in names {
287        let (was, now) = (value(&checkpoint.bindings, name), value(bindings, name));
288        if was == now {
289            continue;
290        }
291        let lands = |v: Option<f64>| {
292            landing_step(v.unwrap_or(f64::INFINITY), f64::from(rate))
293                .is_none_or(|at| at >= checkpoint.at as u64)
294        };
295        if name != RELEASE || !lands(was) || !lands(now) {
296            return Err(EngineError::refused(Diagnostic {
297                code: "engine.binding_not_causal".to_string(),
298                message: format!(
299                    "`{name}` moves from {} to {} at sample {}, and a sample before it can \
300                     hear that",
301                    shown(was),
302                    shown(now),
303                    checkpoint.at
304                ),
305                location: Located::at(target, None),
306                help: "move only release, to a key-up at or past the checkpoint".to_string(),
307            }));
308        }
309    }
310    Ok(())
311}
312
313fn shown(value: Option<f64>) -> String {
314    value.map_or("unbound".to_string(), |v| v.to_string())
315}
316
317fn mismatch(target: &str, what: &str) -> EngineError {
318    EngineError::refused(Diagnostic {
319        code: "engine.checkpoint_mismatch".to_string(),
320        message: format!("this checkpoint was taken of {what}, not of this `{target}` stream"),
321        location: Located::at(target, None),
322        help: "resume a checkpoint on the stream it was taken of".to_string(),
323    })
324}
325
326/// The graph with `streamed` defined as the target called with `bindings`.
327fn bound(graph: &Graph, target: &str, bindings: &[(String, f64)]) -> Result<Graph, EngineError> {
328    if !graph.defines(target) {
329        return Err(EngineError::UnknownNode(target.to_string()));
330    }
331    let mut named = String::new();
332    for (name, value) in bindings {
333        let word = name.starts_with(|c: char| c.is_ascii_alphabetic() || c == '_')
334            && name.chars().all(|c| c.is_ascii_alphanumeric() || c == '_');
335        if !word || !value.is_finite() {
336            return Err(refusal(target, format!("`{name}` bound to {value}")));
337        }
338        named.push_str(&format!(", {name}={value}"));
339    }
340    let call = sva_ast::parse_expr(&format!("@{target}(t{named})"))
341        .map_err(|d| refusal(target, d.message))?;
342    let mut wrapped = graph.clone();
343    if !wrapped.define(STREAMED, call) {
344        return Err(refusal(
345            target,
346            format!("this composition already has a node named `{STREAMED}`"),
347        ));
348    }
349    Ok(wrapped)
350}
351
352fn refusal(target: &str, what: String) -> EngineError {
353    EngineError::refused(Diagnostic {
354        code: "engine.no_stream".to_string(),
355        message: format!("`{target}` opens no stream: {what}"),
356        location: Located::at(target, None),
357        help: "name a node the composition defines, and bind each name to a finite number"
358            .to_string(),
359    })
360}