use std::cell::RefCell;
use std::collections::{BTreeMap, BTreeSet};
use std::sync::Arc;
use sva_ast::{Expr, Graph};
use sva_formula::{Hash, NodeId};
use sva_samples::{Buffer, Extent};
use super::drive::Driver;
use super::drive::Work;
use super::end::{Ending, Fading, Heard, under};
use super::terms::{Handle, NOTES, Terms, cut, placed};
use super::value_graph::support::Supports;
use super::value_graph::{Past, ValueGraph};
use super::world::{Plan, Root, STREAMED, Walked, Wanted, World};
use super::{Ends, RenderConfig, range_over};
use crate::cache::{Backend, CacheStats, Counters, Memory, Recording, Stored, Tier};
use crate::error::{Diagnostic, EngineError, Located};
use crate::recent::Recent;
#[cfg(test)]
mod rebuilt;
pub const LATEST: usize = 256;
#[derive(Clone, Debug, PartialEq)]
pub struct StreamConfig {
pub block: usize,
pub channels: Option<usize>,
pub render: RenderConfig,
}
pub struct Stream {
config: StreamConfig,
world: World,
driver: Driver,
expr: Expr,
terms: Terms,
width: usize,
heard: BTreeMap<Handle, Heard>,
gain: Gained,
fading: Vec<Fading>,
faded: f64,
ending: Option<Option<i64>>,
treated_as_silent_from_sample: Option<i64>,
generation: u64,
live: bool,
dropped: Recent<String>,
late: usize,
built: Built,
demands: usize,
}
#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
pub struct Counts {
pub dropped: usize,
pub late: usize,
pub terms: usize,
pub built: Built,
pub demands: usize,
pub tier: Counters,
}
#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
pub struct Built {
pub parsed: usize,
pub instances: usize,
pub visited: usize,
pub typed: usize,
pub values: usize,
pub copied: usize,
pub lookups: usize,
}
impl Stream {
pub async fn open<B: Backend>(
graph: &Graph,
target: &Expr,
config: StreamConfig,
tier: &Tier<B>,
) -> Result<Stream, EngineError> {
let render = blocked(&config)?.render.clone();
if graph.defines(STREAMED) {
return Err(refusal(format!(
"this composition already has a node named `{STREAMED}`"
)));
}
let world = World::over(graph, render.rate, !graph.defines(NOTES));
let memory = tier.memory().clone();
let recording = Recording::over(&memory).latest(LATEST);
let value_graph = ValueGraph::new(&render.profile);
let memo = (memory, recording);
let driver = Driver::new(value_graph, Extent::new(0, 0), config.block, &render, memo);
let mut stream = Stream {
world,
driver,
expr: target.clone(),
terms: Terms::default(),
width: 0,
heard: BTreeMap::new(),
gain: Gained::default(),
fading: Vec::new(),
faded: 0.0,
ending: None,
treated_as_silent_from_sample: None,
generation: 0,
live: false,
dropped: Recent::keeping(LATEST),
late: 0,
built: Built::default(),
demands: 0,
config,
};
let opening = Prospect {
target: target.clone(),
terms: Terms::default(),
term: None,
from: None,
answer: Changed::Edited,
landing: None,
parsed: 1,
};
let mut local = Local::default();
let round = tier.begin();
loop {
match stream.attempt(&opening, &mut local, (tier.memory(), round))? {
Attempt::Landed(_) => break,
Attempt::Asks(keys) => local.look(keys, tier, round).await,
Attempt::Reads(wants) => local.read(wants, tier).await,
Attempt::Moved => unreachable!("nothing plays a stream before it opens"),
}
}
let mut needs = stream.needs();
while !needs.is_empty() {
let fetched = tier.fetch(&needs).await;
for (key, parts) in fetched.handed {
stream.driver.value_graph.took(key, &parts);
}
needs = fetched.left;
}
Ok(stream)
}
pub fn graph(&self) -> &Graph {
&self.world.graph
}
fn prospect(&self, change: Change) -> Result<Prospect, Changed> {
let mut landing = None;
let prospect = |terms, term: Option<(Handle, Expr)>, from: Option<Graph>, answer| {
let roots = |(handle, term): &(Handle, Expr)| sva_ast::reads_of(&handle.node(), term);
Prospect {
target: self.expr.clone(),
terms,
from: from.map(|graph| (graph, term.as_ref().map(roots).unwrap_or_default())),
term,
answer,
landing: None,
parsed: 1,
}
};
Ok(match change {
Change::Target(graph, target) => {
let roots = sva_ast::reads_of(STREAMED, &target);
Prospect {
target,
from: Some((graph, roots)),
..prospect(self.terms.clone(), None, None, Changed::Edited)
}
}
Change::Add(graph, term, at) => {
let term = match at {
Placed::Written => term,
Placed::Landing => placed(&term, *landing.insert(self.driver.at)),
};
let (terms, handle) = self.terms.added();
Prospect {
landing,
..prospect(
terms,
Some((handle, term)),
Some(graph),
Changed::Added(handle),
)
}
}
Change::Replace(handle, graph, term, at) => {
let landed = self.terms.landed(handle).ok_or(Changed::Held(false))?;
let term = match at {
Placed::Written => term,
Placed::Landing => placed(&term, landed),
};
let terms = self.terms.clone();
prospect(
terms,
Some((handle, term)),
Some(graph),
Changed::Held(true),
)
}
Change::Remove(handle) => {
let at = self.driver.at;
let terms = self.terms.removed(handle).ok_or(Changed::Held(false))?;
let held = self.world.graph.expr(&handle.node());
let term = cut(held.expect("a held term's node"), at);
Prospect {
parsed: 0,
..prospect(terms, Some((handle, term)), None, Changed::Held(true))
}
}
})
}
fn attempt(
&mut self,
prospect: &Prospect,
local: &mut Local,
(memory, round): (&Memory, u64),
) -> Result<Attempt, EngineError> {
if prospect.landing.is_some_and(|at| at != self.driver.at) {
return Ok(Attempt::Moved);
}
let wanted = Wanted {
root: Root::Streamed(&prospect.target),
terms: &prospect.terms,
term: prospect.term.as_ref().map(|(h, e)| (*h, e)),
from: prospect.from.as_ref().map(|(g, roots)| (g, roots.clone())),
whole: false,
};
let seen = &mut self.driver.recording;
let mut found = |node: &str, key: Hash| memory.answered((key, round), (node, &mut *seen));
let walked = self.world.plan(&wanted, &self.config.render, &mut found)?;
let mut plan = match walked {
Walked::Asks(keys) => return Ok(Attempt::Asks(keys)),
Walked::Planned(plan) => plan,
};
let built = self.built(&mut plan);
let (root, range, treated_as_silent_from_sample) = match built {
Ok(held) => held,
Err(e) => {
let freed = self.world.abort();
self.driver.value_graph.abort(&freed);
return Err(e);
}
};
let from = self.driver.at.max(range.start);
let future = Extent::new(from, self.last(range.end).max(from));
let value_graph = &mut self.driver.value_graph;
let short = value_graph.short((root, future, Past::Stored));
let opened = short
.into_iter()
.try_for_each(|at| value_graph.read_on(&self.world.typing, at));
if let Err(e) = opened {
let freed = self.world.abort();
self.driver.value_graph.abort(&freed);
return Err(e);
}
let asked_range = ahead(from, self.config.render.rate);
let wants = self.driver.value_graph.needs_made(root, asked_range);
let wants: Vec<(Hash, Extent)> = wants
.into_iter()
.filter(|(key, over)| !local.holds(*key, *over))
.collect();
if !wants.is_empty() {
let freed = self.world.abort();
self.driver.value_graph.abort(&freed);
return Ok(Attempt::Reads(wants));
}
self.treated_as_silent_from_sample = treated_as_silent_from_sample;
Ok(Attempt::Landed(self.land(
prospect,
(plan, root, range),
local,
)))
}
fn built(&mut self, plan: &mut Plan) -> Result<(usize, Extent, Option<i64>), EngineError> {
let typing = &mut self.world.typing;
let id = typing
.id(STREAMED)
.ok_or_else(|| EngineError::UnknownNode(STREAMED.to_string()))?;
let hits: BTreeMap<NodeId, Arc<Stored>> = plan
.stored
.iter()
.filter_map(|(path, stored)| Some((typing.id(path)?, Arc::clone(stored))))
.collect();
let value_graph = &mut self.driver.value_graph;
let root = value_graph.grow(typing, id, &hits)?;
let plays = value_graph.values[root].width;
let mut render = self.config.render.clone();
match self.width {
0 if self.config.channels == Some(0) => {
return Err(refusal("a stream of no channels".to_string()));
}
0 => widens(plays, self.config.channels.unwrap_or(plays))?,
width => {
widens(plays, width)?;
render.range.start = Some(self.driver.start);
}
}
let supports = Supports::over(typing, Some(&value_graph.supports));
let ending = Ending::new(typing, &render.profile, &supports);
let end = match (render.range.end, self.live) {
(None, false) => ending.of(id),
_ => ending.exact(id),
};
let range = range_over((&render, STREAMED), end.support, Ends::Pulled)?;
Ok((root, range, end.treated_as_silent_from_sample))
}
fn land(
&mut self,
prospect: &Prospect,
(mut plan, root, range): (Plan, usize, Extent),
local: &mut Local,
) -> Changed {
let freed = self.world.commit(std::mem::take(&mut plan.found));
let now = self.driver.at;
let carried = self
.driver
.value_graph
.settled(root, &freed, (now, self.live));
self.driver.value_graph.offers(&self.world.typing, range);
for (key, parts) in &local.fetched {
self.driver.value_graph.took(*key, parts);
}
for at in &carried.silent {
self.dropped
.push(self.driver.value_graph.values[*at].name.clone());
}
match self.width {
0 => {
let plays = self.driver.value_graph.values[root].width;
self.width = self.config.channels.unwrap_or(plays);
(self.driver.start, self.driver.at) = (range.start, range.start);
}
_ => self.late += usize::from(self.driver.at > local.issued),
}
let last = self.last(range.end);
self.driver.bound(last);
self.expr = prospect.target.clone();
self.terms = prospect.terms.clone();
if let Changed::Added(handle) = prospect.answer {
self.terms.land(handle, self.driver.at);
}
self.hear();
self.built = Built {
parsed: prospect.parsed + plan.adopted,
instances: plan.named,
visited: plan.visited,
typed: self.world.typing.lowered().len(),
values: self.driver.value_graph.built,
copied: carried.taken,
lookups: local.lookups,
};
self.generation += 1;
self.retire_terms_below_silence_threshold();
prospect.answer
}
fn hear(&mut self) {
let typing = &self.world.typing;
let supports = Supports::over(typing, Some(&self.driver.value_graph.supports));
let ending = Ending::new(typing, &self.config.render.profile, &supports);
let ends = typing.id(STREAMED).zip(typing.id(NOTES));
let key = ends.and_then(|(root, notes)| ending.gain_key(root, notes));
let gain = match ends {
_ if key.is_some() && key == self.gain.of => self.gain.gain,
Some((root, notes)) => ending.gain(root, notes),
None => None,
};
let moved = gain != std::mem::replace(&mut self.gain, Gained { gain, of: key }).gain;
let lowered: BTreeSet<&str> = typing.lowered().iter().map(String::as_str).collect();
for handle in self.terms.handles() {
let node = handle.node();
if !moved && !lowered.contains(node.as_str()) && self.heard.contains_key(&handle) {
continue;
}
let id = typing.id(&node).expect("a sounding term is typed");
self.heard.insert(handle, ending.heard(id, gain));
}
self.ending = None;
}
fn needs(&self) -> Vec<(Hash, Extent)> {
let next = ahead(self.driver.at, self.config.render.rate);
self.driver.value_graph.needs(next)
}
pub fn read(&mut self, at: i64, n: usize) -> Result<Option<Buffer>, EngineError> {
let now = self.driver.at;
if at < now {
return Err(refused(
"engine.stream_behind",
format!("sample {at} is before sample {now}, where the stream stands"),
"read from the stream's position or later",
));
}
if n == 0 {
return Err(refused(
"engine.empty_read",
format!("a read of no samples at sample {at}"),
"read one sample or more",
));
}
match self.live {
true if at > now => {
for silenced in self.driver.skip(at)? {
self.dropped
.push(self.driver.value_graph.values[silenced].name.clone());
}
}
_ => {
let block = self.config.block;
while self.driver.at < at {
let step = block.min((at - self.driver.at) as usize);
self.reads_on(step)?;
if !self.driver.pulled(step)? {
return Ok(None);
}
self.retire_terms_below_silence_threshold();
}
}
}
self.reads_on(n)?;
let block = self.driver.read(n)?;
self.retire_terms_below_silence_threshold();
Ok(block.map(|b| widened(b, self.width)))
}
fn reads_on(&mut self, n: usize) -> Result<(), EngineError> {
let at = self.driver.at;
let range = Extent::new(self.driver.start, self.driver.last());
let asked_range = Extent::new(at, at.saturating_add(n as i64).min(range.end).max(at));
let value_graph = &mut self.driver.value_graph;
if asked_range.is_empty() {
return Ok(());
}
for short in value_graph.short((value_graph.root, asked_range, Past::Held)) {
value_graph.read_on(&self.world.typing, short)?;
let landed = value_graph.landed(short);
let carried = value_graph.settled(value_graph.root, &[], (landed, self.live));
for silent in carried.silent {
self.dropped.push(value_graph.values[silent].name.clone());
}
value_graph.offers(&self.world.typing, range);
}
Ok(())
}
pub fn go_live(&mut self) {
self.live = true;
let last = self.last(self.driver.last());
self.driver.bound(last);
}
fn last(&self, range_end: i64) -> i64 {
match self.live {
true => self.config.render.range.end.unwrap_or(i64::MAX),
false => range_end,
}
}
pub fn dropped(&self) -> Vec<&str> {
self.dropped.iter().map(String::as_str).collect()
}
pub fn counts(&self) -> Counts {
Counts {
dropped: self.dropped.made(),
late: self.late,
terms: self.terms.count(),
built: self.built,
demands: self.demands,
tier: self.driver.recording.since(&self.driver.memory),
}
}
fn retire_terms_below_silence_threshold(&mut self) {
let now = self.driver.at;
let heard = &self.heard;
let ending = *self
.ending
.get_or_insert_with(|| heard.values().map(|h| h.support.end).min());
if ending.is_none_or(|end| end > now) {
return;
}
let (value_graph, tys) = (&self.driver.value_graph, &self.world.typing);
let last = self.driver.last();
let asked = match tys.id(NOTES).and_then(|notes| value_graph.of(notes)) {
Some(notes) if now < last => {
self.demands += 1;
let needs = value_graph.demand(Extent::new(now, last));
needs[notes].hold.iter().next().map(|asked| asked.start)
}
Some(_) => None,
None => Some(i64::MIN),
};
let silence_threshold = self.config.render.profile.silence_threshold_amplitude();
let gain = self.gain.gain.unwrap_or(f64::INFINITY);
if let Some(from) = asked {
let let_go = silence_threshold * 2f64.powi(-30);
let mut faded = self.faded;
self.fading.retain(|f| match f.from(from) {
b if b <= let_go => {
faded = (faded + b) * (1.0 + f64::EPSILON);
false
}
_ => true,
});
self.faded = faded;
}
let mut gone = BTreeSet::new();
for (handle, h) in &self.heard {
let end = h.support.end;
if end > now || asked.is_some_and(|from| end > from) {
continue;
}
let (Some(from), Some(fading)) = (asked, &h.fading) else {
gone.insert(*handle);
continue;
};
let own = fading.from(from);
if own > 0.0 {
let left: f64 = self.fading.iter().map(|f| f.from(from)).sum();
let ops = self.fading.len() as f64 + 3.0;
let left = (self.faded + left + own) * (1.0 + ops * f64::EPSILON);
if !under(gain, left, silence_threshold) {
continue;
}
self.fading.push(fading.clone());
}
gone.insert(*handle);
}
let heard = &self.heard;
let support = |handle: Handle| heard.get(&handle).map(|h| h.support);
let went = self
.terms
.retire(&|handle| gone.contains(&handle), &support);
if !went.is_empty() {
self.generation += 1;
}
for handle in went {
self.heard.remove(&handle);
self.ending = None;
}
}
pub fn evaluated(&self, node: &str) -> Vec<sva_samples::Extent> {
let value_graph = &self.driver.value_graph;
self.world
.typing
.id(node)
.and_then(|id| value_graph.of(id))
.map_or(Vec::new(), |at| value_graph.values[at].evaluated.clone())
}
pub fn cutting_below_silence_threshold(&self) -> sva_samples::CuttingBelowSilenceThreshold {
sva_samples::CuttingBelowSilenceThreshold {
silence_threshold_dbfs: self.config.render.profile.silence_threshold_dbfs,
treated_as_silent_from_sample: self
.treated_as_silent_from_sample
.map(|at| (STREAMED.to_string(), at))
.into_iter()
.collect(),
}
}
pub fn landed(&self, handle: Handle) -> Option<i64> {
self.terms.landed(handle)
}
pub fn position(&self) -> i64 {
self.driver.at
}
pub fn work(&self) -> Work {
self.driver.work
}
pub fn stats(&self) -> CacheStats {
self.driver.recording.stats(&self.driver.memory)
}
pub fn held_bytes(&self) -> usize {
self.driver.value_graph.bytes()
}
pub fn end(&self) -> Option<i64> {
self.driver.end()
}
pub fn width(&self) -> usize {
self.width
}
pub fn config(&self) -> &StreamConfig {
&self.config
}
}
#[derive(Clone, Copy, Default)]
struct Gained {
gain: Option<f64>,
of: Option<Hash>,
}
pub enum Change {
Target(Graph, Expr),
Add(Graph, Expr, Placed),
Replace(Handle, Graph, Expr, Placed),
Remove(Handle),
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub enum Placed {
Written,
Landing,
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub enum Changed {
Edited,
Added(Handle),
Held(bool),
}
struct Prospect {
target: Expr,
terms: Terms,
term: Option<(Handle, Expr)>,
from: Option<(Graph, Vec<String>)>,
answer: Changed,
landing: Option<i64>,
parsed: usize,
}
enum Attempt {
Landed(Changed),
Asks(Vec<Hash>),
Reads(Vec<(Hash, Extent)>),
Moved,
}
#[derive(Default)]
struct Local {
fetched: Vec<(Hash, Vec<Arc<Buffer>>)>,
asked: Vec<(Hash, Extent)>,
unread: BTreeSet<Hash>,
lookups: usize,
issued: i64,
}
impl Local {
async fn look<B: Backend>(&mut self, keys: Vec<Hash>, tier: &Tier<B>, round: u64) {
for key in keys {
self.lookups += 1;
tier.lookup(key, round).await;
}
}
async fn read<B: Backend>(&mut self, wants: Vec<(Hash, Extent)>, tier: &Tier<B>) {
let fetched = tier.fetch(&wants).await;
for (key, over) in wants {
if fetched.left.contains(&(key, over)) {
continue;
}
self.asked.push((key, over));
if !fetched.handed.iter().any(|(held, _)| *held == key) {
self.unread.insert(key);
}
}
self.fetched.extend(fetched.handed);
}
fn holds(&self, key: Hash, over: Extent) -> bool {
self.unread.contains(&key)
|| self
.asked
.iter()
.any(|(k, e)| *k == key && e.intersect(over) == over)
}
}
pub async fn change<E: From<EngineError>, B: Backend>(
stream: &RefCell<Stream>,
mut build: impl FnMut(&Stream) -> Result<Change, E>,
tier: &Tier<B>,
) -> Result<Changed, E> {
let mut local = Local {
issued: stream.borrow().driver.at,
..Local::default()
};
let round = tier.begin();
stream.borrow_mut().driver.recording.begin();
loop {
let change = build(&stream.borrow())?;
let prospect = match stream.borrow().prospect(change) {
Ok(prospect) => prospect,
Err(answer) => return Ok(answer),
};
let generation = stream.borrow().generation;
loop {
if stream.borrow().generation != generation {
break;
}
let attempt =
stream
.borrow_mut()
.attempt(&prospect, &mut local, (tier.memory(), round))?;
match attempt {
Attempt::Landed(answer) => return Ok(answer),
Attempt::Moved => break,
Attempt::Asks(keys) => local.look(keys, tier, round).await,
Attempt::Reads(wants) => local.read(wants, tier).await,
}
}
}
}
pub async fn fetch<B: Backend>(stream: &RefCell<Stream>, tier: &Tier<B>) {
let needs = stream.borrow().needs();
let fetched = tier.fetch(&needs).await;
let mut stream = stream.borrow_mut();
for (key, parts) in fetched.handed {
stream.driver.value_graph.took(key, &parts);
}
}
fn ahead(at: i64, rate: u32) -> Extent {
Extent::new(at, at.saturating_add(i64::from(rate)))
}
fn widens(plays: usize, width: usize) -> Result<(), EngineError> {
match plays == width || plays == 1 {
true => Ok(()),
false => Err(refused(
"engine.stream_width",
format!("this plays {plays} channel(s), and the stream plays {width}"),
"play as many channels as the stream, or one, or open a new stream for it",
)),
}
}
fn blocked(config: &StreamConfig) -> Result<&StreamConfig, EngineError> {
match config.block {
0 => Err(refusal("a block of no samples".to_string())),
_ => Ok(config),
}
}
pub(super) fn refusal(what: String) -> EngineError {
refused(
"engine.no_stream",
format!("this target opens no stream: {what}"),
"stream an expression over the nodes the composition defines",
)
}
fn refused(code: &str, message: String, help: &str) -> EngineError {
EngineError::refused(Diagnostic {
code: code.to_string(),
message,
location: Located::at(STREAMED, None),
help: help.to_string(),
})
}
fn widened(mut block: Buffer, width: usize) -> Buffer {
if block.planes.len() == 1 {
let copies = vec![block.planes[0].clone(); width.saturating_sub(1)];
block.planes.extend(copies);
}
block
}