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; (Change) -> Changed
2
3use std::cell::RefCell;
4use std::collections::{BTreeMap, BTreeSet};
5use std::sync::Arc;
6
7use sva_ast::{Expr, Graph};
8use sva_formula::{Hash, NodeId};
9use sva_samples::{Buffer, Extent};
10
11use super::drive::{Block, Driver};
12use super::frontier::{Frontier, Known};
13use super::table::{Table, edit};
14use super::terms::{Handle, NOTES, Terms};
15use super::{Ends, Render, RenderConfig, range_of, through};
16use crate::cache::{Cache, CacheStats, Lookup, Outcome, Recording, Stored, Through};
17use crate::error::{Diagnostic, EngineError, Located};
18use crate::flops::Work;
19use crate::instantiate;
20use crate::schedule;
21use crate::typing;
22
23pub const STREAMED: &str = "streamed";
24
25#[derive(Clone, Debug, PartialEq)]
26pub struct StreamConfig {
27    pub block: usize,
28    pub render: RenderConfig,
29}
30
31/// A target rendered block by block off one table. The target may read `@notes`, the sum of
32/// the terms added under handles. A node the store answers plays from its samples. A hit is
33/// looked up once; a node reading the stream's own note sum never, as no store holds one.
34pub struct Stream {
35    config: StreamConfig,
36    shell: Render,
37    driver: Driver,
38    graph: Graph,
39    expr: Expr,
40    terms: Terms,
41    met: Met,
42    generation: u64,
43    live: bool,
44    dropped: Vec<String>,
45}
46
47/// Each hit met: a header holds until its store moves the samples; any holder may fill a miss.
48#[derive(Default)]
49struct Met {
50    epoch: u64,
51    known: Known,
52}
53
54impl Met {
55    fn over(&mut self, store: &impl Through) -> &Known {
56        if store.epoch() != self.epoch {
57            self.known.clear();
58            self.epoch = store.epoch();
59        }
60        &self.known
61    }
62
63    /// Kept where the store has not changed since `epoch`.
64    fn noted(&mut self, key: Hash, found: &Option<Arc<Stored>>, epoch: u64, store: &impl Through) {
65        self.over(store);
66        if epoch == self.epoch {
67            self.known.insert(key, found.clone());
68        }
69    }
70
71    fn hits(&mut self, store: &impl Through) -> Vec<(Hash, Option<Arc<Stored>>)> {
72        let hits = self.over(store).iter().filter(|(_, found)| found.is_some());
73        hits.map(|(key, found)| (*key, found.clone())).collect()
74    }
75}
76
77struct Shelled {
78    shell: Render,
79    range: Extent,
80    table: Table,
81    hits: Vec<Lookup>,
82}
83
84impl Stream {
85    pub async fn open(
86        graph: &Graph,
87        target: &Expr,
88        config: StreamConfig,
89        cache: Option<&Cache>,
90        store: &impl Through,
91    ) -> Result<Stream, EngineError> {
92        let mut terms = Terms::default();
93        let (mut known, epoch) = (Known::new(), store.epoch());
94        let render = &blocked(&config)?.render;
95        let walk = async |found: &mut Frontier<'_>| found.walked(&mut known, store).await;
96        let mut found = shelled(graph, target, &mut terms, render, walk).await?;
97        let mut met = Met::default();
98        for (key, hit) in known.into_iter().filter(|(_, found)| found.is_some()) {
99            met.noted(key, &hit, epoch, store);
100        }
101        let next = ahead(found.range.start, render.rate);
102        through::load(&mut found.table, store, next, &BTreeSet::new()).await;
103        let mut recording = Recording::over(cache, config.render.cache_policy);
104        recording.found(found.hits);
105        let driver = Driver::new(
106            found.table,
107            found.range,
108            config.block,
109            &config.render,
110            recording,
111        );
112        Ok(Stream {
113            graph: graph.clone(),
114            expr: target.clone(),
115            terms,
116            met,
117            generation: 0,
118            live: false,
119            dropped: Vec::new(),
120            driver,
121            config,
122            shell: found.shell,
123        })
124    }
125
126    pub fn exprs(&self) -> impl Iterator<Item = &Expr> {
127        std::iter::once(&self.expr).chain(self.terms.exprs())
128    }
129
130    fn prospect(&self, change: Change) -> Result<Prospect, Changed> {
131        let mut render = self.config.render.clone();
132        render.range.start = Some(self.driver.start);
133        let (graph, target, terms, answer) = match change {
134            Change::Target(graph, target) => (graph, target, self.terms.clone(), Changed::Edited),
135            Change::Add(graph, term) => {
136                let (terms, handle) = self.terms.added(term);
137                (graph, self.expr.clone(), terms, Changed::Added(handle))
138            }
139            Change::Replace(handle, graph, term) => {
140                let terms = self.terms.replaced(handle, term);
141                let terms = terms.ok_or(Changed::Held(false))?;
142                (graph, self.expr.clone(), terms, Changed::Held(true))
143            }
144            Change::Remove(handle) => {
145                let at = self.driver.at as f64 / f64::from(self.shell.rate());
146                let terms = self.terms.removed(handle, at).ok_or(Changed::Held(false))?;
147                (
148                    self.graph.clone(),
149                    self.expr.clone(),
150                    terms,
151                    Changed::Held(true),
152                )
153            }
154        };
155        Ok(Prospect {
156            graph,
157            target,
158            terms,
159            answer,
160            render,
161            generation: self.generation,
162        })
163    }
164
165    /// The stored samples `table` asks over the next second that it brings in and nothing
166    /// holds; `None` where the stream changed since `generation`.
167    fn wanting(
168        &self,
169        generation: u64,
170        table: &Table,
171        unread: &BTreeSet<Hash>,
172    ) -> Option<Vec<(Arc<Stored>, Extent)>> {
173        if generation != self.generation {
174            return None;
175        }
176        let sounding = self.driver.table.stored_keys();
177        let wants = table.wants(ahead(self.driver.at, self.config.render.rate));
178        let brought = wants.into_iter();
179        let brought =
180            brought.filter(|(s, _)| !sounding.contains(&s.key) && !unread.contains(&s.key));
181        Some(brought.collect())
182    }
183
184    /// Values of the same identity carry on; a changed stateful one takes its predecessor's
185    /// state; the rest start now.
186    fn apply(&mut self, prospect: Prospect, shelled: Shelled, fetched: &[(Hash, Vec<Buffer>)]) {
187        let Shelled {
188            shell,
189            range,
190            mut table,
191            hits,
192        } = shelled;
193        let old = std::mem::replace(&mut self.driver.table, Table::empty());
194        let dropped = edit::carried(&mut table, old, self.driver.at, self.live);
195        for (key, samples) in fetched {
196            table.took(*key, samples);
197        }
198        self.dropped
199            .extend(dropped.into_iter().map(|at| table.values[at].name.clone()));
200        self.driver.replace(table, range.end);
201        self.driver.recording.found(hits);
202        self.shell = shell;
203        self.graph = prospect.graph;
204        self.expr = prospect.target;
205        self.terms = prospect.terms;
206        self.generation += 1;
207        self.prune();
208    }
209
210    /// Loads the stored samples the next second reads; a block reading one not loaded has it
211    /// computed, or, live, started silent where not ready.
212    pub async fn fetch(&mut self, store: &impl Through) {
213        let next = ahead(self.driver.at, self.config.render.rate);
214        through::load(&mut self.driver.table, store, next, &BTreeSet::new()).await;
215    }
216
217    /// What `fetch` reads, for a caller reading it as the stream plays.
218    pub fn wanted(&self) -> Vec<(Arc<Stored>, Extent)> {
219        let next = ahead(self.driver.at, self.config.render.rate);
220        self.driver.table.wants(next)
221    }
222
223    pub fn took(&mut self, key: Hash, samples: &[Buffer]) {
224        self.driver.table.took(key, samples);
225    }
226
227    /// The next block, cut where the stream ends; `None` from there on.
228    pub fn next_block(&mut self) -> Result<Option<Block>, EngineError> {
229        let block = self.driver.next_block()?;
230        self.prune();
231        Ok(block)
232    }
233
234    /// An edited node with no state there starts silent, never computing its past.
235    pub fn go_live(&mut self) {
236        self.live = true;
237    }
238
239    /// Each node a live edit started silent.
240    pub fn dropped(&self) -> &[String] {
241        &self.dropped
242    }
243
244    /// Retires every term whose value no later block reads.
245    fn prune(&mut self) {
246        let (table, shell) = (&self.driver.table, &self.shell);
247        let future = (self.driver.at < self.driver.last())
248            .then(|| table.demand(sva_samples::Extent::new(self.driver.at, self.driver.last())));
249        let gone = |id| {
250            table.of(id).is_some_and(|at| {
251                !table.values[at].evaluated.is_empty()
252                    && future
253                        .as_ref()
254                        .is_none_or(|needs| needs[at].hold.is_empty())
255            })
256        };
257        let named = |leaf| crate::refs::identity(&shell.tys, leaf).ok();
258        if self.terms.prune(&gone, &named) {
259            self.generation += 1;
260        }
261    }
262
263    /// Every segment of its own clock the stream computed of `node`'s value, in order.
264    pub fn evaluated(&self, node: &str) -> Vec<sva_samples::Extent> {
265        let table = &self.driver.table;
266        self.shell
267            .tys
268            .id(node)
269            .and_then(|id| table.of(id))
270            .map_or(Vec::new(), |at| table.values[at].evaluated.clone())
271    }
272
273    pub fn position(&self) -> i64 {
274        self.driver.at
275    }
276
277    pub fn work(&self) -> Work {
278        self.driver.work
279    }
280
281    pub fn stats(&self) -> CacheStats {
282        self.driver.recording.stats()
283    }
284
285    /// The bytes its values hold, samples and state.
286    pub fn held_bytes(&self) -> usize {
287        self.driver.table.bytes()
288    }
289
290    pub fn end(&self) -> Option<i64> {
291        self.driver.end()
292    }
293
294    pub fn width(&self) -> usize {
295        self.driver.table.values[self.driver.table.root].width
296    }
297
298    pub fn config(&self) -> &StreamConfig {
299        &self.config
300    }
301}
302
303/// An edit, each with the graph that reaches it and all the stream plays.
304pub enum Change {
305    Target(Graph, Expr),
306    Add(Graph, Expr),
307    Replace(Handle, Graph, Expr),
308    /// A sounding term is cut where the stream stands as the edit is built; what played stays.
309    Remove(Handle),
310}
311
312/// A replace or remove answers whether the stream still held its handle.
313#[derive(Clone, Copy, Debug, PartialEq, Eq)]
314pub enum Changed {
315    Edited,
316    Added(Handle),
317    Held(bool),
318}
319
320struct Prospect {
321    graph: Graph,
322    target: Expr,
323    terms: Terms,
324    answer: Changed,
325    render: RenderConfig,
326    generation: u64,
327}
328
329/// `build`'s change, holding the stream only to apply it, between two blocks, once the store
330/// answered what it brings in: an exact stream pulled only once its edit is done plays it where
331/// it was issued. `build` runs again whenever the stream changed under it.
332pub async fn change<E: From<EngineError>>(
333    stream: &RefCell<Stream>,
334    mut build: impl FnMut(&Stream) -> Result<Change, E>,
335    store: &impl Through,
336) -> Result<Changed, E> {
337    let mut local = Met::default();
338    let mut fetched: Vec<(Hash, Vec<Buffer>)> = Vec::new();
339    let (mut unread, mut asked) = (BTreeSet::new(), Vec::<(Hash, Extent)>::new());
340    loop {
341        let change = build(&stream.borrow())?;
342        let mut prospect = match stream.borrow().prospect(change) {
343            Ok(prospect) => prospect,
344            Err(answer) => return Ok(answer),
345        };
346        for (key, hit) in stream.borrow_mut().met.hits(store) {
347            local.noted(key, &hit, store.epoch(), store);
348        }
349        let walk = async |found: &mut Frontier<'_>| {
350            while let Some(key) = found.walk(local.over(store)) {
351                let epoch = store.epoch();
352                let found = store.lookup(key).await.map(Arc::new);
353                if found.is_some() {
354                    stream.borrow_mut().met.noted(key, &found, epoch, store);
355                }
356                local.noted(key, &found, epoch, store);
357            }
358        };
359        let (graph, target) = (&prospect.graph, &prospect.target);
360        let config = &prospect.render;
361        let mut shelled = shelled(graph, target, &mut prospect.terms, config, walk).await?;
362        let generation = prospect.generation;
363        loop {
364            for (key, samples) in &fetched {
365                shelled.table.took(*key, samples);
366            }
367            let wants = stream.borrow().wanting(generation, &shelled.table, &unread);
368            let Some(wants) = wants else {
369                break;
370            };
371            if wants.is_empty() {
372                let answer = prospect.answer;
373                stream.borrow_mut().apply(prospect, shelled, &fetched);
374                return Ok(answer);
375            }
376            for (stored, over) in wants {
377                let again = asked
378                    .iter()
379                    .any(|(key, e)| *key == stored.key && !e.intersect(over).is_empty());
380                asked.push((stored.key, over));
381                let read = match again {
382                    false => store.read(&stored, over).await,
383                    true => None,
384                };
385                match read {
386                    Some(samples) => fetched.push((stored.key, samples)),
387                    None => {
388                        unread.insert(stored.key);
389                    }
390                }
391            }
392        }
393    }
394}
395
396fn ahead(at: i64, rate: u32) -> Extent {
397    Extent::new(at, at.saturating_add(i64::from(rate)))
398}
399
400fn blocked(config: &StreamConfig) -> Result<&StreamConfig, EngineError> {
401    match config.block {
402        0 => Err(refusal("a block of no samples".to_string())),
403        _ => Ok(config),
404    }
405}
406
407/// The graph with `streamed` defined as `target` and `notes` as `terms`. A composition's own
408/// `notes` stands while there is no term.
409fn wrapped(graph: &Graph, target: &Expr, terms: &Terms) -> Result<Graph, EngineError> {
410    let mut wrapped = graph.clone();
411    if !terms.is_empty() && graph.defines(NOTES) {
412        return Err(EngineError::refused(Diagnostic {
413            code: "engine.no_stream".to_string(),
414            message: format!(
415                "a term is added to `@{NOTES}`, and this composition defines its own `{NOTES}`"
416            ),
417            location: Located::at(NOTES, None),
418            help: format!("rename the composition's `{NOTES}`, or play it without adding terms"),
419        }));
420    }
421    let own = terms.is_empty() && graph.defines(NOTES);
422    let sum = (!own).then(|| (NOTES, terms.sum()));
423    for (name, body) in std::iter::once((STREAMED, target.clone())).chain(sum) {
424        if !wrapped.define(name, body) {
425            return Err(refusal(format!(
426                "this composition already has a node named `{name}`"
427            )));
428        }
429    }
430    Ok(wrapped)
431}
432
433/// The wrapped graph typed whole, each node the store answers standing on its samples.
434async fn shelled(
435    graph: &Graph,
436    target: &Expr,
437    terms: &mut Terms,
438    config: &RenderConfig,
439    walk: impl AsyncFnOnce(&mut Frontier<'_>),
440) -> Result<Shelled, EngineError> {
441    let wrapped = wrapped(graph, target, terms)?;
442    let instances = instantiate::instantiate(&wrapped, STREAMED, config.rate)?;
443    let root = instances.instance_of(STREAMED)?;
444    let order = schedule::schedule_from(&instances, std::slice::from_ref(&root))?;
445    let keys = through::keys(&wrapped, &instances, &order, config);
446    let mut found = Frontier::from((&instances, &order), &keys, &root, (config, true));
447    if !terms.is_empty() {
448        found.unstored(reading(&order, &instances.instance_of(NOTES)?));
449    }
450    walk(&mut found).await;
451    let mut tys = typing::infer_over(&instances, &order.within(&found.visited), &BTreeMap::new())?;
452    let id = tys
453        .id(&root)
454        .ok_or_else(|| EngineError::UnknownNode(root.clone()))?;
455    terms.typed(&instances, &mut tys);
456    let prefixes: BTreeMap<NodeId, Arc<Stored>> = std::mem::take(&mut found.stored)
457        .into_iter()
458        .map(|(path, stored)| (tys.id(&path).expect("a walked node is typed"), stored))
459        .collect();
460    let schedule = schedule::plan(&tys, id, &[]);
461    let shell = Render::shell(tys, id, config.clone(), schedule);
462    let range = range_of(&shell, Ends::Pulled)?;
463    let table = Table::prefixed(&shell.tys, shell.root, &shell.config.profile, &prefixes)?;
464    let hits = found.lookups.into_iter();
465    let hits = hits.filter(|l| l.outcome == Outcome::Hit).collect();
466    Ok(Shelled {
467        shell,
468        range,
469        table,
470        hits,
471    })
472}
473
474/// `path` and every node reading it, however far down.
475fn reading(order: &schedule::Order, path: &str) -> BTreeSet<String> {
476    let mut out = BTreeSet::from([path.to_string()]);
477    loop {
478        let more: Vec<String> = order
479            .groups
480            .iter()
481            .flatten()
482            .filter(|node| !out.contains(*node))
483            .filter(|node| order.deps(node).iter().any(|read| out.contains(read)))
484            .cloned()
485            .collect();
486        if more.is_empty() {
487            return out;
488        }
489        out.extend(more);
490    }
491}
492
493fn refusal(what: String) -> EngineError {
494    EngineError::refused(Diagnostic {
495        code: "engine.no_stream".to_string(),
496        message: format!("this target opens no stream: {what}"),
497        location: Located::at(STREAMED, None),
498        help: "stream an expression over the nodes the composition defines".to_string(),
499    })
500}