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