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