use std::ffi::{c_char, c_void};
use std::sync::Arc;
use std::time::Duration;
use bytes::Bytes;
use tokio::sync::oneshot;
use crate::ffi::OnStatus;
use crate::{Error, Id, NonZeroSlab, Shared, State, ffi};
#[repr(C)]
#[allow(non_camel_case_types)]
#[derive(Clone, Copy, Debug)]
pub enum moq_audio_sample_format {
MOQ_AUDIO_SAMPLE_FORMAT_U8 = 0,
MOQ_AUDIO_SAMPLE_FORMAT_S16 = 1,
MOQ_AUDIO_SAMPLE_FORMAT_S32 = 2,
MOQ_AUDIO_SAMPLE_FORMAT_F32 = 3,
MOQ_AUDIO_SAMPLE_FORMAT_U8_PLANAR = 4,
MOQ_AUDIO_SAMPLE_FORMAT_S16_PLANAR = 5,
MOQ_AUDIO_SAMPLE_FORMAT_S32_PLANAR = 6,
MOQ_AUDIO_SAMPLE_FORMAT_F32_PLANAR = 7,
}
fn audio_format_from_u32(value: u32) -> Result<moq_audio::Format, Error> {
use moq_audio::Format;
Ok(match value {
v if v == moq_audio_sample_format::MOQ_AUDIO_SAMPLE_FORMAT_U8 as u32 => Format::U8,
v if v == moq_audio_sample_format::MOQ_AUDIO_SAMPLE_FORMAT_S16 as u32 => Format::S16,
v if v == moq_audio_sample_format::MOQ_AUDIO_SAMPLE_FORMAT_S32 as u32 => Format::S32,
v if v == moq_audio_sample_format::MOQ_AUDIO_SAMPLE_FORMAT_F32 as u32 => Format::F32,
v if v == moq_audio_sample_format::MOQ_AUDIO_SAMPLE_FORMAT_U8_PLANAR as u32 => Format::U8Planar,
v if v == moq_audio_sample_format::MOQ_AUDIO_SAMPLE_FORMAT_S16_PLANAR as u32 => Format::S16Planar,
v if v == moq_audio_sample_format::MOQ_AUDIO_SAMPLE_FORMAT_S32_PLANAR as u32 => Format::S32Planar,
v if v == moq_audio_sample_format::MOQ_AUDIO_SAMPLE_FORMAT_F32_PLANAR as u32 => Format::F32Planar,
_ => return Err(Error::InvalidCode),
})
}
#[repr(C)]
#[allow(non_camel_case_types)]
pub struct moq_audio_encoder_input {
pub format: u32,
pub sample_rate: u32,
pub channels: u32,
}
#[repr(C)]
#[allow(non_camel_case_types)]
pub struct moq_audio_encoder_output {
pub codec: *const c_char,
pub codec_len: usize,
pub sample_rate: u32,
pub channels: u32,
pub bitrate: u32,
pub frame_duration_us: u32,
}
#[repr(C)]
#[allow(non_camel_case_types)]
pub struct moq_audio_decoder_output {
pub format: u32,
pub sample_rate: u32,
pub channels: u32,
pub max_age_us: u64,
}
#[repr(C)]
#[allow(non_camel_case_types)]
pub struct moq_audio_frame {
pub timestamp_us: u64,
pub data: *const u8,
pub data_size: usize,
}
pub(crate) struct AudioEncoder {
producer: moq_audio::encode::Producer<moq_mux::catalog::hang::Extra>,
reservation: Option<Arc<moq_net::bandwidth::Reservation>>,
}
type AudioProducer = Shared<AudioEncoder>;
#[derive(Default)]
pub struct Audio {
producers: NonZeroSlab<AudioProducer>,
consumer_tasks: NonZeroSlab<Option<AudioTaskEntry>>,
frames: NonZeroSlab<moq_audio::Frame>,
}
struct AudioTaskEntry {
close: Option<oneshot::Sender<()>>,
callback: OnStatus,
}
impl Audio {
pub fn publish(
&mut self,
broadcast: &mut moq_net::broadcast::Producer,
catalog: moq_mux::catalog::Producer<moq_mux::catalog::hang::Extra>,
input: moq_audio::encode::Input,
options: moq_audio::encode::Options,
reserve: bool,
) -> Result<Id, Error> {
let producer = moq_audio::encode::Producer::new(broadcast, catalog, input, &options)?;
let reservation = reserve.then(|| Arc::new(options.bandwidth.reserve(&producer.demand(), producer.bitrate())));
self.producers
.insert(Shared::new(AudioEncoder { producer, reservation }))
}
pub(crate) fn reservation(&self, id: Id) -> Result<Option<Arc<moq_net::bandwidth::Reservation>>, Error> {
Ok(self
.producer(id)?
.lock()
.as_ref()
.ok_or(Error::MediaNotFound)?
.reservation
.clone())
}
pub(crate) fn demand(&self, id: Id) -> Result<moq_net::track::Demand, Error> {
Ok(self
.producer(id)?
.lock()
.as_ref()
.ok_or(Error::MediaNotFound)?
.producer
.demand())
}
pub(crate) fn producer(&self, id: Id) -> Result<AudioProducer, Error> {
self.producers.get(id).cloned().ok_or(Error::MediaNotFound)
}
pub(crate) fn remove(&mut self, id: Id) -> Result<AudioProducer, Error> {
self.producers.remove(id).ok_or(Error::MediaNotFound)
}
pub fn consume(
&mut self,
broadcast: &moq_net::broadcast::Consumer,
catalog: &hang::catalog::AudioConfig,
name: &str,
config: moq_audio::decode::Options,
on_frame: OnStatus,
) -> Result<Id, Error> {
let broadcast = broadcast.clone();
let catalog = catalog.clone();
let name = name.to_string();
let channel = oneshot::channel();
let entry = AudioTaskEntry {
close: Some(channel.0),
callback: on_frame,
};
let id = self.consumer_tasks.insert(Some(entry))?;
tokio::spawn(async move {
let res = async move {
let consumer = moq_audio::decode::Consumer::new(&broadcast, &catalog, name, config).await?;
Self::run(on_frame, consumer, channel.1).await
}
.await;
let entry = State::lock().audio.consumer_tasks.remove(id).flatten();
if let Some(entry) = entry {
entry.callback.call(res);
}
});
Ok(id)
}
async fn run(
callback: OnStatus,
mut consumer: moq_audio::decode::Consumer,
mut close: oneshot::Receiver<()>,
) -> Result<(), Error> {
loop {
let frame = tokio::select! {
biased;
_ = &mut close => return Ok(()),
frame = consumer.read() => match frame {
Ok(Some(frame)) => frame,
Ok(None) => return Ok(()),
Err(moq_audio::Error::Decode(err)) => {
tracing::warn!(%err, "dropping an audio frame");
continue;
}
Err(err) => return Err(err.into()),
},
};
let frame_id = State::lock().audio.frames.insert(frame)?;
callback.call(Ok(frame_id));
}
}
pub fn consume_close(&mut self, id: Id) -> Result<(), Error> {
self.consumer_tasks
.get_mut(id)
.and_then(|entry| entry.as_mut())
.ok_or(Error::TrackNotFound)?
.close
.take()
.ok_or(Error::TrackNotFound)?;
Ok(())
}
pub fn frame_info(&self, id: Id, dst: &mut moq_audio_frame) -> Result<(), Error> {
let frame = self.frames.get(id).ok_or(Error::FrameNotFound)?;
*dst = moq_audio_frame {
timestamp_us: u64::try_from(frame.timestamp.as_micros()).unwrap_or(u64::MAX),
data: frame.data.as_ptr(),
data_size: frame.data.len(),
};
Ok(())
}
pub fn frame_free(&mut self, id: Id) -> Result<(), Error> {
self.frames.remove(id).ok_or(Error::FrameNotFound)?;
Ok(())
}
}
#[unsafe(no_mangle)]
pub unsafe extern "C" fn moq_encode_audio(
broadcast: u32,
name: *const c_char,
name_len: usize,
input: *const moq_audio_encoder_input,
output: *const moq_audio_encoder_output,
bandwidth: u32,
) -> i32 {
ffi::enter(move || {
let broadcast = ffi::parse_id(broadcast)?;
let name = unsafe { ffi::parse_str(name, name_len)? }.to_string();
let raw_input = unsafe { input.as_ref() }.ok_or(Error::InvalidPointer)?;
let raw_output = unsafe { output.as_ref() }.ok_or(Error::InvalidPointer)?;
let codec_str = unsafe { ffi::parse_str(raw_output.codec, raw_output.codec_len)? };
let layout = moq_audio::Layout::from_channels(raw_input.channels)?;
let mut encoder_input = moq_audio::encode::Input::new(raw_input.sample_rate, layout);
encoder_input.format = audio_format_from_u32(raw_input.format)?;
let mut options = moq_audio::encode::Options::default();
options.track = Some(name);
let codec = codec_str
.parse()
.map_err(|_| Error::UnknownFormat(codec_str.to_string()))?;
options.settings = moq_audio::encode::Settings::from_input(codec, &encoder_input);
if let Some(sample_rate) = zeroable(raw_output.sample_rate) {
options.settings.sample_rate = sample_rate;
}
if let Some(channels) = zeroable(raw_output.channels) {
options.settings.layout = moq_audio::Layout::from_channels(channels)?;
}
options.settings.bitrate =
zeroable(raw_output.bitrate).map(|bps| moq_net::bandwidth::Rate::from_bps(bps.into()));
if let Some(micros) = zeroable(raw_output.frame_duration_us) {
options.settings.frame_duration = Duration::from_micros(micros.into());
}
let bandwidth = ffi::parse_id_optional(bandwidth)?;
let mut state = State::lock();
if let Some(id) = bandwidth {
options.bandwidth = state.bandwidth.allocator(id)?;
}
let State { publish, audio, .. } = &mut *state;
let (broadcast_producer, catalog) = publish.pair_mut(broadcast)?;
audio.publish(
broadcast_producer,
catalog.clone(),
encoder_input,
options,
bandwidth.is_some(),
)
})
}
#[unsafe(no_mangle)]
pub extern "C" fn moq_encode_audio_reservation(producer: u32) -> i32 {
ffi::enter(move || {
let producer = ffi::parse_id(producer)?;
let mut state = State::lock();
match state.audio.reservation(producer)? {
Some(reservation) => Ok(i32::from(state.bandwidth.hold(reservation)?)),
None => Ok(0),
}
})
}
#[unsafe(no_mangle)]
pub unsafe extern "C" fn moq_encode_audio_demand(
producer: u32,
on_demand: crate::moq_status_callback,
user_data: *mut c_void,
) -> i32 {
ffi::enter(move || {
let producer = ffi::parse_id(producer)?;
let on_demand = unsafe { OnStatus::new(user_data, on_demand)? };
let mut state = State::lock();
let demand = state.audio.demand(producer)?;
state.publish.demand(demand, on_demand)
})
}
fn zeroable(value: u32) -> Option<u32> {
(value != 0).then_some(value)
}
#[unsafe(no_mangle)]
pub unsafe extern "C" fn moq_encode_audio_frame(producer: u32, frame: *const moq_audio_frame) -> i32 {
ffi::enter(move || {
let producer = ffi::parse_id(producer)?;
let frame = unsafe { frame.as_ref() }.ok_or(Error::InvalidPointer)?;
let data = unsafe { ffi::parse_slice(frame.data, frame.data_size)? };
let timestamp = moq_net::Timestamp::from_micros(frame.timestamp_us).map_err(moq_audio::Error::from)?;
let owned = moq_audio::Frame::new(Bytes::copy_from_slice(data), timestamp);
let producer = State::lock().audio.producer(producer)?;
producer
.lock()
.as_mut()
.ok_or(Error::MediaNotFound)?
.producer
.write(&owned)?;
Ok(())
})
}
#[unsafe(no_mangle)]
pub extern "C" fn moq_encode_audio_finish(producer: u32) -> i32 {
ffi::enter(move || {
let producer = ffi::parse_id(producer)?;
let producer = State::lock().audio.remove(producer)?;
let mut producer = producer.take().ok_or(Error::MediaNotFound)?.producer;
producer.finish()?;
Ok(())
})
}
#[unsafe(no_mangle)]
pub unsafe extern "C" fn moq_decode_audio(
catalog: u32,
index: u32,
output: *const moq_audio_decoder_output,
on_frame: crate::moq_status_callback,
user_data: *mut c_void,
) -> i32 {
ffi::enter(move || {
let catalog = ffi::parse_id(catalog)?;
let raw = unsafe { output.as_ref() }.ok_or(Error::InvalidPointer)?;
let mut config = moq_audio::decode::Options::default();
config.start = moq_audio::decode::Start::Latest;
config.output.format = audio_format_from_u32(raw.format)?;
config.output.sample_rate = zeroable(raw.sample_rate);
config.output.layout = zeroable(raw.channels)
.map(moq_audio::Layout::from_channels)
.transpose()?;
config.max_age = Duration::from_micros(raw.max_age_us);
let on_frame = unsafe { OnStatus::new(user_data, on_frame)? };
let mut state = State::lock();
let (broadcast, audio_cfg, name) = state.consume.audio_rendition(catalog, index as usize)?;
let State { audio, .. } = &mut *state;
audio.consume(&broadcast, &audio_cfg, &name, config, on_frame)
})
}
#[unsafe(no_mangle)]
pub extern "C" fn moq_decode_audio_cancel(consumer: u32) -> i32 {
ffi::enter(move || {
let consumer = ffi::parse_id(consumer)?;
State::lock().audio.consume_close(consumer)
})
}
#[unsafe(no_mangle)]
pub unsafe extern "C" fn moq_decode_audio_frame(id: u32, dst: *mut moq_audio_frame) -> i32 {
ffi::enter(move || {
let id = ffi::parse_id(id)?;
let dst = unsafe { dst.as_mut() }.ok_or(Error::InvalidPointer)?;
State::lock().audio.frame_info(id, dst)
})
}
#[unsafe(no_mangle)]
pub extern "C" fn moq_decode_audio_frame_free(id: u32) -> i32 {
ffi::enter(move || {
let id = ffi::parse_id(id)?;
State::lock().audio.frame_free(id)
})
}