1use 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
27pub const LATEST: usize = 256;
29
30#[derive(Clone, Debug, PartialEq)]
31pub struct StreamConfig {
32 pub block: usize,
33 pub channels: Option<usize>,
35 pub render: RenderConfig,
36}
37
38pub 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 pub late: usize,
62 pub terms: usize,
63}
64
65#[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 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 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 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 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 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 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 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 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 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 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 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
421pub enum Change {
423 Target(Graph, Expr),
424 Add(Graph, Expr, Placed),
425 Replace(Handle, Graph, Expr, Placed),
426 Remove(Handle),
428}
429
430#[derive(Clone, Copy, Debug, PartialEq, Eq)]
432pub enum Placed {
433 Written,
434 Landing,
435}
436
437#[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
455pub 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
548fn 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
575async 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
623fn 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}