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::{Table, edit};
14use super::terms::{Handle, NOTES, Terms};
15use super::{Ends, Render, RenderConfig, range_of, through};
16use crate::cache::{Cache, CacheStats, Lookup, Outcome, Recording, Stored, Through};
17use crate::error::{Diagnostic, EngineError, Located};
18use crate::flops::Work;
19use crate::instantiate;
20use crate::schedule;
21use crate::typing;
22
23pub const STREAMED: &str = "streamed";
24
25#[derive(Clone, Debug, PartialEq)]
26pub struct StreamConfig {
27 pub block: usize,
28 pub render: RenderConfig,
29}
30
31pub struct Stream {
35 config: StreamConfig,
36 shell: Render,
37 driver: Driver,
38 graph: Graph,
39 expr: Expr,
40 terms: Terms,
41 met: Met,
42 generation: u64,
43 live: bool,
44 dropped: Vec<String>,
45}
46
47#[derive(Default)]
49struct Met {
50 epoch: u64,
51 known: Known,
52}
53
54impl Met {
55 fn over(&mut self, store: &impl Through) -> &Known {
56 if store.epoch() != self.epoch {
57 self.known.clear();
58 self.epoch = store.epoch();
59 }
60 &self.known
61 }
62
63 fn noted(&mut self, key: Hash, found: &Option<Arc<Stored>>, epoch: u64, store: &impl Through) {
65 self.over(store);
66 if epoch == self.epoch {
67 self.known.insert(key, found.clone());
68 }
69 }
70
71 fn hits(&mut self, store: &impl Through) -> Vec<(Hash, Option<Arc<Stored>>)> {
72 let hits = self.over(store).iter().filter(|(_, found)| found.is_some());
73 hits.map(|(key, found)| (*key, found.clone())).collect()
74 }
75}
76
77struct Shelled {
78 shell: Render,
79 range: Extent,
80 table: Table,
81 hits: Vec<Lookup>,
82}
83
84impl Stream {
85 pub async fn open(
86 graph: &Graph,
87 target: &Expr,
88 config: StreamConfig,
89 cache: Option<&Cache>,
90 store: &impl Through,
91 ) -> Result<Stream, EngineError> {
92 let mut terms = Terms::default();
93 let (mut known, epoch) = (Known::new(), store.epoch());
94 let render = &blocked(&config)?.render;
95 let walk = async |found: &mut Frontier<'_>| found.walked(&mut known, store).await;
96 let mut found = shelled(graph, target, &mut terms, render, walk).await?;
97 let mut met = Met::default();
98 for (key, hit) in known.into_iter().filter(|(_, found)| found.is_some()) {
99 met.noted(key, &hit, epoch, store);
100 }
101 let next = ahead(found.range.start, render.rate);
102 through::load(&mut found.table, store, next, &BTreeSet::new()).await;
103 let mut recording = Recording::over(cache, config.render.cache_policy);
104 recording.found(found.hits);
105 let driver = Driver::new(
106 found.table,
107 found.range,
108 config.block,
109 &config.render,
110 recording,
111 );
112 Ok(Stream {
113 graph: graph.clone(),
114 expr: target.clone(),
115 terms,
116 met,
117 generation: 0,
118 live: false,
119 dropped: Vec::new(),
120 driver,
121 config,
122 shell: found.shell,
123 })
124 }
125
126 pub fn exprs(&self) -> impl Iterator<Item = &Expr> {
127 std::iter::once(&self.expr).chain(self.terms.exprs())
128 }
129
130 fn prospect(&self, change: Change) -> Result<Prospect, Changed> {
131 let mut render = self.config.render.clone();
132 render.range.start = Some(self.driver.start);
133 let (graph, target, terms, answer) = match change {
134 Change::Target(graph, target) => (graph, target, self.terms.clone(), Changed::Edited),
135 Change::Add(graph, term) => {
136 let (terms, handle) = self.terms.added(term);
137 (graph, self.expr.clone(), terms, Changed::Added(handle))
138 }
139 Change::Replace(handle, graph, term) => {
140 let terms = self.terms.replaced(handle, term);
141 let terms = terms.ok_or(Changed::Held(false))?;
142 (graph, self.expr.clone(), terms, Changed::Held(true))
143 }
144 Change::Remove(handle) => {
145 let at = self.driver.at as f64 / f64::from(self.shell.rate());
146 let terms = self.terms.removed(handle, at).ok_or(Changed::Held(false))?;
147 (
148 self.graph.clone(),
149 self.expr.clone(),
150 terms,
151 Changed::Held(true),
152 )
153 }
154 };
155 Ok(Prospect {
156 graph,
157 target,
158 terms,
159 answer,
160 render,
161 generation: self.generation,
162 })
163 }
164
165 fn wanting(
168 &self,
169 generation: u64,
170 table: &Table,
171 unread: &BTreeSet<Hash>,
172 ) -> Option<Vec<(Arc<Stored>, Extent)>> {
173 if generation != self.generation {
174 return None;
175 }
176 let sounding = self.driver.table.stored_keys();
177 let wants = table.wants(ahead(self.driver.at, self.config.render.rate));
178 let brought = wants.into_iter();
179 let brought =
180 brought.filter(|(s, _)| !sounding.contains(&s.key) && !unread.contains(&s.key));
181 Some(brought.collect())
182 }
183
184 fn apply(&mut self, prospect: Prospect, shelled: Shelled, fetched: &[(Hash, Vec<Buffer>)]) {
187 let Shelled {
188 shell,
189 range,
190 mut table,
191 hits,
192 } = shelled;
193 let old = std::mem::replace(&mut self.driver.table, Table::empty());
194 let dropped = edit::carried(&mut table, old, self.driver.at, self.live);
195 for (key, samples) in fetched {
196 table.took(*key, samples);
197 }
198 self.dropped
199 .extend(dropped.into_iter().map(|at| table.values[at].name.clone()));
200 self.driver.replace(table, range.end);
201 self.driver.recording.found(hits);
202 self.shell = shell;
203 self.graph = prospect.graph;
204 self.expr = prospect.target;
205 self.terms = prospect.terms;
206 self.generation += 1;
207 self.prune();
208 }
209
210 pub async fn fetch(&mut self, store: &impl Through) {
213 let next = ahead(self.driver.at, self.config.render.rate);
214 through::load(&mut self.driver.table, store, next, &BTreeSet::new()).await;
215 }
216
217 pub fn wanted(&self) -> Vec<(Arc<Stored>, Extent)> {
219 let next = ahead(self.driver.at, self.config.render.rate);
220 self.driver.table.wants(next)
221 }
222
223 pub fn took(&mut self, key: Hash, samples: &[Buffer]) {
224 self.driver.table.took(key, samples);
225 }
226
227 pub fn next_block(&mut self) -> Result<Option<Block>, EngineError> {
229 let block = self.driver.next_block()?;
230 self.prune();
231 Ok(block)
232 }
233
234 pub fn go_live(&mut self) {
236 self.live = true;
237 }
238
239 pub fn dropped(&self) -> &[String] {
241 &self.dropped
242 }
243
244 fn prune(&mut self) {
246 let (table, shell) = (&self.driver.table, &self.shell);
247 let future = (self.driver.at < self.driver.last())
248 .then(|| table.demand(sva_samples::Extent::new(self.driver.at, self.driver.last())));
249 let gone = |id| {
250 table.of(id).is_some_and(|at| {
251 !table.values[at].evaluated.is_empty()
252 && future
253 .as_ref()
254 .is_none_or(|needs| needs[at].hold.is_empty())
255 })
256 };
257 let named = |leaf| crate::refs::identity(&shell.tys, leaf).ok();
258 if self.terms.prune(&gone, &named) {
259 self.generation += 1;
260 }
261 }
262
263 pub fn evaluated(&self, node: &str) -> Vec<sva_samples::Extent> {
265 let table = &self.driver.table;
266 self.shell
267 .tys
268 .id(node)
269 .and_then(|id| table.of(id))
270 .map_or(Vec::new(), |at| table.values[at].evaluated.clone())
271 }
272
273 pub fn position(&self) -> i64 {
274 self.driver.at
275 }
276
277 pub fn work(&self) -> Work {
278 self.driver.work
279 }
280
281 pub fn stats(&self) -> CacheStats {
282 self.driver.recording.stats()
283 }
284
285 pub fn held_bytes(&self) -> usize {
287 self.driver.table.bytes()
288 }
289
290 pub fn end(&self) -> Option<i64> {
291 self.driver.end()
292 }
293
294 pub fn width(&self) -> usize {
295 self.driver.table.values[self.driver.table.root].width
296 }
297
298 pub fn config(&self) -> &StreamConfig {
299 &self.config
300 }
301}
302
303pub enum Change {
305 Target(Graph, Expr),
306 Add(Graph, Expr),
307 Replace(Handle, Graph, Expr),
308 Remove(Handle),
310}
311
312#[derive(Clone, Copy, Debug, PartialEq, Eq)]
314pub enum Changed {
315 Edited,
316 Added(Handle),
317 Held(bool),
318}
319
320struct Prospect {
321 graph: Graph,
322 target: Expr,
323 terms: Terms,
324 answer: Changed,
325 render: RenderConfig,
326 generation: u64,
327}
328
329pub async fn change<E: From<EngineError>>(
333 stream: &RefCell<Stream>,
334 mut build: impl FnMut(&Stream) -> Result<Change, E>,
335 store: &impl Through,
336) -> Result<Changed, E> {
337 let mut local = Met::default();
338 let mut fetched: Vec<(Hash, Vec<Buffer>)> = Vec::new();
339 let (mut unread, mut asked) = (BTreeSet::new(), Vec::<(Hash, Extent)>::new());
340 loop {
341 let change = build(&stream.borrow())?;
342 let mut prospect = match stream.borrow().prospect(change) {
343 Ok(prospect) => prospect,
344 Err(answer) => return Ok(answer),
345 };
346 for (key, hit) in stream.borrow_mut().met.hits(store) {
347 local.noted(key, &hit, store.epoch(), store);
348 }
349 let walk = async |found: &mut Frontier<'_>| {
350 while let Some(key) = found.walk(local.over(store)) {
351 let epoch = store.epoch();
352 let found = store.lookup(key).await.map(Arc::new);
353 if found.is_some() {
354 stream.borrow_mut().met.noted(key, &found, epoch, store);
355 }
356 local.noted(key, &found, epoch, store);
357 }
358 };
359 let (graph, target) = (&prospect.graph, &prospect.target);
360 let config = &prospect.render;
361 let mut shelled = shelled(graph, target, &mut prospect.terms, config, walk).await?;
362 let generation = prospect.generation;
363 loop {
364 for (key, samples) in &fetched {
365 shelled.table.took(*key, samples);
366 }
367 let wants = stream.borrow().wanting(generation, &shelled.table, &unread);
368 let Some(wants) = wants else {
369 break;
370 };
371 if wants.is_empty() {
372 let answer = prospect.answer;
373 stream.borrow_mut().apply(prospect, shelled, &fetched);
374 return Ok(answer);
375 }
376 for (stored, over) in wants {
377 let again = asked
378 .iter()
379 .any(|(key, e)| *key == stored.key && !e.intersect(over).is_empty());
380 asked.push((stored.key, over));
381 let read = match again {
382 false => store.read(&stored, over).await,
383 true => None,
384 };
385 match read {
386 Some(samples) => fetched.push((stored.key, samples)),
387 None => {
388 unread.insert(stored.key);
389 }
390 }
391 }
392 }
393 }
394}
395
396fn ahead(at: i64, rate: u32) -> Extent {
397 Extent::new(at, at.saturating_add(i64::from(rate)))
398}
399
400fn blocked(config: &StreamConfig) -> Result<&StreamConfig, EngineError> {
401 match config.block {
402 0 => Err(refusal("a block of no samples".to_string())),
403 _ => Ok(config),
404 }
405}
406
407fn wrapped(graph: &Graph, target: &Expr, terms: &Terms) -> Result<Graph, EngineError> {
410 let mut wrapped = graph.clone();
411 if !terms.is_empty() && graph.defines(NOTES) {
412 return Err(EngineError::refused(Diagnostic {
413 code: "engine.no_stream".to_string(),
414 message: format!(
415 "a term is added to `@{NOTES}`, and this composition defines its own `{NOTES}`"
416 ),
417 location: Located::at(NOTES, None),
418 help: format!("rename the composition's `{NOTES}`, or play it without adding terms"),
419 }));
420 }
421 let own = terms.is_empty() && graph.defines(NOTES);
422 let sum = (!own).then(|| (NOTES, terms.sum()));
423 for (name, body) in std::iter::once((STREAMED, target.clone())).chain(sum) {
424 if !wrapped.define(name, body) {
425 return Err(refusal(format!(
426 "this composition already has a node named `{name}`"
427 )));
428 }
429 }
430 Ok(wrapped)
431}
432
433async fn shelled(
435 graph: &Graph,
436 target: &Expr,
437 terms: &mut Terms,
438 config: &RenderConfig,
439 walk: impl AsyncFnOnce(&mut Frontier<'_>),
440) -> Result<Shelled, EngineError> {
441 let wrapped = wrapped(graph, target, terms)?;
442 let instances = instantiate::instantiate(&wrapped, STREAMED, config.rate)?;
443 let root = instances.instance_of(STREAMED)?;
444 let order = schedule::schedule_from(&instances, std::slice::from_ref(&root))?;
445 let keys = through::keys(&wrapped, &instances, &order, config);
446 let mut found = Frontier::from((&instances, &order), &keys, &root, (config, true));
447 if !terms.is_empty() {
448 found.unstored(reading(&order, &instances.instance_of(NOTES)?));
449 }
450 walk(&mut found).await;
451 let mut tys = typing::infer_over(&instances, &order.within(&found.visited), &BTreeMap::new())?;
452 let id = tys
453 .id(&root)
454 .ok_or_else(|| EngineError::UnknownNode(root.clone()))?;
455 terms.typed(&instances, &mut tys);
456 let prefixes: BTreeMap<NodeId, Arc<Stored>> = std::mem::take(&mut found.stored)
457 .into_iter()
458 .map(|(path, stored)| (tys.id(&path).expect("a walked node is typed"), stored))
459 .collect();
460 let schedule = schedule::plan(&tys, id, &[]);
461 let shell = Render::shell(tys, id, config.clone(), schedule);
462 let range = range_of(&shell, Ends::Pulled)?;
463 let table = Table::prefixed(&shell.tys, shell.root, &shell.config.profile, &prefixes)?;
464 let hits = found.lookups.into_iter();
465 let hits = hits.filter(|l| l.outcome == Outcome::Hit).collect();
466 Ok(Shelled {
467 shell,
468 range,
469 table,
470 hits,
471 })
472}
473
474fn reading(order: &schedule::Order, path: &str) -> BTreeSet<String> {
476 let mut out = BTreeSet::from([path.to_string()]);
477 loop {
478 let more: Vec<String> = order
479 .groups
480 .iter()
481 .flatten()
482 .filter(|node| !out.contains(*node))
483 .filter(|node| order.deps(node).iter().any(|read| out.contains(read)))
484 .cloned()
485 .collect();
486 if more.is_empty() {
487 return out;
488 }
489 out.extend(more);
490 }
491}
492
493fn refusal(what: String) -> EngineError {
494 EngineError::refused(Diagnostic {
495 code: "engine.no_stream".to_string(),
496 message: format!("this target opens no stream: {what}"),
497 location: Located::at(STREAMED, None),
498 help: "stream an expression over the nodes the composition defines".to_string(),
499 })
500}