use std::collections::VecDeque;
use std::time::{Duration, Instant};
use bytes::{Bytes, BytesMut};
use moq_net::Timestamp;
use v4l::v4l_sys::{
V4L2_CID_MPEG_VIDEO_BITRATE, V4L2_CID_MPEG_VIDEO_BITRATE_MODE, V4L2_CID_MPEG_VIDEO_FORCE_KEY_FRAME,
V4L2_CID_MPEG_VIDEO_GOP_SIZE, V4L2_CID_MPEG_VIDEO_H264_LEVEL, V4L2_CID_MPEG_VIDEO_H264_PROFILE,
V4L2_CID_MPEG_VIDEO_HEADER_MODE, V4L2_CID_MPEG_VIDEO_PREPEND_SPSPPS_TO_IDR, V4L2_CID_MPEG_VIDEO_REPEAT_SEQ_HEADER,
V4L2_ENC_CMD_START, V4L2_ENC_CMD_STOP, v4l2_mpeg_video_bitrate_mode_V4L2_MPEG_VIDEO_BITRATE_MODE_CBR,
v4l2_mpeg_video_h264_level, v4l2_mpeg_video_h264_level_V4L2_MPEG_VIDEO_H264_LEVEL_1_0,
v4l2_mpeg_video_h264_level_V4L2_MPEG_VIDEO_H264_LEVEL_1_1,
v4l2_mpeg_video_h264_level_V4L2_MPEG_VIDEO_H264_LEVEL_1_2,
v4l2_mpeg_video_h264_level_V4L2_MPEG_VIDEO_H264_LEVEL_1_3,
v4l2_mpeg_video_h264_level_V4L2_MPEG_VIDEO_H264_LEVEL_2_0,
v4l2_mpeg_video_h264_level_V4L2_MPEG_VIDEO_H264_LEVEL_2_1,
v4l2_mpeg_video_h264_level_V4L2_MPEG_VIDEO_H264_LEVEL_2_2,
v4l2_mpeg_video_h264_level_V4L2_MPEG_VIDEO_H264_LEVEL_3_0,
v4l2_mpeg_video_h264_level_V4L2_MPEG_VIDEO_H264_LEVEL_3_1,
v4l2_mpeg_video_h264_level_V4L2_MPEG_VIDEO_H264_LEVEL_3_2,
v4l2_mpeg_video_h264_level_V4L2_MPEG_VIDEO_H264_LEVEL_4_0,
v4l2_mpeg_video_h264_level_V4L2_MPEG_VIDEO_H264_LEVEL_4_1,
v4l2_mpeg_video_h264_level_V4L2_MPEG_VIDEO_H264_LEVEL_4_2,
v4l2_mpeg_video_h264_level_V4L2_MPEG_VIDEO_H264_LEVEL_5_0,
v4l2_mpeg_video_h264_level_V4L2_MPEG_VIDEO_H264_LEVEL_5_1,
v4l2_mpeg_video_h264_profile_V4L2_MPEG_VIDEO_H264_PROFILE_CONSTRAINED_BASELINE,
v4l2_mpeg_video_header_mode_V4L2_MPEG_VIDEO_HEADER_MODE_JOINED_WITH_1ST_FRAME,
};
use super::super::encoder::Config;
use super::{Backend, Encoded};
use crate::v4l2::{self, Dequeue, Device, Dir, Planes, Queue, Rect, Request, Role};
use crate::{Error, Frame, Size};
pub(crate) const NAME: &str = "v4l2";
const ROLE: Role = Role {
env: "MOQ_V4L2_ENCODER",
input: v4l2::RAW,
output: &[v4l2::H264],
};
const RAW_BUFFERS: u32 = 4;
const CODED_BUFFERS: u32 = 8;
const BUFFER_TIMEOUT: Duration = Duration::from_millis(500);
const FLUSH_TIMEOUT: Duration = Duration::from_millis(500);
const POLL_INTERVAL: Duration = Duration::from_millis(5);
const CARRY_LIMIT: usize = 64 * 1024;
pub(crate) struct V4l2 {
device: Device,
raw: Queue,
coded: Queue,
planes: Planes,
size: Size,
pending: Pending,
drainable: bool,
keyframes: bool,
}
impl V4l2 {
pub(crate) fn open(config: &Config) -> Result<Box<dyn Backend>, Error> {
let size = config.size();
size.validate("V4L2 encode of")?;
let device = v4l2::open(&ROLE)?;
device.set_control(
V4L2_CID_MPEG_VIDEO_H264_PROFILE,
v4l2_mpeg_video_h264_profile_V4L2_MPEG_VIDEO_H264_PROFILE_CONSTRAINED_BASELINE as i32,
)?;
device.set_control(V4L2_CID_MPEG_VIDEO_H264_LEVEL, h264_level(config)?)?;
device.set_control(V4L2_CID_MPEG_VIDEO_GOP_SIZE, config.gop as i32)?;
set_bitrate(&device, config.resolved_bitrate())?;
device.try_control(
V4L2_CID_MPEG_VIDEO_BITRATE_MODE,
v4l2_mpeg_video_bitrate_mode_V4L2_MPEG_VIDEO_BITRATE_MODE_CBR as i32,
);
device.try_control(
V4L2_CID_MPEG_VIDEO_HEADER_MODE,
v4l2_mpeg_video_header_mode_V4L2_MPEG_VIDEO_HEADER_MODE_JOINED_WITH_1ST_FRAME as i32,
);
let repeat = device.try_control(V4L2_CID_MPEG_VIDEO_REPEAT_SEQ_HEADER, 1)
| device.try_control(V4L2_CID_MPEG_VIDEO_PREPEND_SPSPPS_TO_IDR, 1);
if !repeat {
tracing::warn!(
encoder = NAME,
device = %device.path().display(),
"driver repeats no parameter sets; subscribers can only join at the first keyframe"
);
}
let coded = device.set_format(
Dir::Capture,
&Request {
pixelformat: v4l2::H264,
size,
sizeimage: Some(coded_size(size)),
color: None,
},
)?;
if coded.pixelformat != v4l2::H264 {
return Err(Error::Codec(anyhow::anyhow!(
"V4L2 encoder answered an H264 request with {}",
v4l2::name(coded.pixelformat)
)));
}
let raw = device.set_format(
Dir::Output,
&Request {
pixelformat: v4l2::NV12,
size,
sizeimage: None,
color: Some(config.resolved_color()),
},
)?;
if let Err(err) = device.set_framerate(Dir::Output, config.framerate) {
tracing::debug!(encoder = NAME, %err, "driver does not take a framerate");
}
let planes = Planes::new(&raw, Rect::whole(size))?;
tracing::info!(
encoder = NAME,
device = %device.path().display(),
format = v4l2::name(raw.pixelformat),
stride = raw.planes[0].stride,
width = size.width,
height = size.height,
"opened H.264 encoder"
);
let raw = Queue::alloc(&device, Dir::Output, raw, RAW_BUFFERS)?;
let coded = Queue::alloc(&device, Dir::Capture, coded, CODED_BUFFERS)?;
Ok(Box::new(Self {
device,
raw,
coded,
planes,
size,
pending: Pending::default(),
drainable: true,
keyframes: true,
}))
}
fn start(&mut self) -> Result<(), Error> {
self.raw.stream_on(&self.device)?;
while let Some(index) = self.coded.take_free() {
self.coded.queue(&self.device, index, &[], Duration::ZERO)?;
}
self.coded.stream_on(&self.device)
}
fn free_buffer(&mut self) -> Result<u32, Error> {
let deadline = Instant::now() + BUFFER_TIMEOUT;
loop {
self.reclaim()?;
if let Some(index) = self.raw.take_free() {
return Ok(index);
}
if Instant::now() >= deadline {
return Err(Error::Codec(anyhow::anyhow!(
"V4L2 encoder held every input buffer for {BUFFER_TIMEOUT:?}"
)));
}
self.device.wait(POLL_INTERVAL);
}
}
fn reclaim(&mut self) -> Result<(), Error> {
while let Some(buffer) = self.raw.dequeue(&self.device)?.buffer() {
if buffer.failed() {
tracing::warn!(encoder = NAME, buffer = buffer.index, "V4L2 encoder dropped a frame");
self.pending.dropped(buffer.timestamp);
}
self.raw.reclaim(buffer.index);
}
Ok(())
}
fn drain(&mut self, out: &mut Vec<Encoded>) -> Result<bool, Error> {
self.reclaim()?;
loop {
let buffer = match self.coded.dequeue(&self.device)? {
Dequeue::Buffer(buffer) => buffer,
Dequeue::Empty => return Ok(false),
Dequeue::Ended => return Ok(true),
};
let payload = access_unit(self.coded.payload(&buffer, 0), buffer.written(0));
self.coded.queue(&self.device, buffer.index, &[], Duration::ZERO)?;
if buffer.failed() {
tracing::warn!(
encoder = NAME,
buffer = buffer.index,
bytes = buffer.bytesused[0],
"V4L2 encoder flagged an access unit bad"
);
self.pending.dropped(buffer.timestamp);
} else if !payload.is_empty()
&& let Some((payload, timestamp)) = self.pending.matched(buffer.timestamp, payload)
{
out.push(Encoded::new(payload, timestamp));
}
if buffer.last() {
return Ok(true);
}
}
}
fn drain_tail(&mut self) -> Result<Vec<Encoded>, Error> {
if !self.raw.streaming() || !self.coded.streaming() {
return Ok(Vec::new());
}
if self.drainable
&& let Err(err) = self.device.encoder_cmd(V4L2_ENC_CMD_STOP)
{
tracing::warn!(
encoder = NAME,
%err,
"driver takes no encoder command; a group boundary can only drain what it has finished"
);
self.drainable = false;
}
match self.drainable {
true => self.drain_to_last(),
false => self.wait_out(),
}
}
fn drain_to_last(&mut self) -> Result<Vec<Encoded>, Error> {
let mut out = Vec::new();
let mut deadline = Instant::now() + FLUSH_TIMEOUT;
loop {
let before = out.len();
if self.drain(&mut out)? {
break;
}
if out.len() > before {
deadline = Instant::now() + FLUSH_TIMEOUT;
} else if Instant::now() >= deadline {
return Err(Error::Codec(anyhow::anyhow!(
"V4L2 encoder did not finish its drain within {FLUSH_TIMEOUT:?}, holding {} frame(s)",
self.pending.len()
)));
}
self.device.wait(POLL_INTERVAL);
}
if !self.pending.is_empty() {
tracing::debug!(
encoder = NAME,
frames = self.pending.len(),
"V4L2 encoder ended its drain still owing access units"
);
self.pending.forget();
}
self.device.encoder_cmd(V4L2_ENC_CMD_START)?;
Ok(out)
}
fn wait_out(&mut self) -> Result<Vec<Encoded>, Error> {
let mut out = Vec::new();
let mut deadline = Instant::now() + FLUSH_TIMEOUT;
loop {
let before = out.len();
self.drain(&mut out)?;
if self.pending.is_empty() {
return Ok(out);
}
if out.len() > before {
deadline = Instant::now() + FLUSH_TIMEOUT;
} else if Instant::now() >= deadline {
return Err(Error::Codec(anyhow::anyhow!(
"V4L2 encoder held {} frame(s) for {FLUSH_TIMEOUT:?} without encoding them",
self.pending.len()
)));
}
self.device.wait(POLL_INTERVAL);
}
}
}
impl Backend for V4l2 {
fn encode(&mut self, frame: &Frame, keyframe: bool) -> Result<Vec<Encoded>, Error> {
if frame.size() != self.size {
return Err(Error::Codec(anyhow::anyhow!(
"V4L2 encoder opened for {} was given a {} frame",
self.size,
frame.size()
)));
}
let i420 = frame.surface.to_i420()?;
let index = self.free_buffer()?;
self.planes.write(&mut self.raw, index, &i420)?;
if keyframe && self.keyframes {
self.keyframes = self.device.try_control(V4L2_CID_MPEG_VIDEO_FORCE_KEY_FRAME, 0);
if !self.keyframes {
tracing::warn!(
encoder = NAME,
device = %self.device.path().display(),
"driver takes no keyframe request; groups fall on the encoder's own GOP boundary"
);
}
}
let key = key(frame.timestamp);
let bytesused: Vec<u32> = self.raw.format().planes.iter().map(|plane| plane.sizeimage).collect();
self.raw.queue(&self.device, index, &bytesused, key)?;
self.pending.queued(key, frame.timestamp);
if !self.raw.streaming() {
self.start()?;
}
let mut out = Vec::new();
self.drain(&mut out)?;
Ok(out)
}
fn flush(&mut self) -> Result<Vec<Encoded>, Error> {
self.drain_tail()
}
fn finish(&mut self) -> Result<Vec<Encoded>, Error> {
self.drain_tail()
}
fn set_bitrate(&mut self, bitrate: u64) -> Result<(), Error> {
set_bitrate(&self.device, bitrate).map_err(|err| {
tracing::debug!(encoder = NAME, %err, "driver refused a bitrate change");
Error::BitrateUnsupported(NAME)
})
}
fn name(&self) -> &str {
NAME
}
}
fn key(timestamp: Timestamp) -> Duration {
Duration::from_micros(timestamp.as_micros() as u64)
}
#[derive(Debug, Default)]
struct Pending {
frames: VecDeque<(Duration, Timestamp)>,
header: Option<Bytes>,
unstamped: bool,
}
impl Pending {
fn queued(&mut self, key: Duration, timestamp: Timestamp) {
self.frames.push_back((key, timestamp));
}
fn len(&self) -> usize {
self.frames.len()
}
fn is_empty(&self) -> bool {
self.frames.is_empty()
}
fn matched(&mut self, key: Duration, payload: Bytes) -> Option<(Bytes, Timestamp)> {
if !has_picture(&payload) {
self.carry(payload);
return None;
}
let Some(at) = self.frames.iter().position(|(frame, _)| *frame == key) else {
if !self.unstamped {
tracing::warn!(
encoder = NAME,
"V4L2 encoder answered a picture that matches no frame; the driver is not copying timestamps"
);
self.unstamped = true;
}
return None;
};
let (_, timestamp) = self.frames[at];
self.frames.drain(..=at);
let payload = match self.header.take() {
Some(header) => join(&header, &payload),
None => payload,
};
Some((payload, timestamp))
}
fn carry(&mut self, payload: Bytes) {
let held = self.header.as_ref().map_or(0, Bytes::len);
if held + payload.len() > CARRY_LIMIT {
tracing::debug!(
encoder = NAME,
held,
"V4L2 encoder discarded coded bytes that hold no picture"
);
self.header = None;
return;
}
self.header = Some(match self.header.take() {
Some(header) => join(&header, &payload),
None => payload,
});
}
fn dropped(&mut self, key: Duration) {
if let Some(at) = self.frames.iter().position(|(frame, _)| *frame == key) {
self.frames.remove(at);
}
}
fn forget(&mut self) {
self.frames.clear();
}
}
fn join(header: &[u8], payload: &[u8]) -> Bytes {
let mut joined = BytesMut::with_capacity(header.len() + payload.len());
joined.extend_from_slice(header);
joined.extend_from_slice(payload);
joined.freeze()
}
fn set_bitrate(device: &Device, bitrate: u64) -> Result<(), Error> {
device.set_control(V4L2_CID_MPEG_VIDEO_BITRATE, bitrate.min(i32::MAX as u64) as i32)
}
fn has_picture(annexb: &[u8]) -> bool {
annexb
.windows(4)
.any(|bytes| matches!(bytes, [0, 0, 1, header] if (1..=5).contains(&(header & 0x1f))))
}
struct Level {
code: v4l2_mpeg_video_h264_level,
per_second: u32,
per_frame: u32,
kbps: u32,
}
const LEVEL_5_2: v4l2_mpeg_video_h264_level = 16;
const LEVEL_6_0: v4l2_mpeg_video_h264_level = 17;
const LEVEL_6_1: v4l2_mpeg_video_h264_level = 18;
const LEVEL_6_2: v4l2_mpeg_video_h264_level = 19;
const LEVELS: &[Level] = &[
Level {
code: v4l2_mpeg_video_h264_level_V4L2_MPEG_VIDEO_H264_LEVEL_1_0,
per_second: 1_485,
per_frame: 99,
kbps: 64,
},
Level {
code: v4l2_mpeg_video_h264_level_V4L2_MPEG_VIDEO_H264_LEVEL_1_1,
per_second: 3_000,
per_frame: 396,
kbps: 192,
},
Level {
code: v4l2_mpeg_video_h264_level_V4L2_MPEG_VIDEO_H264_LEVEL_1_2,
per_second: 6_000,
per_frame: 396,
kbps: 384,
},
Level {
code: v4l2_mpeg_video_h264_level_V4L2_MPEG_VIDEO_H264_LEVEL_1_3,
per_second: 11_880,
per_frame: 396,
kbps: 768,
},
Level {
code: v4l2_mpeg_video_h264_level_V4L2_MPEG_VIDEO_H264_LEVEL_2_0,
per_second: 11_880,
per_frame: 396,
kbps: 2_000,
},
Level {
code: v4l2_mpeg_video_h264_level_V4L2_MPEG_VIDEO_H264_LEVEL_2_1,
per_second: 19_800,
per_frame: 792,
kbps: 4_000,
},
Level {
code: v4l2_mpeg_video_h264_level_V4L2_MPEG_VIDEO_H264_LEVEL_2_2,
per_second: 20_250,
per_frame: 1_620,
kbps: 4_000,
},
Level {
code: v4l2_mpeg_video_h264_level_V4L2_MPEG_VIDEO_H264_LEVEL_3_0,
per_second: 40_500,
per_frame: 1_620,
kbps: 10_000,
},
Level {
code: v4l2_mpeg_video_h264_level_V4L2_MPEG_VIDEO_H264_LEVEL_3_1,
per_second: 108_000,
per_frame: 3_600,
kbps: 14_000,
},
Level {
code: v4l2_mpeg_video_h264_level_V4L2_MPEG_VIDEO_H264_LEVEL_3_2,
per_second: 216_000,
per_frame: 5_120,
kbps: 20_000,
},
Level {
code: v4l2_mpeg_video_h264_level_V4L2_MPEG_VIDEO_H264_LEVEL_4_0,
per_second: 245_760,
per_frame: 8_192,
kbps: 20_000,
},
Level {
code: v4l2_mpeg_video_h264_level_V4L2_MPEG_VIDEO_H264_LEVEL_4_1,
per_second: 245_760,
per_frame: 8_192,
kbps: 50_000,
},
Level {
code: v4l2_mpeg_video_h264_level_V4L2_MPEG_VIDEO_H264_LEVEL_4_2,
per_second: 522_240,
per_frame: 8_704,
kbps: 50_000,
},
Level {
code: v4l2_mpeg_video_h264_level_V4L2_MPEG_VIDEO_H264_LEVEL_5_0,
per_second: 589_824,
per_frame: 22_080,
kbps: 135_000,
},
Level {
code: v4l2_mpeg_video_h264_level_V4L2_MPEG_VIDEO_H264_LEVEL_5_1,
per_second: 983_040,
per_frame: 36_864,
kbps: 240_000,
},
Level {
code: LEVEL_5_2,
per_second: 2_073_600,
per_frame: 36_864,
kbps: 240_000,
},
Level {
code: LEVEL_6_0,
per_second: 4_177_920,
per_frame: 139_264,
kbps: 240_000,
},
Level {
code: LEVEL_6_1,
per_second: 8_355_840,
per_frame: 139_264,
kbps: 480_000,
},
Level {
code: LEVEL_6_2,
per_second: 16_711_680,
per_frame: 139_264,
kbps: 800_000,
},
];
fn h264_level(config: &Config) -> Result<i32, Error> {
let size = config.size();
let per_frame = size.width.div_ceil(16) * size.height.div_ceil(16);
let per_second = per_frame as u64 * config.framerate as u64;
let kbps = config.resolved_bitrate().div_ceil(1_000);
LEVELS
.iter()
.find(|level| {
per_frame <= level.per_frame && per_second <= level.per_second as u64 && kbps <= level.kbps as u64
})
.map(|level| level.code as i32)
.ok_or_else(|| {
Error::Codec(anyhow::anyhow!(
"{size} at {}fps and {kbps} kbit/s exceeds H.264 level 6.2",
config.framerate
))
})
}
fn coded_size(size: Size) -> u32 {
const MIN: u32 = 256 * 1024;
const MAX: u32 = 4 * 1024 * 1024;
(size.pixels().min(u32::MAX as u64) as u32 / 2).clamp(MIN, MAX)
}
fn access_unit(buffer: &[u8], bytesused: u32) -> Bytes {
let used = (bytesused as usize).min(buffer.len());
let end = buffer[..used]
.iter()
.rposition(|byte| *byte != 0)
.map_or(0, |at| at + 1);
Bytes::copy_from_slice(&buffer[..end])
}
#[cfg(test)]
mod tests {
use super::*;
fn config(width: u32, height: u32, framerate: u32) -> Config {
Config::new(width, height, framerate)
}
fn micros(micros: u64) -> Timestamp {
Timestamp::from_micros(micros).unwrap()
}
fn queue(pending: &mut Pending, timestamp: Timestamp) {
pending.queued(key(timestamp), timestamp);
}
#[test]
fn the_level_clears_the_resolution() {
assert_eq!(h264_level(&config(320, 240, 30)).unwrap(), 4); assert_eq!(h264_level(&config(640, 480, 30)).unwrap(), 8); assert_eq!(h264_level(&config(1280, 720, 30)).unwrap(), 9); assert_eq!(h264_level(&config(1920, 1080, 30)).unwrap(), 11); assert_eq!(h264_level(&config(3840, 2160, 30)).unwrap(), 15); }
#[test]
fn the_level_clears_the_framerate() {
assert_eq!(h264_level(&config(1920, 1080, 60)).unwrap(), 13); assert_eq!(h264_level(&config(1280, 720, 60)).unwrap(), 10); assert_eq!(h264_level(&config(320, 240, 60)).unwrap(), 6); }
#[test]
fn the_level_clears_the_bitrate() {
let mut config = config(1280, 720, 30);
config.bitrate = Some(14_000_000);
assert_eq!(h264_level(&config).unwrap(), 9); config.bitrate = Some(14_000_001);
assert_eq!(h264_level(&config).unwrap(), 10); }
#[test]
fn the_level_runs_to_the_end_of_the_menu() {
assert_eq!(h264_level(&config(3840, 2160, 60)).unwrap(), 16); assert_eq!(h264_level(&config(7680, 4320, 60)).unwrap(), 18); assert_eq!(h264_level(&config(7680, 4320, 120)).unwrap(), 19); assert!(h264_level(&config(7680, 4320, 240)).is_err());
}
#[test]
fn a_picture_is_recognized_by_its_vcl_nal() {
assert!(has_picture(&[0, 0, 0, 1, 0x65]));
assert!(has_picture(&[0, 0, 1, 0x41]));
assert!(!has_picture(&[0, 0, 0, 1, 0x67, 0x42, 0, 0, 0, 1, 0x68, 0xce]));
assert!(!has_picture(&[0, 0, 0, 1, 0x06, 0, 0, 0, 1, 0x09]));
assert!(!has_picture(&[]));
assert!(!has_picture(&[0x65]));
}
#[test]
fn the_coded_buffer_is_bounded() {
assert_eq!(coded_size(Size::new(160, 120)), 256 * 1024);
assert_eq!(coded_size(Size::new(1920, 1080)), 1920 * 1080 / 2);
assert_eq!(coded_size(Size::new(7680, 4320)), 4 * 1024 * 1024);
}
#[test]
fn an_access_unit_is_matched_to_the_frame_it_came_from() {
let mut pending = Pending::default();
for at in [0, 33_000, 66_000] {
queue(&mut pending, micros(at));
}
assert_eq!(pending.len(), 3);
let payload = Bytes::from_static(&[0, 0, 0, 1, 0x65]);
assert_eq!(
pending.matched(Duration::from_micros(0), payload.clone()),
Some((payload.clone(), micros(0)))
);
assert_eq!(pending.len(), 2);
assert!(pending.matched(Duration::from_micros(66_000), payload).is_some());
assert!(pending.is_empty());
}
#[test]
fn separate_parameter_sets_join_the_next_access_unit() {
let mut pending = Pending::default();
queue(&mut pending, micros(500));
assert_eq!(
pending.matched(Duration::from_micros(0), Bytes::from_static(&[0, 0, 0, 1, 0x67])),
None
);
assert_eq!(pending.len(), 1);
let (joined, _) = pending
.matched(Duration::from_micros(500), Bytes::from_static(&[0, 0, 0, 1, 0x65]))
.unwrap();
assert_eq!(&joined[..], &[0, 0, 0, 1, 0x67, 0, 0, 0, 1, 0x65]);
assert!(pending.is_empty());
let next = Bytes::from_static(&[0, 0, 0, 1, 0x41]);
queue(&mut pending, micros(533));
assert_eq!(
pending.matched(Duration::from_micros(533), next.clone()),
Some((next, micros(533)))
);
}
#[test]
fn the_answer_carries_the_frame_timestamp_unchanged() {
let ninety_khz = moq_net::Timescale::new(90_000).unwrap();
let timestamp = Timestamp::new(3003, ninety_khz).unwrap();
assert_ne!(Timestamp::from_micros(timestamp.as_micros() as u64).unwrap(), timestamp);
let mut pending = Pending::default();
queue(&mut pending, timestamp);
let (_, answered) = pending
.matched(key(timestamp), Bytes::from_static(&[0, 0, 0, 1, 0x65]))
.unwrap();
assert_eq!(answered, timestamp);
}
#[test]
fn parameter_sets_stamped_like_the_first_frame_still_join_it() {
let mut pending = Pending::default();
queue(&mut pending, micros(0));
assert_eq!(
pending.matched(
Duration::ZERO,
Bytes::from_static(&[0, 0, 0, 1, 0x67, 0, 0, 0, 1, 0x68])
),
None
);
assert_eq!(pending.len(), 1);
let (joined, _) = pending
.matched(Duration::ZERO, Bytes::from_static(&[0, 0, 0, 1, 0x65]))
.unwrap();
assert_eq!(&joined[..], &[0, 0, 0, 1, 0x67, 0, 0, 0, 1, 0x68, 0, 0, 0, 1, 0x65]);
assert!(pending.is_empty());
}
#[test]
fn unmatched_coded_buffers_stop_accumulating() {
let mut pending = Pending::default();
queue(&mut pending, micros(1));
let payload = Bytes::from(vec![0u8; 8 * 1024]);
for _ in 0..64 {
assert_eq!(pending.matched(Duration::from_micros(0), payload.clone()), None);
assert!(pending.header.as_ref().is_none_or(|held| held.len() <= CARRY_LIMIT));
}
assert_eq!(pending.len(), 1);
}
#[test]
fn an_unmatched_picture_is_not_carried() {
let mut pending = Pending::default();
queue(&mut pending, micros(500));
let stray = Bytes::from_static(&[0, 0, 0, 1, 0x65, 0xaa]);
assert_eq!(pending.matched(Duration::from_micros(999), stray), None);
assert!(pending.unstamped);
assert!(pending.header.is_none());
assert_eq!(pending.len(), 1);
let payload = Bytes::from_static(&[0, 0, 0, 1, 0x65]);
assert_eq!(
pending.matched(Duration::from_micros(500), payload.clone()),
Some((payload, micros(500)))
);
}
#[test]
fn a_dropped_frame_leaves_the_rest_matchable() {
let mut pending = Pending::default();
for at in [10, 20, 30] {
queue(&mut pending, micros(at));
}
pending.dropped(Duration::from_micros(20));
assert_eq!(pending.len(), 2);
pending.dropped(Duration::from_micros(999));
assert_eq!(pending.len(), 2);
let payload = Bytes::from_static(&[0, 0, 0, 1, 0x65]);
assert_eq!(
pending.matched(Duration::from_micros(10), payload.clone()),
Some((payload.clone(), micros(10)))
);
assert_eq!(
pending.matched(Duration::from_micros(30), payload.clone()),
Some((payload, micros(30)))
);
assert!(pending.is_empty());
}
#[test]
fn padding_is_not_part_of_the_access_unit() {
let buffer = [0, 0, 0, 1, 0x65, 0x88, 0, 0, 0, 0];
assert_eq!(&access_unit(&buffer, 10)[..], &[0, 0, 0, 1, 0x65, 0x88]);
assert!(access_unit(&[0; 8], 8).is_empty());
assert_eq!(access_unit(&buffer, 64).len(), 6);
}
}