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