1use std::cell::{Ref, 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::{Beside, Table, edit};
15use super::terms::{Handle, NOTES, Terms, placed};
16use super::{Ends, Render, RenderConfig, range_over, 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
25#[cfg(test)]
26mod rebuilt;
27
28pub const STREAMED: &str = "streamed";
29
30pub const LATEST: usize = 256;
32
33#[derive(Clone, Debug, PartialEq)]
34pub struct StreamConfig {
35 pub block: usize,
36 pub channels: Option<usize>,
38 pub render: RenderConfig,
39}
40
41pub struct Stream {
45 config: StreamConfig,
46 played: Played,
47 driver: Driver,
48 graph: Graph,
49 expr: Expr,
50 terms: Terms,
51 width: usize,
52 supports: BTreeMap<Handle, Extent>,
53 met: Met,
54 generation: u64,
55 live: bool,
56 dropped: Recent<String>,
57 late: usize,
58 built: Built,
59}
60
61struct Played {
64 shell: Render,
65 keys: BTreeMap<String, Hash>,
66 reads: BTreeMap<String, Vec<String>>,
67}
68
69#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
70pub struct Counts {
71 pub dropped: usize,
72 pub late: usize,
74 pub terms: usize,
75 pub built: Built,
77}
78
79#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
81pub struct Built {
82 pub instances: usize,
83 pub typed: usize,
84 pub values: usize,
85 pub lookups: usize,
86}
87
88#[derive(Default)]
90struct Met {
91 epoch: u64,
92 known: Known,
93}
94
95impl Met {
96 fn over(&mut self, store: &impl Through) -> &Known {
97 if store.epoch() != self.epoch {
98 self.known.clear();
99 self.epoch = store.epoch();
100 }
101 &self.known
102 }
103
104 fn noted(&mut self, key: Hash, found: &Option<Arc<Stored>>, epoch: u64, store: &impl Through) {
106 self.over(store);
107 if epoch == self.epoch {
108 self.known.insert(key, found.clone());
109 }
110 }
111
112 fn entries(&mut self, store: &impl Through) -> Vec<(Hash, Option<Arc<Stored>>)> {
113 let entries = self.over(store).iter();
114 entries.map(|(key, found)| (*key, found.clone())).collect()
115 }
116}
117
118struct Shelled {
119 played: Played,
120 range: Extent,
121 table: Table,
122 hits: Vec<Lookup>,
123 built: Built,
124}
125
126impl Stream {
127 pub async fn open(
128 graph: &Graph,
129 target: &Expr,
130 config: StreamConfig,
131 cache: Option<&Cache>,
132 store: &impl Through,
133 ) -> Result<Stream, EngineError> {
134 let terms = Terms::default();
135 let (mut known, epoch) = (Known::new(), store.epoch());
136 let render = &blocked(&config)?.render;
137 let lookups = std::cell::Cell::new(0);
138 let walk = async |found: &mut Frontier<'_>| {
139 lookups.set(found.walked(&mut known, store).await);
140 };
141 let mut found = shelled(graph, target, &terms, render, (walk, || None)).await?;
142 found.built.lookups = lookups.get();
143 let width = config.channels.unwrap_or(found.width());
144 if width == 0 {
145 return Err(refusal("a stream of no channels".to_string()));
146 }
147 widens(found.width(), width)?;
148 let mut met = Met::default();
149 for (key, found) in known {
150 met.noted(key, &found, epoch, store);
151 }
152 let next = ahead(found.range.start, render.rate);
153 through::load(&mut found.table, store, next, &BTreeSet::new()).await;
154 let mut recording = Recording::over(cache, config.render.cache_policy).latest(LATEST);
155 recording.found(found.hits);
156 let driver = Driver::new(
157 found.table,
158 found.range,
159 config.block,
160 &config.render,
161 recording,
162 );
163 Ok(Stream {
164 graph: graph.clone(),
165 expr: target.clone(),
166 terms,
167 width,
168 supports: BTreeMap::new(),
169 met,
170 generation: 0,
171 live: false,
172 dropped: Recent::keeping(LATEST),
173 late: 0,
174 built: found.built,
175 driver,
176 config,
177 played: found.played,
178 })
179 }
180
181 pub fn exprs(&self) -> impl Iterator<Item = &Expr> {
182 std::iter::once(&self.expr).chain(self.terms.exprs())
183 }
184
185 fn prospect(&self, change: Change) -> Result<Prospect, Changed> {
186 let mut render = self.config.render.clone();
187 render.range.start = Some(self.driver.start);
188 let mut landing = None;
189 let (graph, target, terms, answer) = match change {
190 Change::Target(graph, target) => (graph, target, self.terms.clone(), Changed::Edited),
191 Change::Add(graph, term, at) => {
192 let term = match at {
193 Placed::Written => term,
194 Placed::Landing => placed(&term, *landing.insert(self.driver.at)),
195 };
196 let (terms, handle) = self.terms.added(term);
197 (graph, self.expr.clone(), terms, Changed::Added(handle))
198 }
199 Change::Replace(handle, graph, term, at) => {
200 let terms = self.terms.replaced(handle, |landed| match at {
201 Placed::Written => term,
202 Placed::Landing => placed(&term, landed),
203 });
204 let terms = terms.ok_or(Changed::Held(false))?;
205 (graph, self.expr.clone(), terms, Changed::Held(true))
206 }
207 Change::Remove(handle) => {
208 let at = self.driver.at as f64 / f64::from(self.played.shell.rate());
209 let terms = self.terms.removed(handle, at);
210 let terms = terms.ok_or(Changed::Held(false))?;
211 (
212 self.graph.clone(),
213 self.expr.clone(),
214 terms,
215 Changed::Held(true),
216 )
217 }
218 };
219 Ok(Prospect {
220 graph,
221 target,
222 terms,
223 answer,
224 render,
225 generation: self.generation,
226 landing,
227 })
228 }
229
230 fn wanting(
233 &self,
234 (generation, landing): (u64, Option<i64>),
235 table: &Table,
236 unread: &BTreeSet<Hash>,
237 ) -> Option<Vec<(Arc<Stored>, Extent)>> {
238 if generation != self.generation || landing.is_some_and(|at| at != self.driver.at) {
239 return None;
240 }
241 let sounding = self.driver.table.stored_keys();
242 let wants = table.wants(ahead(self.driver.at, self.config.render.rate));
243 let brought = wants.into_iter();
244 let brought =
245 brought.filter(|(s, _)| !sounding.contains(&s.key) && !unread.contains(&s.key));
246 Some(brought.collect())
247 }
248
249 fn apply(
252 &mut self,
253 (prospect, issued): (Prospect, i64),
254 shelled: Shelled,
255 fetched: &[(Hash, Vec<Buffer>)],
256 ) {
257 let Shelled {
258 played,
259 range,
260 mut table,
261 hits,
262 built,
263 } = shelled;
264 self.built = built;
265 let old = std::mem::replace(&mut self.driver.table, Table::empty());
266 let dropped = edit::carried(&mut table, old, self.driver.at, self.live);
267 for (key, samples) in fetched {
268 table.took(*key, samples);
269 }
270 for at in dropped {
271 self.dropped.push(table.values[at].name.clone());
272 }
273 self.late += usize::from(self.driver.at > issued);
274 let last = self.last(range.end);
275 self.driver.replace(table, last);
276 self.driver.recording.found(hits);
277 self.played = played;
278 self.graph = prospect.graph;
279 self.expr = prospect.target;
280 self.terms = prospect.terms;
281 if let Changed::Added(handle) = prospect.answer {
282 self.terms.land(handle, self.driver.at);
283 }
284 let (tys, profile) = (&self.played.shell.tys, &self.played.shell.config.profile);
285 let supports = Supports::over(tys, profile, Some(&self.driver.table.supports));
286 let support = |handle: Handle| Some((handle, supports.of(tys.id(&handle.node())?)));
287 self.supports = self.terms.handles().filter_map(support).collect();
288 self.generation += 1;
289 self.prune();
290 }
291
292 pub async fn fetch(&mut self, store: &impl Through) {
295 let next = ahead(self.driver.at, self.config.render.rate);
296 through::load(&mut self.driver.table, store, next, &BTreeSet::new()).await;
297 }
298
299 pub fn wanted(&self) -> Vec<(Arc<Stored>, Extent)> {
301 let next = ahead(self.driver.at, self.config.render.rate);
302 self.driver.table.wants(next)
303 }
304
305 pub fn took(&mut self, key: Hash, samples: &[Buffer]) {
306 self.driver.table.took(key, samples);
307 }
308
309 pub fn read(&mut self, at: i64, n: usize) -> Result<Option<Block>, EngineError> {
313 let now = self.driver.at;
314 if at < now {
315 return Err(refused(
316 "engine.stream_behind",
317 format!("sample {at} is before sample {now}, where the stream stands"),
318 "read from the stream's position or later",
319 ));
320 }
321 if n == 0 {
322 return Err(refused(
323 "engine.empty_read",
324 format!("a read of no samples at sample {at}"),
325 "read one sample or more",
326 ));
327 }
328 match self.live {
329 true if at > now => {
330 for silenced in self.driver.skip(at)? {
331 self.dropped
332 .push(self.driver.table.values[silenced].name.clone());
333 }
334 }
335 _ => {
336 let block = self.config.block;
337 while self.driver.at < at {
338 let step = block.min((at - self.driver.at) as usize);
339 if !self.driver.pulled(step)? {
340 return Ok(None);
341 }
342 self.prune();
343 }
344 }
345 }
346 let block = self.driver.read(n)?;
347 self.prune();
348 Ok(block.map(|b| b.widened(self.width)))
349 }
350
351 pub fn go_live(&mut self) {
355 self.live = true;
356 let last = self.last(self.driver.last());
357 self.driver.bound(last);
358 }
359
360 fn last(&self, range_end: i64) -> i64 {
361 match self.live {
362 true => self.config.render.range.end.unwrap_or(i64::MAX),
363 false => range_end,
364 }
365 }
366
367 pub fn dropped(&self) -> Vec<&str> {
369 self.dropped.iter().map(String::as_str).collect()
370 }
371
372 pub fn counts(&self) -> Counts {
373 Counts {
374 dropped: self.dropped.made(),
375 late: self.late,
376 terms: self.terms.count(),
377 built: self.built,
378 }
379 }
380
381 fn prune(&mut self) {
384 let (table, tys) = (&self.driver.table, &self.played.shell.tys);
385 let (now, last) = (self.driver.at, self.driver.last());
386 let asked = match tys.id(NOTES).and_then(|notes| table.of(notes)) {
387 Some(notes) if now < last => {
388 let needs = table.demand(Extent::new(now, last));
389 needs[notes].hold.iter().next().map(|asked| asked.start)
390 }
391 Some(_) => None,
392 None => Some(i64::MIN),
393 };
394 let supports = &self.supports;
395 let gone = |handle: Handle| {
396 let support = supports.get(&handle);
397 support.is_some_and(|s| s.end <= now && asked.is_none_or(|from| s.end <= from))
398 };
399 let support = |handle: Handle| supports.get(&handle).copied();
400 if self.terms.prune(&gone, &support) {
401 self.generation += 1;
402 }
403 }
404
405 pub fn evaluated(&self, node: &str) -> Vec<sva_samples::Extent> {
407 let table = &self.driver.table;
408 self.played
409 .shell
410 .tys
411 .id(node)
412 .and_then(|id| table.of(id))
413 .map_or(Vec::new(), |at| table.values[at].evaluated.clone())
414 }
415
416 pub fn pruned(&self) -> sva_samples::Pruned {
417 self.driver.table.pruned()
418 }
419
420 pub fn landed(&self, handle: Handle) -> Option<i64> {
421 self.terms.landed(handle)
422 }
423
424 pub fn position(&self) -> i64 {
425 self.driver.at
426 }
427
428 pub fn work(&self) -> Work {
429 self.driver.work
430 }
431
432 pub fn stats(&self) -> CacheStats {
433 self.driver.recording.stats()
434 }
435
436 pub fn held_bytes(&self) -> usize {
438 self.driver.table.bytes()
439 }
440
441 pub fn end(&self) -> Option<i64> {
442 self.driver.end()
443 }
444
445 pub fn width(&self) -> usize {
446 self.width
447 }
448
449 pub fn config(&self) -> &StreamConfig {
450 &self.config
451 }
452}
453
454pub enum Change {
456 Target(Graph, Expr),
457 Add(Graph, Expr, Placed),
458 Replace(Handle, Graph, Expr, Placed),
459 Remove(Handle),
461}
462
463#[derive(Clone, Copy, Debug, PartialEq, Eq)]
465pub enum Placed {
466 Written,
467 Landing,
468}
469
470#[derive(Clone, Copy, Debug, PartialEq, Eq)]
472pub enum Changed {
473 Edited,
474 Added(Handle),
475 Held(bool),
476}
477
478struct Prospect {
479 graph: Graph,
480 target: Expr,
481 terms: Terms,
482 answer: Changed,
483 render: RenderConfig,
484 generation: u64,
485 landing: Option<i64>,
486}
487
488pub async fn change<E: From<EngineError>>(
492 stream: &RefCell<Stream>,
493 mut build: impl FnMut(&Stream) -> Result<Change, E>,
494 store: &impl Through,
495) -> Result<Changed, E> {
496 let mut local = Met::default();
497 let mut fetched: Vec<(Hash, Vec<Buffer>)> = Vec::new();
498 let (mut unread, mut asked) = (BTreeSet::new(), Vec::<(Hash, Extent)>::new());
499 let issued = stream.borrow().driver.at;
500 let lookups = std::cell::Cell::new(0);
501 loop {
502 let change = build(&stream.borrow())?;
503 let prospect = match stream.borrow().prospect(change) {
504 Ok(prospect) => prospect,
505 Err(answer) => return Ok(answer),
506 };
507 for (key, found) in stream.borrow_mut().met.entries(store) {
508 local.noted(key, &found, store.epoch(), store);
509 }
510 let walk = async |found: &mut Frontier<'_>| {
511 while let Some(key) = found.walk(local.over(store)) {
512 let epoch = store.epoch();
513 lookups.set(lookups.get() + 1);
514 let found = store.lookup(key).await.map(Arc::new);
515 stream.borrow_mut().met.noted(key, &found, epoch, store);
516 local.noted(key, &found, epoch, store);
517 }
518 };
519 let (graph, target) = (&prospect.graph, &prospect.target);
520 let config = &prospect.render;
521 let prior = || Some(stream.borrow());
522 let mut shelled = shelled(graph, target, &prospect.terms, config, (walk, prior)).await?;
523 widens(shelled.width(), stream.borrow().width())?;
524 let standing = (prospect.generation, prospect.landing);
525 loop {
526 for (key, samples) in &fetched {
527 shelled.table.took(*key, samples);
528 }
529 let wants = stream.borrow().wanting(standing, &shelled.table, &unread);
530 let Some(wants) = wants else {
531 break;
532 };
533 if wants.is_empty() {
534 let answer = prospect.answer;
535 shelled.built.lookups = lookups.get();
536 stream
537 .borrow_mut()
538 .apply((prospect, issued), shelled, &fetched);
539 return Ok(answer);
540 }
541 for (stored, over) in wants {
542 let again = asked
543 .iter()
544 .any(|(key, e)| *key == stored.key && !e.intersect(over).is_empty());
545 asked.push((stored.key, over));
546 let read = match again {
547 false => store.read(&stored, over).await,
548 true => None,
549 };
550 match read {
551 Some(samples) => fetched.push((stored.key, samples)),
552 None => {
553 unread.insert(stored.key);
554 }
555 }
556 }
557 }
558 }
559}
560
561fn ahead(at: i64, rate: u32) -> Extent {
562 Extent::new(at, at.saturating_add(i64::from(rate)))
563}
564
565fn widens(plays: usize, width: usize) -> Result<(), EngineError> {
566 match plays == width || plays == 1 {
567 true => Ok(()),
568 false => Err(refused(
569 "engine.stream_width",
570 format!("this plays {plays} channel(s), and the stream plays {width}"),
571 "play as many channels as the stream, or one, or open a new stream for it",
572 )),
573 }
574}
575
576fn blocked(config: &StreamConfig) -> Result<&StreamConfig, EngineError> {
577 match config.block {
578 0 => Err(refusal("a block of no samples".to_string())),
579 _ => Ok(config),
580 }
581}
582
583fn wrapped(graph: &Graph, target: &Expr, terms: &Terms) -> Result<Graph, EngineError> {
586 let mut wrapped = graph.clone();
587 if !terms.is_empty() && graph.defines(NOTES) {
588 return Err(EngineError::refused(Diagnostic {
589 code: "engine.no_stream".to_string(),
590 message: format!(
591 "a term is added to `@{NOTES}`, and this composition defines its own `{NOTES}`"
592 ),
593 location: Located::at(NOTES, None),
594 help: format!("rename the composition's `{NOTES}`, or play it without adding terms"),
595 }));
596 }
597 let own = terms.is_empty() && graph.defines(NOTES);
598 let sum = (!own).then(|| (NOTES.to_string(), terms.sum()));
599 let defined = std::iter::once((STREAMED.to_string(), target.clone()));
600 for (name, body) in defined.chain(sum).chain(terms.nodes()) {
601 if !wrapped.define(&name, body) {
602 return Err(refusal(format!(
603 "this composition already has a node named `{name}`"
604 )));
605 }
606 }
607 Ok(wrapped)
608}
609
610async fn shelled<'s>(
612 graph: &Graph,
613 target: &Expr,
614 terms: &Terms,
615 config: &RenderConfig,
616 (walk, prior): (
617 impl AsyncFnOnce(&mut Frontier<'_>),
618 impl Fn() -> Option<Ref<'s, Stream>>,
619 ),
620) -> Result<Shelled, EngineError> {
621 let wrapped = wrapped(graph, target, terms)?;
622 let instances = instantiate::instantiate(&wrapped, STREAMED, config.rate)?;
623 let named = instances.paths().count();
624 let root = instances.instance_of(STREAMED)?;
625 let order = schedule::schedule_from(&instances, std::slice::from_ref(&root))?;
626 let keys = through::keys(&wrapped, &instances, &order, config);
627 let mut found = Frontier::from((&instances, &order), &keys, &root, (config, true));
628 if !terms.is_empty() {
629 found.unstored(reading(&order, &instances.instance_of(NOTES)?));
630 found.unstored(terms.handles().map(Handle::node));
631 }
632 if let Some(prior) = prior() {
633 found.asking(anew(
634 (&prior.played.keys, &prior.played.reads),
635 &keys,
636 &order,
637 ));
638 }
639 walk(&mut found).await;
640 let within = order.within(&found.visited);
641 let prior = prior();
642 let (mut tys, carried) = match &prior {
643 Some(prior) => {
644 let (held, now) = (&prior.played.keys, &keys);
645 let same = |path: &str| held.get(path) == now.get(path);
646 let typing = &prior.played.shell.tys;
647 let prior_typing = typing::Prior {
648 typing,
649 same: &same,
650 };
651 let (tys, carried) = typing::infer_beside(&instances, &within, &prior_typing)?;
652 (tys, Some(carried))
653 }
654 None => (
655 typing::infer_over(&instances, &within, &BTreeMap::new())?,
656 None,
657 ),
658 };
659 let id = tys
660 .id(&root)
661 .ok_or_else(|| EngineError::UnknownNode(root.clone()))?;
662 terms.name(&mut tys);
663 let prefixes: BTreeMap<NodeId, Arc<Stored>> = std::mem::take(&mut found.stored)
664 .into_iter()
665 .map(|(path, stored)| (tys.id(&path).expect("a walked node is typed"), stored))
666 .collect();
667 let schedule = schedule::plan(&tys, id, &[]);
668 let shell = Render::shell(tys, id, config.clone(), schedule);
669 let profile = &shell.config.profile;
670 let beside = match (&prior, &carried) {
671 (Some(prior), Some(carried)) => Beside::of(&prior.driver.table, &carried.nodes),
672 _ => Beside::default(),
673 };
674 let table = Table::prefixed(&shell.tys, shell.root, profile, (&prefixes, beside))?;
675 drop(prior);
676 let support = Supports::over(&shell.tys, profile, Some(&table.supports)).of(shell.root);
677 let range = range_over(&shell, support, Ends::Pulled)?;
678 let hits = std::mem::take(&mut found.lookups).into_iter();
679 let hits = hits.filter(|l| l.outcome == Outcome::Hit).collect();
680 drop(found);
681 let built = Built {
682 instances: named,
683 typed: shell.tys.lowered().len(),
684 values: table.built,
685 lookups: 0,
686 };
687 let reads = order.into_deps();
688 Ok(Shelled {
689 played: Played { shell, keys, reads },
690 range,
691 table,
692 hits,
693 built,
694 })
695}
696
697fn anew(
700 (prior, reads): (&BTreeMap<String, Hash>, &BTreeMap<String, Vec<String>>),
701 keys: &BTreeMap<String, Hash>,
702 order: &schedule::Order,
703) -> BTreeSet<Hash> {
704 let changed = keys
705 .iter()
706 .filter(|(path, key)| prior.get(*path) != Some(key));
707 let mut out = BTreeSet::new();
708 for (path, _) in changed {
709 let before = reads.get(path).map_or(&[][..], Vec::as_slice);
710 let new = order
711 .deps(path)
712 .iter()
713 .filter(|read| !before.contains(read));
714 out.extend(new.filter_map(|read| keys.get(read)));
715 }
716 out
717}
718
719impl Shelled {
720 fn width(&self) -> usize {
721 self.table.values[self.table.root].width
722 }
723}
724
725fn reading(order: &schedule::Order, path: &str) -> BTreeSet<String> {
727 let mut out = BTreeSet::from([path.to_string()]);
728 loop {
729 let more: Vec<String> = order
730 .groups
731 .iter()
732 .flatten()
733 .filter(|node| !out.contains(*node))
734 .filter(|node| order.deps(node).iter().any(|read| out.contains(read)))
735 .cloned()
736 .collect();
737 if more.is_empty() {
738 return out;
739 }
740 out.extend(more);
741 }
742}
743
744fn refusal(what: String) -> EngineError {
745 refused(
746 "engine.no_stream",
747 format!("this target opens no stream: {what}"),
748 "stream an expression over the nodes the composition defines",
749 )
750}
751
752fn refused(code: &str, message: String, help: &str) -> EngineError {
753 EngineError::refused(Diagnostic {
754 code: code.to_string(),
755 message,
756 location: Located::at(STREAMED, None),
757 help: help.to_string(),
758 })
759}