1use std::marker::PhantomData;
9use std::rc::Rc;
10
11use super::Encoded;
12use super::backend::{self, Backend};
13use crate::{Color, Error, Frame, Rate, Size};
14
15#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
23#[non_exhaustive]
24pub enum Codec {
25 #[default]
28 H264,
29 H265,
31}
32
33#[derive(Clone, Debug, Default, PartialEq, Eq)]
36#[non_exhaustive]
37pub enum Kind {
38 #[default]
41 Auto,
42 Hardware,
44 Software,
46 Named(String),
49}
50
51#[derive(Clone, Copy, Debug, PartialEq, Eq)]
58#[non_exhaustive]
59pub enum Gop {
60 Keyframe {
64 interval: u32,
66 },
67}
68
69impl Gop {
70 pub fn keyframe_every(interval: std::time::Duration, framerate: Rate) -> Self {
73 Self::Keyframe {
74 interval: framerate.frames(interval).max(1),
75 }
76 }
77
78 pub(crate) fn validate(&self) -> Result<(), Error> {
80 match *self {
81 Self::Keyframe { interval: 0 } => Err(Error::InvalidGop(
82 "a keyframe interval of 0 frames opens no group".to_owned(),
83 )),
84 Self::Keyframe { .. } => Ok(()),
85 }
86 }
87}
88
89#[derive(Clone, Debug)]
95#[non_exhaustive]
96pub struct Config {
97 pub width: u32,
98 pub height: u32,
99 pub framerate: Rate,
100 pub bitrate: Option<moq_net::bandwidth::Rate>,
103 pub gop: Gop,
106 pub codec: Codec,
108 pub kind: Kind,
109 pub color: Option<Color>,
117}
118
119impl Config {
120 pub fn new(width: u32, height: u32, framerate: Rate) -> Self {
123 Self {
124 width,
125 height,
126 framerate,
127 bitrate: None,
128 gop: Gop::keyframe_every(std::time::Duration::from_secs(2), framerate),
129 codec: Codec::default(),
130 kind: Kind::Auto,
131 color: None,
132 }
133 }
134
135 pub fn size(&self) -> Size {
137 Size::new(self.width, self.height)
138 }
139
140 pub async fn probe(&self) -> Result<hang::catalog::VideoConfig, Error> {
161 let mut sink = super::Sink::open(self).await?;
164
165 let size = self.size();
167 let i420 = crate::I420::new(size, vec![0x80u8; crate::I420::len(size)?])?;
168 let frame = Frame::new(crate::Surface::I420(i420), moq_net::Timestamp::from_micros(0)?);
169
170 let mut encoded = sink.encode(frame).await?;
173 if encoded.is_empty() {
175 encoded = sink.flush().await?;
176 }
177
178 let annexb: Vec<u8> = encoded.iter().flat_map(|frame| frame.payload.iter().copied()).collect();
179 let parsed = match self.codec {
180 Codec::H264 => moq_mux::codec::h264::config(&annexb),
181 Codec::H265 => moq_mux::codec::h265::config(&annexb),
182 };
183 let mut rendition = parsed.map_err(|err| {
184 Error::Codec(anyhow::anyhow!(
185 "{} emitted no usable parameter sets: {err}",
186 sink.name()
187 ))
188 })?;
189
190 rendition.bitrate.get_or_insert(self.resolved_bitrate().as_bps());
193 rendition.framerate.get_or_insert(self.framerate.as_f64());
194 Ok(rendition)
195 }
196
197 pub(crate) fn resolved_color(&self) -> Color {
202 self.color.unwrap_or_else(|| Color::infer(self.size()))
203 }
204
205 pub(crate) fn resolved_bitrate(&self) -> moq_net::bandwidth::Rate {
207 self.bitrate
208 .unwrap_or_else(|| default_bitrate(self.size(), self.framerate))
209 }
210}
211
212pub(crate) fn default_bitrate(size: Size, framerate: Rate) -> moq_net::bandwidth::Rate {
215 moq_net::bandwidth::Rate::from_bps((size.pixels() as f64 * framerate.as_f64() * 0.07) as u64)
216}
217
218pub struct Encoder {
231 backend: Box<dyn Backend>,
232 codec: Codec,
233 size: Size,
234 bitrate: moq_net::bandwidth::Rate,
235 color: Color,
238 pending_cut: bool,
242 _thread_bound: PhantomData<Rc<()>>,
244}
245
246impl Encoder {
247 pub fn new(config: &Config) -> Result<Self, Error> {
249 let size = config.size();
251 size.validate("encoder")?;
252 size.validate_encodable("encoder", config.framerate)?;
253 config.gop.validate()?;
254
255 let backend = backend::open(config)?;
256 Ok(Self {
257 backend,
258 codec: config.codec,
259 size,
260 bitrate: config.resolved_bitrate(),
261 color: config.resolved_color(),
262 pending_cut: false,
263 _thread_bound: PhantomData,
264 })
265 }
266
267 pub fn name(&self) -> &str {
269 self.backend.name()
270 }
271
272 pub fn size(&self) -> Size {
274 self.size
275 }
276
277 pub fn bitrate(&self) -> moq_net::bandwidth::Rate {
280 self.bitrate
281 }
282
283 pub fn set_bitrate(&mut self, bitrate: moq_net::bandwidth::Rate) -> Result<(), Error> {
297 if bitrate == self.bitrate {
298 return Ok(());
299 }
300 self.backend.set_bitrate(bitrate.as_bps())?;
303 self.bitrate = bitrate;
306 Ok(())
307 }
308
309 pub fn codec(&self) -> Codec {
312 self.codec
313 }
314
315 pub fn cut(&mut self) -> Result<(), Error> {
336 if !self.backend.can_cut() {
337 return Err(Error::CutUnsupported(self.backend.name()));
338 }
339 self.pending_cut = true;
340 Ok(())
341 }
342
343 pub fn encode(&mut self, frame: &Frame) -> Result<Vec<Encoded>, Error> {
359 let size = frame.size();
362 if size != self.size {
363 return Err(Error::Codec(anyhow::anyhow!(
364 "frame {size} does not match encoder {}",
365 self.size
366 )));
367 }
368 if let Some(color) = frame.surface.color()
378 && color != self.color
379 {
380 static WARN_ONCE: std::sync::Once = std::sync::Once::new();
381 WARN_ONCE.call_once(|| {
382 tracing::warn!(
383 frame = ?color,
384 encoder = ?self.color,
385 "frame color space differs from the one written into the bitstream; set encode::Config::color"
386 );
387 });
388 }
389 let encoded = self.backend.encode(frame, self.pending_cut)?;
390 self.pending_cut = false;
394 Ok(encoded)
395 }
396
397 pub fn flush(&mut self) -> Result<Vec<Encoded>, Error> {
412 self.backend.flush()
413 }
414
415 pub fn finish(mut self) -> Result<Vec<Encoded>, Error> {
422 self.backend.finish()
423 }
424}
425
426#[cfg(test)]
427mod tests {
428 #![cfg_attr(not(feature = "openh264"), allow(dead_code, unused_imports))]
429
430 use super::*;
431
432 use crate::{I420, Surface};
433
434 fn gray_rgba(width: u32, height: u32) -> Vec<u8> {
436 vec![0x80u8; width as usize * height as usize * 4]
437 }
438
439 fn gray_frame(width: u32, height: u32, index: u64) -> Frame {
442 let surface = Surface::rgba(&gray_rgba(width, height), Size::new(width, height)).unwrap();
443 Frame::new(surface, at(index))
444 }
445
446 fn at(index: u64) -> moq_net::Timestamp {
448 moq_net::Timestamp::from_micros(index * 33_333).unwrap()
449 }
450
451 fn payloads(frames: &[Encoded]) -> Vec<bytes::Bytes> {
453 frames.iter().map(|f| f.payload.clone()).collect()
454 }
455
456 #[test]
457 #[cfg(feature = "openh264")]
458 fn software_encoder_emits_annexb() {
459 let config = Config {
460 kind: Kind::Software,
461 ..Config::new(320, 240, crate::Rate::new(30, 1).unwrap())
462 };
463 let mut encoder = Encoder::new(&config).expect("openh264 is vendored, always available");
464 assert_eq!(encoder.name(), "openh264");
465
466 let mut frames = Vec::new();
467 for i in 0..30 {
468 if i == 0 {
469 encoder.cut().unwrap();
470 }
471 frames.extend(encoder.encode(&gray_frame(320, 240, i)).unwrap());
472 }
473 frames.extend(encoder.finish().unwrap());
474
475 assert!(!frames.is_empty(), "encoder produced no packets");
476
477 let micros: Vec<u128> = frames.iter().map(|f| f.timestamp.as_micros()).collect();
480 assert!(
481 micros.windows(2).all(|w| w[0] < w[1]),
482 "encoded timestamps not strictly increasing: {micros:?}"
483 );
484 assert!(
485 micros.iter().all(|&t| t % 33_333 == 0 && t < 30 * 33_333),
486 "encoded timestamp outside the fed set: {micros:?}"
487 );
488
489 let packets = payloads(&frames);
492 let first = &packets[0];
493 let has_start_code = first.starts_with(&[0, 0, 0, 1]) || first.starts_with(&[0, 0, 1]);
494 assert!(
495 has_start_code,
496 "first packet is not Annex-B: {:02x?}",
497 &first[..first.len().min(8)]
498 );
499 }
500
501 #[test]
503 #[cfg(feature = "openh264")]
504 fn encode_rgba_surface_emits_annexb() {
505 let config = Config {
506 kind: Kind::Software,
507 ..Config::new(320, 240, crate::Rate::new(30, 1).unwrap())
508 };
509 let mut encoder = Encoder::new(&config).unwrap();
510
511 let mut frames = encoder.encode(&gray_frame(320, 240, 0)).unwrap();
512 frames.extend(encoder.finish().unwrap());
513 assert!(!frames.is_empty());
514 let packets = payloads(&frames);
515 assert!(packets[0].starts_with(&[0, 0, 0, 1]) || packets[0].starts_with(&[0, 0, 1]));
516 }
517
518 #[test]
520 #[cfg(feature = "openh264")]
521 fn encode_i420_surface_emits_annexb() {
522 let config = Config {
523 kind: Kind::Software,
524 ..Config::new(320, 240, crate::Rate::new(30, 1).unwrap())
525 };
526 let mut encoder = Encoder::new(&config).unwrap();
527
528 let size = Size::new(320, 240);
530 let i420 = I420::new(size, vec![0x80u8; I420::len(size).unwrap()]).unwrap();
531 let frame = Frame::new(Surface::I420(i420), at(0));
532 let mut frames = encoder.encode(&frame).unwrap();
533 frames.extend(encoder.finish().unwrap());
534 assert!(!frames.is_empty());
535 let packets = payloads(&frames);
536 assert!(packets[0].starts_with(&[0, 0, 0, 1]) || packets[0].starts_with(&[0, 0, 1]));
537 }
538
539 #[test]
542 fn encode_rejects_dimension_mismatch() {
543 let Ok(mut encoder) = Encoder::new(&Config::new(320, 240, crate::Rate::new(30, 1).unwrap())) else {
544 return;
545 };
546 assert!(matches!(encoder.encode(&gray_frame(640, 480, 0)), Err(Error::Codec(_))));
547 }
548
549 #[test]
553 fn encode_rejects_transposed_frame() {
554 let Ok(mut encoder) = Encoder::new(&Config::new(320, 240, crate::Rate::new(30, 1).unwrap())) else {
555 return;
556 };
557
558 let transposed = gray_frame(240, 320, 0);
559 assert_eq!(
560 gray_rgba(240, 320).len(),
561 gray_rgba(320, 240).len(),
562 "the byte counts must collide"
563 );
564 assert!(matches!(encoder.encode(&transposed), Err(Error::Codec(_))));
565 }
566
567 #[test]
568 fn unknown_named_encoder_errors() {
569 let config = Config {
570 kind: Kind::Named("definitely_not_a_codec".into()),
571 ..Config::new(320, 240, crate::Rate::new(30, 1).unwrap())
572 };
573 let err = Encoder::new(&config).err().expect("an unknown name cannot open");
577 assert!(
578 matches!(&err, Error::UnknownEncoder { name, .. } if name == "definitely_not_a_codec"),
579 "unexpected error: {err:?}",
580 );
581 }
582
583 #[cfg(target_os = "macos")]
587 #[test]
588 fn videotoolbox_emits_annexb_keyframe() {
589 let config = Config {
590 kind: Kind::Named("videotoolbox".into()),
591 ..Config::new(320, 240, crate::Rate::new(30, 1).unwrap())
592 };
593 let mut encoder = Encoder::new(&config).expect("videotoolbox is available on macOS");
594 assert_eq!(encoder.name(), "videotoolbox");
595
596 let mut frames = Vec::new();
597 for i in 0..10 {
598 if i == 0 {
599 encoder.cut().unwrap();
600 }
601 frames.extend(encoder.encode(&gray_frame(320, 240, i)).unwrap());
602 }
603 frames.extend(encoder.finish().unwrap());
604
605 assert!(!frames.is_empty(), "encoder produced no packets");
606 let micros: Vec<u128> = frames.iter().map(|f| f.timestamp.as_micros()).collect();
609 assert!(
610 micros.windows(2).all(|w| w[0] < w[1]),
611 "encoded timestamps not strictly increasing: {micros:?}"
612 );
613
614 let packets = payloads(&frames);
615 let first = &packets[0];
616 assert!(
617 first.starts_with(&[0, 0, 0, 1]) || first.starts_with(&[0, 0, 1]),
618 "first packet is not Annex-B"
619 );
620
621 let types = nal_types(first);
624 assert!(types.contains(&7), "no SPS in first packet: {types:?}");
625 assert!(types.contains(&8), "no PPS in first packet: {types:?}");
626 assert!(types.contains(&5), "first packet is not an IDR: {types:?}");
627 }
628
629 #[cfg(target_os = "macos")]
633 #[test]
634 fn videotoolbox_emits_annexb_keyframe_h265() {
635 let config = Config {
636 codec: Codec::H265,
637 kind: Kind::Named("videotoolbox".into()),
638 ..Config::new(320, 240, crate::Rate::new(30, 1).unwrap())
639 };
640 let mut encoder = Encoder::new(&config).expect("videotoolbox HEVC is available on macOS");
641 assert_eq!(encoder.name(), "videotoolbox");
642 assert_eq!(encoder.codec(), Codec::H265);
643
644 let mut frames = Vec::new();
645 for i in 0..10 {
646 if i == 0 {
647 encoder.cut().unwrap();
648 }
649 frames.extend(encoder.encode(&gray_frame(320, 240, i)).unwrap());
650 }
651 frames.extend(encoder.finish().unwrap());
652
653 assert!(!frames.is_empty(), "encoder produced no packets");
654 let packets = payloads(&frames);
655 let first = &packets[0];
656 assert!(
657 first.starts_with(&[0, 0, 0, 1]) || first.starts_with(&[0, 0, 1]),
658 "first packet is not Annex-B"
659 );
660
661 let types = hevc_nal_types(first);
664 assert!(types.contains(&32), "no VPS in first packet: {types:?}");
665 assert!(types.contains(&33), "no SPS in first packet: {types:?}");
666 assert!(types.contains(&34), "no PPS in first packet: {types:?}");
667 assert!(
668 types.iter().any(|t| (16..=23).contains(t)),
669 "first packet is not an IRAP: {types:?}"
670 );
671 }
672
673 #[cfg(target_os = "macos")]
675 fn hevc_nal_types(annexb: &[u8]) -> Vec<u8> {
676 let mut types = Vec::new();
677 let mut i = 0;
678 while i + 3 < annexb.len() {
679 if annexb[i..i + 3] == [0, 0, 1] {
680 types.push((annexb[i + 3] >> 1) & 0x3f);
681 i += 3;
682 } else {
683 i += 1;
684 }
685 }
686 types
687 }
688
689 #[cfg(target_os = "macos")]
692 #[test]
693 fn videotoolbox_encodes_surface_zero_copy() {
694 let config = Config {
695 kind: Kind::Named("videotoolbox".into()),
696 ..Config::new(320, 240, crate::Rate::new(30, 1).unwrap())
697 };
698 let mut encoder = Encoder::new(&config).unwrap();
699
700 let mut frames = Vec::new();
701 for i in 0..10 {
702 if i == 0 {
703 encoder.cut().unwrap();
704 }
705 let frame = Frame::new(Surface::PixelBuffer(nv12_surface(320, 240)), at(i));
706 frames.extend(encoder.encode(&frame).unwrap());
707 }
708 frames.extend(encoder.finish().unwrap());
709
710 assert!(!frames.is_empty());
711 let packets = payloads(&frames);
712 let types = nal_types(&packets[0]);
713 assert!(
714 types.contains(&7) && types.contains(&8) && types.contains(&5),
715 "no IDR: {types:?}"
716 );
717 }
718
719 #[cfg(all(target_os = "macos", feature = "openh264"))]
722 #[test]
723 fn openh264_downloads_surface() {
724 let config = Config {
725 kind: Kind::Software,
726 ..Config::new(320, 240, crate::Rate::new(30, 1).unwrap())
727 };
728 let mut encoder = Encoder::new(&config).unwrap();
729
730 encoder.cut().unwrap();
731 let frame = Frame::new(Surface::PixelBuffer(nv12_surface(320, 240)), at(0));
732 let mut frames = encoder.encode(&frame).unwrap();
733 frames.extend(encoder.finish().unwrap());
734
735 assert!(!frames.is_empty());
736 let packets = payloads(&frames);
737 assert!(packets[0].starts_with(&[0, 0, 0, 1]) || packets[0].starts_with(&[0, 0, 1]));
738 }
739
740 #[cfg(target_os = "macos")]
743 fn nv12_surface(width: u32, height: u32) -> crate::frame::macos::PixelBuffer {
744 use std::ptr::{self, NonNull};
745
746 use objc2_core_foundation::CFRetained;
747 use objc2_core_video::{
748 CVPixelBuffer, CVPixelBufferCreate, CVPixelBufferGetBaseAddressOfPlane, CVPixelBufferGetBytesPerRowOfPlane,
749 CVPixelBufferLockBaseAddress, CVPixelBufferLockFlags, CVPixelBufferUnlockBaseAddress,
750 kCVPixelFormatType_420YpCbCr8BiPlanarVideoRange,
751 };
752
753 let mut raw: *mut CVPixelBuffer = ptr::null_mut();
754 let status = unsafe {
755 CVPixelBufferCreate(
756 None,
757 width as usize,
758 height as usize,
759 kCVPixelFormatType_420YpCbCr8BiPlanarVideoRange,
760 None,
761 NonNull::new(&mut raw).unwrap(),
762 )
763 };
764 assert_eq!(status, 0, "CVPixelBufferCreate failed");
765 let buffer = unsafe { CFRetained::from_raw(NonNull::new(raw).unwrap()) };
766
767 let flags = CVPixelBufferLockFlags(0);
768 assert_eq!(unsafe { CVPixelBufferLockBaseAddress(&buffer, flags) }, 0);
769 for (plane, rows) in [(0usize, height as usize), (1usize, height as usize / 2)] {
770 let base = CVPixelBufferGetBaseAddressOfPlane(&buffer, plane) as *mut u8;
771 let stride = CVPixelBufferGetBytesPerRowOfPlane(&buffer, plane);
772 unsafe { ptr::write_bytes(base, 128, stride * rows) };
773 }
774 unsafe { CVPixelBufferUnlockBaseAddress(&buffer, flags) };
775
776 crate::frame::macos::PixelBuffer::new(buffer, width, height)
777 }
778
779 fn nal_types(annexb: &[u8]) -> Vec<u8> {
782 let mut types = Vec::new();
783 let mut i = 0;
784 while i + 3 < annexb.len() {
785 if annexb[i..i + 3] == [0, 0, 1] {
786 types.push(annexb[i + 3] & 0x1f);
787 i += 3;
788 } else {
789 i += 1;
790 }
791 }
792 types
793 }
794
795 #[cfg(target_os = "windows")]
799 #[test]
800 #[ignore]
801 fn mediafoundation_cpu_rgba() {
802 let config = Config {
803 kind: Kind::Named("mediafoundation".into()),
804 ..Config::new(640, 480, crate::Rate::new(30, 1).unwrap())
805 };
806 let mut encoder = Encoder::new(&config).expect("hardware H.264 encoder available");
807 assert_eq!(encoder.name(), "mediafoundation");
808
809 let mut frames = Vec::new();
810 for i in 0..30 {
811 if i == 0 {
812 encoder.cut().unwrap();
813 }
814 frames.extend(encoder.encode(&gray_frame(640, 480, i)).unwrap());
815 }
816 frames.extend(encoder.finish().unwrap());
817
818 assert!(!frames.is_empty(), "encoder produced no packets");
819 let micros: Vec<u128> = frames.iter().map(|f| f.timestamp.as_micros()).collect();
822 assert!(
823 micros.windows(2).all(|w| w[0] < w[1]),
824 "encoded timestamps not strictly increasing: {micros:?}"
825 );
826 assert!(
827 micros.iter().all(|&t| t % 33_333 == 0 && t < 30 * 33_333),
828 "encoded timestamp outside the fed set: {micros:?}"
829 );
830
831 let packets = payloads(&frames);
832 let types = nal_types(&packets[0]);
833 assert!(types.contains(&7), "no SPS in first packet: {types:?}");
834 assert!(types.contains(&8), "no PPS in first packet: {types:?}");
835 assert!(types.contains(&5), "first packet is not an IDR: {types:?}");
836 }
837
838 #[cfg(all(target_os = "windows", feature = "capture"))]
842 #[tokio::test]
843 #[ignore]
844 async fn mediafoundation_camera_texture() {
845 let mut camera = crate::capture::open(&crate::capture::Config::default())
846 .await
847 .expect("open default camera");
848 let (w, h) = (camera.width(), camera.height());
849
850 let config = Config {
851 kind: Kind::Named("mediafoundation".into()),
852 ..Config::new(w, h, camera.framerate().unwrap_or(crate::Rate::integer(30)))
853 };
854 let mut encoder = Encoder::new(&config).expect("hardware H.264 encoder available");
855
856 let mut frames = Vec::new();
857 let mut textures = 0;
858 for i in 0..30 {
859 let frame = camera
860 .read()
861 .await
862 .expect("read camera frame")
863 .expect("frame, not end of stream");
864 if matches!(frame.surface, Surface::Texture(_)) {
865 textures += 1;
866 }
867 if i == 0 {
868 encoder.cut().unwrap();
869 }
870 frames.extend(encoder.encode(&frame).unwrap());
871 }
872 frames.extend(encoder.finish().unwrap());
873
874 assert!(textures > 0, "capture never produced a GPU texture");
877 assert!(!frames.is_empty(), "encoder produced no packets");
878 let packets = payloads(&frames);
879 let types = nal_types(&packets[0]);
880 assert!(
881 types.contains(&7) && types.contains(&8) && types.contains(&5),
882 "no IDR: {types:?}"
883 );
884 }
885
886 #[test]
890 #[cfg(feature = "openh264")]
891 fn set_bitrate_retunes_software_encoder() {
892 let config = Config {
893 kind: Kind::Software,
894 ..Config::new(320, 240, crate::Rate::new(30, 1).unwrap())
895 };
896 let mut encoder = Encoder::new(&config).unwrap();
897
898 let opened = encoder.bitrate();
899 assert_eq!(opened, config.resolved_bitrate());
900
901 encoder.encode(&gray_frame(320, 240, 0)).unwrap();
903
904 let halved = opened.scaled(0.5);
905 encoder.set_bitrate(halved).unwrap();
906 assert_eq!(encoder.bitrate(), halved);
907
908 let frames = encoder.encode(&gray_frame(320, 240, 1)).unwrap();
910 assert!(!frames.is_empty(), "encoder produced nothing after a retune");
911 let packets = payloads(&frames);
912 assert!(packets[0].starts_with(&[0, 0, 0, 1]) || packets[0].starts_with(&[0, 0, 1]));
913 }
914
915 #[test]
919 #[cfg(feature = "openh264")]
920 fn set_bitrate_before_the_first_frame_is_deferred() {
921 let config = Config {
922 kind: Kind::Software,
923 ..Config::new(320, 240, crate::Rate::new(30, 1).unwrap())
924 };
925 let mut encoder = Encoder::new(&config).unwrap();
926
927 let halved = encoder.bitrate().scaled(0.5);
928 encoder.set_bitrate(halved).expect("a retune before the first frame");
929 assert_eq!(encoder.bitrate(), halved);
930
931 let frames = encoder.encode(&gray_frame(320, 240, 0)).unwrap();
933 assert!(!frames.is_empty());
934
935 encoder.set_bitrate(halved.scaled(0.5)).unwrap();
937 assert!(encoder.encode(&gray_frame(320, 240, 1)).is_ok());
938 }
939
940 #[test]
943 #[cfg(feature = "openh264")]
944 fn set_bitrate_to_current_is_a_noop() {
945 let config = Config {
946 kind: Kind::Software,
947 ..Config::new(320, 240, crate::Rate::new(30, 1).unwrap())
948 };
949 let mut encoder = Encoder::new(&config).unwrap();
950
951 let opened = encoder.bitrate();
952 encoder.set_bitrate(opened).unwrap();
953 assert_eq!(encoder.bitrate(), opened);
954 }
955
956 #[test]
957 fn default_bitrate_scales_with_resolution() {
958 let small = Config::new(320, 240, crate::Rate::new(30, 1).unwrap()).resolved_bitrate();
959 let large = Config::new(1920, 1080, crate::Rate::new(30, 1).unwrap()).resolved_bitrate();
960 assert!(large > small);
961 assert!(small > moq_net::bandwidth::Rate::ZERO);
962 }
963
964 struct Delayed {
969 pending: Option<Encoded>,
970 }
971
972 impl Backend for Delayed {
973 fn encode(&mut self, frame: &Frame, _cut: bool) -> Result<Vec<Encoded>, Error> {
974 let payload = bytes::Bytes::from(frame.timestamp.as_micros().to_string());
975 let previous = self.pending.replace(Encoded::new(payload, frame.timestamp));
976 Ok(previous.into_iter().collect())
977 }
978
979 fn flush(&mut self) -> Result<Vec<Encoded>, Error> {
980 Ok(self.pending.take().into_iter().collect())
981 }
982
983 fn finish(&mut self) -> Result<Vec<Encoded>, Error> {
984 self.flush()
985 }
986
987 fn set_bitrate(&mut self, _bitrate: u64) -> Result<(), Error> {
988 Ok(())
989 }
990
991 fn can_cut(&self) -> bool {
992 true
993 }
994
995 fn name(&self) -> &'static str {
996 "delayed"
997 }
998 }
999
1000 fn encoder_with(backend: Box<dyn Backend>, config: &Config) -> Encoder {
1003 Encoder {
1004 backend,
1005 codec: config.codec,
1006 size: config.size(),
1007 bitrate: config.resolved_bitrate(),
1008 color: config.resolved_color(),
1009 pending_cut: false,
1010 _thread_bound: PhantomData,
1011 }
1012 }
1013
1014 struct Recorder {
1019 log: std::sync::Arc<std::sync::Mutex<Vec<bool>>>,
1020 cuts: bool,
1021 }
1022
1023 impl Backend for Recorder {
1024 fn encode(&mut self, _frame: &Frame, cut: bool) -> Result<Vec<Encoded>, Error> {
1025 self.log.lock().unwrap().push(cut);
1026 Ok(Vec::new())
1027 }
1028
1029 fn flush(&mut self) -> Result<Vec<Encoded>, Error> {
1030 Ok(Vec::new())
1031 }
1032
1033 fn finish(&mut self) -> Result<Vec<Encoded>, Error> {
1034 Ok(Vec::new())
1035 }
1036
1037 fn set_bitrate(&mut self, _bitrate: u64) -> Result<(), Error> {
1038 Ok(())
1039 }
1040
1041 fn can_cut(&self) -> bool {
1042 self.cuts
1043 }
1044
1045 fn name(&self) -> &'static str {
1046 "recorder"
1047 }
1048 }
1049
1050 #[test]
1055 fn a_cut_waits_for_the_next_frame_then_clears() {
1056 let log = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
1057 let config = Config::new(320, 240, crate::Rate::new(30, 1).unwrap());
1058 let backend = Recorder {
1059 log: log.clone(),
1060 cuts: true,
1061 };
1062 let mut encoder = encoder_with(Box::new(backend), &config);
1063
1064 encoder.encode(&gray_frame(320, 240, 0)).unwrap();
1066
1067 encoder.cut().unwrap();
1069 encoder.cut().unwrap();
1070 encoder.encode(&gray_frame(320, 240, 1)).unwrap();
1071
1072 encoder.encode(&gray_frame(320, 240, 2)).unwrap();
1074
1075 assert_eq!(*log.lock().unwrap(), vec![false, true, false]);
1076 }
1077
1078 #[test]
1082 fn every_cut_reaches_the_codec() {
1083 let log = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
1084 let config = Config::new(320, 240, crate::Rate::new(30, 1).unwrap());
1085 let backend = Recorder {
1086 log: log.clone(),
1087 cuts: true,
1088 };
1089 let mut encoder = encoder_with(Box::new(backend), &config);
1090
1091 for i in 0..6 {
1092 if i % 3 == 0 {
1093 encoder.cut().unwrap();
1094 }
1095 encoder.encode(&gray_frame(320, 240, i)).unwrap();
1096 }
1097
1098 assert_eq!(*log.lock().unwrap(), vec![true, false, false, true, false, false]);
1099 }
1100
1101 #[test]
1107 fn a_backend_that_cannot_cut_refuses_up_front() {
1108 let log = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
1109 let config = Config::new(320, 240, crate::Rate::new(30, 1).unwrap());
1110 let backend = Recorder {
1111 log: log.clone(),
1112 cuts: false,
1113 };
1114 let mut encoder = encoder_with(Box::new(backend), &config);
1115
1116 let err = encoder.cut().expect_err("the backend cannot cut");
1117 assert!(
1118 matches!(err, Error::CutUnsupported("recorder")),
1119 "unexpected error: {err:?}"
1120 );
1121 assert!(!encoder.pending_cut, "a refused cut must not be queued");
1122
1123 assert!(err.to_string().contains("recorder"), "{err}");
1125
1126 encoder.encode(&gray_frame(320, 240, 0)).unwrap();
1129 assert_eq!(*log.lock().unwrap(), vec![false]);
1130 }
1131
1132 #[test]
1137 fn a_zero_keyframe_interval_is_refused() {
1138 let config = Config {
1139 gop: Gop::Keyframe { interval: 0 },
1140 ..Config::new(320, 240, crate::Rate::new(30, 1).unwrap())
1141 };
1142 let Err(err) = Encoder::new(&config) else {
1143 panic!("an empty group cannot open");
1144 };
1145 assert!(matches!(err, Error::InvalidGop(_)), "unexpected error: {err:?}");
1146 }
1147
1148 #[test]
1152 fn a_sub_frame_keyframe_window_rounds_up_to_one_frame() {
1153 let framerate = crate::Rate::new(30, 1).unwrap();
1154 assert_eq!(
1155 Gop::keyframe_every(std::time::Duration::from_millis(1), framerate),
1156 Gop::Keyframe { interval: 1 }
1157 );
1158 assert_eq!(
1159 Gop::keyframe_every(std::time::Duration::from_secs(2), framerate),
1160 Gop::Keyframe { interval: 60 }
1161 );
1162 assert!(
1163 Gop::keyframe_every(std::time::Duration::ZERO, framerate)
1164 .validate()
1165 .is_ok()
1166 );
1167 }
1168
1169 #[test]
1174 #[cfg(feature = "openh264")]
1175 fn a_mid_stream_cut_emits_an_idr() {
1176 let config = Config {
1177 kind: Kind::Software,
1178 ..Config::new(320, 240, crate::Rate::new(30, 1).unwrap())
1179 };
1180 let mut encoder = Encoder::new(&Config {
1183 gop: Gop::Keyframe { interval: 1000 },
1184 ..config
1185 })
1186 .unwrap();
1187
1188 let mut per_frame = Vec::new();
1189 for i in 0..6 {
1190 if i == 3 {
1192 encoder.cut().unwrap();
1193 }
1194 let encoded = encoder.encode(&gray_frame(320, 240, i)).unwrap();
1195 let joined: Vec<u8> = encoded.iter().flat_map(|f| f.payload.iter()).copied().collect();
1196 per_frame.push(nal_types(&joined));
1197 }
1198
1199 let asked = &per_frame[3];
1202 assert!(asked.contains(&5), "the requested frame is not an IDR: {asked:?}");
1203 assert!(asked.contains(&7), "no SPS with the requested IDR: {asked:?}");
1204 assert!(asked.contains(&8), "no PPS with the requested IDR: {asked:?}");
1205
1206 for i in [1, 2, 4, 5] {
1210 assert!(
1211 !per_frame[i].contains(&5),
1212 "frame {i} was keyed without being asked: {:?}",
1213 per_frame[i]
1214 );
1215 }
1216 }
1217
1218 struct Failing;
1220
1221 impl Backend for Failing {
1222 fn encode(&mut self, _frame: &Frame, _cut: bool) -> Result<Vec<Encoded>, Error> {
1223 Err(Error::Codec(anyhow::anyhow!("no")))
1224 }
1225
1226 fn flush(&mut self) -> Result<Vec<Encoded>, Error> {
1227 Ok(Vec::new())
1228 }
1229
1230 fn finish(&mut self) -> Result<Vec<Encoded>, Error> {
1231 Ok(Vec::new())
1232 }
1233
1234 fn set_bitrate(&mut self, _bitrate: u64) -> Result<(), Error> {
1235 Ok(())
1236 }
1237
1238 fn can_cut(&self) -> bool {
1239 true
1240 }
1241
1242 fn name(&self) -> &'static str {
1243 "failing"
1244 }
1245 }
1246
1247 #[test]
1252 fn a_failed_encode_keeps_the_cut() {
1253 let config = Config::new(320, 240, crate::Rate::new(30, 1).unwrap());
1254
1255 let mut encoder = encoder_with(Box::new(Failing), &config);
1256 encoder.cut().unwrap();
1257 assert!(encoder.encode(&gray_frame(320, 240, 0)).is_err());
1258 assert!(encoder.pending_cut, "the backend error swallowed the request");
1259
1260 let mut encoder = encoder_with(Box::new(Delayed { pending: None }), &config);
1262 encoder.cut().unwrap();
1263 assert!(encoder.encode(&gray_frame(640, 480, 0)).is_err());
1264 assert!(encoder.pending_cut, "the size check swallowed the request");
1265
1266 encoder.encode(&gray_frame(320, 240, 1)).unwrap();
1268 assert!(!encoder.pending_cut);
1269 }
1270
1271 #[test]
1277 fn a_buffering_backend_keeps_each_frames_timestamp() {
1278 let config = Config::new(320, 240, crate::Rate::new(30, 1).unwrap());
1279 let mut encoder = encoder_with(Box::new(Delayed { pending: None }), &config);
1280
1281 for i in 0..5 {
1285 let encoded = encoder.encode(&gray_frame(320, 240, i)).unwrap();
1286 if i == 0 {
1287 assert!(encoded.is_empty(), "the first frame is still buffered");
1288 continue;
1289 }
1290 assert_eq!(encoded.len(), 1);
1291 assert_eq!(encoded[0].timestamp, at(i - 1));
1292 assert_eq!(&encoded[0].payload[..], at(i - 1).as_micros().to_string().as_bytes());
1295 }
1296
1297 let tail = encoder.finish().unwrap();
1300 assert_eq!(tail.len(), 1);
1301 assert_eq!(tail[0].timestamp, at(4));
1302 assert!(
1303 encoder_with(Box::new(Delayed { pending: None }), &config)
1304 .finish()
1305 .unwrap()
1306 .is_empty()
1307 );
1308 }
1309
1310 #[test]
1318 fn a_flush_empties_a_pipelined_backend_and_leaves_it_running() {
1319 let config = Config::new(320, 240, crate::Rate::new(30, 1).unwrap());
1320 let mut encoder = encoder_with(Box::new(Delayed { pending: None }), &config);
1321
1322 let mut group = Vec::new();
1323 for i in 0..3 {
1324 group.extend(encoder.encode(&gray_frame(320, 240, i)).unwrap());
1325 }
1326 group.extend(encoder.flush().unwrap());
1327
1328 let times: Vec<_> = group.iter().map(|packet| packet.timestamp).collect();
1330 assert_eq!(times, vec![at(0), at(1), at(2)], "the group lost or reordered frames");
1331
1332 let mut next = encoder.encode(&gray_frame(320, 240, 3)).unwrap();
1335 assert!(next.is_empty(), "frame 3 is buffered, so nothing comes back yet");
1336 next.extend(encoder.finish().unwrap());
1337 let times: Vec<_> = next.iter().map(|packet| packet.timestamp).collect();
1338 assert_eq!(times, vec![at(3)], "the flush left something behind");
1339 }
1340
1341 #[cfg(target_os = "macos")]
1346 #[test]
1347 fn videotoolbox_sps_declares_the_color_space() {
1348 use super::backend::test_util::{BT601_DESCRIBED, BT709_DESCRIBED, declared_color};
1349
1350 for (size, described) in [
1351 (Size::new(640, 480), BT601_DESCRIBED),
1352 (Size::new(1920, 1080), BT709_DESCRIBED),
1353 ] {
1354 let config = Config {
1355 kind: Kind::Named("videotoolbox".into()),
1356 ..Config::new(size.width, size.height, crate::Rate::new(30, 1).unwrap())
1357 };
1358 let mut encoder = Encoder::new(&config).expect("videotoolbox is available on macOS");
1359
1360 let rgba = [255u8, 0, 0, 255].repeat(size.pixels() as usize);
1361 let surface = crate::frame::Surface::rgba(&rgba, size).unwrap();
1362 encoder.cut().unwrap();
1363 let frames = encoder
1364 .encode(&Frame::new(surface, moq_net::Timestamp::from_micros(0).unwrap()))
1365 .unwrap();
1366
1367 let keyframe = frames.first().expect("a keyframe");
1368 assert_eq!(declared_color(&keyframe.payload), Some(described), "{size} SPS");
1369 }
1370 }
1371
1372 #[test]
1377 #[cfg(feature = "openh264")]
1378 fn config_color_pins_the_space_a_resize_carried() {
1379 use crate::Color;
1380
1381 let big = Size::new(1280, 720);
1382 let small = Size::new(640, 480);
1383
1384 let rgba = vec![0x80u8; big.pixels() as usize * 4];
1386 let frame = Frame::new(
1387 crate::frame::Surface::rgba(&rgba, big).unwrap(),
1388 moq_net::Timestamp::from_micros(0).unwrap(),
1389 );
1390 let scaled = frame.resize(small, &crate::resize::Config::default()).unwrap();
1391 assert_eq!(
1392 scaled.surface.color(),
1393 Some(Color::Bt709Limited),
1394 "resize keeps the space"
1395 );
1396
1397 let config = Config {
1401 kind: Kind::Software,
1402 ..Config::new(small.width, small.height, crate::Rate::new(30, 1).unwrap())
1403 };
1404 let mut encoder = Encoder::new(&config).unwrap();
1405 encoder.cut().unwrap();
1406 let frames = encoder.encode(&scaled).expect("a mismatch warns rather than fails");
1407 use super::backend::test_util::{BT601_DESCRIBED, BT709_DESCRIBED, declared_color};
1408 assert_eq!(
1409 declared_color(&frames.first().expect("a keyframe").payload),
1410 Some(BT601_DESCRIBED),
1411 "the inferred label is the wrong one, which is the case Config::color covers"
1412 );
1413
1414 let config = Config {
1416 kind: Kind::Software,
1417 color: Some(Color::Bt709Limited),
1418 ..Config::new(small.width, small.height, crate::Rate::new(30, 1).unwrap())
1419 };
1420 let mut encoder = Encoder::new(&config).unwrap();
1421 encoder.cut().unwrap();
1422 let frames = encoder.encode(&scaled).expect("a declared space encodes");
1423
1424 let keyframe = frames.first().expect("a keyframe");
1425 assert_eq!(declared_color(&keyframe.payload), Some(BT709_DESCRIBED));
1426 }
1427}