use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::Arc;
use cpal::Sample;
use crossbeam_channel::Receiver;
use super::{AudioRenderQuantum, NodeIndex};
use crate::buffer::{AudioBuffer, AudioBufferOptions};
use crate::message::ControlMessage;
use crate::{SampleRate, RENDER_QUANTUM_SIZE};
use super::graph::Graph;
pub(crate) struct RenderThread {
graph: Graph,
sample_rate: SampleRate,
channels: usize,
frames_played: Arc<AtomicU64>,
receiver: Receiver<ControlMessage>,
buffer_offset: Option<(usize, AudioRenderQuantum)>,
}
unsafe impl Send for RenderThread {}
impl RenderThread {
pub fn new(
sample_rate: SampleRate,
channels: usize,
receiver: Receiver<ControlMessage>,
frames_played: Arc<AtomicU64>,
) -> Self {
Self {
graph: Graph::new(),
sample_rate,
channels,
frames_played,
receiver,
buffer_offset: None,
}
}
fn handle_control_messages(&mut self) {
for msg in self.receiver.try_iter() {
use ControlMessage::*;
match msg {
RegisterNode {
id,
node,
inputs,
outputs,
channel_config,
} => {
self.graph
.add_node(NodeIndex(id), node, inputs, outputs, channel_config);
}
ConnectNode {
from,
to,
output,
input,
} => {
self.graph
.add_edge((NodeIndex(from), output), (NodeIndex(to), input));
}
DisconnectNode { from, to } => {
self.graph.remove_edge(NodeIndex(from), NodeIndex(to));
}
DisconnectAll { from } => {
self.graph.remove_edges_from(NodeIndex(from));
}
FreeWhenFinished { id } => {
self.graph.mark_free_when_finished(NodeIndex(id));
}
AudioParamEvent { to, event } => {
to.send(event).expect("Audioparam disappeared unexpectedly")
}
}
}
}
pub fn render_audiobuffer(&mut self, length: usize) -> AudioBuffer {
debug_assert_eq!(length % RENDER_QUANTUM_SIZE, 0);
let options = AudioBufferOptions {
number_of_channels: self.channels,
length: 0,
sample_rate: self.sample_rate,
};
let mut buf = AudioBuffer::new(options);
for _ in 0..length / RENDER_QUANTUM_SIZE {
self.handle_control_messages();
let timestamp =
self.frames_played
.fetch_add(RENDER_QUANTUM_SIZE as u64, Ordering::SeqCst) as f64
/ self.sample_rate.0 as f64;
let rendered = self.graph.render(timestamp, self.sample_rate);
buf.extend_alloc(rendered);
}
buf
}
#[allow(dead_code)]
pub fn render<S: Sample>(&mut self, mut buffer: &mut [S]) {
if let Some((offset, prev_rendered)) = self.buffer_offset.take() {
let leftover_len = (RENDER_QUANTUM_SIZE - offset) * self.channels;
let (first, next) = buffer.split_at_mut(leftover_len.min(buffer.len()));
for i in 0..self.channels {
let output = first.iter_mut().skip(i).step_by(self.channels);
let channel = prev_rendered.channel_data(i)[offset..].iter();
for (sample, input) in output.zip(channel) {
let value = Sample::from::<f32>(input);
*sample = value;
}
}
if next.is_empty() {
self.buffer_offset = Some((offset + first.len() / self.channels, prev_rendered));
return;
}
buffer = next;
}
let chunk_size = RENDER_QUANTUM_SIZE * self.channels as usize;
for data in buffer.chunks_mut(chunk_size) {
self.handle_control_messages();
let timestamp =
self.frames_played
.fetch_add(RENDER_QUANTUM_SIZE as u64, Ordering::SeqCst) as f64
/ self.sample_rate.0 as f64;
let rendered = self.graph.render(timestamp, self.sample_rate);
for i in 0..self.channels {
let output = data.iter_mut().skip(i).step_by(self.channels);
let channel = rendered.channel_data(i).iter();
for (sample, input) in output.zip(channel) {
let value = Sample::from::<f32>(input);
*sample = value;
}
}
if data.len() != chunk_size {
let channel_offset = data.len() / self.channels;
debug_assert!(channel_offset < RENDER_QUANTUM_SIZE);
self.buffer_offset = Some((channel_offset, rendered.clone()));
}
}
}
}