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