use super::{Input, Options, Producer};
use crate::capture;
use crate::{Error, Format, Frame};
pub async fn publish_capture(
mut broadcast: moq_net::broadcast::Producer,
catalog: moq_mux::catalog::Producer,
capture: capture::Config,
encode: Options,
clock: moq_mux::Clock,
) -> Result<(), Error> {
let (sample_rate, channels) = capture::format(&capture).await?;
let input = Input {
format: Format::F32,
sample_rate,
channels,
};
let mut producer = Producer::new(&mut broadcast, catalog, input, &encode)?;
let track = producer.track().clone();
let result = capture_loop(&mut producer, &track, &capture, &clock).await;
if let Err(err) = producer.finish() {
tracing::debug!(error = %err, "audio track finish after capture ended");
}
result
}
async fn capture_loop(
producer: &mut Producer,
track: &moq_net::track::Producer,
config: &capture::Config,
clock: &moq_mux::Clock,
) -> Result<(), Error> {
loop {
if let Err(err) = track.used().await {
log_track_ended(err);
return Ok(());
}
let mut input = capture::open(config).await?;
loop {
let samples = tokio::select! {
biased;
res = track.unused() => {
if let Err(err) = res {
log_track_ended(err);
return Ok(());
}
break; }
samples = input.read() => samples,
};
let Some(samples) = samples else { break };
if samples.gap {
producer.reset_epoch();
}
producer.write(&frame(samples.data, clock.micros())?)?;
}
drop(input);
producer.reset_epoch();
tracing::info!("no listeners: released audio capture");
}
}
fn log_track_ended(err: moq_net::Error) {
if matches!(err, moq_net::Error::Dropped | moq_net::Error::Closed) {
tracing::debug!("audio track no longer announced; stopping capture");
} else {
tracing::warn!(error = %err, "audio track aborted; stopping capture");
}
}
fn frame(samples: Vec<f32>, timestamp_us: u64) -> Result<Frame, Error> {
let mut bytes = Vec::with_capacity(samples.len() * size_of::<f32>());
for sample in &samples {
bytes.extend_from_slice(&sample.to_le_bytes());
}
Ok(Frame {
timestamp: moq_net::Timestamp::from_micros(timestamp_us)?,
data: bytes.into(),
})
}