1use sva_ast::{Expr, Graph};
4
5use super::drive::{Block, Driver};
6use super::table::{Table, edit};
7use super::terms::{Handle, NOTES, Terms};
8use super::{Ends, Render, RenderConfig, prepared, range_of};
9use crate::cache::{Cache, CacheStats, Recording};
10use crate::error::{Diagnostic, EngineError, Located};
11use crate::flops::Work;
12use crate::schedule;
13
14pub const STREAMED: &str = "streamed";
15
16#[derive(Clone, Debug, PartialEq)]
17pub struct StreamConfig {
18 pub block: usize,
19 pub render: RenderConfig,
20}
21
22pub struct Stream {
25 config: StreamConfig,
26 shell: Render,
27 driver: Driver,
28 graph: Graph,
29 expr: Expr,
30 terms: Terms,
31 live: bool,
32 dropped: Vec<String>,
33}
34
35impl Stream {
36 pub fn open(
37 graph: &Graph,
38 target: &Expr,
39 config: StreamConfig,
40 cache: Option<&Cache>,
41 ) -> Result<Stream, EngineError> {
42 if config.block == 0 {
43 return Err(refusal("a block of no samples".to_string()));
44 }
45 let mut terms = Terms::default();
46 let (shell, range) = shelled(graph, target, &mut terms, &config.render)?;
47 let table = Table::build(&shell.tys, shell.root, &[], &shell.config.profile)?;
48 let recording = Recording::over(cache, config.render.cache_policy);
49 let driver = Driver::new(table, range, config.block, &config.render, recording);
50 Ok(Stream {
51 graph: graph.clone(),
52 expr: target.clone(),
53 terms,
54 live: false,
55 dropped: Vec::new(),
56 driver,
57 config,
58 shell,
59 })
60 }
61
62 pub fn edit(&mut self, graph: &Graph, target: &Expr) -> Result<(), EngineError> {
63 self.rebuilt(graph, target.clone(), self.terms.clone())
64 }
65
66 pub fn add(&mut self, graph: &Graph, term: &Expr) -> Result<Handle, EngineError> {
67 let (terms, handle) = self.terms.added(term.clone());
68 self.rebuilt(graph, self.expr.clone(), terms)?;
69 Ok(handle)
70 }
71
72 pub fn replace(
74 &mut self,
75 graph: &Graph,
76 handle: Handle,
77 term: &Expr,
78 ) -> Result<bool, EngineError> {
79 let Some(terms) = self.terms.replaced(handle, term.clone()) else {
80 return Ok(false);
81 };
82 self.rebuilt(graph, self.expr.clone(), terms)?;
83 Ok(true)
84 }
85
86 pub fn remove(&mut self, handle: Handle) -> Result<bool, EngineError> {
88 let at = self.driver.at as f64 / f64::from(self.shell.rate());
89 let Some(terms) = self.terms.removed(handle, at) else {
90 return Ok(false);
91 };
92 let graph = self.graph.clone();
93 self.rebuilt(&graph, self.expr.clone(), terms)?;
94 Ok(true)
95 }
96
97 pub fn exprs(&self) -> impl Iterator<Item = &Expr> {
98 std::iter::once(&self.expr).chain(self.terms.exprs())
99 }
100
101 fn rebuilt(
104 &mut self,
105 graph: &Graph,
106 target: Expr,
107 mut terms: Terms,
108 ) -> Result<(), EngineError> {
109 let mut render = self.config.render.clone();
110 render.range.start = Some(self.driver.start);
111 let (shell, range) = shelled(graph, &target, &mut terms, &render)?;
112 let mut table = Table::build(&shell.tys, shell.root, &[], &shell.config.profile)?;
113 let old = std::mem::replace(&mut self.driver.table, Table::empty());
114 let dropped = edit::carried(&mut table, old, self.driver.at, self.live);
115 self.dropped
116 .extend(dropped.into_iter().map(|at| table.values[at].name.clone()));
117 self.driver.replace(table, range.end);
118 self.shell = shell;
119 self.graph = graph.clone();
120 self.expr = target;
121 self.terms = terms;
122 self.prune();
123 Ok(())
124 }
125
126 pub fn next_block(&mut self) -> Result<Option<Block>, EngineError> {
128 let block = self.driver.next_block()?;
129 self.prune();
130 Ok(block)
131 }
132
133 pub fn go_live(&mut self) {
135 self.live = true;
136 }
137
138 pub fn dropped(&self) -> &[String] {
140 &self.dropped
141 }
142
143 fn prune(&mut self) {
145 let (table, shell) = (&self.driver.table, &self.shell);
146 let future = (self.driver.at < self.driver.last())
147 .then(|| table.demand(sva_samples::Extent::new(self.driver.at, self.driver.last())));
148 let gone = |id| {
149 table.of(id).is_some_and(|at| {
150 !table.values[at].evaluated.is_empty()
151 && future
152 .as_ref()
153 .is_none_or(|needs| needs[at].hold.is_empty())
154 })
155 };
156 self.terms
157 .prune(&gone, &|leaf| crate::refs::identity(&shell.tys, leaf).ok());
158 }
159
160 pub fn evaluated(&self, node: &str) -> Vec<sva_samples::Extent> {
162 let table = &self.driver.table;
163 self.shell
164 .tys
165 .id(node)
166 .and_then(|id| table.of(id))
167 .map_or(Vec::new(), |at| table.values[at].evaluated.clone())
168 }
169
170 pub fn position(&self) -> i64 {
171 self.driver.at
172 }
173
174 pub fn work(&self) -> Work {
175 self.driver.work
176 }
177
178 pub fn stats(&self) -> CacheStats {
179 self.driver.recording.stats()
180 }
181
182 pub fn held_bytes(&self) -> usize {
184 self.driver.table.bytes()
185 }
186
187 pub fn end(&self) -> Option<i64> {
188 self.driver.end()
189 }
190
191 pub fn width(&self) -> usize {
192 self.driver.table.values[self.driver.table.root].width
193 }
194
195 pub fn config(&self) -> &StreamConfig {
196 &self.config
197 }
198}
199
200fn shelled(
203 graph: &Graph,
204 target: &Expr,
205 terms: &mut Terms,
206 config: &RenderConfig,
207) -> Result<(Render, sva_samples::Extent), EngineError> {
208 let mut wrapped = graph.clone();
209 if !terms.is_empty() && graph.defines(NOTES) {
210 return Err(EngineError::refused(Diagnostic {
211 code: "engine.no_stream".to_string(),
212 message: format!(
213 "a term is added to `@{NOTES}`, and this composition defines its own `{NOTES}`"
214 ),
215 location: Located::at(NOTES, None),
216 help: format!("rename the composition's `{NOTES}`, or play it without adding terms"),
217 }));
218 }
219 let own = terms.is_empty() && graph.defines(NOTES);
220 let sum = (!own).then(|| (NOTES, terms.sum()));
221 for (name, body) in std::iter::once((STREAMED, target.clone())).chain(sum) {
222 if !wrapped.define(name, body) {
223 return Err(refusal(format!(
224 "this composition already has a node named `{name}`"
225 )));
226 }
227 }
228 let mut held = prepared(&wrapped, STREAMED, config.rate)?;
229 terms.typed(&held.instances, &mut held.tys);
230 let schedule = schedule::plan(&held.tys, held.root, &[]);
231 let shell = Render::shell(held.tys, held.root, config.clone(), schedule);
232 let range = range_of(&shell, Ends::Pulled)?;
233 Ok((shell, range))
234}
235
236fn refusal(what: String) -> EngineError {
237 EngineError::refused(Diagnostic {
238 code: "engine.no_stream".to_string(),
239 message: format!("this target opens no stream: {what}"),
240 location: Located::at(STREAMED, None),
241 help: "stream an expression over the nodes the composition defines".to_string(),
242 })
243}