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