use error::Error;
use portaudio::pa;
use portaudio::pa::Sample as PaSample;
use sample::{Sample, Wave};
use settings::{Channels, Settings, Frames, SampleHz};
use std::collections::VecDeque;
use std::marker::PhantomData;
use utils::take_front;
use super::{
BufferFrequency,
CallbackFlags,
CallbackResult,
DeltaTimeSeconds,
MINIMUM_BUFFER_RESERVATION,
SoundStream,
StreamFlags,
StreamParams,
wait_for_stream,
};
pub struct Builder<I, O> {
pub stream_params: SoundStream,
pub input_params: StreamParams<I>,
pub output_params: StreamParams<O>,
}
pub struct BlockingStream<'a, I=Wave, O=Wave>
where
I: Sample + PaSample,
O: Sample + PaSample,
{
input_buffer: VecDeque<I>,
output_buffer: VecDeque<O>,
user_buffer: Vec<O>,
in_channels: Channels,
out_channels: Channels,
sample_hz: SampleHz,
frames: Frames,
last_event: Option<LastEvent>,
stream: pa::Stream<I, O>,
is_closed: bool,
marker: PhantomData<&'a ()>,
}
pub type Callback<I, O> =
Box<FnMut(&[I], Settings, &mut[O], Settings, DeltaTimeSeconds, CallbackFlags) -> CallbackResult>;
pub struct NonBlockingStream<I=Wave, O=Wave>
where
I: Sample + PaSample,
O: Sample + PaSample,
{
stream: pa::Stream<I, O>,
is_closed: bool,
}
#[derive(Debug)]
pub enum Event<'a, I=Wave, O=Wave> where O: 'a {
In(Vec<I>, Settings),
Out(&'a mut [O], Settings),
}
#[derive(Clone, Copy)]
pub enum LastEvent {
In,
Out,
Update,
}
type PaParams = (StreamFlags, pa::StreamParameters, pa::StreamParameters, f64, u32);
impl<I, O> Builder<I, O>
where
I: Sample + PaSample,
O: Sample + PaSample,
{
fn unwrap_params(self) -> Result<PaParams, Error> {
let Builder { stream_params, input_params, output_params } = self;
let SoundStream { maybe_buffer_frequency, maybe_sample_hz, maybe_flags } = stream_params;
let flags = maybe_flags.unwrap_or_else(|| StreamFlags::empty());
let input_params = {
let idx = input_params.idx.unwrap_or_else(|| pa::device::get_default_input());
let info = match pa::device::get_info(idx) {
Ok(info) => info,
Err(err) => return Err(Error::PortAudio(err)),
};
let channels = input_params.channel_count
.map(|n| ::std::cmp::min(n, info.max_input_channels))
.unwrap_or_else(|| ::std::cmp::min(2, info.max_input_channels));
let sample_format = input_params.sample_format();
let suggested_latency = input_params.suggested_latency
.unwrap_or_else(|| info.default_low_input_latency);
pa::StreamParameters {
device: idx,
channel_count: channels,
sample_format: sample_format,
suggested_latency: suggested_latency,
}
};
let output_params = {
let idx = output_params.idx.unwrap_or_else(|| pa::device::get_default_output());
let info = match pa::device::get_info(idx) {
Ok(info) => info,
Err(err) => return Err(Error::PortAudio(err)),
};
let channels = output_params.channel_count
.map(|n| ::std::cmp::min(n, info.max_output_channels))
.unwrap_or_else(|| ::std::cmp::min(2, info.max_output_channels));
let sample_format = output_params.sample_format();
let suggested_latency = output_params.suggested_latency
.unwrap_or_else(|| info.default_low_output_latency);
pa::StreamParameters {
device: idx,
channel_count: channels,
sample_format: sample_format,
suggested_latency: suggested_latency,
}
};
let sample_hz = match maybe_sample_hz {
Some(sample_hz) => sample_hz,
None => match pa::device::get_info(input_params.device) {
Ok(info) => info.default_sample_rate,
Err(err) => return Err(Error::PortAudio(err)),
},
};
let frames = match maybe_buffer_frequency {
Some(BufferFrequency::Frames(frames)) => frames as u32,
Some(BufferFrequency::Hz(hz)) => (sample_hz as f32 / hz).round() as u32,
None => 0,
};
Ok((flags, input_params, output_params, sample_hz, frames))
}
#[inline]
pub fn run_callback(self, mut callback: Callback<I, O>)
-> Result<NonBlockingStream<I, O>, Error>
where I: 'static,
O: 'static,
{
try!(pa::initialize().map_err(|err| Error::PortAudio(err)));
let (flags, input_params, output_params, sample_hz, frames) = try!(self.unwrap_params());
let in_channels = input_params.channel_count;
let out_channels = output_params.channel_count;
let mut stream = pa::Stream::new();
let mut maybe_last_time = None;
let f = Box::new(move |input: &[I],
output: &mut[O],
frames: u32,
time_info: &pa::StreamCallbackTimeInfo,
flags: pa::StreamCallbackFlags| -> pa::StreamCallbackResult {
let in_settings = Settings {
sample_hz: sample_hz as u32,
frames: frames as u16,
channels: in_channels as u16,
};
let out_settings = Settings { channels: out_channels as u16, ..in_settings };
let dt = time_info.current_time - maybe_last_time.unwrap_or(time_info.current_time);
maybe_last_time = Some(time_info.current_time);
match callback(input, in_settings, output, out_settings, dt, flags) {
CallbackResult::Continue => pa::StreamCallbackResult::Continue,
CallbackResult::Complete => pa::StreamCallbackResult::Complete,
CallbackResult::Abort => pa::StreamCallbackResult::Abort,
}
});
try!(stream.open(Some(&input_params), Some(&output_params), sample_hz, frames, flags, Some(f))
.map_err(|err| Error::PortAudio(err)));
try!(stream.start().map_err(|err| Error::PortAudio(err)));
Ok(NonBlockingStream { stream: stream, is_closed: false })
}
#[inline]
pub fn run<'a>(self) -> Result<BlockingStream<'a, I, O>, Error>
where I: 'static,
O: 'static,
{
try!(pa::initialize().map_err(|err| Error::PortAudio(err)));
let (flags, input_params, output_params, sample_hz, frames) = try!(self.unwrap_params());
let mut stream = pa::Stream::new();
try!(stream.open(Some(&input_params), Some(&output_params), sample_hz, frames, flags, None)
.map_err(|err| Error::PortAudio(err)));
try!(stream.start().map_err(|err| Error::PortAudio(err)));
let in_channels = input_params.channel_count;
let double_input_buffer_len = (frames as usize * in_channels as usize) * 2;
let input_buffer_len = ::std::cmp::max(double_input_buffer_len, MINIMUM_BUFFER_RESERVATION);
let out_channels = output_params.channel_count;
let double_output_buffer_len = (frames as usize * out_channels as usize) * 2;
let output_buffer_len = ::std::cmp::max(double_output_buffer_len, MINIMUM_BUFFER_RESERVATION);
Ok(BlockingStream {
stream: stream,
input_buffer: VecDeque::with_capacity(input_buffer_len),
output_buffer: VecDeque::with_capacity(output_buffer_len),
user_buffer: Vec::with_capacity(frames as usize * out_channels as usize),
frames: frames as u16,
in_channels: in_channels as u16,
out_channels: out_channels as u16,
sample_hz: sample_hz as u32,
last_event: None,
is_closed: false,
marker: PhantomData,
})
}
}
impl<I, O> NonBlockingStream<I, O>
where
I: Sample + PaSample,
O: Sample + PaSample,
{
pub fn close(&mut self) -> Result<(), Error> {
self.is_closed = true;
try!(self.stream.close().map_err(|err| Error::PortAudio(err)));
try!(pa::terminate().map_err(|err| Error::PortAudio(err)));
Ok(())
}
pub fn is_active(&self) -> Result<bool, Error> {
self.stream.is_active().map_err(|err| Error::PortAudio(err))
}
}
impl<I, O> Drop for NonBlockingStream<I, O>
where
I: Sample + PaSample,
O: Sample + PaSample,
{
fn drop(&mut self) {
if !self.is_closed {
if let Err(err) = self.close() {
println!("An error occurred while closing NonBlockingStream: {}", err);
}
}
}
}
impl<'a, I, O> BlockingStream<'a, I, O>
where
I: Sample + PaSample,
O: Sample + PaSample,
{
pub fn close(&mut self) -> Result<(), Error> {
self.is_closed = true;
try!(self.stream.close().map_err(|err| Error::PortAudio(err)));
try!(pa::terminate().map_err(|err| Error::PortAudio(err)));
Ok(())
}
}
impl<'a, I, O> Drop for BlockingStream<'a, I, O>
where
I: Sample + PaSample,
O: Sample + PaSample,
{
fn drop(&mut self) {
if !self.is_closed {
if let Err(err) = self.close() {
println!("An error occurred while closing BlockingStream: {}", err);
}
}
}
}
impl<'a, I, O> Iterator for BlockingStream<'a, I, O>
where
I: Sample + PaSample + 'a,
O: Sample + PaSample + 'a,
{
type Item = Event<'a, I, O>;
fn next(&mut self) -> Option<Event<'a, I, O>> {
let BlockingStream {
ref mut stream,
ref mut input_buffer,
ref mut output_buffer,
ref mut user_buffer,
ref mut last_event,
ref frames,
ref in_channels,
ref out_channels,
ref sample_hz,
..
} = *self;
let input_settings = Settings { channels: *in_channels, frames: *frames, sample_hz: *sample_hz };
let target_input_buffer_size = input_settings.buffer_size();
let output_settings = Settings { channels: *out_channels, frames: *frames, sample_hz: *sample_hz };
let target_output_buffer_size = output_settings.buffer_size();
if let Some(LastEvent::Out) = *last_event {
if user_buffer.len() > 0 {
output_buffer.extend(user_buffer.iter().map(|&sample| sample));
user_buffer.clear();
}
if input_buffer.len() >= target_input_buffer_size {
let event_buffer = take_front(input_buffer, input_settings.buffer_size());
*last_event = Some(LastEvent::In);
return Some(Event::In(event_buffer, input_settings));
}
}
loop {
use std::error::Error as StdError;
let available_in_frames = match wait_for_stream(|| stream.get_stream_read_available()) {
Ok(frames) => frames,
Err(err) => {
println!("An error occurred while requesting the number of available \
frames for reading from the input stream: {}. BlockingStream will \
now exit the event loop.", StdError::description(&err));
return None;
},
};
if available_in_frames > 0 {
match stream.read(available_in_frames) {
Ok(input_samples) => input_buffer.extend(input_samples.into_iter()),
Err(err) => {
println!("An error occurred while reading from the input stream: {}. \
BlockingStream will now exit the event loop.",
StdError::description(&err));
return None;
},
}
}
let available_out_frames = match wait_for_stream(|| stream.get_stream_write_available()) {
Ok(frames) => frames,
Err(err) => {
println!("An error occurred while requesting the number of available \
frames for writing from the output stream: {}. BlockingStream will \
now exit the event loop.", StdError::description(&err));
return None;
},
};
let output_buffer_frames = (output_buffer.len() / *out_channels as usize) as u32;
if available_out_frames > 0 && output_buffer_frames > 0 {
let (write_buffer, write_frames) = if output_buffer_frames >= available_out_frames {
let out_samples = (available_out_frames * *out_channels as u32) as usize;
let write_buffer = take_front(output_buffer, out_samples);
(write_buffer, available_out_frames)
}
else {
let len = output_buffer.len();
let write_buffer = take_front(output_buffer, len);
(write_buffer, output_buffer_frames)
};
if let Err(err) = stream.write(write_buffer, write_frames) {
println!("An error occurred while writing to the output stream: {}. \
BlockingStream will now exit the event loop.",
StdError::description(&err));
return None
}
}
if output_buffer.len() <= output_buffer.capacity() - target_output_buffer_size {
use std::iter::repeat;
let start = user_buffer.len();
user_buffer.extend(repeat(O::zero()).take(output_settings.buffer_size()));
let slice = unsafe { ::std::mem::transmute(&mut user_buffer[start..]) };
*last_event = Some(LastEvent::Out);
return Some(Event::Out(slice, output_settings));
}
else if input_buffer.len() >= target_input_buffer_size {
let event_buffer = take_front(input_buffer, input_settings.buffer_size());
*last_event = Some(LastEvent::In);
return Some(Event::In(event_buffer, input_settings));
}
*last_event = None;
}
}
}