Skip to main content

sva_engine/render/
stream.rs

1// Concern: opens a target as a stream over memory, 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::Driver;
12use super::drive::Work;
13use super::end::{Ending, Fading, Heard, under};
14use super::terms::{Handle, NOTES, Terms, cut, placed};
15use super::value_graph::support::Supports;
16use super::value_graph::{Past, ValueGraph};
17use super::world::{Plan, Root, STREAMED, Walked, Wanted, World};
18use super::{Ends, RenderConfig, range_over};
19use crate::cache::{Backend, CacheStats, Counters, Memory, Recording, Stored, Tier};
20use crate::error::{Diagnostic, EngineError, Located};
21use crate::recent::Recent;
22
23#[cfg(test)]
24mod rebuilt;
25
26/// The lookups, and the names started silent, a stream keeps of all it made.
27pub const LATEST: usize = 256;
28
29#[derive(Clone, Debug, PartialEq)]
30pub struct StreamConfig {
31    pub block: usize,
32    /// The target's own at open where `None`; a mono one plays in each, a wider is refused.
33    pub channels: Option<usize>,
34    pub render: RenderConfig,
35}
36
37/// A target rendered block by block off one value graph, reading `@notes`, the sum of the terms
38/// added. A change builds only what it changed or newly reads; the rest stays as it stood.
39pub struct Stream {
40    config: StreamConfig,
41    world: World,
42    driver: Driver,
43    expr: Expr,
44    terms: Terms,
45    width: usize,
46    /// Each sounding term as heard at the root, through `gain` from the note sum.
47    heard: BTreeMap<Handle, Heard>,
48    gain: Gained,
49    /// Each term retired before its support ended.
50    fading: Vec<Fading>,
51    faded: f64,
52    /// The first sounding term's end, which `heard` evicts.
53    ending: Option<Option<i64>>,
54    /// The sample the stream's root is treated as silent from.
55    treated_as_silent_from_sample: Option<i64>,
56    generation: u64,
57    live: bool,
58    dropped: Recent<String>,
59    late: usize,
60    built: Built,
61    demands: usize,
62}
63
64#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
65pub struct Counts {
66    pub dropped: usize,
67    /// Edits that landed past where issued.
68    pub late: usize,
69    pub terms: usize,
70    /// The latest change's, or the open's.
71    pub built: Built,
72    /// The reads that worked out what the root asks of the note sum.
73    pub demands: usize,
74    pub tier: Counters,
75}
76
77/// What one change did anew, not what it carried over.
78#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
79pub struct Built {
80    /// Texts it took in: its expression and each node it newly reads.
81    pub parsed: usize,
82    pub instances: usize,
83    pub visited: usize,
84    pub typed: usize,
85    pub values: usize,
86    /// Values it made taking an old one's state.
87    pub copied: usize,
88    pub lookups: usize,
89}
90
91impl Stream {
92    pub async fn open<B: Backend>(
93        graph: &Graph,
94        target: &Expr,
95        config: StreamConfig,
96        tier: &Tier<B>,
97    ) -> Result<Stream, EngineError> {
98        let render = blocked(&config)?.render.clone();
99        if graph.defines(STREAMED) {
100            return Err(refusal(format!(
101                "this composition already has a node named `{STREAMED}`"
102            )));
103        }
104        let world = World::over(graph, render.rate, !graph.defines(NOTES));
105        let memory = tier.memory().clone();
106        let recording = Recording::over(&memory).latest(LATEST);
107        let value_graph = ValueGraph::new(&render.profile);
108        let memo = (memory, recording);
109        let driver = Driver::new(value_graph, Extent::new(0, 0), config.block, &render, memo);
110        let mut stream = Stream {
111            world,
112            driver,
113            expr: target.clone(),
114            terms: Terms::default(),
115            width: 0,
116            heard: BTreeMap::new(),
117            gain: Gained::default(),
118            fading: Vec::new(),
119            faded: 0.0,
120            ending: None,
121            treated_as_silent_from_sample: None,
122            generation: 0,
123            live: false,
124            dropped: Recent::keeping(LATEST),
125            late: 0,
126            built: Built::default(),
127            demands: 0,
128            config,
129        };
130        let opening = Prospect {
131            target: target.clone(),
132            terms: Terms::default(),
133            term: None,
134            from: None,
135            answer: Changed::Edited,
136            landing: None,
137            parsed: 1,
138        };
139        let mut local = Local::default();
140        let round = tier.begin();
141        loop {
142            match stream.attempt(&opening, &mut local, (tier.memory(), round))? {
143                Attempt::Landed(_) => break,
144                Attempt::Asks(keys) => local.look(keys, tier, round).await,
145                Attempt::Reads(wants) => local.read(wants, tier).await,
146                Attempt::Moved => unreachable!("nothing plays a stream before it opens"),
147            }
148        }
149        let mut needs = stream.needs();
150        while !needs.is_empty() {
151            let fetched = tier.fetch(&needs).await;
152            for (key, parts) in fetched.handed {
153                stream.driver.value_graph.took(key, &parts);
154            }
155            needs = fetched.left;
156        }
157        Ok(stream)
158    }
159
160    pub fn graph(&self) -> &Graph {
161        &self.world.graph
162    }
163
164    fn prospect(&self, change: Change) -> Result<Prospect, Changed> {
165        let mut landing = None;
166        let prospect = |terms, term: Option<(Handle, Expr)>, from: Option<Graph>, answer| {
167            let roots = |(handle, term): &(Handle, Expr)| sva_ast::reads_of(&handle.node(), term);
168            Prospect {
169                target: self.expr.clone(),
170                terms,
171                from: from.map(|graph| (graph, term.as_ref().map(roots).unwrap_or_default())),
172                term,
173                answer,
174                landing: None,
175                parsed: 1,
176            }
177        };
178        Ok(match change {
179            Change::Target(graph, target) => {
180                let roots = sva_ast::reads_of(STREAMED, &target);
181                Prospect {
182                    target,
183                    from: Some((graph, roots)),
184                    ..prospect(self.terms.clone(), None, None, Changed::Edited)
185                }
186            }
187            Change::Add(graph, term, at) => {
188                let term = match at {
189                    Placed::Written => term,
190                    Placed::Landing => placed(&term, *landing.insert(self.driver.at)),
191                };
192                let (terms, handle) = self.terms.added();
193                Prospect {
194                    landing,
195                    ..prospect(
196                        terms,
197                        Some((handle, term)),
198                        Some(graph),
199                        Changed::Added(handle),
200                    )
201                }
202            }
203            Change::Replace(handle, graph, term, at) => {
204                let landed = self.terms.landed(handle).ok_or(Changed::Held(false))?;
205                let term = match at {
206                    Placed::Written => term,
207                    Placed::Landing => placed(&term, landed),
208                };
209                let terms = self.terms.clone();
210                prospect(
211                    terms,
212                    Some((handle, term)),
213                    Some(graph),
214                    Changed::Held(true),
215                )
216            }
217            Change::Remove(handle) => {
218                let at = self.driver.at;
219                let terms = self.terms.removed(handle).ok_or(Changed::Held(false))?;
220                let held = self.world.graph.expr(&handle.node());
221                let term = cut(held.expect("a held term's node"), at);
222                Prospect {
223                    parsed: 0,
224                    ..prospect(terms, Some((handle, term)), None, Changed::Held(true))
225                }
226            }
227        })
228    }
229
230    /// One try at `prospect`: planned, typed and built beside what plays, then held, or undone
231    /// where memory must answer first.
232    fn attempt(
233        &mut self,
234        prospect: &Prospect,
235        local: &mut Local,
236        (memory, round): (&Memory, u64),
237    ) -> Result<Attempt, EngineError> {
238        if prospect.landing.is_some_and(|at| at != self.driver.at) {
239            return Ok(Attempt::Moved);
240        }
241        let wanted = Wanted {
242            root: Root::Streamed(&prospect.target),
243            terms: &prospect.terms,
244            term: prospect.term.as_ref().map(|(h, e)| (*h, e)),
245            from: prospect.from.as_ref().map(|(g, roots)| (g, roots.clone())),
246            whole: false,
247        };
248        let seen = &mut self.driver.recording;
249        let mut found = |node: &str, key: Hash| memory.answered((key, round), (node, &mut *seen));
250        let walked = self.world.plan(&wanted, &self.config.render, &mut found)?;
251        let mut plan = match walked {
252            Walked::Asks(keys) => return Ok(Attempt::Asks(keys)),
253            Walked::Planned(plan) => plan,
254        };
255        let built = self.built(&mut plan);
256        let (root, range, treated_as_silent_from_sample) = match built {
257            Ok(held) => held,
258            Err(e) => {
259                let freed = self.world.abort();
260                self.driver.value_graph.abort(&freed);
261                return Err(e);
262            }
263        };
264        let from = self.driver.at.max(range.start);
265        let future = Extent::new(from, self.last(range.end).max(from));
266        let value_graph = &mut self.driver.value_graph;
267        let short = value_graph.short((root, future, Past::Stored));
268        let opened = short
269            .into_iter()
270            .try_for_each(|at| value_graph.read_on(&self.world.typing, at));
271        if let Err(e) = opened {
272            let freed = self.world.abort();
273            self.driver.value_graph.abort(&freed);
274            return Err(e);
275        }
276        let asked_range = ahead(from, self.config.render.rate);
277        let wants = self.driver.value_graph.needs_made(root, asked_range);
278        let wants: Vec<(Hash, Extent)> = wants
279            .into_iter()
280            .filter(|(key, over)| !local.holds(*key, *over))
281            .collect();
282        if !wants.is_empty() {
283            let freed = self.world.abort();
284            self.driver.value_graph.abort(&freed);
285            return Ok(Attempt::Reads(wants));
286        }
287        self.treated_as_silent_from_sample = treated_as_silent_from_sample;
288        Ok(Attempt::Landed(self.land(
289            prospect,
290            (plan, root, range),
291            local,
292        )))
293    }
294
295    /// The plan typed and built: the root's value, its range, and the sample it is treated as
296    /// silent from.
297    fn built(&mut self, plan: &mut Plan) -> Result<(usize, Extent, Option<i64>), EngineError> {
298        let typing = &mut self.world.typing;
299        let id = typing
300            .id(STREAMED)
301            .ok_or_else(|| EngineError::UnknownNode(STREAMED.to_string()))?;
302        let hits: BTreeMap<NodeId, Arc<Stored>> = plan
303            .stored
304            .iter()
305            .filter_map(|(path, stored)| Some((typing.id(path)?, Arc::clone(stored))))
306            .collect();
307        let value_graph = &mut self.driver.value_graph;
308        let root = value_graph.grow(typing, id, &hits)?;
309        let plays = value_graph.values[root].width;
310        let mut render = self.config.render.clone();
311        match self.width {
312            0 if self.config.channels == Some(0) => {
313                return Err(refusal("a stream of no channels".to_string()));
314            }
315            0 => widens(plays, self.config.channels.unwrap_or(plays))?,
316            width => {
317                widens(plays, width)?;
318                render.range.start = Some(self.driver.start);
319            }
320        }
321        let supports = Supports::over(typing, Some(&value_graph.supports));
322        let ending = Ending::new(typing, &render.profile, &supports);
323        let end = match (render.range.end, self.live) {
324            (None, false) => ending.of(id),
325            _ => ending.exact(id),
326        };
327        let range = range_over((&render, STREAMED), end.support, Ends::Pulled)?;
328        Ok((root, range, end.treated_as_silent_from_sample))
329    }
330
331    /// The change held, each value it made carrying on what it continues.
332    fn land(
333        &mut self,
334        prospect: &Prospect,
335        (mut plan, root, range): (Plan, usize, Extent),
336        local: &mut Local,
337    ) -> Changed {
338        let freed = self.world.commit(std::mem::take(&mut plan.found));
339        let now = self.driver.at;
340        let carried = self
341            .driver
342            .value_graph
343            .settled(root, &freed, (now, self.live));
344        self.driver.value_graph.offers(&self.world.typing, range);
345        for (key, parts) in &local.fetched {
346            self.driver.value_graph.took(*key, parts);
347        }
348        for at in &carried.silent {
349            self.dropped
350                .push(self.driver.value_graph.values[*at].name.clone());
351        }
352        match self.width {
353            0 => {
354                let plays = self.driver.value_graph.values[root].width;
355                self.width = self.config.channels.unwrap_or(plays);
356                (self.driver.start, self.driver.at) = (range.start, range.start);
357            }
358            _ => self.late += usize::from(self.driver.at > local.issued),
359        }
360        let last = self.last(range.end);
361        self.driver.bound(last);
362        self.expr = prospect.target.clone();
363        self.terms = prospect.terms.clone();
364        if let Changed::Added(handle) = prospect.answer {
365            self.terms.land(handle, self.driver.at);
366        }
367        self.hear();
368        self.built = Built {
369            parsed: prospect.parsed + plan.adopted,
370            instances: plan.named,
371            visited: plan.visited,
372            typed: self.world.typing.lowered().len(),
373            values: self.driver.value_graph.built,
374            copied: carried.taken,
375            lookups: local.lookups,
376        };
377        self.generation += 1;
378        self.retire_terms_below_silence_threshold();
379        prospect.answer
380    }
381
382    /// A term is its own sound, judged at the root through the note sum's gain to it, bounded
383    /// again only where the key of what it is a function of moved.
384    fn hear(&mut self) {
385        let typing = &self.world.typing;
386        let supports = Supports::over(typing, Some(&self.driver.value_graph.supports));
387        let ending = Ending::new(typing, &self.config.render.profile, &supports);
388        let ends = typing.id(STREAMED).zip(typing.id(NOTES));
389        let key = ends.and_then(|(root, notes)| ending.gain_key(root, notes));
390        let gain = match ends {
391            _ if key.is_some() && key == self.gain.of => self.gain.gain,
392            Some((root, notes)) => ending.gain(root, notes),
393            None => None,
394        };
395        let moved = gain != std::mem::replace(&mut self.gain, Gained { gain, of: key }).gain;
396        let lowered: BTreeSet<&str> = typing.lowered().iter().map(String::as_str).collect();
397        for handle in self.terms.handles() {
398            let node = handle.node();
399            if !moved && !lowered.contains(node.as_str()) && self.heard.contains_key(&handle) {
400                continue;
401            }
402            let id = typing.id(&node).expect("a sounding term is typed");
403            self.heard.insert(handle, ending.heard(id, gain));
404        }
405        self.ending = None;
406    }
407
408    fn needs(&self) -> Vec<(Hash, Extent)> {
409        let next = ahead(self.driver.at, self.config.render.rate);
410        self.driver.value_graph.needs(next)
411    }
412
413    /// `n` samples from `at`, `None` past the end. An `at` behind is refused; one ahead skips
414    /// there, computing through the span, or, live, as `go_live` says.
415    pub fn read(&mut self, at: i64, n: usize) -> Result<Option<Buffer>, EngineError> {
416        let now = self.driver.at;
417        if at < now {
418            return Err(refused(
419                "engine.stream_behind",
420                format!("sample {at} is before sample {now}, where the stream stands"),
421                "read from the stream's position or later",
422            ));
423        }
424        if n == 0 {
425            return Err(refused(
426                "engine.empty_read",
427                format!("a read of no samples at sample {at}"),
428                "read one sample or more",
429            ));
430        }
431        match self.live {
432            true if at > now => {
433                for silenced in self.driver.skip(at)? {
434                    self.dropped
435                        .push(self.driver.value_graph.values[silenced].name.clone());
436                }
437            }
438            _ => {
439                let block = self.config.block;
440                while self.driver.at < at {
441                    let step = block.min((at - self.driver.at) as usize);
442                    self.reads_on(step)?;
443                    if !self.driver.pulled(step)? {
444                        return Ok(None);
445                    }
446                    self.retire_terms_below_silence_threshold();
447                }
448            }
449        }
450        self.reads_on(n)?;
451        let block = self.driver.read(n)?;
452        self.retire_terms_below_silence_threshold();
453        Ok(block.map(|b| widened(b, self.width)))
454    }
455
456    /// Each node memory answered that the next `n` samples ask past what it holds reads on from
457    /// its own value, made as the change that made it would have: live, a stateful one already
458    /// sounding then starts silent.
459    fn reads_on(&mut self, n: usize) -> Result<(), EngineError> {
460        let at = self.driver.at;
461        let range = Extent::new(self.driver.start, self.driver.last());
462        let asked_range = Extent::new(at, at.saturating_add(n as i64).min(range.end).max(at));
463        let value_graph = &mut self.driver.value_graph;
464        if asked_range.is_empty() {
465            return Ok(());
466        }
467        for short in value_graph.short((value_graph.root, asked_range, Past::Held)) {
468            value_graph.read_on(&self.world.typing, short)?;
469            let landed = value_graph.landed(short);
470            let carried = value_graph.settled(value_graph.root, &[], (landed, self.live));
471            for silent in carried.silent {
472                self.dropped.push(value_graph.values[silent].name.clone());
473            }
474            value_graph.offers(&self.world.typing, range);
475        }
476        Ok(())
477    }
478
479    /// An edited node with no state, or one a read skips past, starts silent, never computing
480    /// its past, named in `dropped`; a formula reads on exactly.
481    pub fn go_live(&mut self) {
482        self.live = true;
483        let last = self.last(self.driver.last());
484        self.driver.bound(last);
485    }
486
487    fn last(&self, range_end: i64) -> i64 {
488        match self.live {
489            true => self.config.render.range.end.unwrap_or(i64::MAX),
490            false => range_end,
491        }
492    }
493
494    /// The latest nodes a live edit started silent, of `counts().dropped`.
495    pub fn dropped(&self) -> Vec<&str> {
496        self.dropped.iter().map(String::as_str).collect()
497    }
498
499    pub fn counts(&self) -> Counts {
500        Counts {
501            dropped: self.dropped.made(),
502            late: self.late,
503            terms: self.terms.count(),
504            built: self.built,
505            demands: self.demands,
506            tier: self.driver.recording.since(&self.driver.memory),
507        }
508    }
509
510    /// Retires every term silent at the root by now and from the first sample of `notes` the
511    /// root asks, while all it retired early, summed through the gain, stay under the silence
512    /// threshold.
513    fn retire_terms_below_silence_threshold(&mut self) {
514        let now = self.driver.at;
515        let heard = &self.heard;
516        let ending = *self
517            .ending
518            .get_or_insert_with(|| heard.values().map(|h| h.support.end).min());
519        if ending.is_none_or(|end| end > now) {
520            return;
521        }
522        let (value_graph, tys) = (&self.driver.value_graph, &self.world.typing);
523        let last = self.driver.last();
524        let asked = match tys.id(NOTES).and_then(|notes| value_graph.of(notes)) {
525            Some(notes) if now < last => {
526                self.demands += 1;
527                let needs = value_graph.demand(Extent::new(now, last));
528                needs[notes].hold.iter().next().map(|asked| asked.start)
529            }
530            Some(_) => None,
531            None => Some(i64::MIN),
532        };
533        let silence_threshold = self.config.render.profile.silence_threshold_amplitude();
534        let gain = self.gain.gain.unwrap_or(f64::INFINITY);
535        if let Some(from) = asked {
536            // Far under the silence threshold, a bound joins one sum, so the list stays short.
537            let let_go = silence_threshold * 2f64.powi(-30);
538            let mut faded = self.faded;
539            self.fading.retain(|f| match f.from(from) {
540                b if b <= let_go => {
541                    faded = (faded + b) * (1.0 + f64::EPSILON);
542                    false
543                }
544                _ => true,
545            });
546            self.faded = faded;
547        }
548        let mut gone = BTreeSet::new();
549        for (handle, h) in &self.heard {
550            let end = h.support.end;
551            if end > now || asked.is_some_and(|from| end > from) {
552                continue;
553            }
554            let (Some(from), Some(fading)) = (asked, &h.fading) else {
555                gone.insert(*handle);
556                continue;
557            };
558            let own = fading.from(from);
559            if own > 0.0 {
560                let left: f64 = self.fading.iter().map(|f| f.from(from)).sum();
561                let ops = self.fading.len() as f64 + 3.0;
562                let left = (self.faded + left + own) * (1.0 + ops * f64::EPSILON);
563                if !under(gain, left, silence_threshold) {
564                    continue;
565                }
566                self.fading.push(fading.clone());
567            }
568            gone.insert(*handle);
569        }
570        let heard = &self.heard;
571        let support = |handle: Handle| heard.get(&handle).map(|h| h.support);
572        let went = self
573            .terms
574            .retire(&|handle| gone.contains(&handle), &support);
575        if !went.is_empty() {
576            self.generation += 1;
577        }
578        for handle in went {
579            self.heard.remove(&handle);
580            self.ending = None;
581        }
582    }
583
584    /// Every segment of its own clock the stream computed of `node`'s value, in order.
585    pub fn evaluated(&self, node: &str) -> Vec<sva_samples::Extent> {
586        let value_graph = &self.driver.value_graph;
587        self.world
588            .typing
589            .id(node)
590            .and_then(|id| value_graph.of(id))
591            .map_or(Vec::new(), |at| value_graph.values[at].evaluated.clone())
592    }
593
594    pub fn cutting_below_silence_threshold(&self) -> sva_samples::CuttingBelowSilenceThreshold {
595        sva_samples::CuttingBelowSilenceThreshold {
596            silence_threshold_dbfs: self.config.render.profile.silence_threshold_dbfs,
597            treated_as_silent_from_sample: self
598                .treated_as_silent_from_sample
599                .map(|at| (STREAMED.to_string(), at))
600                .into_iter()
601                .collect(),
602        }
603    }
604
605    pub fn landed(&self, handle: Handle) -> Option<i64> {
606        self.terms.landed(handle)
607    }
608
609    pub fn position(&self) -> i64 {
610        self.driver.at
611    }
612
613    pub fn work(&self) -> Work {
614        self.driver.work
615    }
616
617    pub fn stats(&self) -> CacheStats {
618        self.driver.recording.stats(&self.driver.memory)
619    }
620
621    /// The bytes its values hold, samples and state.
622    pub fn held_bytes(&self) -> usize {
623        self.driver.value_graph.bytes()
624    }
625
626    pub fn end(&self) -> Option<i64> {
627        self.driver.end()
628    }
629
630    pub fn width(&self) -> usize {
631        self.width
632    }
633
634    pub fn config(&self) -> &StreamConfig {
635        &self.config
636    }
637}
638
639#[derive(Clone, Copy, Default)]
640struct Gained {
641    gain: Option<f64>,
642    of: Option<Hash>,
643}
644
645/// An edit, each with the graph that reaches it and all the stream plays.
646pub enum Change {
647    Target(Graph, Expr),
648    Add(Graph, Expr, Placed),
649    Replace(Handle, Graph, Expr, Placed),
650    /// A sounding term is cut where the stream stands as the edit is built; what played stays.
651    Remove(Handle),
652}
653
654/// Where a term's sample 0 sits: the stream's own, or the sample its add lands at.
655#[derive(Clone, Copy, Debug, PartialEq, Eq)]
656pub enum Placed {
657    Written,
658    Landing,
659}
660
661/// A replace or remove answers whether the stream still held its handle.
662#[derive(Clone, Copy, Debug, PartialEq, Eq)]
663pub enum Changed {
664    Edited,
665    Added(Handle),
666    Held(bool),
667}
668
669/// What a change wants played, and what it answers once it lands.
670struct Prospect {
671    target: Expr,
672    terms: Terms,
673    term: Option<(Handle, Expr)>,
674    from: Option<(Graph, Vec<String>)>,
675    answer: Changed,
676    landing: Option<i64>,
677    /// Its own expression, where it brought one.
678    parsed: usize,
679}
680
681/// What one try at a change came to.
682enum Attempt {
683    Landed(Changed),
684    Asks(Vec<Hash>),
685    Reads(Vec<(Hash, Extent)>),
686    /// The stream left the sample it was placed at.
687    Moved,
688}
689
690/// What one change has looked up and read, across its tries.
691#[derive(Default)]
692struct Local {
693    fetched: Vec<(Hash, Vec<Arc<Buffer>>)>,
694    asked: Vec<(Hash, Extent)>,
695    unread: BTreeSet<Hash>,
696    lookups: usize,
697    issued: i64,
698}
699
700impl Local {
701    async fn look<B: Backend>(&mut self, keys: Vec<Hash>, tier: &Tier<B>, round: u64) {
702        for key in keys {
703            self.lookups += 1;
704            tier.lookup(key, round).await;
705        }
706    }
707
708    /// What `wants` asks, off memory; one a fetch's reads did not reach is asked again.
709    async fn read<B: Backend>(&mut self, wants: Vec<(Hash, Extent)>, tier: &Tier<B>) {
710        let fetched = tier.fetch(&wants).await;
711        for (key, over) in wants {
712            if fetched.left.contains(&(key, over)) {
713                continue;
714            }
715            self.asked.push((key, over));
716            if !fetched.handed.iter().any(|(held, _)| *held == key) {
717                self.unread.insert(key);
718            }
719        }
720        self.fetched.extend(fetched.handed);
721    }
722
723    /// Whether this change already read `key` over `over`, or could not.
724    fn holds(&self, key: Hash, over: Extent) -> bool {
725        self.unread.contains(&key)
726            || self
727                .asked
728                .iter()
729                .any(|(k, e)| *k == key && e.intersect(over) == over)
730    }
731}
732
733/// `build`'s change, the stream held only while each try runs, between two blocks: an exact
734/// stream pulled once its edit is done plays it where issued. `build` runs again whenever the
735/// stream changed under it.
736pub async fn change<E: From<EngineError>, B: Backend>(
737    stream: &RefCell<Stream>,
738    mut build: impl FnMut(&Stream) -> Result<Change, E>,
739    tier: &Tier<B>,
740) -> Result<Changed, E> {
741    let mut local = Local {
742        issued: stream.borrow().driver.at,
743        ..Local::default()
744    };
745    let round = tier.begin();
746    stream.borrow_mut().driver.recording.begin();
747    loop {
748        let change = build(&stream.borrow())?;
749        let prospect = match stream.borrow().prospect(change) {
750            Ok(prospect) => prospect,
751            Err(answer) => return Ok(answer),
752        };
753        let generation = stream.borrow().generation;
754        loop {
755            if stream.borrow().generation != generation {
756                break;
757            }
758            let attempt =
759                stream
760                    .borrow_mut()
761                    .attempt(&prospect, &mut local, (tier.memory(), round))?;
762            match attempt {
763                Attempt::Landed(answer) => return Ok(answer),
764                Attempt::Moved => break,
765                Attempt::Asks(keys) => local.look(keys, tier, round).await,
766                Attempt::Reads(wants) => local.read(wants, tier).await,
767            }
768        }
769    }
770}
771
772/// What the next second of `stream` reads off memory, as far as one fetch reaches; a block
773/// reading samples not ready computes them, or, live, starts them silent.
774pub async fn fetch<B: Backend>(stream: &RefCell<Stream>, tier: &Tier<B>) {
775    let needs = stream.borrow().needs();
776    let fetched = tier.fetch(&needs).await;
777    let mut stream = stream.borrow_mut();
778    for (key, parts) in fetched.handed {
779        stream.driver.value_graph.took(key, &parts);
780    }
781}
782
783fn ahead(at: i64, rate: u32) -> Extent {
784    Extent::new(at, at.saturating_add(i64::from(rate)))
785}
786
787fn widens(plays: usize, width: usize) -> Result<(), EngineError> {
788    match plays == width || plays == 1 {
789        true => Ok(()),
790        false => Err(refused(
791            "engine.stream_width",
792            format!("this plays {plays} channel(s), and the stream plays {width}"),
793            "play as many channels as the stream, or one, or open a new stream for it",
794        )),
795    }
796}
797
798fn blocked(config: &StreamConfig) -> Result<&StreamConfig, EngineError> {
799    match config.block {
800        0 => Err(refusal("a block of no samples".to_string())),
801        _ => Ok(config),
802    }
803}
804
805pub(super) fn refusal(what: String) -> EngineError {
806    refused(
807        "engine.no_stream",
808        format!("this target opens no stream: {what}"),
809        "stream an expression over the nodes the composition defines",
810    )
811}
812
813fn refused(code: &str, message: String, help: &str) -> EngineError {
814    EngineError::refused(Diagnostic {
815        code: code.to_string(),
816        message,
817        location: Located::at(STREAMED, None),
818        help: help.to_string(),
819    })
820}
821
822/// A mono block in each of `width` channels; any other as it is.
823fn widened(mut block: Buffer, width: usize) -> Buffer {
824    if block.planes.len() == 1 {
825        let copies = vec![block.planes[0].clone(); width.saturating_sub(1)];
826        block.planes.extend(copies);
827    }
828    block
829}