1mod live;
4mod node;
5
6use std::cell::RefCell;
7use std::collections::BTreeMap;
8
9use sva_ast::Graph;
10use sva_samples::Tape;
11use sva_samples::physics::chaigne_askenfelt::landing_step;
12
13use super::silent::{Silent, bound_from, not_silent_by};
14use super::{Render, RenderConfig, prepared};
15use crate::error::{Diagnostic, EngineError, Located};
16use crate::instantiate::RELEASE;
17use crate::schedule;
18use node::Streamed;
19
20pub const STREAMED: &str = "streamed";
21
22#[derive(Clone, Copy, Debug, PartialEq)]
23pub struct StreamConfig {
24 pub rate: u32,
25 pub block: usize,
26 pub silent: Option<Silent>,
28}
29
30pub struct Stream {
33 graph: Graph,
34 target: String,
35 bindings: Vec<(String, f64)>,
36 config: StreamConfig,
37 shell: Render,
38 nodes: Vec<Streamed>,
39 root: usize,
40 at: usize,
41 silent: Option<(Silent, usize)>,
43 bound: f64,
44 heard: RefCell<BTreeMap<(sva_formula::NodeId, u64), sva_samples::Buffer>>,
45 end: Option<usize>,
46}
47
48pub struct Block<'a> {
50 tape: &'a Tape,
51 start: usize,
52}
53
54impl Block<'_> {
55 pub fn start(&self) -> usize {
56 self.start
57 }
58
59 pub fn len(&self) -> usize {
60 self.tape.end() - self.start
61 }
62
63 pub fn is_empty(&self) -> bool {
64 self.len() == 0
65 }
66
67 pub fn width(&self) -> usize {
68 self.tape.width()
69 }
70
71 pub fn plane(&self, c: usize) -> &[f64] {
72 self.tape.since(c, self.start)
73 }
74}
75
76impl Stream {
77 pub fn open(
79 graph: &Graph,
80 target: &str,
81 bindings: &[(String, f64)],
82 config: StreamConfig,
83 ) -> Result<Stream, EngineError> {
84 let mut stream = Stream::opened_at(graph, target, bindings, config, 0)?;
85 stream.settle()?;
86 Ok(stream)
87 }
88
89 fn opened_at(
91 graph: &Graph,
92 target: &str,
93 bindings: &[(String, f64)],
94 config: StreamConfig,
95 from: usize,
96 ) -> Result<Stream, EngineError> {
97 if config.block == 0 {
98 return Err(refusal(target, "a block of no samples".to_string()));
99 }
100 let wrapped = bound(graph, target, bindings)?;
101 let held = prepared(&wrapped, STREAMED)?;
102 let rate = f64::from(config.rate);
103 let silent = config
104 .silent
105 .map(|silent| (silent, from + (silent.max_secs * rate).ceil() as usize));
106 let render_config = RenderConfig::seconds(config.rate, 1.0);
108 let schedule = schedule::plan(&held.tys, &held.order, held.root, &[]);
109 let shell = Render {
110 root: held.root,
111 tys: held.tys,
112 buffers: BTreeMap::new(),
113 frames: BTreeMap::new(),
114 symbolic: BTreeMap::new(),
115 labels: BTreeMap::new(),
116 traces: Vec::new(),
117 config: render_config,
118 schedule,
119 bindings: BTreeMap::new(),
120 cache_stats: None,
121 };
122 let nodes = node::built(&shell, config.block)?;
123 let root = nodes
124 .iter()
125 .position(|n| n.id == shell.root)
126 .ok_or_else(|| EngineError::UnknownNode(shell.tys.name(shell.root).to_string()))?;
127 Ok(Stream {
128 graph: graph.clone(),
129 target: target.to_string(),
130 bindings: bindings.to_vec(),
131 config,
132 shell,
133 nodes,
134 root,
135 at: 0,
136 silent,
137 bound: f64::INFINITY,
138 heard: RefCell::new(BTreeMap::new()),
139 end: None,
140 })
141 }
142
143 fn settle(&mut self) -> Result<(), EngineError> {
146 let (Some((silent, _)), None) = (self.silent, self.end) else {
147 return Ok(());
148 };
149 let live = live::View::of(&self.nodes, self.at);
150 let (tys, root) = (&self.shell.tys, self.shell.root);
151 let config = &self.shell.config;
152 self.bound = bound_from(tys, root, config, silent, &live, self.at, &self.heard)?;
153 if self.bound < silent.threshold() {
154 self.end = Some(self.at.max(1));
155 }
156 Ok(())
157 }
158
159 fn in_time(&self) -> Result<(), EngineError> {
161 match (self.silent, self.end) {
162 (Some((silent, limit)), None) if self.at >= limit => {
163 let max_secs = limit as f64 / f64::from(self.config.rate);
164 let (tys, root) = (&self.shell.tys, self.shell.root);
165 Err(not_silent_by(
166 tys,
167 root,
168 self.bound,
169 Silent { max_secs, ..silent },
170 ))
171 }
172 _ => Ok(()),
173 }
174 }
175
176 pub fn next_block(&mut self) -> Result<Option<Block<'_>>, EngineError> {
178 self.in_time()?;
179 let from = self.at;
180 let to = match self.end {
181 Some(end) if from >= end => return Ok(None),
182 Some(end) => end.min(from + self.config.block),
183 None => from + self.config.block,
184 };
185 for at in 0..self.nodes.len() {
186 let (done, rest) = self.nodes.split_at_mut(at);
187 rest[0].run(&self.shell, done, from, to)?;
188 }
189 self.at = to;
190 self.settle()?;
191 Ok(Some(Block {
192 tape: &self.nodes[self.root].tape,
193 start: from,
194 }))
195 }
196
197 pub fn position(&self) -> usize {
198 self.at
199 }
200
201 pub fn end(&self) -> Option<usize> {
203 self.end
204 }
205
206 pub fn width(&self) -> usize {
207 self.nodes[self.root].width
208 }
209
210 pub fn config(&self) -> StreamConfig {
211 self.config
212 }
213
214 pub fn checkpoint(&self) -> Checkpoint {
215 Checkpoint {
216 at: self.at,
217 target: self.target.clone(),
218 bindings: self.bindings.clone(),
219 rate: self.config.rate,
220 block: self.config.block,
221 nodes: self.nodes.iter().map(Streamed::held).collect(),
222 }
223 }
224
225 pub fn resume(
229 &self,
230 checkpoint: &Checkpoint,
231 bindings: &[(String, f64)],
232 silent: Option<Silent>,
233 ) -> Result<Stream, EngineError> {
234 let (rate, block) = (self.config.rate, self.config.block);
235 if checkpoint.target != self.target || (checkpoint.rate, checkpoint.block) != (rate, block)
236 {
237 return Err(mismatch(&self.target, "another stream"));
238 }
239 causal(checkpoint, bindings, rate, &self.target)?;
240 let config = StreamConfig {
241 rate,
242 block,
243 silent,
244 };
245 let mut resumed =
246 Stream::opened_at(&self.graph, &self.target, bindings, config, checkpoint.at)?;
247 if resumed.nodes.len() != checkpoint.nodes.len() {
248 return Err(mismatch(&self.target, "a graph of another shape"));
249 }
250 for (node, held) in resumed.nodes.iter_mut().zip(&checkpoint.nodes) {
251 node.resume(&resumed.shell, held, checkpoint.at)?;
252 }
253 resumed.at = checkpoint.at;
254 resumed.settle()?;
255 Ok(resumed)
256 }
257}
258
259#[derive(Clone)]
261pub struct Checkpoint {
262 at: usize,
263 target: String,
264 bindings: Vec<(String, f64)>,
265 rate: u32,
266 block: usize,
267 nodes: Vec<Option<node::NodeState>>,
268}
269
270impl Checkpoint {
271 pub fn position(&self) -> usize {
272 self.at
273 }
274}
275
276fn causal(
278 checkpoint: &Checkpoint,
279 bindings: &[(String, f64)],
280 rate: u32,
281 target: &str,
282) -> Result<(), EngineError> {
283 let value =
284 |set: &[(String, f64)], name: &str| set.iter().find(|(n, _)| n == name).map(|(_, v)| *v);
285 let names = checkpoint.bindings.iter().chain(bindings).map(|(n, _)| n);
286 for name in names {
287 let (was, now) = (value(&checkpoint.bindings, name), value(bindings, name));
288 if was == now {
289 continue;
290 }
291 let lands = |v: Option<f64>| {
292 landing_step(v.unwrap_or(f64::INFINITY), f64::from(rate))
293 .is_none_or(|at| at >= checkpoint.at as u64)
294 };
295 if name != RELEASE || !lands(was) || !lands(now) {
296 return Err(EngineError::refused(Diagnostic {
297 code: "engine.binding_not_causal".to_string(),
298 message: format!(
299 "`{name}` moves from {} to {} at sample {}, and a sample before it can \
300 hear that",
301 shown(was),
302 shown(now),
303 checkpoint.at
304 ),
305 location: Located::at(target, None),
306 help: "move only release, to a key-up at or past the checkpoint".to_string(),
307 }));
308 }
309 }
310 Ok(())
311}
312
313fn shown(value: Option<f64>) -> String {
314 value.map_or("unbound".to_string(), |v| v.to_string())
315}
316
317fn mismatch(target: &str, what: &str) -> EngineError {
318 EngineError::refused(Diagnostic {
319 code: "engine.checkpoint_mismatch".to_string(),
320 message: format!("this checkpoint was taken of {what}, not of this `{target}` stream"),
321 location: Located::at(target, None),
322 help: "resume a checkpoint on the stream it was taken of".to_string(),
323 })
324}
325
326fn bound(graph: &Graph, target: &str, bindings: &[(String, f64)]) -> Result<Graph, EngineError> {
328 if !graph.defines(target) {
329 return Err(EngineError::UnknownNode(target.to_string()));
330 }
331 let mut named = String::new();
332 for (name, value) in bindings {
333 let word = name.starts_with(|c: char| c.is_ascii_alphabetic() || c == '_')
334 && name.chars().all(|c| c.is_ascii_alphanumeric() || c == '_');
335 if !word || !value.is_finite() {
336 return Err(refusal(target, format!("`{name}` bound to {value}")));
337 }
338 named.push_str(&format!(", {name}={value}"));
339 }
340 let call = sva_ast::parse_expr(&format!("@{target}(t{named})"))
341 .map_err(|d| refusal(target, d.message))?;
342 let mut wrapped = graph.clone();
343 if !wrapped.define(STREAMED, call) {
344 return Err(refusal(
345 target,
346 format!("this composition already has a node named `{STREAMED}`"),
347 ));
348 }
349 Ok(wrapped)
350}
351
352fn refusal(target: &str, what: String) -> EngineError {
353 EngineError::refused(Diagnostic {
354 code: "engine.no_stream".to_string(),
355 message: format!("`{target}` opens no stream: {what}"),
356 location: Located::at(target, None),
357 help: "name a node the composition defines, and bind each name to a finite number"
358 .to_string(),
359 })
360}