moq_audio/decode/
consumer.rs1use bytes::Bytes;
4
5use super::decoder::{Config, Decoder};
6use crate::resample::{Resampler, remix};
7use crate::{Error, Frame};
8
9pub struct Consumer {
16 decoder: Decoder,
17 track: moq_mux::container::Consumer<moq_mux::container::legacy::Wire>,
18 resampler: Option<Resampler>,
19 config: Config,
20 resolved_sample_rate: u32,
21 resolved_channels: u32,
22}
23
24impl Consumer {
25 pub async fn new(
28 broadcast: &moq_net::broadcast::Consumer,
29 catalog: &hang::catalog::AudioConfig,
30 name: impl Into<String>,
31 config: Config,
32 ) -> Result<Self, Error> {
33 let decoder = Decoder::new(catalog)?;
34 let sample_rate = config.sample_rate.unwrap_or_else(|| decoder.sample_rate());
35 let channels = config.channels.unwrap_or_else(|| decoder.channel_count());
36 crate::opus::validate_channels(channels)?;
37
38 let resampler = if sample_rate == decoder.sample_rate() {
39 None
40 } else {
41 let chunk_frames = (decoder.sample_rate() as usize * 20) / 1000;
42 Some(Resampler::new(
43 decoder.sample_rate(),
44 sample_rate,
45 decoder.channel_count(),
46 chunk_frames,
47 )?)
48 };
49
50 let name = name.into();
51 let track = broadcast
52 .track(&name)?
53 .subscribe(moq_net::track::Subscription::default().with_priority(hang::catalog::PRIORITY.audio))
54 .await?;
55 let mut track = moq_mux::container::Consumer::new(track, moq_mux::container::legacy::Wire);
56 if let Some(latency) = config.latency_max {
57 track = track.with_latency(latency);
58 }
59
60 Ok(Self {
61 decoder,
62 track,
63 resampler,
64 config,
65 resolved_sample_rate: sample_rate,
66 resolved_channels: channels,
67 })
68 }
69
70 pub fn config(&self) -> &Config {
72 &self.config
73 }
74
75 pub fn sample_rate(&self) -> u32 {
78 self.resolved_sample_rate
79 }
80
81 pub fn channels(&self) -> u32 {
84 self.resolved_channels
85 }
86
87 pub async fn read(&mut self) -> Result<Option<Frame>, Error> {
89 let Some(mux_frame) = self.track.read().await? else {
90 return Ok(None);
91 };
92
93 let decoded = self.decoder.decode(&mux_frame.payload)?;
94 let pcm = match self.resampler.as_mut() {
95 Some(r) => r.process(&decoded)?,
96 None => decoded,
97 };
98 let pcm = if self.decoder.channel_count() == self.resolved_channels {
99 pcm
100 } else {
101 remix(&pcm, self.decoder.channel_count(), self.resolved_channels)?
102 };
103
104 let bytes = self.config.format.from_interleaved_f32(&pcm, self.resolved_channels)?;
105 Ok(Some(Frame {
106 timestamp: mux_frame.timestamp,
107 data: Bytes::from(bytes),
108 }))
109 }
110}
111
112#[cfg(test)]
113mod tests {
114 use moq_net::Timestamp;
115
116 use super::*;
117 use crate::Format;
118 use crate::encode::{Encoder, Input, Options, Producer};
119
120 #[tokio::test]
121 async fn remixes_mono_stream_to_stereo_output() {
122 let mut broadcast = moq_net::broadcast::Info::new().produce();
123 let catalog = moq_mux::catalog::Producer::new(&mut broadcast).unwrap();
124 let subscriber = broadcast.consume();
125 let input = Input {
126 format: Format::F32,
127 sample_rate: 48_000,
128 channels: 1,
129 };
130 let options = Options {
131 track: Some("audio".to_string()),
132 ..Options::default()
133 };
134 let mut producer = Producer::new(&mut broadcast, catalog, input.clone(), &options).unwrap();
135 let catalog = Encoder::new(&crate::encode::Config::new(input)).unwrap().catalog();
136 let mut consumer = Consumer::new(
137 &subscriber,
138 &catalog,
139 "audio",
140 Config {
141 channels: Some(2),
142 ..Config::new()
143 },
144 )
145 .await
146 .unwrap();
147
148 let samples = vec![0.1f32; 960];
149 let mut data = Vec::with_capacity(samples.len() * size_of::<f32>());
150 for sample in samples {
151 data.extend_from_slice(&sample.to_le_bytes());
152 }
153 producer
154 .write(&Frame {
155 timestamp: Timestamp::ZERO,
156 data: data.into(),
157 })
158 .unwrap();
159
160 let frame = consumer.read().await.unwrap().expect("decoded frame");
161 let samples = Format::F32.as_interleaved_f32(&frame.data, 2).unwrap();
162 assert_eq!(samples.len(), (960 - 312) * 2);
163 for pair in samples.chunks_exact(2) {
164 assert_eq!(pair[0], pair[1]);
165 }
166 }
167}