1use std::time::Duration;
4
5use bytes::Bytes;
6
7use moq_mux::catalog::hang::CatalogExt;
8use moq_mux::container::Frame as MuxFrame;
9use moq_net::Timestamp;
10
11use super::encoder::{Codec, Config, Encoder, Input};
12use crate::resample::Resampler;
13use crate::{Error, Frame};
14
15#[derive(Clone, Debug)]
24#[non_exhaustive]
25pub struct Options {
26 pub track: Option<String>,
30 pub codec: Codec,
32 pub sample_rate: Option<u32>,
35 pub channels: Option<u32>,
38 pub bitrate: Option<u32>,
41 pub fec: bool,
43 pub dtx: bool,
45 pub frame_duration: Duration,
48}
49
50impl Default for Options {
51 fn default() -> Self {
52 Self {
53 track: None,
54 codec: Codec::default(),
55 sample_rate: None,
56 channels: None,
57 bitrate: None,
58 fec: false,
59 dtx: false,
60 frame_duration: Duration::from_millis(20),
61 }
62 }
63}
64
65impl Options {
66 fn config(&self, input: Input) -> Config {
68 Config {
69 input,
70 codec: self.codec,
71 sample_rate: self.sample_rate,
72 channels: self.channels,
73 bitrate: self.bitrate,
74 fec: self.fec,
75 dtx: self.dtx,
76 frame_duration: self.frame_duration,
77 }
78 }
79}
80
81pub struct Producer<E: CatalogExt = ()> {
91 encoder: Encoder,
92 resampler: Option<Resampler>,
93 track: moq_mux::container::Producer<moq_mux::container::legacy::Wire>,
94 rendition: Rendition<E>,
96 pending: Vec<f32>,
97 frames_produced: u64,
99 epoch_us: Option<u64>,
103}
104
105impl<E: CatalogExt> Producer<E> {
106 pub fn new(
109 broadcast: &mut moq_net::broadcast::Producer,
110 catalog: moq_mux::catalog::Producer<E>,
111 input: Input,
112 options: &Options,
113 ) -> Result<Self, Error> {
114 let encoder = Encoder::new(&options.config(input))?;
115 let input = &encoder.config().input;
116
117 let resampler = if input.sample_rate == encoder.codec_rate() {
118 None
119 } else {
120 let chunk_frames =
123 ((input.sample_rate as u128 * encoder.config().frame_duration.as_micros()) / 1_000_000) as usize;
124 Some(Resampler::new(
125 input.sample_rate,
126 encoder.codec_rate(),
127 input.channels,
128 chunk_frames,
129 )?)
130 };
131
132 let track = match &options.track {
133 Some(name) => broadcast.create_track(name.clone(), hang::container::track_info())?,
137 None => moq_mux::import::unique_track(broadcast, &format!(".{}", options.codec))?,
140 };
141 let name = track.name().to_string();
142 let track = catalog.media_producer(track, moq_mux::container::legacy::Wire)?;
143
144 let mut catalog_mut = catalog.clone();
145 let mut config = encoder.catalog();
146 config.timeline = Some(catalog.timeline(&name)?.section());
147 catalog_mut.lock().audio.insert(&name, config)?;
148
149 Ok(Self {
150 encoder,
151 resampler,
152 track,
153 rendition: Rendition { catalog, name },
154 pending: Vec::new(),
155 frames_produced: 0,
156 epoch_us: None,
157 })
158 }
159
160 pub fn track_name(&self) -> &str {
162 &self.rendition.name
163 }
164
165 pub fn track(&self) -> &moq_net::track::Producer {
168 self.track.track()
169 }
170
171 pub fn bitrate(&self) -> u64 {
173 self.encoder.bitrate()
174 }
175
176 pub fn set_bitrate(&mut self, bitrate: u64) -> Result<(), Error> {
178 self.encoder.set_bitrate(bitrate)
179 }
180
181 pub fn reset_epoch(&mut self) {
188 self.epoch_us = None;
189 self.frames_produced = 0;
190 self.pending.clear();
191 }
192
193 pub fn write(&mut self, frame: &Frame) -> Result<(), Error> {
205 let timestamp_us = u64::try_from(frame.timestamp.as_micros())
206 .map_err(|_| Error::Unsupported(format!("frame timestamp {:?} out of range", frame.timestamp)))?;
207 let epoch_us = *self.epoch_us.get_or_insert(timestamp_us);
208
209 let input = &self.encoder.config().input;
210 let (format, channels) = (input.format, input.channels);
211 let pcm = format.as_interleaved_f32(frame.data.as_ref(), channels)?;
212 let pcm: Vec<f32> = match self.resampler.as_mut() {
213 Some(r) => r.process(&pcm)?,
214 None => pcm.into_owned(),
215 };
216
217 self.pending.extend(pcm);
218
219 let frame_samples = self.encoder.frame_size() * self.encoder.codec_channels() as usize;
220 while self.pending.len() >= frame_samples {
221 let chunk: Vec<f32> = self.pending.drain(..frame_samples).collect();
222 let packet = self.encoder.encode(&chunk)?;
223
224 let timestamp = self.timestamp(epoch_us)?;
225 self.frames_produced += self.encoder.frame_size() as u64;
226 self.publish(packet, timestamp)?;
227 }
228
229 Ok(())
230 }
231
232 fn timestamp(&self, epoch_us: u64) -> Result<Timestamp, Error> {
234 let offset_us = (self.frames_produced * 1_000_000) / self.encoder.codec_rate() as u64;
235 Ok(Timestamp::from_micros(epoch_us + offset_us)?)
236 }
237
238 fn publish(&mut self, payload: Bytes, timestamp: Timestamp) -> Result<(), Error> {
239 let mux_frame = MuxFrame {
243 timestamp,
244 payload,
245 keyframe: true,
246 duration: None,
247 };
248 self.track.write(mux_frame)?;
249 self.track.cut(None)?;
252 Ok(())
253 }
254
255 pub fn discontinuity(&mut self) -> Result<(), Error> {
263 self.track.discontinuity()?;
264 Ok(())
265 }
266
267 pub fn finish(mut self) -> Result<(), Error> {
270 let frame_samples = self.encoder.frame_size() * self.encoder.codec_channels() as usize;
271 if !self.pending.is_empty() {
272 self.pending.resize(frame_samples, 0.0);
273 let chunk = std::mem::take(&mut self.pending);
274 let packet = self.encoder.encode(&chunk)?;
275 let timestamp = self.timestamp(self.epoch_us.unwrap_or(0))?;
276 self.publish(packet, timestamp)?;
277 }
278 self.track.finish()?;
279 Ok(())
280 }
281
282 pub fn abort(self, err: moq_net::Error) {
285 self.track.abort(err);
286 }
287}
288
289struct Rendition<E: CatalogExt> {
294 catalog: moq_mux::catalog::Producer<E>,
295 name: String,
296}
297
298impl<E: CatalogExt> Drop for Rendition<E> {
299 fn drop(&mut self) {
300 self.catalog.lock().audio.remove(&self.name);
301 }
302}
303
304#[cfg(test)]
305mod tests {
306 use super::*;
307 use crate::Format;
308
309 fn full_frame(timestamp_us: u64) -> Frame {
312 let mut data = Vec::with_capacity(960 * 4);
313 for _ in 0..960 {
314 data.extend_from_slice(&0.1f32.to_le_bytes());
315 }
316 Frame {
317 timestamp: Timestamp::from_micros(timestamp_us).unwrap(),
318 data: data.into(),
319 }
320 }
321
322 async fn published_pts(frames: &[Frame], reset_before: Option<usize>) -> Vec<u128> {
326 let mut broadcast = moq_net::broadcast::Info::new().produce();
327 let catalog = moq_mux::catalog::Producer::new(&mut broadcast).unwrap();
328 let consumer = broadcast.consume();
329
330 let input = Input {
333 format: Format::F32,
334 sample_rate: 48_000,
335 channels: 1,
336 };
337 let options = Options {
338 track: Some("audio".to_string()),
339 ..Options::default()
340 };
341 let mut producer = Producer::new(&mut broadcast, catalog, input, &options).unwrap();
342
343 let track = consumer.track("audio").unwrap().subscribe(None).await.unwrap();
344 let mut reader = moq_mux::container::Consumer::new(track, moq_mux::container::legacy::Wire);
345
346 let mut pts = Vec::new();
347 for (i, frame) in frames.iter().enumerate() {
348 if reset_before == Some(i) {
349 producer.reset_epoch();
350 }
351 producer.write(frame).unwrap();
352 let read = reader.read().await.unwrap().expect("a packet per full frame");
353 pts.push(read.timestamp.as_micros());
354 }
355 pts
356 }
357
358 #[tokio::test]
359 async fn epoch_anchors_to_first_frame_timestamp() {
360 let pts = published_pts(&[full_frame(1_000_000)], None).await;
363 assert_eq!(pts, vec![1_000_000]);
364 }
365
366 #[tokio::test]
367 async fn pts_advances_by_frame_duration_ignoring_later_timestamps() {
368 let pts = published_pts(&[full_frame(1_000), full_frame(999_999)], None).await;
371 assert_eq!(pts, vec![1_000, 1_000 + 20_000]);
372 }
373
374 #[tokio::test]
375 async fn reset_epoch_reanchors_so_the_gap_lands_in_pts() {
376 let pts = published_pts(&[full_frame(0), full_frame(5_000_000)], Some(1)).await;
379 assert_eq!(pts, vec![0, 5_000_000]);
380 }
381
382 #[tokio::test]
386 async fn default_options_derive_the_track_name() {
387 let mut broadcast = moq_net::broadcast::Info::new().produce();
388 let catalog = moq_mux::catalog::Producer::new(&mut broadcast).unwrap();
389
390 let first = Producer::new(&mut broadcast, catalog.clone(), Input::default(), &Options::default()).unwrap();
391 assert_eq!(first.track_name(), "0.opus");
392
393 let second = Producer::new(&mut broadcast, catalog, Input::default(), &Options::default()).unwrap();
394 assert_eq!(second.track_name(), "1.opus");
395 }
396}