Skip to main content

sva_engine/render/
stream.rs

1// Concern: opens a target as a stream and edits it and its terms as it plays | Non-concern: pulling its blocks, what an edit carries on (edit.rs) | IO: (&Graph, target) -> Stream; (expr) -> Handle, bool
2
3use sva_ast::{Expr, Graph};
4
5use super::drive::{Block, Driver};
6use super::table::{Table, edit};
7use super::terms::{Handle, NOTES, Terms};
8use super::{Ends, Render, RenderConfig, prepared, range_of};
9use crate::cache::{Cache, CacheStats, Recording};
10use crate::error::{Diagnostic, EngineError, Located};
11use crate::flops::Work;
12use crate::schedule;
13
14pub const STREAMED: &str = "streamed";
15
16#[derive(Clone, Debug, PartialEq)]
17pub struct StreamConfig {
18    pub block: usize,
19    pub render: RenderConfig,
20}
21
22/// A target rendered block by block off one table. The target may read `@notes`, the sum of
23/// the terms added under handles; an edit to it or to a term plays from the next block on.
24pub struct Stream {
25    config: StreamConfig,
26    shell: Render,
27    driver: Driver,
28    graph: Graph,
29    expr: Expr,
30    terms: Terms,
31    live: bool,
32    dropped: Vec<String>,
33}
34
35impl Stream {
36    pub fn open(
37        graph: &Graph,
38        target: &Expr,
39        config: StreamConfig,
40        cache: Option<&Cache>,
41    ) -> Result<Stream, EngineError> {
42        if config.block == 0 {
43            return Err(refusal("a block of no samples".to_string()));
44        }
45        let mut terms = Terms::default();
46        let (shell, range) = shelled(graph, target, &mut terms, &config.render)?;
47        let table = Table::build(&shell.tys, shell.root, &[], &shell.config.profile)?;
48        let recording = Recording::over(cache, config.render.cache_policy);
49        let driver = Driver::new(table, range, config.block, &config.render, recording);
50        Ok(Stream {
51            graph: graph.clone(),
52            expr: target.clone(),
53            terms,
54            live: false,
55            dropped: Vec::new(),
56            driver,
57            config,
58            shell,
59        })
60    }
61
62    pub fn edit(&mut self, graph: &Graph, target: &Expr) -> Result<(), EngineError> {
63        self.rebuilt(graph, target.clone(), self.terms.clone())
64    }
65
66    pub fn add(&mut self, graph: &Graph, term: &Expr) -> Result<Handle, EngineError> {
67        let (terms, handle) = self.terms.added(term.clone());
68        self.rebuilt(graph, self.expr.clone(), terms)?;
69        Ok(handle)
70    }
71
72    /// False, and nothing edited, where the stream no longer holds `handle`.
73    pub fn replace(
74        &mut self,
75        graph: &Graph,
76        handle: Handle,
77        term: &Expr,
78    ) -> Result<bool, EngineError> {
79        let Some(terms) = self.terms.replaced(handle, term.clone()) else {
80            return Ok(false);
81        };
82        self.rebuilt(graph, self.expr.clone(), terms)?;
83        Ok(true)
84    }
85
86    /// A sounding term is cut where the stream stands, so what it played stays what it was.
87    pub fn remove(&mut self, handle: Handle) -> Result<bool, EngineError> {
88        let at = self.driver.at as f64 / f64::from(self.shell.rate());
89        let Some(terms) = self.terms.removed(handle, at) else {
90            return Ok(false);
91        };
92        let graph = self.graph.clone();
93        self.rebuilt(&graph, self.expr.clone(), terms)?;
94        Ok(true)
95    }
96
97    pub fn exprs(&self) -> impl Iterator<Item = &Expr> {
98        std::iter::once(&self.expr).chain(self.terms.exprs())
99    }
100
101    /// Values of the same identity carry on; a changed stateful one takes its predecessor's
102    /// state; the rest start now.
103    fn rebuilt(
104        &mut self,
105        graph: &Graph,
106        target: Expr,
107        mut terms: Terms,
108    ) -> Result<(), EngineError> {
109        let mut render = self.config.render.clone();
110        render.range.start = Some(self.driver.start);
111        let (shell, range) = shelled(graph, &target, &mut terms, &render)?;
112        let mut table = Table::build(&shell.tys, shell.root, &[], &shell.config.profile)?;
113        let old = std::mem::replace(&mut self.driver.table, Table::empty());
114        let dropped = edit::carried(&mut table, old, self.driver.at, self.live);
115        self.dropped
116            .extend(dropped.into_iter().map(|at| table.values[at].name.clone()));
117        self.driver.replace(table, range.end);
118        self.shell = shell;
119        self.graph = graph.clone();
120        self.expr = target;
121        self.terms = terms;
122        self.prune();
123        Ok(())
124    }
125
126    /// The next block, cut where the stream ends; `None` from there on.
127    pub fn next_block(&mut self) -> Result<Option<Block>, EngineError> {
128        let block = self.driver.next_block()?;
129        self.prune();
130        Ok(block)
131    }
132
133    /// An edited node with no state there starts silent, never computing its past.
134    pub fn go_live(&mut self) {
135        self.live = true;
136    }
137
138    /// Each node a live edit started silent.
139    pub fn dropped(&self) -> &[String] {
140        &self.dropped
141    }
142
143    /// Retires every term whose value no later block reads.
144    fn prune(&mut self) {
145        let (table, shell) = (&self.driver.table, &self.shell);
146        let future = (self.driver.at < self.driver.last())
147            .then(|| table.demand(sva_samples::Extent::new(self.driver.at, self.driver.last())));
148        let gone = |id| {
149            table.of(id).is_some_and(|at| {
150                !table.values[at].evaluated.is_empty()
151                    && future
152                        .as_ref()
153                        .is_none_or(|needs| needs[at].hold.is_empty())
154            })
155        };
156        self.terms
157            .prune(&gone, &|leaf| crate::refs::identity(&shell.tys, leaf).ok());
158    }
159
160    /// Every segment of its own clock the stream computed of `node`'s value, in order.
161    pub fn evaluated(&self, node: &str) -> Vec<sva_samples::Extent> {
162        let table = &self.driver.table;
163        self.shell
164            .tys
165            .id(node)
166            .and_then(|id| table.of(id))
167            .map_or(Vec::new(), |at| table.values[at].evaluated.clone())
168    }
169
170    pub fn position(&self) -> i64 {
171        self.driver.at
172    }
173
174    pub fn work(&self) -> Work {
175        self.driver.work
176    }
177
178    pub fn stats(&self) -> CacheStats {
179        self.driver.recording.stats()
180    }
181
182    /// The bytes its values hold, samples and state.
183    pub fn held_bytes(&self) -> usize {
184        self.driver.table.bytes()
185    }
186
187    pub fn end(&self) -> Option<i64> {
188        self.driver.end()
189    }
190
191    pub fn width(&self) -> usize {
192        self.driver.table.values[self.driver.table.root].width
193    }
194
195    pub fn config(&self) -> &StreamConfig {
196        &self.config
197    }
198}
199
200/// The graph with `streamed` defined as `target` and `notes` as `terms`, typed, scheduled
201/// and ranged. A composition's own `notes` stands while there is no term.
202fn shelled(
203    graph: &Graph,
204    target: &Expr,
205    terms: &mut Terms,
206    config: &RenderConfig,
207) -> Result<(Render, sva_samples::Extent), EngineError> {
208    let mut wrapped = graph.clone();
209    if !terms.is_empty() && graph.defines(NOTES) {
210        return Err(EngineError::refused(Diagnostic {
211            code: "engine.no_stream".to_string(),
212            message: format!(
213                "a term is added to `@{NOTES}`, and this composition defines its own `{NOTES}`"
214            ),
215            location: Located::at(NOTES, None),
216            help: format!("rename the composition's `{NOTES}`, or play it without adding terms"),
217        }));
218    }
219    let own = terms.is_empty() && graph.defines(NOTES);
220    let sum = (!own).then(|| (NOTES, terms.sum()));
221    for (name, body) in std::iter::once((STREAMED, target.clone())).chain(sum) {
222        if !wrapped.define(name, body) {
223            return Err(refusal(format!(
224                "this composition already has a node named `{name}`"
225            )));
226        }
227    }
228    let mut held = prepared(&wrapped, STREAMED, config.rate)?;
229    terms.typed(&held.instances, &mut held.tys);
230    let schedule = schedule::plan(&held.tys, held.root, &[]);
231    let shell = Render::shell(held.tys, held.root, config.clone(), schedule);
232    let range = range_of(&shell, Ends::Pulled)?;
233    Ok((shell, range))
234}
235
236fn refusal(what: String) -> EngineError {
237    EngineError::refused(Diagnostic {
238        code: "engine.no_stream".to_string(),
239        message: format!("this target opens no stream: {what}"),
240        location: Located::at(STREAMED, None),
241        help: "stream an expression over the nodes the composition defines".to_string(),
242    })
243}