Skip to main content

sva_engine/render/
stream.rs

1// Concern: opens a target as a stream over a store, editing it and its terms live | Non-concern: pulling its blocks, what an edit carries on | IO: (Graph, target) -> Stream; (expr) -> Handle, bool
2
3use std::collections::BTreeMap;
4use std::sync::Arc;
5
6use sva_ast::{Expr, Graph};
7use sva_formula::NodeId;
8use sva_samples::Extent;
9
10use super::drive::{Block, Driver};
11use super::frontier::Frontier;
12use super::table::{Table, edit};
13use super::terms::{Handle, NOTES, Terms};
14use super::{Ends, Render, RenderConfig, range_of, through};
15use crate::cache::{Cache, CacheStats, Lookup, Outcome, Recording, Stored, Through};
16use crate::error::{Diagnostic, EngineError, Located};
17use crate::flops::Work;
18use crate::instantiate;
19use crate::schedule;
20use crate::typing;
21
22pub const STREAMED: &str = "streamed";
23
24#[derive(Clone, Debug, PartialEq)]
25pub struct StreamConfig {
26    pub block: usize,
27    pub render: RenderConfig,
28}
29
30/// A target rendered block by block off one table. The target may read `@notes`, the sum of
31/// the terms added under handles; an edit to it or to a term plays from the next block on.
32/// Opening and each edit look `store` up first: a node it answers plays from its samples, its
33/// own value computing what they miss.
34pub struct Stream {
35    config: StreamConfig,
36    shell: Render,
37    driver: Driver,
38    graph: Graph,
39    expr: Expr,
40    terms: Terms,
41    live: bool,
42    dropped: Vec<String>,
43}
44
45struct Shelled {
46    shell: Render,
47    range: Extent,
48    table: Table,
49    hits: Vec<Lookup>,
50}
51
52impl Stream {
53    pub async fn open(
54        graph: &Graph,
55        target: &Expr,
56        config: StreamConfig,
57        cache: Option<&Cache>,
58        store: &impl Through,
59    ) -> Result<Stream, EngineError> {
60        let mut terms = Terms::default();
61        let render = &blocked(&config)?.render;
62        let mut found = shelled(graph, target, &mut terms, render, store).await?;
63        let next = ahead(found.range.start, render.rate);
64        through::load(&mut found.table, store, next).await;
65        let mut recording = Recording::over(cache, config.render.cache_policy);
66        recording.found(found.hits);
67        let driver = Driver::new(
68            found.table,
69            found.range,
70            config.block,
71            &config.render,
72            recording,
73        );
74        Ok(Stream {
75            graph: graph.clone(),
76            expr: target.clone(),
77            terms,
78            live: false,
79            dropped: Vec::new(),
80            driver,
81            config,
82            shell: found.shell,
83        })
84    }
85
86    pub async fn edit(
87        &mut self,
88        graph: &Graph,
89        target: &Expr,
90        store: &impl Through,
91    ) -> Result<(), EngineError> {
92        let terms = self.terms.clone();
93        self.rebuilt(graph, target.clone(), terms, store).await
94    }
95
96    pub async fn add(
97        &mut self,
98        graph: &Graph,
99        term: &Expr,
100        store: &impl Through,
101    ) -> Result<Handle, EngineError> {
102        let (terms, handle) = self.terms.added(term.clone());
103        self.rebuilt(graph, self.expr.clone(), terms, store).await?;
104        Ok(handle)
105    }
106
107    /// False, and nothing edited, where the stream no longer holds `handle`.
108    pub async fn replace(
109        &mut self,
110        graph: &Graph,
111        (handle, term): (Handle, &Expr),
112        store: &impl Through,
113    ) -> Result<bool, EngineError> {
114        let Some(terms) = self.terms.replaced(handle, term.clone()) else {
115            return Ok(false);
116        };
117        self.rebuilt(graph, self.expr.clone(), terms, store).await?;
118        Ok(true)
119    }
120
121    /// A sounding term is cut where the stream stands, so what it played stays what it was.
122    pub async fn remove(
123        &mut self,
124        handle: Handle,
125        store: &impl Through,
126    ) -> Result<bool, EngineError> {
127        let at = self.driver.at as f64 / f64::from(self.shell.rate());
128        let Some(terms) = self.terms.removed(handle, at) else {
129            return Ok(false);
130        };
131        let graph = self.graph.clone();
132        self.rebuilt(&graph, self.expr.clone(), terms, store)
133            .await?;
134        Ok(true)
135    }
136
137    pub fn exprs(&self) -> impl Iterator<Item = &Expr> {
138        std::iter::once(&self.expr).chain(self.terms.exprs())
139    }
140
141    /// Values of the same identity carry on; a changed stateful one takes its predecessor's
142    /// state; the rest start now.
143    async fn rebuilt(
144        &mut self,
145        graph: &Graph,
146        target: Expr,
147        mut terms: Terms,
148        store: &impl Through,
149    ) -> Result<(), EngineError> {
150        let mut render = self.config.render.clone();
151        render.range.start = Some(self.driver.start);
152        let found = shelled(graph, &target, &mut terms, &render, store).await?;
153        let Shelled {
154            shell,
155            range,
156            mut table,
157            hits,
158        } = found;
159        let old = std::mem::replace(&mut self.driver.table, Table::empty());
160        let dropped = edit::carried(&mut table, old, self.driver.at, self.live);
161        through::load(&mut table, store, ahead(self.driver.at, render.rate)).await;
162        self.dropped
163            .extend(dropped.into_iter().map(|at| table.values[at].name.clone()));
164        self.driver.replace(table, range.end);
165        self.driver.recording.found(hits);
166        self.shell = shell;
167        self.graph = graph.clone();
168        self.expr = target;
169        self.terms = terms;
170        self.prune();
171        Ok(())
172    }
173
174    /// Loads the stored samples the next second reads; a block reading one not loaded has it
175    /// computed, or, live, started silent where not ready.
176    pub async fn fetch(&mut self, store: &impl Through) {
177        let next = ahead(self.driver.at, self.config.render.rate);
178        through::load(&mut self.driver.table, store, next).await;
179    }
180
181    /// What `fetch` reads, for a caller reading it as the stream plays.
182    pub fn wanted(&self) -> Vec<(Arc<Stored>, Extent)> {
183        let next = ahead(self.driver.at, self.config.render.rate);
184        self.driver.table.wants(next)
185    }
186
187    pub fn took(&mut self, key: sva_formula::Hash, samples: Vec<sva_samples::Buffer>) {
188        self.driver.table.took(key, samples);
189    }
190
191    /// The next block, cut where the stream ends; `None` from there on.
192    pub fn next_block(&mut self) -> Result<Option<Block>, EngineError> {
193        let block = self.driver.next_block()?;
194        self.prune();
195        Ok(block)
196    }
197
198    /// An edited node with no state there starts silent, never computing its past.
199    pub fn go_live(&mut self) {
200        self.live = true;
201    }
202
203    /// Each node a live edit started silent.
204    pub fn dropped(&self) -> &[String] {
205        &self.dropped
206    }
207
208    /// Retires every term whose value no later block reads.
209    fn prune(&mut self) {
210        let (table, shell) = (&self.driver.table, &self.shell);
211        let future = (self.driver.at < self.driver.last())
212            .then(|| table.demand(sva_samples::Extent::new(self.driver.at, self.driver.last())));
213        let gone = |id| {
214            table.of(id).is_some_and(|at| {
215                !table.values[at].evaluated.is_empty()
216                    && future
217                        .as_ref()
218                        .is_none_or(|needs| needs[at].hold.is_empty())
219            })
220        };
221        self.terms
222            .prune(&gone, &|leaf| crate::refs::identity(&shell.tys, leaf).ok());
223    }
224
225    /// Every segment of its own clock the stream computed of `node`'s value, in order.
226    pub fn evaluated(&self, node: &str) -> Vec<sva_samples::Extent> {
227        let table = &self.driver.table;
228        self.shell
229            .tys
230            .id(node)
231            .and_then(|id| table.of(id))
232            .map_or(Vec::new(), |at| table.values[at].evaluated.clone())
233    }
234
235    pub fn position(&self) -> i64 {
236        self.driver.at
237    }
238
239    pub fn work(&self) -> Work {
240        self.driver.work
241    }
242
243    pub fn stats(&self) -> CacheStats {
244        self.driver.recording.stats()
245    }
246
247    /// The bytes its values hold, samples and state.
248    pub fn held_bytes(&self) -> usize {
249        self.driver.table.bytes()
250    }
251
252    pub fn end(&self) -> Option<i64> {
253        self.driver.end()
254    }
255
256    pub fn width(&self) -> usize {
257        self.driver.table.values[self.driver.table.root].width
258    }
259
260    pub fn config(&self) -> &StreamConfig {
261        &self.config
262    }
263}
264
265fn ahead(at: i64, rate: u32) -> Extent {
266    Extent::new(at, at.saturating_add(i64::from(rate)))
267}
268
269fn blocked(config: &StreamConfig) -> Result<&StreamConfig, EngineError> {
270    match config.block {
271        0 => Err(refusal("a block of no samples".to_string())),
272        _ => Ok(config),
273    }
274}
275
276/// The graph with `streamed` defined as `target` and `notes` as `terms`. A composition's own
277/// `notes` stands while there is no term.
278fn wrapped(graph: &Graph, target: &Expr, terms: &Terms) -> Result<Graph, EngineError> {
279    let mut wrapped = graph.clone();
280    if !terms.is_empty() && graph.defines(NOTES) {
281        return Err(EngineError::refused(Diagnostic {
282            code: "engine.no_stream".to_string(),
283            message: format!(
284                "a term is added to `@{NOTES}`, and this composition defines its own `{NOTES}`"
285            ),
286            location: Located::at(NOTES, None),
287            help: format!("rename the composition's `{NOTES}`, or play it without adding terms"),
288        }));
289    }
290    let own = terms.is_empty() && graph.defines(NOTES);
291    let sum = (!own).then(|| (NOTES, terms.sum()));
292    for (name, body) in std::iter::once((STREAMED, target.clone())).chain(sum) {
293        if !wrapped.define(name, body) {
294            return Err(refusal(format!(
295                "this composition already has a node named `{name}`"
296            )));
297        }
298    }
299    Ok(wrapped)
300}
301
302/// The wrapped graph typed whole, each node `store` answers standing on its samples.
303async fn shelled(
304    graph: &Graph,
305    target: &Expr,
306    terms: &mut Terms,
307    config: &RenderConfig,
308    store: &impl Through,
309) -> Result<Shelled, EngineError> {
310    let wrapped = wrapped(graph, target, terms)?;
311    let instances = instantiate::instantiate(&wrapped, STREAMED, config.rate)?;
312    let root = instances.instance_of(STREAMED)?;
313    let order = schedule::schedule_from(&instances, std::slice::from_ref(&root))?;
314    let keys = through::keys(&wrapped, &instances, &order, config);
315    let mut found = Frontier::from((&instances, &order), &keys, &root, (config, true));
316    found.walk(store).await;
317    let mut tys = typing::infer_over(&instances, &order.within(&found.visited), &BTreeMap::new())?;
318    let id = tys
319        .id(&root)
320        .ok_or_else(|| EngineError::UnknownNode(root.clone()))?;
321    terms.typed(&instances, &mut tys);
322    let prefixes: BTreeMap<NodeId, Arc<Stored>> = std::mem::take(&mut found.stored)
323        .into_iter()
324        .map(|(path, stored)| (tys.id(&path).expect("a walked node is typed"), stored))
325        .collect();
326    let schedule = schedule::plan(&tys, id, &[]);
327    let shell = Render::shell(tys, id, config.clone(), schedule);
328    let range = range_of(&shell, Ends::Pulled)?;
329    let table = Table::prefixed(&shell.tys, shell.root, &shell.config.profile, &prefixes)?;
330    let hits = found.lookups.into_iter();
331    let hits = hits.filter(|l| l.outcome == Outcome::Hit).collect();
332    Ok(Shelled {
333        shell,
334        range,
335        table,
336        hits,
337    })
338}
339
340fn refusal(what: String) -> EngineError {
341    EngineError::refused(Diagnostic {
342        code: "engine.no_stream".to_string(),
343        message: format!("this target opens no stream: {what}"),
344        location: Located::at(STREAMED, None),
345        help: "stream an expression over the nodes the composition defines".to_string(),
346    })
347}