Skip to main content

moq_audio/decode/
consumer.rs

1//! Subscribe to an encoded audio track and emit raw PCM.
2
3use bytes::Bytes;
4
5use super::decoder::{Config, Decoder};
6use crate::resample::{Resampler, remix};
7use crate::{Error, Frame};
8
9/// Subscribe to a moq-mux audio track and emit decoded PCM in the layout
10/// declared by [`Config`].
11///
12/// The mirror of [`encode::Producer`](crate::encode::Producer): output format /
13/// sample rate / channel count are fixed at construction, and
14/// [`read`](Self::read) returns plain [`Frame`]s.
15pub 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	/// Subscribe to `name` in `broadcast`, using the catalog entry to pick the
26	/// codec.
27	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	/// The config this consumer was built with.
71	pub fn config(&self) -> &Config {
72		&self.config
73	}
74
75	/// Sample rate samples are actually delivered at, which is
76	/// [`Config::sample_rate`] resolved against the catalog.
77	pub fn sample_rate(&self) -> u32 {
78		self.resolved_sample_rate
79	}
80
81	/// Channel count samples are actually delivered at, which is
82	/// [`Config::channels`] resolved against the catalog.
83	pub fn channels(&self) -> u32 {
84		self.resolved_channels
85	}
86
87	/// Read the next decoded PCM frame, or `None` when the track ends.
88	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}