Skip to main content

moq/
audio.rs

1//! Raw-audio import/export via [`moq_audio`].
2//!
3//! Sibling to `moq_publish_media_*` / `moq_consume_audio`
4//! (those handle already-encoded frames). These functions accept and
5//! return raw PCM, with Opus encode/decode happening inside the FFI
6//! boundary.
7//!
8//! Format / sample rate / channel count are fixed at producer or
9//! consumer construction via [`moq_audio_encoder_input`] /
10//! [`moq_audio_encoder_output`] / [`moq_audio_decoder_output`], so
11//! each [`moq_audio_frame`] carries only payload bytes and a
12//! timestamp.
13
14use std::ffi::{c_char, c_void};
15use std::sync::Arc;
16use std::time::Duration;
17
18use bytes::Bytes;
19use tokio::sync::oneshot;
20
21use crate::ffi::OnStatus;
22use crate::{Error, Id, NonZeroSlab, Shared, State, ffi};
23
24// ---- C-visible types ----
25
26/// Raw PCM sample layout, mirroring WebCodecs `AudioData.format`.
27///
28/// The enum is exposed in the C header for readability, but ABI
29/// fields/parameters that carry it are typed `u32`. A C caller
30/// passing an unknown discriminant gets `Error::InvalidCode` instead
31/// of UB.
32///
33/// <https://developer.mozilla.org/en-US/docs/Web/API/AudioData/format>
34#[repr(C)]
35#[allow(non_camel_case_types)]
36#[derive(Clone, Copy, Debug)]
37pub enum moq_audio_sample_format {
38	MOQ_AUDIO_SAMPLE_FORMAT_U8 = 0,
39	MOQ_AUDIO_SAMPLE_FORMAT_S16 = 1,
40	MOQ_AUDIO_SAMPLE_FORMAT_S32 = 2,
41	MOQ_AUDIO_SAMPLE_FORMAT_F32 = 3,
42	MOQ_AUDIO_SAMPLE_FORMAT_U8_PLANAR = 4,
43	MOQ_AUDIO_SAMPLE_FORMAT_S16_PLANAR = 5,
44	MOQ_AUDIO_SAMPLE_FORMAT_S32_PLANAR = 6,
45	MOQ_AUDIO_SAMPLE_FORMAT_F32_PLANAR = 7,
46}
47
48fn audio_format_from_u32(value: u32) -> Result<moq_audio::Format, Error> {
49	use moq_audio::Format;
50	Ok(match value {
51		v if v == moq_audio_sample_format::MOQ_AUDIO_SAMPLE_FORMAT_U8 as u32 => Format::U8,
52		v if v == moq_audio_sample_format::MOQ_AUDIO_SAMPLE_FORMAT_S16 as u32 => Format::S16,
53		v if v == moq_audio_sample_format::MOQ_AUDIO_SAMPLE_FORMAT_S32 as u32 => Format::S32,
54		v if v == moq_audio_sample_format::MOQ_AUDIO_SAMPLE_FORMAT_F32 as u32 => Format::F32,
55		v if v == moq_audio_sample_format::MOQ_AUDIO_SAMPLE_FORMAT_U8_PLANAR as u32 => Format::U8Planar,
56		v if v == moq_audio_sample_format::MOQ_AUDIO_SAMPLE_FORMAT_S16_PLANAR as u32 => Format::S16Planar,
57		v if v == moq_audio_sample_format::MOQ_AUDIO_SAMPLE_FORMAT_S32_PLANAR as u32 => Format::S32Planar,
58		v if v == moq_audio_sample_format::MOQ_AUDIO_SAMPLE_FORMAT_F32_PLANAR as u32 => Format::F32Planar,
59		_ => return Err(Error::InvalidCode),
60	})
61}
62
63/// PCM layout the caller hands to [`moq_encode_audio_frame`].
64#[repr(C)]
65#[allow(non_camel_case_types)]
66pub struct moq_audio_encoder_input {
67	/// `moq_audio_sample_format` discriminant.
68	pub format: u32,
69	pub sample_rate: u32,
70	/// Interleaved channel count, which also names the speaker layout by the
71	/// WAVE convention: 1 mono, 2 stereo, 3 2.1, 4 quad, 5 5.0, 6 5.1, 7 6.1,
72	/// 8 7.1, in front left, front right, center, LFE, back, side order.
73	pub channels: u32,
74}
75
76/// Codec-side configuration. `sample_rate` / `channels` = 0 means
77/// "match the input (snapping the rate up to a libopus-supported
78/// value if necessary)".
79#[repr(C)]
80#[allow(non_camel_case_types)]
81pub struct moq_audio_encoder_output {
82	/// Codec id, UTF-8: "opus", "pcm", or "aac". AAC encodes through the
83	/// platform's encoder, so a host without one refuses it.
84	pub codec: *const c_char,
85	pub codec_len: usize,
86	/// 0 = derive from input.
87	pub sample_rate: u32,
88	/// 0 = derive from input.
89	pub channels: u32,
90	/// 0 = libopus default.
91	pub bitrate: u32,
92	/// Encoded frame duration in microseconds. Opus accepts exactly
93	/// 2500/5000/10000/20000/40000/60000 us. 0 = the codec's default: 20 ms for
94	/// Opus, which matches the JS publish path, and 1024 samples for AAC.
95	pub frame_duration_us: u32,
96}
97
98/// PCM layout the caller wants out of [`moq_decode_audio`].
99#[repr(C)]
100#[allow(non_camel_case_types)]
101pub struct moq_audio_decoder_output {
102	pub format: u32,
103	/// 0 = deliver at the codec's native sample rate.
104	pub sample_rate: u32,
105	/// 0 = deliver at the codec's native channel count. A count names its
106	/// layout as `moq_audio_encoder_input.channels` describes, and the decoder
107	/// remixes to it.
108	pub channels: u32,
109	/// Upper bound on buffering before skipping a stalled group, in
110	/// microseconds. Same congestion-control knob as
111	/// `moq_consume_audio`'s `max_age_us`. 0 = skip
112	/// aggressively (the moq-mux default); set to your playout
113	/// buffer (tens to a few hundred ms) for a softer skip. Named
114	/// `_max` to leave room for a future `min_buffer_us`, a
115	/// jitter-buffer floor rather than a staleness bound.
116	pub max_age_us: u64,
117}
118
119/// One audio frame: payload bytes plus a presentation timestamp.
120///
121/// `data` is owned by the consume slab (see
122/// [`moq_decode_audio_frame_free`]) or borrowed by the publish call
123/// (the publisher copies before returning).
124#[repr(C)]
125#[allow(non_camel_case_types)]
126pub struct moq_audio_frame {
127	pub timestamp_us: u64,
128	pub data: *const u8,
129	pub data_size: usize,
130}
131
132// ---- State extensions (used internally by lib.rs) ----
133
134/// An audio producer, shared so the Opus encode in `write` runs with the global
135/// lock released. See [`Shared`].
136pub(crate) struct AudioEncoder {
137	producer: moq_audio::encode::Producer<moq_mux::catalog::hang::Extra>,
138	reservation: Option<Arc<moq_net::bandwidth::Reservation>>,
139}
140
141type AudioProducer = Shared<AudioEncoder>;
142
143#[derive(Default)]
144pub struct Audio {
145	producers: NonZeroSlab<AudioProducer>,
146	consumer_tasks: NonZeroSlab<Option<AudioTaskEntry>>,
147	frames: NonZeroSlab<moq_audio::Frame>,
148}
149
150/// A spawned task entry: `close` signals shutdown, `callback` delivers status.
151///
152/// `close` is an `Option` so `consume_close` can drop just the sender without
153/// removing the entry. The task delivers one final terminal callback and then
154/// removes itself, so `user_data` stays valid until that callback fires.
155struct AudioTaskEntry {
156	close: Option<oneshot::Sender<()>>,
157	callback: OnStatus,
158}
159
160impl Audio {
161	pub fn publish(
162		&mut self,
163		broadcast: &mut moq_net::broadcast::Producer,
164		catalog: moq_mux::catalog::Producer<moq_mux::catalog::hang::Extra>,
165		input: moq_audio::encode::Input,
166		options: moq_audio::encode::Options,
167		reserve: bool,
168	) -> Result<Id, Error> {
169		let producer = moq_audio::encode::Producer::new(broadcast, catalog, input, &options)?;
170		let reservation = reserve.then(|| Arc::new(options.bandwidth.reserve(&producer.demand(), producer.bitrate())));
171		self.producers
172			.insert(Shared::new(AudioEncoder { producer, reservation }))
173	}
174
175	pub(crate) fn reservation(&self, id: Id) -> Result<Option<Arc<moq_net::bandwidth::Reservation>>, Error> {
176		Ok(self
177			.producer(id)?
178			.lock()
179			.as_ref()
180			.ok_or(Error::MediaNotFound)?
181			.reservation
182			.clone())
183	}
184
185	/// A watch-only handle to the encoded track's subscriber demand.
186	pub(crate) fn demand(&self, id: Id) -> Result<moq_net::track::Demand, Error> {
187		Ok(self
188			.producer(id)?
189			.lock()
190			.as_ref()
191			.ok_or(Error::MediaNotFound)?
192			.producer
193			.demand())
194	}
195
196	/// Resolve a producer handle, so the caller can encode with the global lock
197	/// released.
198	///
199	/// Bind the result before locking it: a temporary [`State`] guard lives to the
200	/// end of the statement that created it, so resolving and locking in one
201	/// expression would put the encode back under the global lock.
202	pub(crate) fn producer(&self, id: Id) -> Result<AudioProducer, Error> {
203		self.producers.get(id).cloned().ok_or(Error::MediaNotFound)
204	}
205
206	/// Resolve a producer and drop its id, so nothing can be published to it after.
207	pub(crate) fn remove(&mut self, id: Id) -> Result<AudioProducer, Error> {
208		self.producers.remove(id).ok_or(Error::MediaNotFound)
209	}
210
211	pub fn consume(
212		&mut self,
213		broadcast: &moq_net::broadcast::Consumer,
214		catalog: &hang::catalog::AudioConfig,
215		name: &str,
216		config: moq_audio::decode::Options,
217		on_frame: OnStatus,
218	) -> Result<Id, Error> {
219		let broadcast = broadcast.clone();
220		let catalog = catalog.clone();
221		let name = name.to_string();
222
223		let channel = oneshot::channel();
224		let entry = AudioTaskEntry {
225			close: Some(channel.0),
226			callback: on_frame,
227		};
228		let id = self.consumer_tasks.insert(Some(entry))?;
229
230		// `decode::Consumer::new` subscribes (blocking on SUBSCRIBE_OK), so run it
231		// inside the task to keep this entrypoint non-blocking.
232		tokio::spawn(async move {
233			let res = async move {
234				let consumer = moq_audio::decode::Consumer::new(&broadcast, &catalog, name, config).await?;
235				Self::run(on_frame, consumer, channel.1).await
236			}
237			.await;
238
239			// Deliver one final terminal callback (code <= 0), then drop the entry.
240			// Pull it out from under the lock so the callback never runs while held.
241			let entry = State::lock().audio.consumer_tasks.remove(id).flatten();
242			if let Some(entry) = entry {
243				entry.callback.call(res);
244			}
245		});
246
247		Ok(id)
248	}
249
250	async fn run(
251		callback: OnStatus,
252		mut consumer: moq_audio::decode::Consumer,
253		mut close: oneshot::Receiver<()>,
254	) -> Result<(), Error> {
255		loop {
256			// `biased` so a pending close always wins over a ready frame.
257			let frame = tokio::select! {
258				biased;
259				_ = &mut close => return Ok(()),
260				frame = consumer.read() => match frame {
261					Ok(Some(frame)) => frame,
262					Ok(None) => return Ok(()),
263					// One packet the codec rejected is that packet's problem: the
264					// decoder stays usable, so drop it and keep the subscription rather
265					// than ending the caller's stream over a single bad frame.
266					Err(moq_audio::Error::Decode(err)) => {
267						tracing::warn!(%err, "dropping an audio frame");
268						continue;
269					}
270					Err(err) => return Err(err.into()),
271				},
272			};
273
274			// Hold the lock only to buffer the frame; release it before the callback.
275			let frame_id = State::lock().audio.frames.insert(frame)?;
276			callback.call(Ok(frame_id));
277		}
278	}
279
280	pub fn consume_close(&mut self, id: Id) -> Result<(), Error> {
281		// Signal shutdown; the task delivers a final callback and removes itself.
282		self.consumer_tasks
283			.get_mut(id)
284			.and_then(|entry| entry.as_mut())
285			.ok_or(Error::TrackNotFound)?
286			.close
287			.take()
288			.ok_or(Error::TrackNotFound)?;
289		Ok(())
290	}
291
292	pub fn frame_info(&self, id: Id, dst: &mut moq_audio_frame) -> Result<(), Error> {
293		let frame = self.frames.get(id).ok_or(Error::FrameNotFound)?;
294		*dst = moq_audio_frame {
295			// The C ABI carries plain microseconds, so flatten the scaled
296			// `Timestamp` here at the boundary. Saturating rather than erroring:
297			// this is a getter on a frame we already decoded, and a u64 overflow
298			// needs a timestamp ~580,000 years out.
299			timestamp_us: u64::try_from(frame.timestamp.as_micros()).unwrap_or(u64::MAX),
300			data: frame.data.as_ptr(),
301			data_size: frame.data.len(),
302		};
303		Ok(())
304	}
305
306	pub fn frame_free(&mut self, id: Id) -> Result<(), Error> {
307		self.frames.remove(id).ok_or(Error::FrameNotFound)?;
308		Ok(())
309	}
310}
311
312// ---- C entry points ----
313
314/// Open an audio track on a broadcast.
315///
316/// The encoder configuration is fixed at construction; subsequent
317/// frame writes pass only payload + timestamp via
318/// [`moq_encode_audio_frame`].
319///
320/// Returns a non-zero handle on success or a negative error code.
321///
322/// # Safety
323/// - `name` must point to `name_len` bytes of UTF-8.
324/// - `input` / `output` must point to fully populated structs.
325/// - `output->codec` must point to `output->codec_len` bytes of UTF-8.
326/// - `bandwidth` is a handle from [`crate::moq_session_bandwidth`], or 0 to leave the
327///   configured bitrate unclaimed.
328#[unsafe(no_mangle)]
329pub unsafe extern "C" fn moq_encode_audio(
330	broadcast: u32,
331	name: *const c_char,
332	name_len: usize,
333	input: *const moq_audio_encoder_input,
334	output: *const moq_audio_encoder_output,
335	bandwidth: u32,
336) -> i32 {
337	ffi::enter(move || {
338		let broadcast = ffi::parse_id(broadcast)?;
339		let name = unsafe { ffi::parse_str(name, name_len)? }.to_string();
340		let raw_input = unsafe { input.as_ref() }.ok_or(Error::InvalidPointer)?;
341		let raw_output = unsafe { output.as_ref() }.ok_or(Error::InvalidPointer)?;
342		let codec_str = unsafe { ffi::parse_str(raw_output.codec, raw_output.codec_len)? };
343
344		let layout = moq_audio::Layout::from_channels(raw_input.channels)?;
345		let mut encoder_input = moq_audio::encode::Input::new(raw_input.sample_rate, layout);
346		encoder_input.format = audio_format_from_u32(raw_input.format)?;
347
348		// The C ABI takes an explicit track name and spells "unset" as 0, so map
349		// both onto the Rust options here rather than leaking either convention.
350		let mut options = moq_audio::encode::Options::default();
351		options.track = Some(name);
352		let codec = codec_str
353			.parse()
354			.map_err(|_| Error::UnknownFormat(codec_str.to_string()))?;
355		options.settings = moq_audio::encode::Settings::from_input(codec, &encoder_input);
356		if let Some(sample_rate) = zeroable(raw_output.sample_rate) {
357			options.settings.sample_rate = sample_rate;
358		}
359		if let Some(channels) = zeroable(raw_output.channels) {
360			options.settings.layout = moq_audio::Layout::from_channels(channels)?;
361		}
362		options.settings.bitrate =
363			zeroable(raw_output.bitrate).map(|bps| moq_net::bandwidth::Rate::from_bps(bps.into()));
364		if let Some(micros) = zeroable(raw_output.frame_duration_us) {
365			options.settings.frame_duration = Duration::from_micros(micros.into());
366		}
367
368		let bandwidth = ffi::parse_id_optional(bandwidth)?;
369		let mut state = State::lock();
370		if let Some(id) = bandwidth {
371			options.bandwidth = state.bandwidth.allocator(id)?;
372		}
373		let State { publish, audio, .. } = &mut *state;
374		let (broadcast_producer, catalog) = publish.pair_mut(broadcast)?;
375
376		audio.publish(
377			broadcast_producer,
378			catalog.clone(),
379			encoder_input,
380			options,
381			bandwidth.is_some(),
382		)
383	})
384}
385
386/// This encoder's bandwidth reservation, or 0 if it was published without one.
387///
388/// Closing the returned handle does not release the encoder's claim; that lasts
389/// until [moq_encode_audio_finish].
390#[unsafe(no_mangle)]
391pub extern "C" fn moq_encode_audio_reservation(producer: u32) -> i32 {
392	ffi::enter(move || {
393		let producer = ffi::parse_id(producer)?;
394		let mut state = State::lock();
395		match state.audio.reservation(producer)? {
396			Some(reservation) => Ok(i32::from(state.bandwidth.hold(reservation)?)),
397			None => Ok(0),
398		}
399	})
400}
401
402/// Watch whether the encoded audio track has subscribers, so the microphone and
403/// encoder run only while someone listens. See [`crate::moq_publish_media_demand`]
404/// for the callback contract.
405///
406/// Returns a non-zero watcher handle on success, or a negative code on failure.
407///
408/// # Safety
409/// - `on_demand` must be non-NULL.
410/// - The caller must keep `user_data` valid until the terminal (`<= 0`) `on_demand` callback.
411#[unsafe(no_mangle)]
412pub unsafe extern "C" fn moq_encode_audio_demand(
413	producer: u32,
414	on_demand: crate::moq_status_callback,
415	user_data: *mut c_void,
416) -> i32 {
417	ffi::enter(move || {
418		let producer = ffi::parse_id(producer)?;
419		let on_demand = unsafe { OnStatus::new(user_data, on_demand)? };
420		let mut state = State::lock();
421		let demand = state.audio.demand(producer)?;
422		state.publish.demand(demand, on_demand)
423	})
424}
425
426/// The C ABI spells an unset `u32` knob as 0, which no field here accepts as a
427/// real value.
428fn zeroable(value: u32) -> Option<u32> {
429	(value != 0).then_some(value)
430}
431
432/// Push one audio frame.
433///
434/// `frame->data` is borrowed for the duration of the call; the
435/// producer copies before returning.
436///
437/// # Safety
438/// - `frame` must point to a valid [`moq_audio_frame`].
439/// - `frame->data` must point to `frame->data_size` bytes.
440#[unsafe(no_mangle)]
441pub unsafe extern "C" fn moq_encode_audio_frame(producer: u32, frame: *const moq_audio_frame) -> i32 {
442	ffi::enter(move || {
443		let producer = ffi::parse_id(producer)?;
444		let frame = unsafe { frame.as_ref() }.ok_or(Error::InvalidPointer)?;
445		let data = unsafe { ffi::parse_slice(frame.data, frame.data_size)? };
446
447		// The C ABI carries plain microseconds; scale them at the boundary.
448		let timestamp = moq_net::Timestamp::from_micros(frame.timestamp_us).map_err(moq_audio::Error::from)?;
449		let owned = moq_audio::Frame::new(Bytes::copy_from_slice(data), timestamp);
450
451		let producer = State::lock().audio.producer(producer)?;
452		producer
453			.lock()
454			.as_mut()
455			.ok_or(Error::MediaNotFound)?
456			.producer
457			.write(&owned)?;
458		Ok(())
459	})
460}
461
462/// Flush any pending samples and finalize an audio producer.
463#[unsafe(no_mangle)]
464pub extern "C" fn moq_encode_audio_finish(producer: u32) -> i32 {
465	ffi::enter(move || {
466		let producer = ffi::parse_id(producer)?;
467		// The id is dropped first, so nothing new queues behind the flush; whatever
468		// is mid-encode still finishes before this takes the producer.
469		let producer = State::lock().audio.remove(producer)?;
470		let mut producer = producer.take().ok_or(Error::MediaNotFound)?.producer;
471		producer.finish()?;
472		Ok(())
473	})
474}
475
476/// Subscribe to an audio track and decode it into PCM.
477///
478/// The catalog `index` identifies which audio rendition to subscribe
479/// to, matching the existing `moq_consume_audio` selection
480/// model. TODO: a future API will pick the right rendition
481/// automatically (ABR).
482///
483/// Returns a non-zero handle on success or a negative error code.
484///
485/// `on_frame` is called with a positive frame ID per frame, then exactly once
486/// more with a terminal code: `0` (closed cleanly) or a negative error. After
487/// the terminal (`<= 0`) callback, `on_frame` is never called again and
488/// `user_data` is never touched again, so release `user_data` there. The
489/// terminal callback fires even after [`moq_decode_audio_cancel`].
490///
491/// Starts at the newest cached group so reopening live playback skips the backlog.
492///
493/// A packet the codec cannot decode is logged and skipped rather than ending
494/// the subscription, so a single bad frame costs that frame and not the stream.
495///
496/// # Safety
497/// - `output` must point to a valid [`moq_audio_decoder_output`].
498/// - `user_data` must stay valid until the terminal (`<= 0`) `on_frame` callback.
499#[unsafe(no_mangle)]
500pub unsafe extern "C" fn moq_decode_audio(
501	catalog: u32,
502	index: u32,
503	output: *const moq_audio_decoder_output,
504	on_frame: crate::moq_status_callback,
505	user_data: *mut c_void,
506) -> i32 {
507	ffi::enter(move || {
508		let catalog = ffi::parse_id(catalog)?;
509		let raw = unsafe { output.as_ref() }.ok_or(Error::InvalidPointer)?;
510
511		let mut config = moq_audio::decode::Options::default();
512		config.start = moq_audio::decode::Start::Latest;
513		config.output.format = audio_format_from_u32(raw.format)?;
514		config.output.sample_rate = zeroable(raw.sample_rate);
515		config.output.layout = zeroable(raw.channels)
516			.map(moq_audio::Layout::from_channels)
517			.transpose()?;
518		config.max_age = Duration::from_micros(raw.max_age_us);
519
520		let on_frame = unsafe { OnStatus::new(user_data, on_frame)? };
521
522		let mut state = State::lock();
523		let (broadcast, audio_cfg, name) = state.consume.audio_rendition(catalog, index as usize)?;
524
525		let State { audio, .. } = &mut *state;
526		audio.consume(&broadcast, &audio_cfg, &name, config, on_frame)
527	})
528}
529
530/// Stop an audio (raw PCM) consumer's background task.
531///
532/// Returns immediately: zero on success, or a negative code if already closed.
533/// Does NOT free `user_data`; the on-frame callback still fires once more with a
534/// terminal `0` (or a negative error), which is where `user_data` should be
535/// released. Frame IDs already delivered to the callback are likewise not freed;
536/// release each with [`moq_decode_audio_frame_free`].
537#[unsafe(no_mangle)]
538pub extern "C" fn moq_decode_audio_cancel(consumer: u32) -> i32 {
539	ffi::enter(move || {
540		let consumer = ffi::parse_id(consumer)?;
541		State::lock().audio.consume_close(consumer)
542	})
543}
544
545/// Copy a delivered frame's metadata into `dst`.
546///
547/// The written `dst->data` pointer remains valid until the same `id`
548/// is released with [`moq_decode_audio_frame_free`].
549///
550/// # Safety
551/// - `dst` must point to a writable [`moq_audio_frame`].
552#[unsafe(no_mangle)]
553pub unsafe extern "C" fn moq_decode_audio_frame(id: u32, dst: *mut moq_audio_frame) -> i32 {
554	ffi::enter(move || {
555		let id = ffi::parse_id(id)?;
556		let dst = unsafe { dst.as_mut() }.ok_or(Error::InvalidPointer)?;
557		State::lock().audio.frame_info(id, dst)
558	})
559}
560
561/// Free a frame previously delivered through the consume callback.
562/// Required for every delivered frame ID; closing the parent consumer
563/// is not enough.
564#[unsafe(no_mangle)]
565pub extern "C" fn moq_decode_audio_frame_free(id: u32) -> i32 {
566	ffi::enter(move || {
567		let id = ffi::parse_id(id)?;
568		State::lock().audio.frame_free(id)
569	})
570}