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::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
28pub const LATEST: usize = 256;
30
31#[derive(Clone, Debug, PartialEq)]
32pub struct StreamConfig {
33 pub block: usize,
34 pub channels: Option<usize>,
36 pub render: RenderConfig,
37}
38
39pub 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 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 pub late: usize,
64 pub terms: usize,
65 pub built: Built,
67 pub demands: usize,
69 pub tier: Counters,
70}
71
72#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
74pub struct Built {
75 pub parsed: usize,
77 pub instances: usize,
78 pub visited: usize,
79 pub typed: usize,
80 pub values: usize,
81 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 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 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 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 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 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 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 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 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 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 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
531pub enum Change {
533 Target(Graph, Expr),
534 Add(Graph, Expr, Placed),
535 Replace(Handle, Graph, Expr, Placed),
536 Remove(Handle),
538}
539
540#[derive(Clone, Copy, Debug, PartialEq, Eq)]
542pub enum Placed {
543 Written,
544 Landing,
545}
546
547#[derive(Clone, Copy, Debug, PartialEq, Eq)]
549pub enum Changed {
550 Edited,
551 Added(Handle),
552 Held(bool),
553}
554
555struct 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 parsed: usize,
565}
566
567enum Attempt {
569 Landed(Changed),
570 Asks(Vec<Hash>),
571 Reads(Vec<(Hash, Extent)>),
572 Moved,
574}
575
576#[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 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 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
619pub 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
655pub 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}