mod decoding;
pub use decoding::MediaDecoder;
mod mic;
pub use mic::Microphone;
#[cfg(not(test))]
pub(crate) use mic::MicrophoneRender;
use std::error::Error;
use crate::buffer::AudioBuffer;
use crate::control::Controller;
use crate::BufferDepletedError;
use crossbeam_channel::{self, Receiver};
pub trait MediaStream:
Iterator<Item = Result<AudioBuffer, Box<dyn Error + Send + Sync>>> + Send + 'static
{
}
impl<M: Iterator<Item = Result<AudioBuffer, Box<dyn Error + Send + Sync>>> + Send + 'static>
MediaStream for M
{
}
pub struct MediaElement {
input: Receiver<Option<Result<AudioBuffer, Box<dyn Error + Send + Sync>>>>,
buffer: Vec<AudioBuffer>,
buffer_complete: bool,
buffer_index: usize,
controller: Controller,
timestamp: f64,
seeking: Option<f64>,
}
impl MediaElement {
pub fn new<S: MediaStream>(input: S) -> Self {
let (sender, receiver) = crossbeam_channel::unbounded();
let fill_buffer = move || {
let _ = sender.send(None); input.map(Some).for_each(|i| {
let _ = sender.send(i);
});
let _ = sender.send(None); };
std::thread::spawn(fill_buffer);
let ping = receiver.recv().expect("buffer channel disconnected");
assert!(ping.is_none());
Self {
input: receiver,
buffer: vec![],
buffer_complete: false,
buffer_index: 0,
controller: Controller::new(),
timestamp: 0.,
seeking: None,
}
}
pub fn controller(&self) -> &Controller {
&self.controller
}
fn load_next(&mut self) -> Option<Result<AudioBuffer, Box<dyn Error + Send + Sync>>> {
if !self.buffer_complete {
let next = match self.input.try_recv() {
Err(_) => return Some(Err(Box::new(BufferDepletedError {}))),
Ok(v) => v,
};
match next {
Some(Err(e)) => {
self.buffer_complete = true;
return Some(Err(e));
}
Some(Ok(data)) => {
self.buffer.push(data.clone());
self.buffer_index += 1;
self.timestamp += data.duration();
return Some(Ok(data));
}
None => {
self.buffer_complete = true;
return None;
}
}
}
None
}
pub fn seek(&mut self, ts: f64) {
if ts == 0. {
self.timestamp = 0.;
self.buffer_index = 0;
return;
}
self.timestamp = 0.;
for (i, buf) in self.buffer.iter().enumerate() {
self.buffer_index = i;
self.timestamp += buf.duration();
if self.timestamp > ts {
return; }
}
loop {
match self.load_next() {
Some(Ok(buf)) => {
self.timestamp += buf.duration();
if self.timestamp > ts {
return; }
}
Some(Err(e)) if e.is::<BufferDepletedError>() => {
self.seeking = Some(ts);
return;
}
_ => {
self.buffer_index += 1;
return;
}
}
}
}
}
impl Iterator for MediaElement {
type Item = Result<AudioBuffer, Box<dyn Error + Send + Sync>>;
fn next(&mut self) -> Option<Self::Item> {
if let Some(seek) = self.controller().should_seek() {
self.seek(seek);
} else if let Some(seek) = self.seeking.take() {
self.seek(seek);
}
if self.seeking.is_some() {
return Some(Err(Box::new(BufferDepletedError {})));
}
if self.controller.loop_() && self.timestamp > self.controller.loop_end() {
self.seek(self.controller.loop_start());
}
if let Some(data) = self.buffer.get(self.buffer_index) {
self.buffer_index += 1;
self.timestamp += data.duration();
return Some(Ok(data.clone()));
}
match self.load_next() {
Some(Ok(data)) => {
return Some(Ok(data));
}
Some(Err(e)) if e.is::<BufferDepletedError>() => {
return Some(Err(e));
}
_ => (), };
if !self.controller.loop_() || self.buffer.is_empty() {
return None;
}
self.seek(self.controller.loop_start());
self.next()
}
}