1use sva_ast::{Expr, Graph};
4use sva_formula::NodeId;
5use sva_samples::Extent;
6
7use super::drive::{Block, Driver};
8use super::frontier::Frontier;
9use super::table::{Table, edit};
10use super::terms::{Handle, NOTES, Terms};
11use super::{Ends, Render, RenderConfig, range_of, through};
12use crate::cache::{Cache, CacheStats, Lookup, Outcome, Recording, Through};
13use crate::error::{Diagnostic, EngineError, Located};
14use crate::flops::Work;
15use crate::instantiate::{self, Instances};
16use crate::schedule;
17use crate::typing::{self, Typing};
18
19pub const STREAMED: &str = "streamed";
20
21#[derive(Clone, Debug, PartialEq)]
22pub struct StreamConfig {
23 pub block: usize,
24 pub render: RenderConfig,
25}
26
27pub struct Stream {
31 config: StreamConfig,
32 shell: Render,
33 driver: Driver,
34 graph: Graph,
35 expr: Expr,
36 terms: Terms,
37 live: bool,
38 dropped: Vec<String>,
39}
40
41struct Shelled {
42 shell: Render,
43 range: Extent,
44 table: Table,
45 hits: Vec<Lookup>,
46}
47
48impl Stream {
49 pub async fn open(
50 graph: &Graph,
51 target: &Expr,
52 config: StreamConfig,
53 cache: Option<&Cache>,
54 store: &impl Through,
55 ) -> Result<Stream, EngineError> {
56 let mut terms = Terms::default();
57 let render = &blocked(&config)?.render;
58 let found = shelled(graph, target, &mut terms, render, None, store).await?;
59 let mut recording = Recording::over(cache, config.render.cache_policy);
60 recording.found(found.hits);
61 let driver = Driver::new(
62 found.table,
63 found.range,
64 config.block,
65 &config.render,
66 recording,
67 );
68 Ok(Stream {
69 graph: graph.clone(),
70 expr: target.clone(),
71 terms,
72 live: false,
73 dropped: Vec::new(),
74 driver,
75 config,
76 shell: found.shell,
77 })
78 }
79
80 pub async fn edit(
81 &mut self,
82 graph: &Graph,
83 target: &Expr,
84 store: &impl Through,
85 ) -> Result<(), EngineError> {
86 let terms = self.terms.clone();
87 self.rebuilt(graph, target.clone(), terms, store).await
88 }
89
90 pub async fn add(
91 &mut self,
92 graph: &Graph,
93 term: &Expr,
94 store: &impl Through,
95 ) -> Result<Handle, EngineError> {
96 let (terms, handle) = self.terms.added(term.clone());
97 self.rebuilt(graph, self.expr.clone(), terms, store).await?;
98 Ok(handle)
99 }
100
101 pub async fn replace(
103 &mut self,
104 graph: &Graph,
105 (handle, term): (Handle, &Expr),
106 store: &impl Through,
107 ) -> Result<bool, EngineError> {
108 let Some(terms) = self.terms.replaced(handle, term.clone()) else {
109 return Ok(false);
110 };
111 self.rebuilt(graph, self.expr.clone(), terms, store).await?;
112 Ok(true)
113 }
114
115 pub async fn remove(
117 &mut self,
118 handle: Handle,
119 store: &impl Through,
120 ) -> Result<bool, EngineError> {
121 let at = self.driver.at as f64 / f64::from(self.shell.rate());
122 let Some(terms) = self.terms.removed(handle, at) else {
123 return Ok(false);
124 };
125 let graph = self.graph.clone();
126 self.rebuilt(&graph, self.expr.clone(), terms, store)
127 .await?;
128 Ok(true)
129 }
130
131 pub fn exprs(&self) -> impl Iterator<Item = &Expr> {
132 std::iter::once(&self.expr).chain(self.terms.exprs())
133 }
134
135 async fn rebuilt(
138 &mut self,
139 graph: &Graph,
140 target: Expr,
141 mut terms: Terms,
142 store: &impl Through,
143 ) -> Result<(), EngineError> {
144 let mut render = self.config.render.clone();
145 render.range.start = Some(self.driver.start);
146 let now = Some(self.driver.at);
147 let found = shelled(graph, &target, &mut terms, &render, now, store).await?;
148 let Shelled {
149 shell,
150 range,
151 mut table,
152 hits,
153 } = found;
154 let old = std::mem::replace(&mut self.driver.table, Table::empty());
155 let dropped = edit::carried(&mut table, old, self.driver.at, self.live);
156 self.dropped
157 .extend(dropped.into_iter().map(|at| table.values[at].name.clone()));
158 self.driver.replace(table, range.end);
159 self.driver.recording.found(hits);
160 self.shell = shell;
161 self.graph = graph.clone();
162 self.expr = target;
163 self.terms = terms;
164 self.prune();
165 Ok(())
166 }
167
168 pub fn next_block(&mut self) -> Result<Option<Block>, EngineError> {
170 let block = self.driver.next_block()?;
171 self.prune();
172 Ok(block)
173 }
174
175 pub fn go_live(&mut self) {
177 self.live = true;
178 }
179
180 pub fn dropped(&self) -> &[String] {
182 &self.dropped
183 }
184
185 fn prune(&mut self) {
187 let (table, shell) = (&self.driver.table, &self.shell);
188 let future = (self.driver.at < self.driver.last())
189 .then(|| table.demand(sva_samples::Extent::new(self.driver.at, self.driver.last())));
190 let gone = |id| {
191 table.of(id).is_some_and(|at| {
192 !table.values[at].evaluated.is_empty()
193 && future
194 .as_ref()
195 .is_none_or(|needs| needs[at].hold.is_empty())
196 })
197 };
198 self.terms
199 .prune(&gone, &|leaf| crate::refs::identity(&shell.tys, leaf).ok());
200 }
201
202 pub fn evaluated(&self, node: &str) -> Vec<sva_samples::Extent> {
204 let table = &self.driver.table;
205 self.shell
206 .tys
207 .id(node)
208 .and_then(|id| table.of(id))
209 .map_or(Vec::new(), |at| table.values[at].evaluated.clone())
210 }
211
212 pub fn position(&self) -> i64 {
213 self.driver.at
214 }
215
216 pub fn work(&self) -> Work {
217 self.driver.work
218 }
219
220 pub fn stats(&self) -> CacheStats {
221 self.driver.recording.stats()
222 }
223
224 pub fn held_bytes(&self) -> usize {
226 self.driver.table.bytes()
227 }
228
229 pub fn end(&self) -> Option<i64> {
230 self.driver.end()
231 }
232
233 pub fn width(&self) -> usize {
234 self.driver.table.values[self.driver.table.root].width
235 }
236
237 pub fn config(&self) -> &StreamConfig {
238 &self.config
239 }
240}
241
242fn blocked(config: &StreamConfig) -> Result<&StreamConfig, EngineError> {
243 match config.block {
244 0 => Err(refusal("a block of no samples".to_string())),
245 _ => Ok(config),
246 }
247}
248
249fn wrapped(graph: &Graph, target: &Expr, terms: &Terms) -> Result<Graph, EngineError> {
252 let mut wrapped = graph.clone();
253 if !terms.is_empty() && graph.defines(NOTES) {
254 return Err(EngineError::refused(Diagnostic {
255 code: "engine.no_stream".to_string(),
256 message: format!(
257 "a term is added to `@{NOTES}`, and this composition defines its own `{NOTES}`"
258 ),
259 location: Located::at(NOTES, None),
260 help: format!("rename the composition's `{NOTES}`, or play it without adding terms"),
261 }));
262 }
263 let own = terms.is_empty() && graph.defines(NOTES);
264 let sum = (!own).then(|| (NOTES, terms.sum()));
265 for (name, body) in std::iter::once((STREAMED, target.clone())).chain(sum) {
266 if !wrapped.define(name, body) {
267 return Err(refusal(format!(
268 "this composition already has a node named `{name}`"
269 )));
270 }
271 }
272 Ok(wrapped)
273}
274
275async fn shelled(
278 graph: &Graph,
279 target: &Expr,
280 terms: &mut Terms,
281 config: &RenderConfig,
282 now: Option<i64>,
283 store: &impl Through,
284) -> Result<Shelled, EngineError> {
285 let wrapped = wrapped(graph, target, terms)?;
286 let instances = instantiate::instantiate(&wrapped, STREAMED, config.rate)?;
287 let root = instances.instance_of(STREAMED)?;
288 let order = schedule::schedule_from(&instances, std::slice::from_ref(&root))?;
289 let keys = through::keys(&wrapped, &instances, &order, config);
290 let mut found = Frontier::from((&instances, &order), &keys, &root, config);
291 loop {
292 found.walk(store).await;
293 let tys = typing::infer_over(&instances, &order.within(&found.visited), &found.stored)?;
294 let id = tys
295 .id(&root)
296 .ok_or_else(|| EngineError::UnknownNode(root.clone()))?;
297 let shelled = typed(&instances, tys, id, terms, config)?;
298 let range = shelled.range;
299 let from = now.map_or(range.start, |now| now.clamp(range.start, range.end));
300 let short = through::short(&shelled.table, Extent::new(from, range.end));
301 if short.is_empty() {
302 let hits = found.lookups.into_iter();
303 let hits = hits
304 .filter(|lookup| lookup.outcome == Outcome::Hit)
305 .collect();
306 return Ok(Shelled { hits, ..shelled });
307 }
308 for path in short {
309 found.reopen(&path);
310 }
311 }
312}
313
314fn typed(
316 instances: &Instances,
317 mut tys: Typing,
318 root: NodeId,
319 terms: &mut Terms,
320 config: &RenderConfig,
321) -> Result<Shelled, EngineError> {
322 terms.typed(instances, &mut tys);
323 let schedule = schedule::plan(&tys, root, &[]);
324 let shell = Render::shell(tys, root, config.clone(), schedule);
325 let range = range_of(&shell, Ends::Pulled)?;
326 let table = Table::build(&shell.tys, shell.root, &[], &shell.config.profile)?;
327 Ok(Shelled {
328 shell,
329 range,
330 table,
331 hits: Vec::new(),
332 })
333}
334
335fn refusal(what: String) -> EngineError {
336 EngineError::refused(Diagnostic {
337 code: "engine.no_stream".to_string(),
338 message: format!("this target opens no stream: {what}"),
339 location: Located::at(STREAMED, None),
340 help: "stream an expression over the nodes the composition defines".to_string(),
341 })
342}