use anyhow::{Context, Result};
pub struct AudioEncoder {
opus: moq_mux::codec::opus::Import,
ffmpeg_encoder: ffmpeg_next::encoder::audio::Encoder,
resampler: Option<ffmpeg_next::software::resampling::Context>,
input_buffer: Vec<i16>,
encode_buffer: Vec<i16>,
frame_size: usize,
frame_count: u64,
input_sample_rate: u32,
epoch: Option<u64>,
}
const OPUS_SAMPLE_RATE: u32 = 48000;
const OPUS_FRAME_SAMPLES: usize = 960;
const CHANNELS: u32 = 2;
const OPUS_BITRATE: usize = 64000;
impl AudioEncoder {
pub fn new(
broadcast: moq_net::BroadcastProducer,
catalog: moq_mux::catalog::hang::Producer,
input_sample_rate: u32,
) -> Result<Self> {
let opus = moq_mux::codec::opus::Import::new(
broadcast,
catalog,
moq_mux::codec::opus::Config {
sample_rate: OPUS_SAMPLE_RATE,
channel_count: CHANNELS,
},
)?;
let codec = ffmpeg_next::encoder::find(ffmpeg_next::codec::Id::OPUS).context("Opus encoder not found")?;
let ctx = ffmpeg_next::codec::Context::new_with_codec(codec);
let mut enc = ctx.encoder().audio()?;
enc.set_rate(OPUS_SAMPLE_RATE as i32);
enc.set_format(ffmpeg_next::format::Sample::I16(
ffmpeg_next::format::sample::Type::Packed,
));
enc.set_channel_layout(ffmpeg_next::ChannelLayout::STEREO);
enc.set_time_base(ffmpeg_next::Rational::new(1, OPUS_SAMPLE_RATE as i32));
enc.set_bit_rate(OPUS_BITRATE);
let ffmpeg_encoder = enc.open()?;
let frame_size = ffmpeg_encoder.frame_size() as usize;
let resampler = if input_sample_rate != OPUS_SAMPLE_RATE {
Some(ffmpeg_next::software::resampling::Context::get(
ffmpeg_next::format::Sample::I16(ffmpeg_next::format::sample::Type::Packed),
ffmpeg_next::ChannelLayout::STEREO,
input_sample_rate,
ffmpeg_next::format::Sample::I16(ffmpeg_next::format::sample::Type::Packed),
ffmpeg_next::ChannelLayout::STEREO,
OPUS_SAMPLE_RATE,
)?)
} else {
None
};
Ok(Self {
opus,
ffmpeg_encoder,
resampler,
input_buffer: Vec::new(),
encode_buffer: Vec::new(),
frame_size: if frame_size > 0 { frame_size } else { OPUS_FRAME_SAMPLES },
frame_count: 0,
input_sample_rate,
epoch: None,
})
}
pub fn track(&self) -> &moq_net::TrackProducer {
self.opus.track()
}
pub fn reset_epoch(&mut self) {
self.epoch = None;
self.frame_count = 0;
self.input_buffer.clear();
self.encode_buffer.clear();
if let Some(resampler) = &mut self.resampler {
let mut flushed = ffmpeg_next::frame::Audio::empty();
let _ = resampler.flush(&mut flushed);
}
}
pub fn push_samples(&mut self, samples: &[u8], elapsed: std::time::Duration) -> Result<()> {
let i16_samples: Vec<i16> = samples.iter().map(|&s| ((s as i16) - 128) * 256).collect();
self.input_buffer.extend_from_slice(&i16_samples);
self.resample()?;
let samples_per_frame = self.frame_size * CHANNELS as usize;
let frame_duration_us = self.frame_size as u64 * 1_000_000 / OPUS_SAMPLE_RATE as u64;
if self.epoch.is_none() && self.encode_buffer.len() >= samples_per_frame {
let buffered_us = self.encode_buffer.len() as u64 * 1_000_000 / (OPUS_SAMPLE_RATE as u64 * CHANNELS as u64);
self.epoch = Some((elapsed.as_micros() as u64).saturating_sub(buffered_us));
}
while self.encode_buffer.len() >= samples_per_frame {
let frame_samples: Vec<i16> = self.encode_buffer.drain(..samples_per_frame).collect();
let ts_micros = self.epoch.unwrap() + self.frame_count * frame_duration_us;
self.encode_frame(&frame_samples, ts_micros)?;
}
Ok(())
}
fn resample(&mut self) -> Result<()> {
let Some(resampler) = &mut self.resampler else {
self.encode_buffer.append(&mut self.input_buffer);
return Ok(());
};
if self.input_buffer.is_empty() {
return Ok(());
}
let nb_samples = self.input_buffer.len() / CHANNELS as usize;
let mut frame = ffmpeg_next::frame::Audio::new(
ffmpeg_next::format::Sample::I16(ffmpeg_next::format::sample::Type::Packed),
nb_samples,
ffmpeg_next::ChannelLayout::STEREO,
);
frame.set_rate(self.input_sample_rate);
let data = frame.data_mut(0);
let bytes: &[u8] =
unsafe { std::slice::from_raw_parts(self.input_buffer.as_ptr() as *const u8, self.input_buffer.len() * 2) };
data[..bytes.len()].copy_from_slice(bytes);
self.input_buffer.clear();
let delay = resampler.delay().map(|d| d.input as u64).unwrap_or(0);
let out_samples =
((nb_samples as u64 + delay) * OPUS_SAMPLE_RATE as u64).div_ceil(self.input_sample_rate as u64);
let mut resampled = ffmpeg_next::frame::Audio::new(
ffmpeg_next::format::Sample::I16(ffmpeg_next::format::sample::Type::Packed),
out_samples as usize,
ffmpeg_next::ChannelLayout::STEREO,
);
resampler.run(&frame, &mut resampled)?;
let out_samples = resampled.samples() * CHANNELS as usize;
let out_data = resampled.data(0);
let out_i16: &[i16] = unsafe { std::slice::from_raw_parts(out_data.as_ptr() as *const i16, out_samples) };
self.encode_buffer.extend_from_slice(out_i16);
Ok(())
}
fn encode_frame(&mut self, samples: &[i16], ts_micros: u64) -> Result<()> {
let mut frame = ffmpeg_next::frame::Audio::new(
ffmpeg_next::format::Sample::I16(ffmpeg_next::format::sample::Type::Packed),
self.frame_size,
ffmpeg_next::ChannelLayout::STEREO,
);
frame.set_rate(OPUS_SAMPLE_RATE);
frame.set_pts(Some(self.frame_count as i64 * self.frame_size as i64));
let data = frame.data_mut(0);
let bytes: &[u8] = unsafe { std::slice::from_raw_parts(samples.as_ptr() as *const u8, samples.len() * 2) };
data[..bytes.len()].copy_from_slice(bytes);
self.ffmpeg_encoder.send_frame(&frame)?;
let mut pkt = ffmpeg_next::Packet::empty();
while self.ffmpeg_encoder.receive_packet(&mut pkt).is_ok() {
if let Some(data) = pkt.data() {
let ts = hang::container::Timestamp::from_micros(ts_micros)?;
self.opus.decode(&mut &*data, Some(ts))?;
}
}
self.frame_count += 1;
Ok(())
}
}