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(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 surface = camera.read().await.expect("read camera frame");
860 if matches!(surface, Some(Surface::Texture(_))) {
861 textures += 1;
862 }
863 if i == 0 {
864 encoder.cut().unwrap();
865 }
866 let surface = surface.expect("frame, not end of stream");
867 frames.extend(encoder.encode(&Frame::new(surface, at(i))).unwrap());
868 }
869 frames.extend(encoder.finish().unwrap());
870
871 assert!(textures > 0, "capture never produced a GPU texture");
874 assert!(!frames.is_empty(), "encoder produced no packets");
875 let packets = payloads(&frames);
876 let types = nal_types(&packets[0]);
877 assert!(
878 types.contains(&7) && types.contains(&8) && types.contains(&5),
879 "no IDR: {types:?}"
880 );
881 }
882
883 #[test]
887 #[cfg(feature = "openh264")]
888 fn set_bitrate_retunes_software_encoder() {
889 let config = Config {
890 kind: Kind::Software,
891 ..Config::new(320, 240, crate::Rate::new(30, 1).unwrap())
892 };
893 let mut encoder = Encoder::new(&config).unwrap();
894
895 let opened = encoder.bitrate();
896 assert_eq!(opened, config.resolved_bitrate());
897
898 encoder.encode(&gray_frame(320, 240, 0)).unwrap();
900
901 let halved = opened.scaled(0.5);
902 encoder.set_bitrate(halved).unwrap();
903 assert_eq!(encoder.bitrate(), halved);
904
905 let frames = encoder.encode(&gray_frame(320, 240, 1)).unwrap();
907 assert!(!frames.is_empty(), "encoder produced nothing after a retune");
908 let packets = payloads(&frames);
909 assert!(packets[0].starts_with(&[0, 0, 0, 1]) || packets[0].starts_with(&[0, 0, 1]));
910 }
911
912 #[test]
916 #[cfg(feature = "openh264")]
917 fn set_bitrate_before_the_first_frame_is_deferred() {
918 let config = Config {
919 kind: Kind::Software,
920 ..Config::new(320, 240, crate::Rate::new(30, 1).unwrap())
921 };
922 let mut encoder = Encoder::new(&config).unwrap();
923
924 let halved = encoder.bitrate().scaled(0.5);
925 encoder.set_bitrate(halved).expect("a retune before the first frame");
926 assert_eq!(encoder.bitrate(), halved);
927
928 let frames = encoder.encode(&gray_frame(320, 240, 0)).unwrap();
930 assert!(!frames.is_empty());
931
932 encoder.set_bitrate(halved.scaled(0.5)).unwrap();
934 assert!(encoder.encode(&gray_frame(320, 240, 1)).is_ok());
935 }
936
937 #[test]
940 #[cfg(feature = "openh264")]
941 fn set_bitrate_to_current_is_a_noop() {
942 let config = Config {
943 kind: Kind::Software,
944 ..Config::new(320, 240, crate::Rate::new(30, 1).unwrap())
945 };
946 let mut encoder = Encoder::new(&config).unwrap();
947
948 let opened = encoder.bitrate();
949 encoder.set_bitrate(opened).unwrap();
950 assert_eq!(encoder.bitrate(), opened);
951 }
952
953 #[test]
954 fn default_bitrate_scales_with_resolution() {
955 let small = Config::new(320, 240, crate::Rate::new(30, 1).unwrap()).resolved_bitrate();
956 let large = Config::new(1920, 1080, crate::Rate::new(30, 1).unwrap()).resolved_bitrate();
957 assert!(large > small);
958 assert!(small > moq_net::bandwidth::Rate::ZERO);
959 }
960
961 struct Delayed {
966 pending: Option<Encoded>,
967 }
968
969 impl Backend for Delayed {
970 fn encode(&mut self, frame: &Frame, _cut: bool) -> Result<Vec<Encoded>, Error> {
971 let payload = bytes::Bytes::from(frame.timestamp.as_micros().to_string());
972 let previous = self.pending.replace(Encoded::new(payload, frame.timestamp));
973 Ok(previous.into_iter().collect())
974 }
975
976 fn flush(&mut self) -> Result<Vec<Encoded>, Error> {
977 Ok(self.pending.take().into_iter().collect())
978 }
979
980 fn finish(&mut self) -> Result<Vec<Encoded>, Error> {
981 self.flush()
982 }
983
984 fn set_bitrate(&mut self, _bitrate: u64) -> Result<(), Error> {
985 Ok(())
986 }
987
988 fn can_cut(&self) -> bool {
989 true
990 }
991
992 fn name(&self) -> &'static str {
993 "delayed"
994 }
995 }
996
997 fn encoder_with(backend: Box<dyn Backend>, config: &Config) -> Encoder {
1000 Encoder {
1001 backend,
1002 codec: config.codec,
1003 size: config.size(),
1004 bitrate: config.resolved_bitrate(),
1005 color: config.resolved_color(),
1006 pending_cut: false,
1007 _thread_bound: PhantomData,
1008 }
1009 }
1010
1011 struct Recorder {
1016 log: std::sync::Arc<std::sync::Mutex<Vec<bool>>>,
1017 cuts: bool,
1018 }
1019
1020 impl Backend for Recorder {
1021 fn encode(&mut self, _frame: &Frame, cut: bool) -> Result<Vec<Encoded>, Error> {
1022 self.log.lock().unwrap().push(cut);
1023 Ok(Vec::new())
1024 }
1025
1026 fn flush(&mut self) -> Result<Vec<Encoded>, Error> {
1027 Ok(Vec::new())
1028 }
1029
1030 fn finish(&mut self) -> Result<Vec<Encoded>, Error> {
1031 Ok(Vec::new())
1032 }
1033
1034 fn set_bitrate(&mut self, _bitrate: u64) -> Result<(), Error> {
1035 Ok(())
1036 }
1037
1038 fn can_cut(&self) -> bool {
1039 self.cuts
1040 }
1041
1042 fn name(&self) -> &'static str {
1043 "recorder"
1044 }
1045 }
1046
1047 #[test]
1052 fn a_cut_waits_for_the_next_frame_then_clears() {
1053 let log = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
1054 let config = Config::new(320, 240, crate::Rate::new(30, 1).unwrap());
1055 let backend = Recorder {
1056 log: log.clone(),
1057 cuts: true,
1058 };
1059 let mut encoder = encoder_with(Box::new(backend), &config);
1060
1061 encoder.encode(&gray_frame(320, 240, 0)).unwrap();
1063
1064 encoder.cut().unwrap();
1066 encoder.cut().unwrap();
1067 encoder.encode(&gray_frame(320, 240, 1)).unwrap();
1068
1069 encoder.encode(&gray_frame(320, 240, 2)).unwrap();
1071
1072 assert_eq!(*log.lock().unwrap(), vec![false, true, false]);
1073 }
1074
1075 #[test]
1079 fn every_cut_reaches_the_codec() {
1080 let log = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
1081 let config = Config::new(320, 240, crate::Rate::new(30, 1).unwrap());
1082 let backend = Recorder {
1083 log: log.clone(),
1084 cuts: true,
1085 };
1086 let mut encoder = encoder_with(Box::new(backend), &config);
1087
1088 for i in 0..6 {
1089 if i % 3 == 0 {
1090 encoder.cut().unwrap();
1091 }
1092 encoder.encode(&gray_frame(320, 240, i)).unwrap();
1093 }
1094
1095 assert_eq!(*log.lock().unwrap(), vec![true, false, false, true, false, false]);
1096 }
1097
1098 #[test]
1104 fn a_backend_that_cannot_cut_refuses_up_front() {
1105 let log = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
1106 let config = Config::new(320, 240, crate::Rate::new(30, 1).unwrap());
1107 let backend = Recorder {
1108 log: log.clone(),
1109 cuts: false,
1110 };
1111 let mut encoder = encoder_with(Box::new(backend), &config);
1112
1113 let err = encoder.cut().expect_err("the backend cannot cut");
1114 assert!(
1115 matches!(err, Error::CutUnsupported("recorder")),
1116 "unexpected error: {err:?}"
1117 );
1118 assert!(!encoder.pending_cut, "a refused cut must not be queued");
1119
1120 assert!(err.to_string().contains("recorder"), "{err}");
1122
1123 encoder.encode(&gray_frame(320, 240, 0)).unwrap();
1126 assert_eq!(*log.lock().unwrap(), vec![false]);
1127 }
1128
1129 #[test]
1134 fn a_zero_keyframe_interval_is_refused() {
1135 let config = Config {
1136 gop: Gop::Keyframe { interval: 0 },
1137 ..Config::new(320, 240, crate::Rate::new(30, 1).unwrap())
1138 };
1139 let Err(err) = Encoder::new(&config) else {
1140 panic!("an empty group cannot open");
1141 };
1142 assert!(matches!(err, Error::InvalidGop(_)), "unexpected error: {err:?}");
1143 }
1144
1145 #[test]
1149 fn a_sub_frame_keyframe_window_rounds_up_to_one_frame() {
1150 let framerate = crate::Rate::new(30, 1).unwrap();
1151 assert_eq!(
1152 Gop::keyframe_every(std::time::Duration::from_millis(1), framerate),
1153 Gop::Keyframe { interval: 1 }
1154 );
1155 assert_eq!(
1156 Gop::keyframe_every(std::time::Duration::from_secs(2), framerate),
1157 Gop::Keyframe { interval: 60 }
1158 );
1159 assert!(
1160 Gop::keyframe_every(std::time::Duration::ZERO, framerate)
1161 .validate()
1162 .is_ok()
1163 );
1164 }
1165
1166 #[test]
1171 #[cfg(feature = "openh264")]
1172 fn a_mid_stream_cut_emits_an_idr() {
1173 let config = Config {
1174 kind: Kind::Software,
1175 ..Config::new(320, 240, crate::Rate::new(30, 1).unwrap())
1176 };
1177 let mut encoder = Encoder::new(&Config {
1180 gop: Gop::Keyframe { interval: 1000 },
1181 ..config
1182 })
1183 .unwrap();
1184
1185 let mut per_frame = Vec::new();
1186 for i in 0..6 {
1187 if i == 3 {
1189 encoder.cut().unwrap();
1190 }
1191 let encoded = encoder.encode(&gray_frame(320, 240, i)).unwrap();
1192 let joined: Vec<u8> = encoded.iter().flat_map(|f| f.payload.iter()).copied().collect();
1193 per_frame.push(nal_types(&joined));
1194 }
1195
1196 let asked = &per_frame[3];
1199 assert!(asked.contains(&5), "the requested frame is not an IDR: {asked:?}");
1200 assert!(asked.contains(&7), "no SPS with the requested IDR: {asked:?}");
1201 assert!(asked.contains(&8), "no PPS with the requested IDR: {asked:?}");
1202
1203 for i in [1, 2, 4, 5] {
1207 assert!(
1208 !per_frame[i].contains(&5),
1209 "frame {i} was keyed without being asked: {:?}",
1210 per_frame[i]
1211 );
1212 }
1213 }
1214
1215 struct Failing;
1217
1218 impl Backend for Failing {
1219 fn encode(&mut self, _frame: &Frame, _cut: bool) -> Result<Vec<Encoded>, Error> {
1220 Err(Error::Codec(anyhow::anyhow!("no")))
1221 }
1222
1223 fn flush(&mut self) -> Result<Vec<Encoded>, Error> {
1224 Ok(Vec::new())
1225 }
1226
1227 fn finish(&mut self) -> Result<Vec<Encoded>, Error> {
1228 Ok(Vec::new())
1229 }
1230
1231 fn set_bitrate(&mut self, _bitrate: u64) -> Result<(), Error> {
1232 Ok(())
1233 }
1234
1235 fn can_cut(&self) -> bool {
1236 true
1237 }
1238
1239 fn name(&self) -> &'static str {
1240 "failing"
1241 }
1242 }
1243
1244 #[test]
1249 fn a_failed_encode_keeps_the_cut() {
1250 let config = Config::new(320, 240, crate::Rate::new(30, 1).unwrap());
1251
1252 let mut encoder = encoder_with(Box::new(Failing), &config);
1253 encoder.cut().unwrap();
1254 assert!(encoder.encode(&gray_frame(320, 240, 0)).is_err());
1255 assert!(encoder.pending_cut, "the backend error swallowed the request");
1256
1257 let mut encoder = encoder_with(Box::new(Delayed { pending: None }), &config);
1259 encoder.cut().unwrap();
1260 assert!(encoder.encode(&gray_frame(640, 480, 0)).is_err());
1261 assert!(encoder.pending_cut, "the size check swallowed the request");
1262
1263 encoder.encode(&gray_frame(320, 240, 1)).unwrap();
1265 assert!(!encoder.pending_cut);
1266 }
1267
1268 #[test]
1274 fn a_buffering_backend_keeps_each_frames_timestamp() {
1275 let config = Config::new(320, 240, crate::Rate::new(30, 1).unwrap());
1276 let mut encoder = encoder_with(Box::new(Delayed { pending: None }), &config);
1277
1278 for i in 0..5 {
1282 let encoded = encoder.encode(&gray_frame(320, 240, i)).unwrap();
1283 if i == 0 {
1284 assert!(encoded.is_empty(), "the first frame is still buffered");
1285 continue;
1286 }
1287 assert_eq!(encoded.len(), 1);
1288 assert_eq!(encoded[0].timestamp, at(i - 1));
1289 assert_eq!(&encoded[0].payload[..], at(i - 1).as_micros().to_string().as_bytes());
1292 }
1293
1294 let tail = encoder.finish().unwrap();
1297 assert_eq!(tail.len(), 1);
1298 assert_eq!(tail[0].timestamp, at(4));
1299 assert!(
1300 encoder_with(Box::new(Delayed { pending: None }), &config)
1301 .finish()
1302 .unwrap()
1303 .is_empty()
1304 );
1305 }
1306
1307 #[test]
1315 fn a_flush_empties_a_pipelined_backend_and_leaves_it_running() {
1316 let config = Config::new(320, 240, crate::Rate::new(30, 1).unwrap());
1317 let mut encoder = encoder_with(Box::new(Delayed { pending: None }), &config);
1318
1319 let mut group = Vec::new();
1320 for i in 0..3 {
1321 group.extend(encoder.encode(&gray_frame(320, 240, i)).unwrap());
1322 }
1323 group.extend(encoder.flush().unwrap());
1324
1325 let times: Vec<_> = group.iter().map(|packet| packet.timestamp).collect();
1327 assert_eq!(times, vec![at(0), at(1), at(2)], "the group lost or reordered frames");
1328
1329 let mut next = encoder.encode(&gray_frame(320, 240, 3)).unwrap();
1332 assert!(next.is_empty(), "frame 3 is buffered, so nothing comes back yet");
1333 next.extend(encoder.finish().unwrap());
1334 let times: Vec<_> = next.iter().map(|packet| packet.timestamp).collect();
1335 assert_eq!(times, vec![at(3)], "the flush left something behind");
1336 }
1337
1338 #[cfg(target_os = "macos")]
1343 #[test]
1344 fn videotoolbox_sps_declares_the_color_space() {
1345 use super::backend::test_util::{BT601_DESCRIBED, BT709_DESCRIBED, declared_color};
1346
1347 for (size, described) in [
1348 (Size::new(640, 480), BT601_DESCRIBED),
1349 (Size::new(1920, 1080), BT709_DESCRIBED),
1350 ] {
1351 let config = Config {
1352 kind: Kind::Named("videotoolbox".into()),
1353 ..Config::new(size.width, size.height, crate::Rate::new(30, 1).unwrap())
1354 };
1355 let mut encoder = Encoder::new(&config).expect("videotoolbox is available on macOS");
1356
1357 let rgba = [255u8, 0, 0, 255].repeat(size.pixels() as usize);
1358 let surface = crate::frame::Surface::rgba(&rgba, size).unwrap();
1359 encoder.cut().unwrap();
1360 let frames = encoder
1361 .encode(&Frame::new(surface, moq_net::Timestamp::from_micros(0).unwrap()))
1362 .unwrap();
1363
1364 let keyframe = frames.first().expect("a keyframe");
1365 assert_eq!(declared_color(&keyframe.payload), Some(described), "{size} SPS");
1366 }
1367 }
1368
1369 #[test]
1374 #[cfg(feature = "openh264")]
1375 fn config_color_pins_the_space_a_resize_carried() {
1376 use crate::Color;
1377
1378 let big = Size::new(1280, 720);
1379 let small = Size::new(640, 480);
1380
1381 let rgba = vec![0x80u8; big.pixels() as usize * 4];
1383 let frame = Frame::new(
1384 crate::frame::Surface::rgba(&rgba, big).unwrap(),
1385 moq_net::Timestamp::from_micros(0).unwrap(),
1386 );
1387 let scaled = frame.resize(small, &crate::resize::Config::default()).unwrap();
1388 assert_eq!(
1389 scaled.surface.color(),
1390 Some(Color::Bt709Limited),
1391 "resize keeps the space"
1392 );
1393
1394 let config = Config {
1398 kind: Kind::Software,
1399 ..Config::new(small.width, small.height, crate::Rate::new(30, 1).unwrap())
1400 };
1401 let mut encoder = Encoder::new(&config).unwrap();
1402 encoder.cut().unwrap();
1403 let frames = encoder.encode(&scaled).expect("a mismatch warns rather than fails");
1404 use super::backend::test_util::{BT601_DESCRIBED, BT709_DESCRIBED, declared_color};
1405 assert_eq!(
1406 declared_color(&frames.first().expect("a keyframe").payload),
1407 Some(BT601_DESCRIBED),
1408 "the inferred label is the wrong one, which is the case Config::color covers"
1409 );
1410
1411 let config = Config {
1413 kind: Kind::Software,
1414 color: Some(Color::Bt709Limited),
1415 ..Config::new(small.width, small.height, crate::Rate::new(30, 1).unwrap())
1416 };
1417 let mut encoder = Encoder::new(&config).unwrap();
1418 encoder.cut().unwrap();
1419 let frames = encoder.encode(&scaled).expect("a declared space encodes");
1420
1421 let keyframe = frames.first().expect("a keyframe");
1422 assert_eq!(declared_color(&keyframe.payload), Some(BT709_DESCRIBED));
1423 }
1424}