1mod live;
4mod node;
5
6use std::collections::BTreeMap;
7
8use sva_ast::Graph;
9use sva_samples::Tape;
10use sva_samples::physics::chaigne_askenfelt::landing_step;
11
12use super::silent::{Kept, Silent, bound_from, not_silent_by};
13use super::{Render, RenderConfig, prepared};
14use crate::error::{Diagnostic, EngineError, Located};
15use crate::flops::Work;
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: Option<f64>,
44 kept: Kept,
45 end: Option<usize>,
46 work: Work,
47}
48
49pub struct Block<'a> {
51 tape: &'a Tape,
52 start: usize,
53}
54
55impl Block<'_> {
56 pub fn start(&self) -> usize {
57 self.start
58 }
59
60 pub fn len(&self) -> usize {
61 self.tape.end() - self.start
62 }
63
64 pub fn is_empty(&self) -> bool {
65 self.len() == 0
66 }
67
68 pub fn width(&self) -> usize {
69 self.tape.width()
70 }
71
72 pub fn plane(&self, c: usize) -> &[f64] {
73 self.tape.since(c, self.start)
74 }
75}
76
77impl Stream {
78 pub fn open(
81 graph: &Graph,
82 target: &str,
83 bindings: &[(String, f64)],
84 config: StreamConfig,
85 ) -> Result<Stream, EngineError> {
86 Stream::opened_at(graph, target, bindings, config, 0)
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 proofs: 0,
122 };
123 let nodes = node::built(&shell, config.block)?;
124 let kept = Kept::new(shell.tys.clone(), shell.config.clone());
125 let root = nodes
126 .iter()
127 .position(|n| n.id == shell.root)
128 .ok_or_else(|| EngineError::UnknownNode(shell.tys.name(shell.root).to_string()))?;
129 Ok(Stream {
130 graph: graph.clone(),
131 target: target.to_string(),
132 bindings: bindings.to_vec(),
133 config,
134 shell,
135 nodes,
136 root,
137 at: 0,
138 silent,
139 bound: None,
140 kept,
141 end: None,
142 work: Work {
143 waves: Some(0),
144 ..Work::default()
145 },
146 })
147 }
148
149 fn settle(&mut self) -> Result<(), EngineError> {
151 let (Some((silent, _)), None) = (self.silent, self.end) else {
152 return Ok(());
153 };
154 let live = live::View::of(&self.nodes, self.at);
155 self.work.proofs += 1;
156 let bound = bound_from(&self.kept, self.shell.root, silent, &live, self.at)?;
157 self.bound = Some(bound);
158 if bound < silent.threshold() {
159 self.end = Some(self.at.max(1));
160 }
161 Ok(())
162 }
163
164 fn in_time(&self) -> Result<(), EngineError> {
166 match (self.silent, self.end) {
167 (Some((silent, limit)), None) if self.at >= limit => {
168 let max_secs = limit as f64 / f64::from(self.config.rate);
169 let (tys, root) = (&self.shell.tys, self.shell.root);
170 Err(not_silent_by(
171 tys,
172 root,
173 self.bound,
174 Silent { max_secs, ..silent },
175 ))
176 }
177 _ => Ok(()),
178 }
179 }
180
181 pub fn next_block(&mut self) -> Result<Option<Block<'_>>, EngineError> {
183 self.in_time()?;
184 let from = self.at;
185 let to = match self.end {
186 Some(end) if from >= end => return Ok(None),
187 Some(end) => end.min(from + self.config.block),
188 None => from + self.config.block,
189 };
190 for at in 0..self.nodes.len() {
191 let (done, rest) = self.nodes.split_at_mut(at);
192 rest[0].run(&self.shell, done, from, to)?;
193 let (priced, waves) = rest[0].work(from, to);
194 self.work.priced_flops += priced;
195 self.work.waves = self.work.waves.zip(waves).map(|(held, more)| held + more);
196 }
197 self.at = to;
198 self.work.samples += (to - from) as u64;
199 self.settle()?;
200 Ok(Some(Block {
201 tape: &self.nodes[self.root].tape,
202 start: from,
203 }))
204 }
205
206 pub fn position(&self) -> usize {
207 self.at
208 }
209
210 pub fn work(&self) -> Work {
212 self.work
213 }
214
215 pub fn end(&self) -> Option<usize> {
217 self.end
218 }
219
220 pub fn width(&self) -> usize {
221 self.nodes[self.root].width
222 }
223
224 pub fn config(&self) -> StreamConfig {
225 self.config
226 }
227
228 pub fn checkpoint(&self) -> Checkpoint {
229 Checkpoint {
230 at: self.at,
231 target: self.target.clone(),
232 bindings: self.bindings.clone(),
233 rate: self.config.rate,
234 block: self.config.block,
235 nodes: self.nodes.iter().map(Streamed::held).collect(),
236 }
237 }
238
239 pub fn resume(
243 &self,
244 checkpoint: &Checkpoint,
245 bindings: &[(String, f64)],
246 silent: Option<Silent>,
247 ) -> Result<Stream, EngineError> {
248 let (rate, block) = (self.config.rate, self.config.block);
249 if checkpoint.target != self.target || (checkpoint.rate, checkpoint.block) != (rate, block)
250 {
251 return Err(mismatch(&self.target, "another stream"));
252 }
253 causal(checkpoint, bindings, rate, &self.target)?;
254 let config = StreamConfig {
255 rate,
256 block,
257 silent,
258 };
259 let mut resumed =
260 Stream::opened_at(&self.graph, &self.target, bindings, config, checkpoint.at)?;
261 if resumed.nodes.len() != checkpoint.nodes.len() {
262 return Err(mismatch(&self.target, "a graph of another shape"));
263 }
264 for (node, held) in resumed.nodes.iter_mut().zip(&checkpoint.nodes) {
265 node.resume(&resumed.shell, held, checkpoint.at)?;
266 }
267 resumed.at = checkpoint.at;
268 resumed.settle()?;
269 Ok(resumed)
270 }
271}
272
273#[derive(Clone)]
275pub struct Checkpoint {
276 at: usize,
277 target: String,
278 bindings: Vec<(String, f64)>,
279 rate: u32,
280 block: usize,
281 nodes: Vec<Option<node::NodeState>>,
282}
283
284impl Checkpoint {
285 pub fn position(&self) -> usize {
286 self.at
287 }
288}
289
290fn causal(
292 checkpoint: &Checkpoint,
293 bindings: &[(String, f64)],
294 rate: u32,
295 target: &str,
296) -> Result<(), EngineError> {
297 let value =
298 |set: &[(String, f64)], name: &str| set.iter().find(|(n, _)| n == name).map(|(_, v)| *v);
299 let names = checkpoint.bindings.iter().chain(bindings).map(|(n, _)| n);
300 for name in names {
301 let (was, now) = (value(&checkpoint.bindings, name), value(bindings, name));
302 if was == now {
303 continue;
304 }
305 let lands = |v: Option<f64>| {
306 landing_step(v.unwrap_or(f64::INFINITY), f64::from(rate))
307 .is_none_or(|at| at >= checkpoint.at as u64)
308 };
309 if name != RELEASE || !lands(was) || !lands(now) {
310 return Err(EngineError::refused(Diagnostic {
311 code: "engine.binding_not_causal".to_string(),
312 message: format!(
313 "`{name}` moves from {} to {} at sample {}, and a sample before it can \
314 hear that",
315 shown(was),
316 shown(now),
317 checkpoint.at
318 ),
319 location: Located::at(target, None),
320 help: "move only release, to a key-up at or past the checkpoint".to_string(),
321 }));
322 }
323 }
324 Ok(())
325}
326
327fn shown(value: Option<f64>) -> String {
328 value.map_or("unbound".to_string(), |v| v.to_string())
329}
330
331fn mismatch(target: &str, what: &str) -> EngineError {
332 EngineError::refused(Diagnostic {
333 code: "engine.checkpoint_mismatch".to_string(),
334 message: format!("this checkpoint was taken of {what}, not of this `{target}` stream"),
335 location: Located::at(target, None),
336 help: "resume a checkpoint on the stream it was taken of".to_string(),
337 })
338}
339
340fn bound(graph: &Graph, target: &str, bindings: &[(String, f64)]) -> Result<Graph, EngineError> {
342 if !graph.defines(target) {
343 return Err(EngineError::UnknownNode(target.to_string()));
344 }
345 let mut named = String::new();
346 for (name, value) in bindings {
347 let word = name.starts_with(|c: char| c.is_ascii_alphabetic() || c == '_')
348 && name.chars().all(|c| c.is_ascii_alphanumeric() || c == '_');
349 if !word || !value.is_finite() {
350 return Err(refusal(target, format!("`{name}` bound to {value}")));
351 }
352 named.push_str(&format!(", {name}={value}"));
353 }
354 let call = sva_ast::parse_expr(&format!("@{target}(t{named})"))
355 .map_err(|d| refusal(target, d.message))?;
356 let mut wrapped = graph.clone();
357 if !wrapped.define(STREAMED, call) {
358 return Err(refusal(
359 target,
360 format!("this composition already has a node named `{STREAMED}`"),
361 ));
362 }
363 Ok(wrapped)
364}
365
366fn refusal(target: &str, what: String) -> EngineError {
367 EngineError::refused(Diagnostic {
368 code: "engine.no_stream".to_string(),
369 message: format!("`{target}` opens no stream: {what}"),
370 location: Located::at(target, None),
371 help: "name a node the composition defines, and bind each name to a finite number"
372 .to_string(),
373 })
374}