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>,
40 pub frame_duration: Duration,
42}
43
44impl Default for Options {
45 fn default() -> Self {
46 Self {
47 track: None,
48 codec: Codec::default(),
49 sample_rate: None,
50 channels: None,
51 bitrate: None,
52 frame_duration: Duration::from_millis(20),
53 }
54 }
55}
56
57impl Options {
58 fn config(&self, input: Input) -> Config {
60 Config {
61 input,
62 codec: self.codec,
63 sample_rate: self.sample_rate,
64 channels: self.channels,
65 bitrate: self.bitrate,
66 frame_duration: self.frame_duration,
67 }
68 }
69}
70
71pub struct Producer<E: CatalogExt = ()> {
81 encoder: Encoder,
82 resampler: Option<Resampler>,
83 track: moq_mux::container::Producer<moq_mux::container::legacy::Wire>,
84 rendition: Rendition<E>,
86 pending: Vec<f32>,
87 frames_produced: u64,
89 epoch_us: Option<u64>,
93}
94
95impl<E: CatalogExt> Producer<E> {
96 pub fn new(
99 broadcast: &mut moq_net::broadcast::Producer,
100 catalog: moq_mux::catalog::Producer<E>,
101 input: Input,
102 options: &Options,
103 ) -> Result<Self, Error> {
104 let encoder = Encoder::new(&options.config(input))?;
105 let input = &encoder.config().input;
106
107 let resampler = if input.sample_rate == encoder.codec_rate() {
108 None
109 } else {
110 let chunk_frames =
113 ((input.sample_rate as u128 * encoder.config().frame_duration.as_micros()) / 1_000_000) as usize;
114 Some(Resampler::new(
115 input.sample_rate,
116 encoder.codec_rate(),
117 input.channels,
118 chunk_frames,
119 )?)
120 };
121
122 let track = match &options.track {
123 Some(name) => broadcast.create_track(name.clone(), hang::container::track_info())?,
127 None => moq_mux::import::unique_track(broadcast, &format!(".{}", options.codec))?,
130 };
131 let name = track.name().to_string();
132 let track = catalog.media_producer(track, moq_mux::container::legacy::Wire)?;
133
134 let mut catalog_mut = catalog.clone();
135 let mut config = encoder.catalog();
136 config.timeline = Some(catalog.timeline(&name)?.section());
137 catalog_mut.lock().audio.insert(&name, config)?;
138
139 Ok(Self {
140 encoder,
141 resampler,
142 track,
143 rendition: Rendition { catalog, name },
144 pending: Vec::new(),
145 frames_produced: 0,
146 epoch_us: None,
147 })
148 }
149
150 pub fn track_name(&self) -> &str {
152 &self.rendition.name
153 }
154
155 pub fn track(&self) -> &moq_net::track::Producer {
158 self.track.track()
159 }
160
161 pub fn reset_epoch(&mut self) {
168 self.epoch_us = None;
169 self.frames_produced = 0;
170 self.pending.clear();
171 }
172
173 pub fn write(&mut self, frame: &Frame) -> Result<(), Error> {
185 let timestamp_us = u64::try_from(frame.timestamp.as_micros())
186 .map_err(|_| Error::Unsupported(format!("frame timestamp {:?} out of range", frame.timestamp)))?;
187 let epoch_us = *self.epoch_us.get_or_insert(timestamp_us);
188
189 let input = &self.encoder.config().input;
190 let (format, channels) = (input.format, input.channels);
191 let pcm = format.as_interleaved_f32(frame.data.as_ref(), channels)?;
192 let pcm: Vec<f32> = match self.resampler.as_mut() {
193 Some(r) => r.process(&pcm)?,
194 None => pcm.into_owned(),
195 };
196
197 self.pending.extend(pcm);
198
199 let frame_samples = self.encoder.frame_size() * self.encoder.codec_channels() as usize;
200 while self.pending.len() >= frame_samples {
201 let chunk: Vec<f32> = self.pending.drain(..frame_samples).collect();
202 let packet = self.encoder.encode(&chunk)?;
203
204 let timestamp = self.timestamp(epoch_us)?;
205 self.frames_produced += self.encoder.frame_size() as u64;
206 self.publish(packet, timestamp)?;
207 }
208
209 Ok(())
210 }
211
212 fn timestamp(&self, epoch_us: u64) -> Result<Timestamp, Error> {
214 let offset_us = (self.frames_produced * 1_000_000) / self.encoder.codec_rate() as u64;
215 Ok(Timestamp::from_micros(epoch_us + offset_us)?)
216 }
217
218 fn publish(&mut self, payload: Bytes, timestamp: Timestamp) -> Result<(), Error> {
219 let mux_frame = MuxFrame {
222 timestamp,
223 payload,
224 keyframe: true,
225 duration: None,
226 };
227 self.track.write(mux_frame)?;
228 self.track.cut(None)?;
231 Ok(())
232 }
233
234 pub fn discontinuity(&mut self) -> Result<(), Error> {
242 self.track.discontinuity()?;
243 Ok(())
244 }
245
246 pub fn finish(mut self) -> Result<(), Error> {
249 let frame_samples = self.encoder.frame_size() * self.encoder.codec_channels() as usize;
250 if !self.pending.is_empty() {
251 self.pending.resize(frame_samples, 0.0);
252 let chunk = std::mem::take(&mut self.pending);
253 let packet = self.encoder.encode(&chunk)?;
254 let timestamp = self.timestamp(self.epoch_us.unwrap_or(0))?;
255 self.publish(packet, timestamp)?;
256 }
257 self.track.finish()?;
258 Ok(())
259 }
260
261 pub fn abort(self, err: moq_net::Error) {
264 self.track.abort(err);
265 }
266}
267
268struct Rendition<E: CatalogExt> {
273 catalog: moq_mux::catalog::Producer<E>,
274 name: String,
275}
276
277impl<E: CatalogExt> Drop for Rendition<E> {
278 fn drop(&mut self) {
279 self.catalog.lock().audio.remove(&self.name);
280 }
281}
282
283#[cfg(test)]
284mod tests {
285 use super::*;
286 use crate::Format;
287
288 fn full_frame(timestamp_us: u64) -> Frame {
291 let mut data = Vec::with_capacity(960 * 4);
292 for _ in 0..960 {
293 data.extend_from_slice(&0.1f32.to_le_bytes());
294 }
295 Frame {
296 timestamp: Timestamp::from_micros(timestamp_us).unwrap(),
297 data: data.into(),
298 }
299 }
300
301 async fn published_pts(frames: &[Frame], reset_before: Option<usize>) -> Vec<u128> {
305 let mut broadcast = moq_net::broadcast::Info::new().produce();
306 let catalog = moq_mux::catalog::Producer::new(&mut broadcast).unwrap();
307 let consumer = broadcast.consume();
308
309 let input = Input {
312 format: Format::F32,
313 sample_rate: 48_000,
314 channels: 1,
315 };
316 let options = Options {
317 track: Some("audio".to_string()),
318 ..Options::default()
319 };
320 let mut producer = Producer::new(&mut broadcast, catalog, input, &options).unwrap();
321
322 let track = consumer.track("audio").unwrap().subscribe(None).await.unwrap();
323 let mut reader = moq_mux::container::Consumer::new(track, moq_mux::container::legacy::Wire);
324
325 let mut pts = Vec::new();
326 for (i, frame) in frames.iter().enumerate() {
327 if reset_before == Some(i) {
328 producer.reset_epoch();
329 }
330 producer.write(frame).unwrap();
331 let read = reader.read().await.unwrap().expect("a packet per full frame");
332 pts.push(read.timestamp.as_micros());
333 }
334 pts
335 }
336
337 #[tokio::test]
338 async fn epoch_anchors_to_first_frame_timestamp() {
339 let pts = published_pts(&[full_frame(1_000_000)], None).await;
342 assert_eq!(pts, vec![1_000_000]);
343 }
344
345 #[tokio::test]
346 async fn pts_advances_by_frame_duration_ignoring_later_timestamps() {
347 let pts = published_pts(&[full_frame(1_000), full_frame(999_999)], None).await;
350 assert_eq!(pts, vec![1_000, 1_000 + 20_000]);
351 }
352
353 #[tokio::test]
354 async fn reset_epoch_reanchors_so_the_gap_lands_in_pts() {
355 let pts = published_pts(&[full_frame(0), full_frame(5_000_000)], Some(1)).await;
358 assert_eq!(pts, vec![0, 5_000_000]);
359 }
360
361 #[tokio::test]
365 async fn default_options_derive_the_track_name() {
366 let mut broadcast = moq_net::broadcast::Info::new().produce();
367 let catalog = moq_mux::catalog::Producer::new(&mut broadcast).unwrap();
368
369 let first = Producer::new(&mut broadcast, catalog.clone(), Input::default(), &Options::default()).unwrap();
370 assert_eq!(first.track_name(), "0.opus");
371
372 let second = Producer::new(&mut broadcast, catalog, Input::default(), &Options::default()).unwrap();
373 assert_eq!(second.track_name(), "1.opus");
374 }
375}