1use std::ffi::{c_char, c_void};
16use std::sync::Arc;
17use std::sync::atomic::{AtomicU64, Ordering};
18use std::time::{Duration, Instant};
19
20use tokio::sync::oneshot;
21
22use crate::ffi::OnStatus;
23use crate::{Error, Id, NonZeroSlab, Shared, State, ffi};
24
25#[repr(C)]
33#[allow(non_camel_case_types)]
34#[derive(Clone, Copy, Debug)]
35pub enum moq_video_pixel_format {
36 MOQ_VIDEO_PIXEL_FORMAT_I420 = 0,
40 MOQ_VIDEO_PIXEL_FORMAT_RGBA = 1,
42}
43
44#[repr(C)]
49#[allow(non_camel_case_types)]
50#[derive(Clone, Copy, Debug)]
51pub enum moq_video_codec {
52 MOQ_VIDEO_CODEC_H264 = 0,
54 MOQ_VIDEO_CODEC_H265 = 1,
56}
57
58#[repr(C)]
60#[allow(non_camel_case_types)]
61#[derive(Clone, Copy, Debug)]
62pub enum moq_video_encoder_kind {
63 MOQ_VIDEO_ENCODER_KIND_AUTO = 0,
65 MOQ_VIDEO_ENCODER_KIND_HARDWARE = 1,
67 MOQ_VIDEO_ENCODER_KIND_SOFTWARE = 2,
69 MOQ_VIDEO_ENCODER_KIND_NAMED = 3,
71}
72
73#[repr(C)]
77#[allow(non_camel_case_types)]
78pub struct moq_video_encoder_input {
79 pub format: u32,
81 pub width: u32,
83 pub height: u32,
85 pub framerate: u32,
88}
89
90#[repr(C)]
93#[allow(non_camel_case_types)]
94pub struct moq_video_encoder_output {
95 pub codec: u32,
97 pub bitrate: u64,
100 pub gop: u32,
103 pub kind: u32,
105 pub encoder: *const c_char,
108 pub encoder_len: usize,
109}
110
111#[repr(C)]
119#[allow(non_camel_case_types)]
120pub struct moq_video_encoder_frame {
121 pub timestamp_us: u64,
123 pub data: *const u8,
124 pub data_size: usize,
125}
126
127#[repr(C)]
138#[allow(non_camel_case_types)]
139pub struct moq_video_decoder_output {
140 pub max_age_us: u64,
145 pub format: u32,
148 pub width: u32,
150 pub height: u32,
152}
153
154#[repr(C)]
168#[allow(non_camel_case_types)]
169pub struct moq_video_frame {
170 pub timestamp_us: u64,
171 pub width: u32,
172 pub height: u32,
173 pub data: *const u8,
174 pub data_size: usize,
175}
176
177#[derive(Default)]
182pub struct Video {
183 producers: NonZeroSlab<Shared<VideoEncoder>>,
184 consumer_tasks: NonZeroSlab<Option<VideoTaskEntry>>,
185 frames: NonZeroSlab<VideoFrame>,
186}
187
188fn block_on<T>(future: impl std::future::Future<Output = T>) -> T {
198 pollster::block_on(future)
199}
200
201fn video_ceiling(
202 config: &moq_video::encode::Config,
203 rendition: &hang::catalog::VideoConfig,
204) -> moq_net::bandwidth::Rate {
205 config
206 .bitrate
207 .or_else(|| rendition.bitrate.map(moq_net::bandwidth::Rate::from_bps))
208 .unwrap_or_else(|| {
209 moq_net::bandwidth::Rate::from_bps(
210 (config.size().pixels() as f64 * config.framerate.as_f64() * 0.07) as u64,
211 )
212 })
213}
214
215async fn follow_reservation(
216 inner: Shared<VideoEncoder>,
217 mut consumer: moq_net::bandwidth::Consumer,
218 ceiling: Arc<AtomicU64>,
219) {
220 use moq_mux::rate::{Control, Policy};
221
222 let mut max = moq_net::bandwidth::Rate::from_bps(ceiling.load(Ordering::SeqCst));
223 let mut control = Control::new(Policy::new(max));
224 loop {
225 let estimate = match consumer.changed().await {
226 Ok(estimate) => estimate,
227 Err(_) => return,
228 };
229 let next = moq_net::bandwidth::Rate::from_bps(ceiling.load(Ordering::SeqCst));
230 if next != max {
231 max = next;
232 control = Control::new(Policy::new(max));
233 }
234 let Some(bitrate) = control.update(estimate, Instant::now()) else {
235 continue;
236 };
237
238 let mut guard = inner.lock();
239 let Some(producer) = guard.as_mut() else {
240 return;
241 };
242 match block_on(producer.encoder.set_bitrate(bitrate)) {
243 Ok(()) => tracing::debug!(bitrate = bitrate.as_bps(), "adjusted encoder bitrate"),
244 Err(moq_video::Error::BitrateUnsupported(name)) => {
245 tracing::warn!(encoder = name, "encoder cannot follow the bandwidth estimate");
246 return;
247 }
248 Err(err) => {
249 tracing::warn!(error = %err, bitrate = bitrate.as_bps(), "failed to adjust encoder bitrate");
250 }
251 }
252 }
253}
254
255pub(crate) struct VideoEncoder {
265 encoder: moq_video::encode::Sink,
266 producer: moq_video::encode::Producer<moq_mux::catalog::hang::Extra>,
267 format: moq_video_pixel_format,
268 size: moq_video::Size,
271 reservation: Option<Arc<moq_net::bandwidth::Reservation>>,
272 follow: Option<oneshot::Sender<()>>,
273 ceiling: Option<Arc<AtomicU64>>,
274}
275
276struct VideoFrame {
280 timestamp_us: u64,
281 width: u32,
282 height: u32,
283 data: bytes::Bytes,
284}
285
286#[derive(Clone, Copy)]
289pub struct DecoderOutput {
290 format: moq_video_pixel_format,
291 size: Option<moq_video::Size>,
292}
293
294fn finalize(
301 mut producer: moq_video::encode::Producer<moq_mux::catalog::hang::Extra>,
302 drained: Result<(), moq_video::Error>,
303) -> Result<(), Error> {
304 match drained {
305 Ok(()) => Ok(producer.finish()?),
306 Err(err) => {
307 producer.abort(moq_net::Error::Transport(err.to_string()));
308 Err(err.into())
309 }
310 }
311}
312
313struct VideoTaskEntry {
320 close: Option<oneshot::Sender<()>>,
321 callback: OnStatus,
322}
323
324impl VideoEncoder {
325 fn publish_frame(&mut self, timestamp_us: u64, data: &[u8]) -> Result<(), Error> {
326 let size = self.size;
329 let surface = match self.format {
330 moq_video_pixel_format::MOQ_VIDEO_PIXEL_FORMAT_I420 => {
331 moq_video::Surface::I420(moq_video::I420::new(size, data.to_vec())?)
332 }
333 moq_video_pixel_format::MOQ_VIDEO_PIXEL_FORMAT_RGBA => moq_video::Surface::rgba(data, size)?,
334 };
335
336 let frame = moq_video::Frame::new(surface, moq_net::Timestamp::from_micros(timestamp_us)?);
337 let encoded = block_on(self.encoder.encode(frame))?;
340 self.producer.publish(&encoded)?;
341 Ok(())
342 }
343
344 fn publish_cut(&mut self) -> Result<(), Error> {
345 Ok(block_on(self.encoder.cut())?)
348 }
349
350 fn publish_bitrate(&mut self, bitrate: u64) -> Result<(), Error> {
351 block_on(self.encoder.set_bitrate(moq_net::bandwidth::Rate::from_bps(bitrate)))?;
352 if let Some(ceiling) = &self.ceiling {
355 ceiling.store(bitrate, Ordering::SeqCst);
356 }
357 if let Some(reservation) = &self.reservation {
358 reservation.update(moq_net::bandwidth::Rate::from_bps(bitrate));
359 }
360 Ok(())
361 }
362
363 fn publish_finish(self) -> Result<(), Error> {
364 let VideoEncoder {
365 encoder, mut producer, ..
366 } = self;
367 let drained = block_on(encoder.finish()).and_then(|encoded| producer.publish(&encoded));
370 finalize(producer, drained)
371 }
372}
373
374impl Video {
375 pub fn publish(
381 &mut self,
382 broadcast: &moq_net::broadcast::Producer,
383 catalog: moq_mux::catalog::Producer<moq_mux::catalog::hang::Extra>,
384 format: moq_video_pixel_format,
385 config: &moq_video::encode::Config,
386 rendition: hang::catalog::VideoConfig,
387 encoder: moq_video::encode::Sink,
388 ) -> Result<Id, Error> {
389 let producer = moq_video::encode::Producer::new(broadcast.clone(), catalog, rendition)?;
390 self.producers.insert(Shared::new(VideoEncoder {
391 encoder,
392 producer,
393 format,
394 size: config.size(),
395 reservation: None,
396 follow: None,
397 ceiling: None,
398 }))
399 }
400
401 pub(crate) fn follow(
402 &self,
403 id: Id,
404 allocator: &moq_net::bandwidth::Allocator,
405 ceiling: moq_net::bandwidth::Rate,
406 ) -> Result<(), Error> {
407 let shared = self.producer(id)?;
408 let max = ceiling;
409 let ceiling = Arc::new(AtomicU64::new(max.as_bps()));
410 let reservation = {
411 let mut guard = shared.lock();
412 let encoder = guard.as_mut().ok_or(Error::MediaNotFound)?;
413 let reservation = Arc::new(allocator.reserve(&encoder.producer.demand(), max));
414 encoder.reservation = Some(reservation.clone());
415 encoder.ceiling = Some(ceiling.clone());
416 reservation
417 };
418 let (close, closed) = oneshot::channel();
419 shared.lock().as_mut().ok_or(Error::MediaNotFound)?.follow = Some(close);
420 let follower = shared.clone();
421 tokio::spawn(async move {
422 tokio::select! {
423 biased;
424 _ = closed => {}
425 _ = follow_reservation(follower, reservation.consumer(), ceiling) => {}
426 }
427 });
428 Ok(())
429 }
430
431 pub(crate) fn reservation(&self, id: Id) -> Result<Option<Arc<moq_net::bandwidth::Reservation>>, Error> {
432 Ok(self
433 .producer(id)?
434 .lock()
435 .as_ref()
436 .ok_or(Error::MediaNotFound)?
437 .reservation
438 .clone())
439 }
440
441 pub(crate) fn demand(&self, id: Id) -> Result<moq_net::track::Demand, Error> {
443 Ok(self
444 .producer(id)?
445 .lock()
446 .as_ref()
447 .ok_or(Error::MediaNotFound)?
448 .producer
449 .demand())
450 }
451
452 pub(crate) fn producer(&self, id: Id) -> Result<Shared<VideoEncoder>, Error> {
459 self.producers.get(id).cloned().ok_or(Error::MediaNotFound)
460 }
461
462 pub(crate) fn remove(&mut self, id: Id) -> Result<Shared<VideoEncoder>, Error> {
464 self.producers.remove(id).ok_or(Error::MediaNotFound)
465 }
466
467 pub fn consume(
468 &mut self,
469 broadcast: &moq_net::broadcast::Consumer,
470 catalog: &hang::catalog::VideoConfig,
471 name: &str,
472 options: moq_video::decode::Options,
473 output: DecoderOutput,
474 on_frame: OnStatus,
475 ) -> Result<Id, Error> {
476 let broadcast = broadcast.clone();
477 let catalog = catalog.clone();
478 let name = name.to_string();
479
480 let channel = oneshot::channel();
481 let entry = VideoTaskEntry {
482 close: Some(channel.0),
483 callback: on_frame,
484 };
485 let id = self.consumer_tasks.insert(Some(entry))?;
486
487 tokio::spawn(async move {
490 let res = async move {
491 let consumer = moq_video::decode::Consumer::new(&broadcast, &catalog, name, options).await?;
492 Self::run(on_frame, consumer, channel.1, output).await
493 }
494 .await;
495
496 let entry = State::lock().video.consumer_tasks.remove(id).flatten();
499 if let Some(entry) = entry {
500 entry.callback.call(res);
501 }
502 });
503
504 Ok(id)
505 }
506
507 async fn run(
508 callback: OnStatus,
509 mut consumer: moq_video::decode::Consumer,
510 mut close: oneshot::Receiver<()>,
511 output: DecoderOutput,
512 ) -> Result<(), Error> {
513 loop {
514 let frame = tokio::select! {
516 biased;
517 _ = &mut close => return Ok(()),
518 frame = consumer.read() => match frame? {
519 Some(frame) => frame,
520 None => return Ok(()),
521 },
522 };
523
524 let mut frame = frame;
530 if let Some(size) = output.size
531 && frame.size() != size
532 {
533 frame = frame.resize(size, &moq_video::resize::Config::default())?;
534 }
535 let size = frame.size();
536 let data = match output.format {
537 moq_video_pixel_format::MOQ_VIDEO_PIXEL_FORMAT_I420 => {
538 bytes::Bytes::from(frame.surface.into_i420()?.into_data())
539 }
540 moq_video_pixel_format::MOQ_VIDEO_PIXEL_FORMAT_RGBA => bytes::Bytes::from(
541 frame
542 .surface
543 .to_rgba(&moq_video::convert::Config::default())?
544 .into_data(),
545 ),
546 };
547 let frame = VideoFrame {
548 timestamp_us: frame.timestamp.as_micros() as u64,
551 width: size.width,
552 height: size.height,
553 data,
554 };
555 let frame_id = State::lock().video.frames.insert(frame)?;
556 callback.call(Ok(frame_id));
557 }
558 }
559
560 pub fn consume_close(&mut self, id: Id) -> Result<(), Error> {
561 self.consumer_tasks
563 .get_mut(id)
564 .and_then(|entry| entry.as_mut())
565 .ok_or(Error::TrackNotFound)?
566 .close
567 .take()
568 .ok_or(Error::TrackNotFound)?;
569 Ok(())
570 }
571
572 pub fn frame_info(&self, id: Id, dst: &mut moq_video_frame) -> Result<(), Error> {
573 let frame = self.frames.get(id).ok_or(Error::FrameNotFound)?;
574 *dst = moq_video_frame {
575 timestamp_us: frame.timestamp_us,
576 width: frame.width,
577 height: frame.height,
578 data: frame.data.as_ptr(),
579 data_size: frame.data.len(),
580 };
581 Ok(())
582 }
583
584 pub fn frame_free(&mut self, id: Id) -> Result<(), Error> {
585 self.frames.remove(id).ok_or(Error::FrameNotFound)?;
586 Ok(())
587 }
588}
589
590fn pixel_format_from_u32(value: u32) -> Result<moq_video_pixel_format, Error> {
593 Ok(match value {
594 v if v == moq_video_pixel_format::MOQ_VIDEO_PIXEL_FORMAT_I420 as u32 => {
595 moq_video_pixel_format::MOQ_VIDEO_PIXEL_FORMAT_I420
596 }
597 v if v == moq_video_pixel_format::MOQ_VIDEO_PIXEL_FORMAT_RGBA as u32 => {
598 moq_video_pixel_format::MOQ_VIDEO_PIXEL_FORMAT_RGBA
599 }
600 _ => return Err(Error::InvalidCode),
601 })
602}
603
604fn decoder_size(width: u32, height: u32) -> Result<Option<moq_video::Size>, Error> {
608 if width == 0 && height == 0 {
609 return Ok(None);
610 }
611 if width == 0 || height == 0 || !width.is_multiple_of(2) || !height.is_multiple_of(2) {
612 return Err(Error::InvalidConfig(format!(
613 "decode size {width}x{height}: use 0x0 for the native size or even non-zero dimensions"
614 )));
615 }
616 let size = moq_video::Size::new(width, height);
617 if size
618 .pixels()
619 .checked_mul(4)
620 .is_none_or(|bytes| usize::try_from(bytes).is_err())
621 {
622 return Err(Error::InvalidConfig(format!(
623 "decode size {width}x{height}: dimensions too large to represent"
624 )));
625 }
626 Ok(Some(size))
627}
628
629fn codec_from_u32(value: u32) -> Result<moq_video::encode::Codec, Error> {
630 use moq_video::encode::Codec;
631 Ok(match value {
632 v if v == moq_video_codec::MOQ_VIDEO_CODEC_H264 as u32 => Codec::H264,
633 v if v == moq_video_codec::MOQ_VIDEO_CODEC_H265 as u32 => Codec::H265,
634 _ => return Err(Error::InvalidCode),
635 })
636}
637
638unsafe fn encoder_kind(output: &moq_video_encoder_output) -> Result<moq_video::encode::Kind, Error> {
642 use moq_video::encode::Kind;
643 Ok(match output.kind {
644 v if v == moq_video_encoder_kind::MOQ_VIDEO_ENCODER_KIND_AUTO as u32 => Kind::Auto,
645 v if v == moq_video_encoder_kind::MOQ_VIDEO_ENCODER_KIND_HARDWARE as u32 => Kind::Hardware,
646 v if v == moq_video_encoder_kind::MOQ_VIDEO_ENCODER_KIND_SOFTWARE as u32 => Kind::Software,
647 v if v == moq_video_encoder_kind::MOQ_VIDEO_ENCODER_KIND_NAMED as u32 => {
648 Kind::Named(unsafe { ffi::parse_str(output.encoder, output.encoder_len)? }.to_string())
649 }
650 _ => return Err(Error::InvalidCode),
651 })
652}
653
654#[unsafe(no_mangle)]
671pub unsafe extern "C" fn moq_encode_video(
672 broadcast: u32,
673 input: *const moq_video_encoder_input,
674 output: *const moq_video_encoder_output,
675 bandwidth: u32,
676) -> i32 {
677 ffi::enter(move || {
678 let broadcast = ffi::parse_id(broadcast)?;
679 let raw_input = unsafe { input.as_ref() }.ok_or(Error::InvalidPointer)?;
680 let raw_output = unsafe { output.as_ref() }.ok_or(Error::InvalidPointer)?;
681
682 let format = pixel_format_from_u32(raw_input.format)?;
683
684 let framerate = moq_video::Rate::new(raw_input.framerate, 1)
685 .map_err(|_| Error::Video(moq_video::Error::InvalidFramerate(raw_input.framerate).into()))?;
686 let mut config = moq_video::encode::Config::new(raw_input.width, raw_input.height, framerate);
687 config.codec = codec_from_u32(raw_output.codec)?;
688 config.kind = unsafe { encoder_kind(raw_output)? };
689 config.bitrate = (raw_output.bitrate != 0).then(|| moq_net::bandwidth::Rate::from_bps(raw_output.bitrate));
692 if raw_output.gop != 0 {
693 config.gop = moq_video::encode::Gop::Keyframe {
694 interval: raw_output.gop,
695 };
696 }
697
698 let rendition = block_on(config.probe())?;
702 let encoder = block_on(moq_video::encode::Sink::open(&config))?;
703
704 let bandwidth = ffi::parse_id_optional(bandwidth)?;
705 let mut state = State::lock();
706 let allocator = bandwidth.map(|id| state.bandwidth.allocator(id)).transpose()?;
707 let State { publish, video, .. } = &mut *state;
708 let (broadcast_producer, catalog) = publish.pair_mut(broadcast)?;
709
710 let id = video.publish(
711 broadcast_producer,
712 catalog.clone(),
713 format,
714 &config,
715 rendition.clone(),
716 encoder,
717 )?;
718 if let Some(allocator) = allocator.as_ref() {
719 video.follow(id, allocator, video_ceiling(&config, &rendition))?;
720 }
721 Ok(id)
722 })
723}
724
725#[unsafe(no_mangle)]
730pub extern "C" fn moq_encode_video_reservation(producer: u32) -> i32 {
731 ffi::enter(move || {
732 let producer = ffi::parse_id(producer)?;
733 let mut state = State::lock();
734 match state.video.reservation(producer)? {
735 Some(reservation) => Ok(i32::from(state.bandwidth.hold(reservation)?)),
736 None => Ok(0),
737 }
738 })
739}
740
741#[unsafe(no_mangle)]
751pub unsafe extern "C" fn moq_encode_video_demand(
752 producer: u32,
753 on_demand: crate::moq_status_callback,
754 user_data: *mut c_void,
755) -> i32 {
756 ffi::enter(move || {
757 let producer = ffi::parse_id(producer)?;
758 let on_demand = unsafe { OnStatus::new(user_data, on_demand)? };
759 let mut state = State::lock();
760 let demand = state.video.demand(producer)?;
761 state.publish.demand(demand, on_demand)
762 })
763}
764
765#[unsafe(no_mangle)]
777pub unsafe extern "C" fn moq_encode_video_frame(producer: u32, frame: *const moq_video_encoder_frame) -> i32 {
778 ffi::enter(move || {
779 let producer = ffi::parse_id(producer)?;
780 let frame = unsafe { frame.as_ref() }.ok_or(Error::InvalidPointer)?;
781 let data = unsafe { ffi::parse_slice(frame.data, frame.data_size)? };
782
783 let producer = State::lock().video.producer(producer)?;
784 producer
785 .lock()
786 .as_mut()
787 .ok_or(Error::MediaNotFound)?
788 .publish_frame(frame.timestamp_us, data)
789 })
790}
791
792#[unsafe(no_mangle)]
808pub extern "C" fn moq_encode_video_cut(producer: u32) -> i32 {
809 ffi::enter(move || {
810 let producer = ffi::parse_id(producer)?;
811 let producer = State::lock().video.producer(producer)?;
812 producer.lock().as_mut().ok_or(Error::MediaNotFound)?.publish_cut()
813 })
814}
815
816#[unsafe(no_mangle)]
832pub extern "C" fn moq_encode_video_bitrate(producer: u32, bitrate: u64) -> i32 {
833 ffi::enter(move || {
834 let producer = ffi::parse_id(producer)?;
835 let producer = State::lock().video.producer(producer)?;
836 producer
837 .lock()
838 .as_mut()
839 .ok_or(Error::MediaNotFound)?
840 .publish_bitrate(bitrate)
841 })
842}
843
844#[unsafe(no_mangle)]
848pub extern "C" fn moq_encode_video_finish(producer: u32) -> i32 {
849 ffi::enter(move || {
850 let producer = ffi::parse_id(producer)?;
851 let producer = State::lock().video.remove(producer)?;
854 producer.take().ok_or(Error::MediaNotFound)?.publish_finish()
855 })
856}
857
858#[unsafe(no_mangle)]
883pub unsafe extern "C" fn moq_decode_video(
884 catalog: u32,
885 index: u32,
886 output: *const moq_video_decoder_output,
887 on_frame: crate::moq_status_callback,
888 user_data: *mut c_void,
889) -> i32 {
890 ffi::enter(move || {
891 let raw = unsafe { output.as_ref() }.ok_or(Error::InvalidPointer)?;
892
893 let format = pixel_format_from_u32(raw.format)?;
897 let size = decoder_size(raw.width, raw.height)?;
898 let catalog = ffi::parse_id(catalog)?;
899
900 let mut options = moq_video::decode::Options::new();
901 options.start = moq_video::decode::Start::Latest;
902 options.max_age = Duration::from_micros(raw.max_age_us);
903 options.decoder.output = moq_video::Output::Cpu;
906 options.decoder.scale_hint = size;
909 let output = DecoderOutput { format, size };
910 let on_frame = unsafe { OnStatus::new(user_data, on_frame)? };
911
912 let mut state = State::lock();
913 let (broadcast, video_cfg, name) = state.consume.video_rendition(catalog, index as usize)?;
914
915 let State { video, .. } = &mut *state;
916 video.consume(&broadcast, &video_cfg, &name, options, output, on_frame)
917 })
918}
919
920#[unsafe(no_mangle)]
928pub extern "C" fn moq_decode_video_cancel(consumer: u32) -> i32 {
929 ffi::enter(move || {
930 let consumer = ffi::parse_id(consumer)?;
931 State::lock().video.consume_close(consumer)
932 })
933}
934
935#[unsafe(no_mangle)]
943pub unsafe extern "C" fn moq_decode_video_frame(id: u32, dst: *mut moq_video_frame) -> i32 {
944 ffi::enter(move || {
945 let id = ffi::parse_id(id)?;
946 let dst = unsafe { dst.as_mut() }.ok_or(Error::InvalidPointer)?;
947 State::lock().video.frame_info(id, dst)
948 })
949}
950
951#[unsafe(no_mangle)]
954pub extern "C" fn moq_decode_video_frame_free(id: u32) -> i32 {
955 ffi::enter(move || {
956 let id = ffi::parse_id(id)?;
957 State::lock().video.frame_free(id)
958 })
959}
960#[cfg(test)]
961mod tests {
962 use super::*;
963
964 async fn track_under_test() -> (
967 moq_video::encode::Producer<moq_mux::catalog::hang::Extra>,
968 moq_net::track::Subscriber,
969 ) {
970 let mut broadcast = moq_net::broadcast::Info::new().produce();
971 let config = moq_mux::catalog::Config::default()
972 .with_catalog(moq_mux::catalog::hang::Catalog::<moq_mux::catalog::hang::Extra>::default());
973 let catalog = moq_mux::catalog::Producer::new(&mut broadcast, config).unwrap();
974 let consumer = broadcast.consume();
975 let rendition = moq_video::encode::Config::new(320, 240, moq_video::Rate::new(30, 1).unwrap())
977 .probe()
978 .await
979 .unwrap();
980 let producer = moq_video::encode::Producer::new(broadcast, catalog, rendition).unwrap();
981
982 let name = producer.demand().name().to_string();
983 let track = consumer.track(&name).unwrap().subscribe(None).await.unwrap();
984 (producer, track)
985 }
986
987 #[tokio::test]
990 async fn a_successful_drain_ends_the_track_cleanly() {
991 let (producer, mut track) = track_under_test().await;
992 finalize(producer, Ok(())).unwrap();
993 assert!(matches!(track.recv_group().await, Ok(None)), "expected a clean end");
994 }
995
996 #[tokio::test]
1000 async fn a_failed_drain_aborts_the_track() {
1001 let (producer, mut track) = track_under_test().await;
1002 let err = moq_video::Error::Codec(anyhow::anyhow!("the codec lost the tail"));
1003 finalize(producer, Err(err)).unwrap_err();
1004
1005 let Err(err) = track.recv_group().await else {
1006 panic!("expected an abort, not a clean end");
1007 };
1008 assert!(
1009 err.to_string().contains("the codec lost the tail"),
1010 "the abort should carry the drain failure: {err}"
1011 );
1012 }
1013}