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::Driver;
12use super::drive::Work;
13use super::end::{Ending, Fading, Heard, under};
14use super::terms::{Handle, NOTES, Terms, cut, placed};
15use super::value_graph::support::Supports;
16use super::value_graph::{Past, ValueGraph};
17use super::world::{Plan, Root, STREAMED, Walked, Wanted, World};
18use super::{Ends, RenderConfig, range_over};
19use crate::cache::{Backend, CacheStats, Counters, Memory, Recording, Stored, Tier};
20use crate::error::{Diagnostic, EngineError, Located};
21use crate::recent::Recent;
22
23#[cfg(test)]
24mod rebuilt;
25
26pub const LATEST: usize = 256;
28
29#[derive(Clone, Debug, PartialEq)]
30pub struct StreamConfig {
31 pub block: usize,
32 pub channels: Option<usize>,
34 pub render: RenderConfig,
35}
36
37pub struct Stream {
40 config: StreamConfig,
41 world: World,
42 driver: Driver,
43 expr: Expr,
44 terms: Terms,
45 width: usize,
46 heard: BTreeMap<Handle, Heard>,
48 gain: Gained,
49 fading: Vec<Fading>,
51 faded: f64,
52 ending: Option<Option<i64>>,
54 treated_as_silent_from_sample: Option<i64>,
56 generation: u64,
57 live: bool,
58 dropped: Recent<String>,
59 late: usize,
60 built: Built,
61 demands: usize,
62}
63
64#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
65pub struct Counts {
66 pub dropped: usize,
67 pub late: usize,
69 pub terms: usize,
70 pub built: Built,
72 pub demands: usize,
74 pub tier: Counters,
75}
76
77#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
79pub struct Built {
80 pub parsed: usize,
82 pub instances: usize,
83 pub visited: usize,
84 pub typed: usize,
85 pub values: usize,
86 pub copied: usize,
88 pub lookups: usize,
89}
90
91impl Stream {
92 pub async fn open<B: Backend>(
93 graph: &Graph,
94 target: &Expr,
95 config: StreamConfig,
96 tier: &Tier<B>,
97 ) -> Result<Stream, EngineError> {
98 let render = blocked(&config)?.render.clone();
99 if graph.defines(STREAMED) {
100 return Err(refusal(format!(
101 "this composition already has a node named `{STREAMED}`"
102 )));
103 }
104 let world = World::over(graph, render.rate, !graph.defines(NOTES));
105 let memory = tier.memory().clone();
106 let recording = Recording::over(&memory).latest(LATEST);
107 let value_graph = ValueGraph::new(&render.profile);
108 let memo = (memory, recording);
109 let driver = Driver::new(value_graph, Extent::new(0, 0), config.block, &render, memo);
110 let mut stream = Stream {
111 world,
112 driver,
113 expr: target.clone(),
114 terms: Terms::default(),
115 width: 0,
116 heard: BTreeMap::new(),
117 gain: Gained::default(),
118 fading: Vec::new(),
119 faded: 0.0,
120 ending: None,
121 treated_as_silent_from_sample: None,
122 generation: 0,
123 live: false,
124 dropped: Recent::keeping(LATEST),
125 late: 0,
126 built: Built::default(),
127 demands: 0,
128 config,
129 };
130 let opening = Prospect {
131 target: target.clone(),
132 terms: Terms::default(),
133 term: None,
134 from: None,
135 answer: Changed::Edited,
136 landing: None,
137 parsed: 1,
138 };
139 let mut local = Local::default();
140 let round = tier.begin();
141 loop {
142 match stream.attempt(&opening, &mut local, (tier.memory(), round))? {
143 Attempt::Landed(_) => break,
144 Attempt::Asks(keys) => local.look(keys, tier, round).await,
145 Attempt::Reads(wants) => local.read(wants, tier).await,
146 Attempt::Moved => unreachable!("nothing plays a stream before it opens"),
147 }
148 }
149 let mut needs = stream.needs();
150 while !needs.is_empty() {
151 let fetched = tier.fetch(&needs).await;
152 for (key, parts) in fetched.handed {
153 stream.driver.value_graph.took(key, &parts);
154 }
155 needs = fetched.left;
156 }
157 Ok(stream)
158 }
159
160 pub fn graph(&self) -> &Graph {
161 &self.world.graph
162 }
163
164 fn prospect(&self, change: Change) -> Result<Prospect, Changed> {
165 let mut landing = None;
166 let prospect = |terms, term: Option<(Handle, Expr)>, from: Option<Graph>, answer| {
167 let roots = |(handle, term): &(Handle, Expr)| sva_ast::reads_of(&handle.node(), term);
168 Prospect {
169 target: self.expr.clone(),
170 terms,
171 from: from.map(|graph| (graph, term.as_ref().map(roots).unwrap_or_default())),
172 term,
173 answer,
174 landing: None,
175 parsed: 1,
176 }
177 };
178 Ok(match change {
179 Change::Target(graph, target) => {
180 let roots = sva_ast::reads_of(STREAMED, &target);
181 Prospect {
182 target,
183 from: Some((graph, roots)),
184 ..prospect(self.terms.clone(), None, None, Changed::Edited)
185 }
186 }
187 Change::Add(graph, term, at) => {
188 let term = match at {
189 Placed::Written => term,
190 Placed::Landing => placed(&term, *landing.insert(self.driver.at)),
191 };
192 let (terms, handle) = self.terms.added();
193 Prospect {
194 landing,
195 ..prospect(
196 terms,
197 Some((handle, term)),
198 Some(graph),
199 Changed::Added(handle),
200 )
201 }
202 }
203 Change::Replace(handle, graph, term, at) => {
204 let landed = self.terms.landed(handle).ok_or(Changed::Held(false))?;
205 let term = match at {
206 Placed::Written => term,
207 Placed::Landing => placed(&term, landed),
208 };
209 let terms = self.terms.clone();
210 prospect(
211 terms,
212 Some((handle, term)),
213 Some(graph),
214 Changed::Held(true),
215 )
216 }
217 Change::Remove(handle) => {
218 let at = self.driver.at;
219 let terms = self.terms.removed(handle).ok_or(Changed::Held(false))?;
220 let held = self.world.graph.expr(&handle.node());
221 let term = cut(held.expect("a held term's node"), at);
222 Prospect {
223 parsed: 0,
224 ..prospect(terms, Some((handle, term)), None, Changed::Held(true))
225 }
226 }
227 })
228 }
229
230 fn attempt(
233 &mut self,
234 prospect: &Prospect,
235 local: &mut Local,
236 (memory, round): (&Memory, u64),
237 ) -> Result<Attempt, EngineError> {
238 if prospect.landing.is_some_and(|at| at != self.driver.at) {
239 return Ok(Attempt::Moved);
240 }
241 let wanted = Wanted {
242 root: Root::Streamed(&prospect.target),
243 terms: &prospect.terms,
244 term: prospect.term.as_ref().map(|(h, e)| (*h, e)),
245 from: prospect.from.as_ref().map(|(g, roots)| (g, roots.clone())),
246 whole: false,
247 };
248 let seen = &mut self.driver.recording;
249 let mut found = |node: &str, key: Hash| memory.answered((key, round), (node, &mut *seen));
250 let walked = self.world.plan(&wanted, &self.config.render, &mut found)?;
251 let mut plan = match walked {
252 Walked::Asks(keys) => return Ok(Attempt::Asks(keys)),
253 Walked::Planned(plan) => plan,
254 };
255 let built = self.built(&mut plan);
256 let (root, range, treated_as_silent_from_sample) = match built {
257 Ok(held) => held,
258 Err(e) => {
259 let freed = self.world.abort();
260 self.driver.value_graph.abort(&freed);
261 return Err(e);
262 }
263 };
264 let from = self.driver.at.max(range.start);
265 let future = Extent::new(from, self.last(range.end).max(from));
266 let value_graph = &mut self.driver.value_graph;
267 let short = value_graph.short((root, future, Past::Stored));
268 let opened = short
269 .into_iter()
270 .try_for_each(|at| value_graph.read_on(&self.world.typing, at));
271 if let Err(e) = opened {
272 let freed = self.world.abort();
273 self.driver.value_graph.abort(&freed);
274 return Err(e);
275 }
276 let asked_range = ahead(from, self.config.render.rate);
277 let wants = self.driver.value_graph.needs_made(root, asked_range);
278 let wants: Vec<(Hash, Extent)> = wants
279 .into_iter()
280 .filter(|(key, over)| !local.holds(*key, *over))
281 .collect();
282 if !wants.is_empty() {
283 let freed = self.world.abort();
284 self.driver.value_graph.abort(&freed);
285 return Ok(Attempt::Reads(wants));
286 }
287 self.treated_as_silent_from_sample = treated_as_silent_from_sample;
288 Ok(Attempt::Landed(self.land(
289 prospect,
290 (plan, root, range),
291 local,
292 )))
293 }
294
295 fn built(&mut self, plan: &mut Plan) -> Result<(usize, Extent, Option<i64>), EngineError> {
298 let typing = &mut self.world.typing;
299 let id = typing
300 .id(STREAMED)
301 .ok_or_else(|| EngineError::UnknownNode(STREAMED.to_string()))?;
302 let hits: BTreeMap<NodeId, Arc<Stored>> = plan
303 .stored
304 .iter()
305 .filter_map(|(path, stored)| Some((typing.id(path)?, Arc::clone(stored))))
306 .collect();
307 let value_graph = &mut self.driver.value_graph;
308 let root = value_graph.grow(typing, id, &hits)?;
309 let plays = value_graph.values[root].width;
310 let mut render = self.config.render.clone();
311 match self.width {
312 0 if self.config.channels == Some(0) => {
313 return Err(refusal("a stream of no channels".to_string()));
314 }
315 0 => widens(plays, self.config.channels.unwrap_or(plays))?,
316 width => {
317 widens(plays, width)?;
318 render.range.start = Some(self.driver.start);
319 }
320 }
321 let supports = Supports::over(typing, Some(&value_graph.supports));
322 let ending = Ending::new(typing, &render.profile, &supports);
323 let end = match (render.range.end, self.live) {
324 (None, false) => ending.of(id),
325 _ => ending.exact(id),
326 };
327 let range = range_over((&render, STREAMED), end.support, Ends::Pulled)?;
328 Ok((root, range, end.treated_as_silent_from_sample))
329 }
330
331 fn land(
333 &mut self,
334 prospect: &Prospect,
335 (mut plan, root, range): (Plan, usize, Extent),
336 local: &mut Local,
337 ) -> Changed {
338 let freed = self.world.commit(std::mem::take(&mut plan.found));
339 let now = self.driver.at;
340 let carried = self
341 .driver
342 .value_graph
343 .settled(root, &freed, (now, self.live));
344 self.driver.value_graph.offers(&self.world.typing, range);
345 for (key, parts) in &local.fetched {
346 self.driver.value_graph.took(*key, parts);
347 }
348 for at in &carried.silent {
349 self.dropped
350 .push(self.driver.value_graph.values[*at].name.clone());
351 }
352 match self.width {
353 0 => {
354 let plays = self.driver.value_graph.values[root].width;
355 self.width = self.config.channels.unwrap_or(plays);
356 (self.driver.start, self.driver.at) = (range.start, range.start);
357 }
358 _ => self.late += usize::from(self.driver.at > local.issued),
359 }
360 let last = self.last(range.end);
361 self.driver.bound(last);
362 self.expr = prospect.target.clone();
363 self.terms = prospect.terms.clone();
364 if let Changed::Added(handle) = prospect.answer {
365 self.terms.land(handle, self.driver.at);
366 }
367 self.hear();
368 self.built = Built {
369 parsed: prospect.parsed + plan.adopted,
370 instances: plan.named,
371 visited: plan.visited,
372 typed: self.world.typing.lowered().len(),
373 values: self.driver.value_graph.built,
374 copied: carried.taken,
375 lookups: local.lookups,
376 };
377 self.generation += 1;
378 self.retire_terms_below_silence_threshold();
379 prospect.answer
380 }
381
382 fn hear(&mut self) {
385 let typing = &self.world.typing;
386 let supports = Supports::over(typing, Some(&self.driver.value_graph.supports));
387 let ending = Ending::new(typing, &self.config.render.profile, &supports);
388 let ends = typing.id(STREAMED).zip(typing.id(NOTES));
389 let key = ends.and_then(|(root, notes)| ending.gain_key(root, notes));
390 let gain = match ends {
391 _ if key.is_some() && key == self.gain.of => self.gain.gain,
392 Some((root, notes)) => ending.gain(root, notes),
393 None => None,
394 };
395 let moved = gain != std::mem::replace(&mut self.gain, Gained { gain, of: key }).gain;
396 let lowered: BTreeSet<&str> = typing.lowered().iter().map(String::as_str).collect();
397 for handle in self.terms.handles() {
398 let node = handle.node();
399 if !moved && !lowered.contains(node.as_str()) && self.heard.contains_key(&handle) {
400 continue;
401 }
402 let id = typing.id(&node).expect("a sounding term is typed");
403 self.heard.insert(handle, ending.heard(id, gain));
404 }
405 self.ending = None;
406 }
407
408 fn needs(&self) -> Vec<(Hash, Extent)> {
409 let next = ahead(self.driver.at, self.config.render.rate);
410 self.driver.value_graph.needs(next)
411 }
412
413 pub fn read(&mut self, at: i64, n: usize) -> Result<Option<Buffer>, EngineError> {
416 let now = self.driver.at;
417 if at < now {
418 return Err(refused(
419 "engine.stream_behind",
420 format!("sample {at} is before sample {now}, where the stream stands"),
421 "read from the stream's position or later",
422 ));
423 }
424 if n == 0 {
425 return Err(refused(
426 "engine.empty_read",
427 format!("a read of no samples at sample {at}"),
428 "read one sample or more",
429 ));
430 }
431 match self.live {
432 true if at > now => {
433 for silenced in self.driver.skip(at)? {
434 self.dropped
435 .push(self.driver.value_graph.values[silenced].name.clone());
436 }
437 }
438 _ => {
439 let block = self.config.block;
440 while self.driver.at < at {
441 let step = block.min((at - self.driver.at) as usize);
442 self.reads_on(step)?;
443 if !self.driver.pulled(step)? {
444 return Ok(None);
445 }
446 self.retire_terms_below_silence_threshold();
447 }
448 }
449 }
450 self.reads_on(n)?;
451 let block = self.driver.read(n)?;
452 self.retire_terms_below_silence_threshold();
453 Ok(block.map(|b| widened(b, self.width)))
454 }
455
456 fn reads_on(&mut self, n: usize) -> Result<(), EngineError> {
460 let at = self.driver.at;
461 let range = Extent::new(self.driver.start, self.driver.last());
462 let asked_range = Extent::new(at, at.saturating_add(n as i64).min(range.end).max(at));
463 let value_graph = &mut self.driver.value_graph;
464 if asked_range.is_empty() {
465 return Ok(());
466 }
467 for short in value_graph.short((value_graph.root, asked_range, Past::Held)) {
468 value_graph.read_on(&self.world.typing, short)?;
469 let landed = value_graph.landed(short);
470 let carried = value_graph.settled(value_graph.root, &[], (landed, self.live));
471 for silent in carried.silent {
472 self.dropped.push(value_graph.values[silent].name.clone());
473 }
474 value_graph.offers(&self.world.typing, range);
475 }
476 Ok(())
477 }
478
479 pub fn go_live(&mut self) {
482 self.live = true;
483 let last = self.last(self.driver.last());
484 self.driver.bound(last);
485 }
486
487 fn last(&self, range_end: i64) -> i64 {
488 match self.live {
489 true => self.config.render.range.end.unwrap_or(i64::MAX),
490 false => range_end,
491 }
492 }
493
494 pub fn dropped(&self) -> Vec<&str> {
496 self.dropped.iter().map(String::as_str).collect()
497 }
498
499 pub fn counts(&self) -> Counts {
500 Counts {
501 dropped: self.dropped.made(),
502 late: self.late,
503 terms: self.terms.count(),
504 built: self.built,
505 demands: self.demands,
506 tier: self.driver.recording.since(&self.driver.memory),
507 }
508 }
509
510 fn retire_terms_below_silence_threshold(&mut self) {
514 let now = self.driver.at;
515 let heard = &self.heard;
516 let ending = *self
517 .ending
518 .get_or_insert_with(|| heard.values().map(|h| h.support.end).min());
519 if ending.is_none_or(|end| end > now) {
520 return;
521 }
522 let (value_graph, tys) = (&self.driver.value_graph, &self.world.typing);
523 let last = self.driver.last();
524 let asked = match tys.id(NOTES).and_then(|notes| value_graph.of(notes)) {
525 Some(notes) if now < last => {
526 self.demands += 1;
527 let needs = value_graph.demand(Extent::new(now, last));
528 needs[notes].hold.iter().next().map(|asked| asked.start)
529 }
530 Some(_) => None,
531 None => Some(i64::MIN),
532 };
533 let silence_threshold = self.config.render.profile.silence_threshold_amplitude();
534 let gain = self.gain.gain.unwrap_or(f64::INFINITY);
535 if let Some(from) = asked {
536 let let_go = silence_threshold * 2f64.powi(-30);
538 let mut faded = self.faded;
539 self.fading.retain(|f| match f.from(from) {
540 b if b <= let_go => {
541 faded = (faded + b) * (1.0 + f64::EPSILON);
542 false
543 }
544 _ => true,
545 });
546 self.faded = faded;
547 }
548 let mut gone = BTreeSet::new();
549 for (handle, h) in &self.heard {
550 let end = h.support.end;
551 if end > now || asked.is_some_and(|from| end > from) {
552 continue;
553 }
554 let (Some(from), Some(fading)) = (asked, &h.fading) else {
555 gone.insert(*handle);
556 continue;
557 };
558 let own = fading.from(from);
559 if own > 0.0 {
560 let left: f64 = self.fading.iter().map(|f| f.from(from)).sum();
561 let ops = self.fading.len() as f64 + 3.0;
562 let left = (self.faded + left + own) * (1.0 + ops * f64::EPSILON);
563 if !under(gain, left, silence_threshold) {
564 continue;
565 }
566 self.fading.push(fading.clone());
567 }
568 gone.insert(*handle);
569 }
570 let heard = &self.heard;
571 let support = |handle: Handle| heard.get(&handle).map(|h| h.support);
572 let went = self
573 .terms
574 .retire(&|handle| gone.contains(&handle), &support);
575 if !went.is_empty() {
576 self.generation += 1;
577 }
578 for handle in went {
579 self.heard.remove(&handle);
580 self.ending = None;
581 }
582 }
583
584 pub fn evaluated(&self, node: &str) -> Vec<sva_samples::Extent> {
586 let value_graph = &self.driver.value_graph;
587 self.world
588 .typing
589 .id(node)
590 .and_then(|id| value_graph.of(id))
591 .map_or(Vec::new(), |at| value_graph.values[at].evaluated.clone())
592 }
593
594 pub fn cutting_below_silence_threshold(&self) -> sva_samples::CuttingBelowSilenceThreshold {
595 sva_samples::CuttingBelowSilenceThreshold {
596 silence_threshold_dbfs: self.config.render.profile.silence_threshold_dbfs,
597 treated_as_silent_from_sample: self
598 .treated_as_silent_from_sample
599 .map(|at| (STREAMED.to_string(), at))
600 .into_iter()
601 .collect(),
602 }
603 }
604
605 pub fn landed(&self, handle: Handle) -> Option<i64> {
606 self.terms.landed(handle)
607 }
608
609 pub fn position(&self) -> i64 {
610 self.driver.at
611 }
612
613 pub fn work(&self) -> Work {
614 self.driver.work
615 }
616
617 pub fn stats(&self) -> CacheStats {
618 self.driver.recording.stats(&self.driver.memory)
619 }
620
621 pub fn held_bytes(&self) -> usize {
623 self.driver.value_graph.bytes()
624 }
625
626 pub fn end(&self) -> Option<i64> {
627 self.driver.end()
628 }
629
630 pub fn width(&self) -> usize {
631 self.width
632 }
633
634 pub fn config(&self) -> &StreamConfig {
635 &self.config
636 }
637}
638
639#[derive(Clone, Copy, Default)]
640struct Gained {
641 gain: Option<f64>,
642 of: Option<Hash>,
643}
644
645pub enum Change {
647 Target(Graph, Expr),
648 Add(Graph, Expr, Placed),
649 Replace(Handle, Graph, Expr, Placed),
650 Remove(Handle),
652}
653
654#[derive(Clone, Copy, Debug, PartialEq, Eq)]
656pub enum Placed {
657 Written,
658 Landing,
659}
660
661#[derive(Clone, Copy, Debug, PartialEq, Eq)]
663pub enum Changed {
664 Edited,
665 Added(Handle),
666 Held(bool),
667}
668
669struct Prospect {
671 target: Expr,
672 terms: Terms,
673 term: Option<(Handle, Expr)>,
674 from: Option<(Graph, Vec<String>)>,
675 answer: Changed,
676 landing: Option<i64>,
677 parsed: usize,
679}
680
681enum Attempt {
683 Landed(Changed),
684 Asks(Vec<Hash>),
685 Reads(Vec<(Hash, Extent)>),
686 Moved,
688}
689
690#[derive(Default)]
692struct Local {
693 fetched: Vec<(Hash, Vec<Arc<Buffer>>)>,
694 asked: Vec<(Hash, Extent)>,
695 unread: BTreeSet<Hash>,
696 lookups: usize,
697 issued: i64,
698}
699
700impl Local {
701 async fn look<B: Backend>(&mut self, keys: Vec<Hash>, tier: &Tier<B>, round: u64) {
702 for key in keys {
703 self.lookups += 1;
704 tier.lookup(key, round).await;
705 }
706 }
707
708 async fn read<B: Backend>(&mut self, wants: Vec<(Hash, Extent)>, tier: &Tier<B>) {
710 let fetched = tier.fetch(&wants).await;
711 for (key, over) in wants {
712 if fetched.left.contains(&(key, over)) {
713 continue;
714 }
715 self.asked.push((key, over));
716 if !fetched.handed.iter().any(|(held, _)| *held == key) {
717 self.unread.insert(key);
718 }
719 }
720 self.fetched.extend(fetched.handed);
721 }
722
723 fn holds(&self, key: Hash, over: Extent) -> bool {
725 self.unread.contains(&key)
726 || self
727 .asked
728 .iter()
729 .any(|(k, e)| *k == key && e.intersect(over) == over)
730 }
731}
732
733pub async fn change<E: From<EngineError>, B: Backend>(
737 stream: &RefCell<Stream>,
738 mut build: impl FnMut(&Stream) -> Result<Change, E>,
739 tier: &Tier<B>,
740) -> Result<Changed, E> {
741 let mut local = Local {
742 issued: stream.borrow().driver.at,
743 ..Local::default()
744 };
745 let round = tier.begin();
746 stream.borrow_mut().driver.recording.begin();
747 loop {
748 let change = build(&stream.borrow())?;
749 let prospect = match stream.borrow().prospect(change) {
750 Ok(prospect) => prospect,
751 Err(answer) => return Ok(answer),
752 };
753 let generation = stream.borrow().generation;
754 loop {
755 if stream.borrow().generation != generation {
756 break;
757 }
758 let attempt =
759 stream
760 .borrow_mut()
761 .attempt(&prospect, &mut local, (tier.memory(), round))?;
762 match attempt {
763 Attempt::Landed(answer) => return Ok(answer),
764 Attempt::Moved => break,
765 Attempt::Asks(keys) => local.look(keys, tier, round).await,
766 Attempt::Reads(wants) => local.read(wants, tier).await,
767 }
768 }
769 }
770}
771
772pub async fn fetch<B: Backend>(stream: &RefCell<Stream>, tier: &Tier<B>) {
775 let needs = stream.borrow().needs();
776 let fetched = tier.fetch(&needs).await;
777 let mut stream = stream.borrow_mut();
778 for (key, parts) in fetched.handed {
779 stream.driver.value_graph.took(key, &parts);
780 }
781}
782
783fn ahead(at: i64, rate: u32) -> Extent {
784 Extent::new(at, at.saturating_add(i64::from(rate)))
785}
786
787fn widens(plays: usize, width: usize) -> Result<(), EngineError> {
788 match plays == width || plays == 1 {
789 true => Ok(()),
790 false => Err(refused(
791 "engine.stream_width",
792 format!("this plays {plays} channel(s), and the stream plays {width}"),
793 "play as many channels as the stream, or one, or open a new stream for it",
794 )),
795 }
796}
797
798fn blocked(config: &StreamConfig) -> Result<&StreamConfig, EngineError> {
799 match config.block {
800 0 => Err(refusal("a block of no samples".to_string())),
801 _ => Ok(config),
802 }
803}
804
805pub(super) fn refusal(what: String) -> EngineError {
806 refused(
807 "engine.no_stream",
808 format!("this target opens no stream: {what}"),
809 "stream an expression over the nodes the composition defines",
810 )
811}
812
813fn refused(code: &str, message: String, help: &str) -> EngineError {
814 EngineError::refused(Diagnostic {
815 code: code.to_string(),
816 message,
817 location: Located::at(STREAMED, None),
818 help: help.to_string(),
819 })
820}
821
822fn widened(mut block: Buffer, width: usize) -> Buffer {
824 if block.planes.len() == 1 {
825 let copies = vec![block.planes[0].clone(); width.saturating_sub(1)];
826 block.planes.extend(copies);
827 }
828 block
829}