use super::*;
use crate::ffi::ReturnCode;
use std::ffi::{c_char, c_void};
use std::sync::mpsc;
use std::time::Duration;
const TIMEOUT: Duration = Duration::from_secs(10);
fn id(raw: i32) -> u32 {
assert!(raw > 0, "expected positive id, got {raw}");
raw as u32
}
fn publish_broadcast(origin: u32, path: &[u8]) -> u32 {
let broadcast = id(unsafe { moq_origin_create_broadcast(origin, path.as_ptr() as *const c_char, path.len()) });
assert_eq!(unsafe { moq_publish_announce(broadcast, std::ptr::null()) }, 0);
broadcast
}
fn request_broadcast(origin: u32, path: &[u8]) -> u32 {
let deadline = std::time::Instant::now() + TIMEOUT;
loop {
let cb = Callback::new();
let _task = id(unsafe {
moq_origin_request(
origin,
path.as_ptr() as *const c_char,
path.len(),
Some(channel_callback),
cb.ptr,
)
});
let code = cb.recv();
if code > 0 {
cb.recv_terminal();
return code as u32;
}
assert!(
std::time::Instant::now() < deadline,
"timed out requesting broadcast: {code}"
);
std::thread::sleep(Duration::from_millis(10));
}
}
fn request_once(origin: u32, path: &[u8]) -> i32 {
let cb = Callback::new();
let _task = id(unsafe {
moq_origin_request(
origin,
path.as_ptr() as *const c_char,
path.len(),
Some(channel_callback),
cb.ptr,
)
});
let code = cb.recv();
if code > 0 {
cb.recv_terminal();
}
code
}
struct Guard<F: FnOnce()>(Option<F>);
impl<F: FnOnce()> Drop for Guard<F> {
fn drop(&mut self) {
if let Some(f) = self.0.take() {
f();
}
}
}
struct Callback {
rx: mpsc::Receiver<i32>,
ptr: *mut c_void,
}
impl Callback {
fn new() -> Self {
let (tx, rx) = mpsc::channel();
let ptr = Box::into_raw(Box::new(tx)) as *mut c_void;
Self { rx, ptr }
}
fn recv(&self) -> i32 {
self.rx.recv_timeout(TIMEOUT).expect("callback timed out")
}
fn recv_terminal(&self) -> i32 {
let code = self.recv();
assert!(code <= 0, "expected terminal code <= 0, got {code}");
code
}
fn recv_catalog_terminal(&self) -> i32 {
loop {
let code = self.recv();
if code <= 0 {
return code;
}
assert_eq!(moq_consume_catalog_free(id(code)), 0);
}
}
}
impl Drop for Callback {
fn drop(&mut self) {
unsafe { drop(Box::from_raw(self.ptr as *mut mpsc::Sender<i32>)) };
}
}
extern "C" fn channel_callback(user_data: *mut c_void, code: i32) {
let tx = unsafe { &*(user_data as *const mpsc::Sender<i32>) };
let _ = tx.send(code);
}
extern "C" fn ignore_callback(_user_data: *mut c_void, _code: i32) {}
fn opus_head() -> Vec<u8> {
let mut head = Vec::with_capacity(19);
head.extend_from_slice(b"OpusHead");
head.push(1); head.push(2); head.extend_from_slice(&0u16.to_le_bytes()); head.extend_from_slice(&48000u32.to_le_bytes()); head.extend_from_slice(&0u16.to_le_bytes()); head.push(0); head
}
fn h264_init() -> Vec<u8> {
let mut init = Vec::new();
init.extend_from_slice(&[0x00, 0x00, 0x00, 0x01]);
init.extend_from_slice(&[
0x67, 0x64, 0x00, 0x1f, 0xac, 0x24, 0x84, 0x01, 0x40, 0x16, 0xec, 0x04, 0x40, 0x00, 0x00, 0x03, 0x00, 0x40,
0x00, 0x00, 0x0c, 0x23, 0xc6, 0x0c, 0x92,
]);
init.extend_from_slice(&[0x00, 0x00, 0x00, 0x01]);
init.extend_from_slice(&[0x68, 0xee, 0x32, 0xc8, 0xb0]);
init
}
fn label_ptr(label: Option<&[u8]>) -> (*const c_char, usize) {
label
.map(|label| (label.as_ptr() as *const c_char, label.len()))
.unwrap_or((std::ptr::null(), 0))
}
fn init_ptr(init: &[u8]) -> *const u8 {
if init.is_empty() {
std::ptr::null()
} else {
init.as_ptr()
}
}
fn publish_audio(broadcast: u32, format: moq_audio_format, init: &[u8], label: Option<&[u8]>) -> i32 {
let (label, label_len) = label_ptr(label);
let config = moq_audio_init {
format: format as u32,
init: init_ptr(init),
init_len: init.len(),
label,
label_len,
};
unsafe { moq_publish_audio(broadcast, &config) }
}
fn publish_video(broadcast: u32, format: moq_video_format, init: &[u8], label: Option<&[u8]>) -> i32 {
let (label, label_len) = label_ptr(label);
let config = moq_video_init {
format: format as u32,
init: init_ptr(init),
init_len: init.len(),
label,
label_len,
hint: moq_video_hint::default(),
};
unsafe { moq_publish_video(broadcast, &config) }
}
fn publish_container(broadcast: u32, format: moq_container_format, init: &[u8]) -> i32 {
let config = moq_container_init {
format: format as u32,
init: init_ptr(init),
init_len: init.len(),
};
unsafe { moq_publish_container(broadcast, &config) }
}
#[test]
fn origin_lifecycle() {
let origin = id(moq_origin_create());
assert_eq!(moq_origin_close(origin), 0, "moq_origin_close should succeed");
assert!(moq_origin_close(origin) < 0, "double-close should fail");
}
#[test]
fn last_error_reports_reason() {
assert!(moq_origin_close(9999) < 0);
let ptr = moq_error();
assert!(!ptr.is_null(), "expected a recorded error message");
let msg = unsafe { std::ffi::CStr::from_ptr(ptr) }.to_str().unwrap();
assert_eq!(msg, "origin not found");
}
#[test]
fn last_error_set_before_callback() {
use crate::Error;
use crate::ffi::OnStatus;
extern "C" fn capture(user_data: *mut c_void, code: i32) {
assert!(code < 0, "expected a negative status, got {code}");
let slot = unsafe { &mut *(user_data as *mut Option<String>) };
let ptr = moq_error();
*slot = (!ptr.is_null()).then(|| unsafe { std::ffi::CStr::from_ptr(ptr) }.to_str().unwrap().to_owned());
}
let mut captured: Option<String> = None;
let cb = unsafe { OnStatus::new(&mut captured as *mut _ as *mut c_void, Some(capture)) }.unwrap();
cb.call(Err::<(), Error>(Error::OriginNotFound));
assert_eq!(captured.as_deref(), Some("origin not found"));
}
#[test]
fn status_callback_refuses_null() {
let callback = unsafe { crate::ffi::OnStatus::new(std::ptr::null_mut(), None) };
assert!(matches!(callback, Err(Error::InvalidPointer)));
}
#[test]
fn last_error_protocol_is_none_for_a_local_failure() {
assert!(moq_origin_close(9999) < 0);
let mut out = moq_protocol_error {
scope: 99,
code: 99,
kind: 99,
};
assert!(unsafe { moq_error_protocol(&mut out) } < 0);
assert_eq!(out.scope, 99, "a non-protocol error must not write the record");
}
#[test]
fn last_error_protocol_captures_a_stream_app_code() {
use crate::Error;
use crate::ffi::OnStatus;
extern "C" fn capture(user_data: *mut c_void, code: i32) {
assert!(code < 0, "expected a negative status, got {code}");
let slot = unsafe { &mut *(user_data as *mut Option<moq_protocol_error>) };
let mut out = moq_protocol_error {
scope: 0,
code: 0,
kind: 0,
};
assert_eq!(unsafe { moq_error_protocol(&mut out) }, 0);
*slot = Some(out);
}
let mut captured: Option<moq_protocol_error> = None;
let cb = unsafe { OnStatus::new(&mut captured as *mut _ as *mut c_void, Some(capture)) }.unwrap();
cb.call(Err::<(), Error>(Error::Moq(moq_net::StreamError::App(404).into())));
let protocol = captured.expect("expected a protocol error");
assert_eq!(protocol.scope, moq_error_scope::MOQ_ERROR_SCOPE_STREAM as u32);
assert_eq!(protocol.code, 64 + 404);
assert_eq!(protocol.kind, moq_protocol_kind::MOQ_PROTOCOL_KIND_APP as u32);
}
#[test]
fn last_error_protocol_captures_session_known_app_and_unknown() {
use crate::Error;
use crate::ffi::OnStatus;
extern "C" fn capture(user_data: *mut c_void, code: i32) {
assert!(code < 0, "expected a negative status, got {code}");
let slot = unsafe { &mut *(user_data as *mut Option<moq_protocol_error>) };
let mut out = moq_protocol_error {
scope: 0,
code: 0,
kind: 0,
};
assert_eq!(unsafe { moq_error_protocol(&mut out) }, 0);
*slot = Some(out);
}
for (err, scope, code, kind) in [
(
moq_net::Error::from(moq_net::SessionError::Cancel),
moq_error_scope::MOQ_ERROR_SCOPE_SESSION as u32,
0,
moq_protocol_kind::MOQ_PROTOCOL_KIND_CANCEL as u32,
),
(
moq_net::Error::from(moq_net::StreamError::Internal),
moq_error_scope::MOQ_ERROR_SCOPE_STREAM as u32,
0,
moq_protocol_kind::MOQ_PROTOCOL_KIND_INTERNAL as u32,
),
(
moq_net::Error::from(moq_net::SessionError::Unauthorized),
moq_error_scope::MOQ_ERROR_SCOPE_SESSION as u32,
0x2,
moq_protocol_kind::MOQ_PROTOCOL_KIND_UNAUTHORIZED as u32,
),
(
moq_net::Error::from(moq_net::SessionError::App(7)),
moq_error_scope::MOQ_ERROR_SCOPE_SESSION as u32,
64 + 7,
moq_protocol_kind::MOQ_PROTOCOL_KIND_APP as u32,
),
(
moq_net::Error::from(moq_net::SessionError::Unknown(0x1f)),
moq_error_scope::MOQ_ERROR_SCOPE_SESSION as u32,
0x1f,
moq_protocol_kind::MOQ_PROTOCOL_KIND_UNKNOWN as u32,
),
(
moq_net::Error::from(moq_net::StreamError::App(404)),
moq_error_scope::MOQ_ERROR_SCOPE_STREAM as u32,
64 + 404,
moq_protocol_kind::MOQ_PROTOCOL_KIND_APP as u32,
),
] {
let mut captured: Option<moq_protocol_error> = None;
let cb = unsafe { OnStatus::new(&mut captured as *mut _ as *mut c_void, Some(capture)) }.unwrap();
cb.call(Err::<(), Error>(Error::Moq(err)));
let protocol = captured.expect("expected a protocol error");
assert_eq!(protocol.scope, scope);
assert_eq!(protocol.code, code);
assert_eq!(protocol.kind, kind);
}
}
#[test]
fn publish_media_lifecycle() {
let origin = id(moq_origin_create());
let broadcast = publish_broadcast(origin, b"publish-media-lifecycle");
let _guard = Guard(Some(|| {
moq_publish_finish(broadcast);
}));
let init = opus_head();
let format = moq_audio_format::MOQ_AUDIO_FORMAT_OPUS;
let media = id(publish_audio(broadcast, format, &init, None));
let payload = b"opus frame";
let ret = unsafe { moq_publish_media_frame(media, payload.as_ptr(), payload.len(), 1000) };
assert_eq!(ret, 0, "moq_publish_media_frame should succeed");
assert_eq!(moq_publish_media_finish(media), 0);
assert_eq!(moq_publish_finish(broadcast), 0);
}
#[test]
fn publish_media_rejects_a_null_config() {
let origin = id(moq_origin_create());
let broadcast = publish_broadcast(origin, b"publish-media-null-config");
assert!(unsafe { moq_publish_audio(broadcast, std::ptr::null()) } < 0);
assert_eq!(moq_publish_finish(broadcast), 0);
assert_eq!(moq_origin_close(origin), 0);
}
#[test]
fn container_and_media_handles_are_not_interchangeable() {
let origin = id(moq_origin_create());
let broadcast = publish_broadcast(origin, b"publish-container-handles");
let container = id(publish_container(
broadcast,
moq_container_format::MOQ_CONTAINER_FORMAT_FMP4,
&[],
));
let media = id(publish_audio(
broadcast,
moq_audio_format::MOQ_AUDIO_FORMAT_OPUS,
&opus_head(),
None,
));
assert!(
moq_publish_media_finish(container) < 0,
"a container is not a media track"
);
assert!(
moq_publish_container_finish(media) < 0,
"a media track is not a container"
);
assert_eq!(moq_publish_container_finish(container), 0);
assert_eq!(moq_publish_media_finish(media), 0);
assert_eq!(moq_publish_finish(broadcast), 0);
assert_eq!(moq_origin_close(origin), 0);
}
fn borrowed_string(ptr: *const c_char, len: usize) -> Option<String> {
if ptr.is_null() {
return None;
}
Some(
unsafe { std::str::from_utf8(std::slice::from_raw_parts(ptr.cast::<u8>(), len)) }
.unwrap()
.to_string(),
)
}
#[test]
fn publish_media_labels_config_without_naming_track() {
let origin = id(moq_origin_create());
let path = b"labeled-track";
let broadcast = publish_broadcast(origin, path);
let init = opus_head();
let format = moq_audio_format::MOQ_AUDIO_FORMAT_OPUS;
let label = b"English";
let media1 = id(publish_audio(broadcast, format, &init, Some(label)));
let consume = request_broadcast(origin, path);
let catalog_cb = Callback::new();
let catalog_task = id(unsafe { moq_consume_catalog(consume, Some(channel_callback), catalog_cb.ptr) });
let catalog_id1 = id(catalog_cb.recv());
let mut audio_cfg = moq_audio_config {
name: std::ptr::null(),
name_len: 0,
label: std::ptr::null(),
label_len: 0,
codec: std::ptr::null(),
codec_len: 0,
description: std::ptr::null(),
description_len: 0,
sample_rate: 0,
channel_count: 0,
container: moq_container::default(),
};
assert_eq!(unsafe { moq_consume_audio_config(catalog_id1, 0, &mut audio_cfg) }, 0);
assert_eq!(
borrowed_string(audio_cfg.name, audio_cfg.name_len).as_deref(),
Some("0.opus")
);
assert_eq!(
borrowed_string(audio_cfg.label, audio_cfg.label_len).as_deref(),
Some("English")
);
let media2 = id(publish_audio(broadcast, format, &init, Some(label)));
let catalog_id2 = id(catalog_cb.recv());
assert_eq!(unsafe { moq_consume_audio_config(catalog_id2, 0, &mut audio_cfg) }, 0);
assert_eq!(
borrowed_string(audio_cfg.name, audio_cfg.name_len).as_deref(),
Some("0.opus")
);
assert_eq!(
borrowed_string(audio_cfg.label, audio_cfg.label_len).as_deref(),
Some("English")
);
assert_eq!(unsafe { moq_consume_audio_config(catalog_id2, 1, &mut audio_cfg) }, 0);
assert_eq!(
borrowed_string(audio_cfg.name, audio_cfg.name_len).as_deref(),
Some("1.opus")
);
assert_eq!(
borrowed_string(audio_cfg.label, audio_cfg.label_len).as_deref(),
Some("English")
);
assert_eq!(moq_consume_catalog_free(catalog_id1), 0);
assert_eq!(moq_consume_catalog_free(catalog_id2), 0);
assert_eq!(moq_consume_catalog_cancel(catalog_task), 0);
assert_eq!(catalog_cb.recv_terminal(), 0, "catalog close delivers terminal 0");
assert_eq!(moq_consume_close(consume), 0);
assert_eq!(moq_publish_media_finish(media1), 0);
assert_eq!(moq_publish_media_finish(media2), 0);
assert_eq!(moq_publish_finish(broadcast), 0);
assert_eq!(moq_origin_close(origin), 0);
}
#[test]
fn session_snapshot_error_preserves_destination() {
let mut dst: moq_connection_snapshot = unsafe { std::mem::zeroed() };
dst.stats.rtt_us = 123;
dst.protocol.len = 456;
assert!(unsafe { moq_session_snapshot(0, &mut dst) } < 0);
assert_eq!(dst.stats.rtt_us, 123);
assert_eq!(dst.protocol.len, 456);
}
#[test]
fn publish_catalog_config_invalid_broadcast() {
let name = "video";
let codec = "vp8";
let video = moq_video_config {
name: name.as_ptr() as *const c_char,
name_len: name.len(),
label: std::ptr::null(),
label_len: 0,
codec: codec.as_ptr() as *const c_char,
codec_len: codec.len(),
description: std::ptr::null(),
description_len: 0,
coded_width: 0,
coded_height: 0,
container: moq_container::default(),
};
assert!(unsafe { moq_publish_video_config(0, &video) } < 0);
assert!(unsafe { moq_publish_video_properties(0, &moq_video_properties::default()) } < 0);
let audio_codec = "opus";
let audio = moq_audio_config {
name: name.as_ptr() as *const c_char,
name_len: name.len(),
label: std::ptr::null(),
label_len: 0,
codec: audio_codec.as_ptr() as *const c_char,
codec_len: audio_codec.len(),
description: std::ptr::null(),
description_len: 0,
sample_rate: 48000,
channel_count: 2,
container: moq_container::default(),
};
assert!(unsafe { moq_publish_audio_config(0, &audio) } < 0);
assert!(unsafe { moq_publish_video_remove(0, name.as_ptr() as *const c_char, name.len()) } < 0);
assert!(unsafe { moq_publish_audio_remove(0, name.as_ptr() as *const c_char, name.len()) } < 0);
}
#[test]
fn publish_media_owns_its_rendition_before_the_first_keyframe() {
let origin = id(moq_origin_create());
let broadcast = publish_broadcast(origin, b"media-owns-its-rendition");
let media = id(publish_video(
broadcast,
moq_video_format::MOQ_VIDEO_FORMAT_AVC3,
&[],
None,
));
let name = "0.avc3";
let codec = "avc1.42c01e";
let video = moq_video_config {
name: name.as_ptr() as *const c_char,
name_len: name.len(),
label: std::ptr::null(),
label_len: 0,
codec: codec.as_ptr() as *const c_char,
codec_len: codec.len(),
description: std::ptr::null(),
description_len: 0,
coded_width: 0,
coded_height: 0,
container: moq_container::default(),
};
assert_eq!(
unsafe { moq_publish_video_config(broadcast, &video) },
-18,
"the media track owns 0.avc3 even before its config resolves"
);
assert_eq!(
unsafe { moq_publish_video_remove(broadcast, name.as_ptr() as *const c_char, name.len()) },
0,
"removing a name the caller never authored is a no-op"
);
assert_eq!(
unsafe { moq_publish_video_config(broadcast, &video) },
-18,
"and the no-op left the media track's rendition alone"
);
assert_eq!(moq_publish_media_finish(media), 0);
assert_eq!(
unsafe { moq_publish_video_config(broadcast, &video) },
0,
"finishing the media track releases its rendition name"
);
assert_eq!(moq_publish_finish(broadcast), 0);
assert_eq!(moq_origin_close(origin), 0);
}
#[test]
fn publish_video_config_replaces_its_own_rendition() {
let origin = id(moq_origin_create());
let broadcast = publish_broadcast(origin, b"catalog-config-replace");
let name = "authored";
let codec = "vp8";
let mut video = moq_video_config {
name: name.as_ptr() as *const c_char,
name_len: name.len(),
label: std::ptr::null(),
label_len: 0,
codec: codec.as_ptr() as *const c_char,
codec_len: codec.len(),
description: std::ptr::null(),
description_len: 0,
coded_width: 640,
coded_height: 360,
container: moq_container::default(),
};
assert_eq!(unsafe { moq_publish_video_config(broadcast, &video) }, 0);
video.coded_width = 1920;
video.coded_height = 1080;
assert_eq!(
unsafe { moq_publish_video_config(broadcast, &video) },
0,
"a caller can refine a rendition it owns"
);
assert_eq!(
unsafe { moq_publish_video_remove(broadcast, name.as_ptr() as *const c_char, name.len()) },
0
);
assert_eq!(
unsafe { moq_publish_video_config(broadcast, &video) },
0,
"the name is free once the caller removes its rendition"
);
assert_eq!(moq_publish_finish(broadcast), 0);
assert_eq!(moq_origin_close(origin), 0);
}
#[test]
fn publish_catalog_config_null_pointer() {
let origin = id(moq_origin_create());
let broadcast = publish_broadcast(origin, b"publish-catalog-config-null-pointer");
assert_eq!(
unsafe { moq_publish_video_config(broadcast, std::ptr::null()) },
-6,
"null config should return InvalidPointer (-6)"
);
assert_eq!(
unsafe { moq_publish_video_properties(broadcast, std::ptr::null()) },
-6,
"null properties should return InvalidPointer (-6)"
);
assert_eq!(
unsafe { moq_publish_audio_config(broadcast, std::ptr::null()) },
-6,
"null config should return InvalidPointer (-6)"
);
assert_eq!(moq_publish_finish(broadcast), 0);
}
#[test]
fn publish_catalog_roundtrip() {
let origin = id(moq_origin_create());
let path = b"catalog-producer";
let broadcast = publish_broadcast(origin, path);
let video_name = "video";
let video_label = "Main camera";
let video_codec = "vp8";
let width: u32 = 1920;
let height: u32 = 1080;
let description: &[u8] = &[0x01, 0x02, 0x03];
let video = moq_video_config {
name: video_name.as_ptr() as *const c_char,
name_len: video_name.len(),
label: video_label.as_ptr() as *const c_char,
label_len: video_label.len(),
codec: video_codec.as_ptr() as *const c_char,
codec_len: video_codec.len(),
description: description.as_ptr(),
description_len: description.len(),
coded_width: width,
coded_height: height,
container: moq_container::default(),
};
assert_eq!(unsafe { moq_publish_video_config(broadcast, &video) }, 0);
let stalled_video_name = "video-stalled";
let stalled_video = moq_video_config {
name: stalled_video_name.as_ptr() as *const c_char,
name_len: stalled_video_name.len(),
label: std::ptr::null(),
label_len: 0,
codec: video_codec.as_ptr() as *const c_char,
codec_len: video_codec.len(),
description: description.as_ptr(),
description_len: description.len(),
coded_width: width,
coded_height: height,
container: moq_container::default(),
};
assert_eq!(unsafe { moq_publish_video_config(broadcast, &stalled_video) }, 0);
{
let mut state = State::lock();
let (_, catalog) = state.publish.pair_mut(Id::try_from(broadcast).unwrap()).unwrap();
catalog
.modify()
.unwrap()
.video
.renditions
.get_mut(stalled_video_name)
.unwrap()
.stalled = Some(true);
}
let properties = moq_video_properties {
display_width: 1080,
display_height: 1920,
has_display: true,
rotation: 315.0,
has_rotation: true,
flip: true,
has_flip: true,
};
assert_eq!(unsafe { moq_publish_video_properties(broadcast, &properties) }, 0);
let audio_name = "audio";
let audio_label = "English";
let audio_codec = "opus";
let audio = moq_audio_config {
name: audio_name.as_ptr() as *const c_char,
name_len: audio_name.len(),
label: audio_label.as_ptr() as *const c_char,
label_len: audio_label.len(),
codec: audio_codec.as_ptr() as *const c_char,
codec_len: audio_codec.len(),
description: std::ptr::null(),
description_len: 0,
sample_rate: 48000,
channel_count: 2,
container: moq_container::default(),
};
assert_eq!(unsafe { moq_publish_audio_config(broadcast, &audio) }, 0);
let consume = request_broadcast(origin, path);
let catalog_cb = Callback::new();
let catalog_task = id(unsafe { moq_consume_catalog(consume, Some(channel_callback), catalog_cb.ptr) });
let catalog_id = id(catalog_cb.recv());
let mut video_cfg = moq_video_config {
name: std::ptr::null(),
name_len: 0,
label: std::ptr::null(),
label_len: 0,
codec: std::ptr::null(),
codec_len: 0,
description: std::ptr::null(),
description_len: 0,
coded_width: 0,
coded_height: 0,
container: moq_container::default(),
};
assert_eq!(unsafe { moq_consume_video_config(catalog_id, 0, &mut video_cfg) }, 0);
let codec = unsafe {
std::str::from_utf8(std::slice::from_raw_parts(
video_cfg.codec.cast::<u8>(),
video_cfg.codec_len,
))
}
.unwrap();
assert_eq!(codec, "vp8");
assert_eq!(
borrowed_string(video_cfg.label, video_cfg.label_len).as_deref(),
Some("Main camera")
);
assert_eq!(video_cfg.coded_width, 1920);
assert_eq!(video_cfg.coded_height, 1080);
let mut stalled = std::mem::MaybeUninit::<bool>::uninit();
assert_eq!(
unsafe { moq_consume_video_stalled(catalog_id, 0, stalled.as_mut_ptr()) },
0
);
assert!(!unsafe { stalled.assume_init() });
let mut stalled = std::mem::MaybeUninit::<bool>::uninit();
assert_eq!(
unsafe { moq_consume_video_stalled(catalog_id, 1, stalled.as_mut_ptr()) },
0
);
assert!(unsafe { stalled.assume_init() });
assert_eq!(
unsafe { moq_consume_video_stalled(catalog_id, 0, std::ptr::null_mut()) },
-6,
"null stalled pointer should return InvalidPointer (-6)"
);
assert_eq!(
unsafe {
moq_publish_video_remove(
broadcast,
stalled_video_name.as_ptr() as *const c_char,
stalled_video_name.len(),
)
},
0
);
let active_catalog_id = id(catalog_cb.recv());
let mut properties = moq_video_properties::default();
assert_eq!(unsafe { moq_consume_video_properties(catalog_id, &mut properties) }, 0);
assert!(properties.has_display);
assert_eq!(properties.display_width, 1080);
assert_eq!(properties.display_height, 1920);
assert!(properties.has_rotation);
assert_eq!(properties.rotation, 0.0);
assert!(properties.has_flip);
assert!(properties.flip);
let mut audio_cfg = moq_audio_config {
name: std::ptr::null(),
name_len: 0,
label: std::ptr::null(),
label_len: 0,
codec: std::ptr::null(),
codec_len: 0,
description: std::ptr::null(),
description_len: 0,
sample_rate: 0,
channel_count: 0,
container: moq_container::default(),
};
assert_eq!(unsafe { moq_consume_audio_config(catalog_id, 0, &mut audio_cfg) }, 0);
assert_eq!(
borrowed_string(audio_cfg.label, audio_cfg.label_len).as_deref(),
Some("English")
);
assert_eq!(audio_cfg.sample_rate, 48000);
assert_eq!(audio_cfg.channel_count, 2);
assert_eq!(
unsafe { moq_publish_video_remove(broadcast, video_name.as_ptr() as *const c_char, video_name.len()) },
0
);
let catalog_id2 = id(catalog_cb.recv());
assert!(
unsafe { moq_consume_video_config(catalog_id2, 0, &mut video_cfg) } < 0,
"video rendition should be gone after remove"
);
assert_eq!(unsafe { moq_consume_audio_config(catalog_id2, 0, &mut audio_cfg) }, 0);
assert_eq!(moq_consume_catalog_free(catalog_id), 0);
assert_eq!(moq_consume_catalog_free(active_catalog_id), 0);
assert_eq!(moq_consume_catalog_free(catalog_id2), 0);
assert_eq!(moq_consume_catalog_cancel(catalog_task), 0);
assert_eq!(catalog_cb.recv_terminal(), 0, "catalog close delivers terminal 0");
assert_eq!(moq_consume_close(consume), 0);
assert_eq!(moq_publish_finish(broadcast), 0);
assert_eq!(moq_origin_close(origin), 0);
}
#[test]
fn a_half_specified_coded_size_round_trips() {
let origin = id(moq_origin_create());
let path = b"half-coded-size";
let broadcast = publish_broadcast(origin, path);
let name = "video";
let codec = "vp8";
let mut video = moq_video_config {
name: name.as_ptr() as *const c_char,
name_len: name.len(),
label: std::ptr::null(),
label_len: 0,
codec: codec.as_ptr() as *const c_char,
codec_len: codec.len(),
description: std::ptr::null(),
description_len: 0,
coded_width: 1920,
coded_height: 0, container: moq_container::default(),
};
assert_eq!(unsafe { moq_publish_video_config(broadcast, &video) }, 0);
let consume = request_broadcast(origin, path);
let catalog_cb = Callback::new();
let catalog_task = id(unsafe { moq_consume_catalog(consume, Some(channel_callback), catalog_cb.ptr) });
let catalog = id(catalog_cb.recv());
let mut read = moq_video_config {
name: std::ptr::null(),
name_len: 0,
label: std::ptr::null(),
label_len: 0,
codec: std::ptr::null(),
codec_len: 0,
description: std::ptr::null(),
description_len: 0,
coded_width: 0,
coded_height: 0,
container: moq_container::default(),
};
assert_eq!(unsafe { moq_consume_video_config(catalog, 0, &mut read) }, 0);
assert_eq!(read.coded_width, 1920);
assert_eq!(read.coded_height, 0, "the absent height must not come back invented");
let forward = publish_broadcast(origin, b"half-coded-size-forwarded");
video.coded_width = read.coded_width;
video.coded_height = read.coded_height;
assert_eq!(unsafe { moq_publish_video_config(forward, &video) }, 0);
let forwarded = request_broadcast(origin, b"half-coded-size-forwarded");
let forwarded_cb = Callback::new();
let forwarded_task = id(unsafe { moq_consume_catalog(forwarded, Some(channel_callback), forwarded_cb.ptr) });
let forwarded_catalog = id(forwarded_cb.recv());
assert_eq!(unsafe { moq_consume_video_config(forwarded_catalog, 0, &mut read) }, 0);
assert_eq!(read.coded_width, 1920);
assert_eq!(read.coded_height, 0, "forwarding must not invent the absent height");
assert_eq!(moq_consume_catalog_free(forwarded_catalog), 0);
assert_eq!(moq_consume_catalog_cancel(forwarded_task), 0);
assert_eq!(forwarded_cb.recv_catalog_terminal(), 0);
assert_eq!(moq_consume_close(forwarded), 0);
assert_eq!(moq_publish_finish(forward), 0);
assert_eq!(moq_consume_catalog_free(catalog), 0);
assert_eq!(moq_consume_catalog_cancel(catalog_task), 0);
assert_eq!(catalog_cb.recv_catalog_terminal(), 0);
assert_eq!(moq_consume_close(consume), 0);
assert_eq!(moq_publish_finish(broadcast), 0);
assert_eq!(moq_origin_close(origin), 0);
}
#[test]
fn raw_loc_video_uses_the_declared_catalog_container() {
let origin = id(moq_origin_create());
let path = b"raw-loc-video";
let broadcast = publish_broadcast(origin, path);
let name = b"video";
let track =
id(unsafe { moq_publish_track(broadcast, name.as_ptr() as *const c_char, name.len(), std::ptr::null()) });
let codec = b"vp8";
let video = moq_video_config {
name: name.as_ptr() as *const c_char,
name_len: name.len(),
label: std::ptr::null(),
label_len: 0,
codec: codec.as_ptr() as *const c_char,
codec_len: codec.len(),
description: std::ptr::null(),
description_len: 0,
coded_width: 0,
coded_height: 0,
container: moq_container {
kind: moq_container_kind::MOQ_CONTAINER_KIND_LOC as u32,
init: std::ptr::null(),
init_len: 0,
},
};
assert_eq!(unsafe { moq_publish_video_config(broadcast, &video) }, 0);
let consume = request_broadcast(origin, path);
let catalog_cb = Callback::new();
let catalog_task = id(unsafe { moq_consume_catalog(consume, Some(channel_callback), catalog_cb.ptr) });
let catalog = id(catalog_cb.recv());
let mut video_cfg = moq_video_config {
name: std::ptr::null(),
name_len: 0,
label: std::ptr::null(),
label_len: 0,
codec: std::ptr::null(),
codec_len: 0,
description: std::ptr::null(),
description_len: 0,
coded_width: 0,
coded_height: 0,
container: moq_container::default(),
};
assert_eq!(unsafe { moq_consume_video_config(catalog, 0, &mut video_cfg) }, 0);
assert_eq!(
video_cfg.container.kind,
moq_container_kind::MOQ_CONTAINER_KIND_LOC as u32
);
assert!(video_cfg.container.init.is_null());
let frame_cb = Callback::new();
let consumer = id(unsafe { moq_consume_video(catalog, 0, 10_000_000, Some(channel_callback), frame_cb.ptr) });
let timestamp_us = 42_000;
let payload = b"codec frame";
let loc = moq_loc::encode(timestamp_us, payload).unwrap();
let group = id(moq_publish_track_group(track));
assert_eq!(
unsafe { moq_publish_group_frame(group, loc.as_ptr(), loc.len(), timestamp_us) },
0
);
assert_eq!(moq_publish_group_finish(group), 0);
let frame_id = id(frame_cb.recv());
let mut frame = moq_frame {
payload: std::ptr::null(),
payload_size: 0,
timestamp_us: 0,
keyframe: false,
};
assert_eq!(unsafe { moq_consume_frame(frame_id, &mut frame) }, 0);
assert_eq!(frame.timestamp_us, timestamp_us);
assert!(frame.keyframe);
assert_eq!(
unsafe { std::slice::from_raw_parts(frame.payload, frame.payload_size) },
payload
);
assert_eq!(moq_consume_frame_free(frame_id), 0);
assert_eq!(moq_consume_video_cancel(consumer), 0);
assert_eq!(frame_cb.recv_terminal(), 0);
assert_eq!(moq_consume_catalog_free(catalog), 0);
assert_eq!(moq_consume_catalog_cancel(catalog_task), 0);
assert_eq!(catalog_cb.recv_terminal(), 0);
assert_eq!(moq_consume_close(consume), 0);
assert_eq!(moq_publish_track_finish(track), 0);
assert_eq!(moq_publish_finish(broadcast), 0);
assert_eq!(moq_origin_close(origin), 0);
}
#[test]
fn cmaf_catalog_container_carries_its_init_segment() {
let origin = id(moq_origin_create());
let path = b"cmaf-container";
let broadcast = publish_broadcast(origin, path);
let name = "audio";
let codec = "opus";
let init: &[u8] = &[0x00, 0x01, 0x02, 0x03];
let audio = moq_audio_config {
name: name.as_ptr() as *const c_char,
name_len: name.len(),
label: std::ptr::null(),
label_len: 0,
codec: codec.as_ptr() as *const c_char,
codec_len: codec.len(),
description: std::ptr::null(),
description_len: 0,
sample_rate: 48000,
channel_count: 2,
container: moq_container {
kind: moq_container_kind::MOQ_CONTAINER_KIND_CMAF as u32,
init: init.as_ptr(),
init_len: init.len(),
},
};
assert_eq!(unsafe { moq_publish_audio_config(broadcast, &audio) }, 0);
let consume = request_broadcast(origin, path);
let catalog_cb = Callback::new();
let catalog_task = id(unsafe { moq_consume_catalog(consume, Some(channel_callback), catalog_cb.ptr) });
let catalog = id(catalog_cb.recv());
let mut audio_cfg = moq_audio_config {
name: std::ptr::null(),
name_len: 0,
label: std::ptr::null(),
label_len: 0,
codec: std::ptr::null(),
codec_len: 0,
description: std::ptr::null(),
description_len: 0,
sample_rate: 0,
channel_count: 0,
container: moq_container::default(),
};
assert_eq!(unsafe { moq_consume_audio_config(catalog, 0, &mut audio_cfg) }, 0);
assert_eq!(
audio_cfg.container.kind,
moq_container_kind::MOQ_CONTAINER_KIND_CMAF as u32
);
assert_eq!(
unsafe { std::slice::from_raw_parts(audio_cfg.container.init, audio_cfg.container.init_len) },
init
);
assert_eq!(moq_consume_catalog_free(catalog), 0);
assert_eq!(moq_consume_catalog_cancel(catalog_task), 0);
assert_eq!(catalog_cb.recv_terminal(), 0);
assert_eq!(moq_consume_close(consume), 0);
assert_eq!(moq_publish_finish(broadcast), 0);
assert_eq!(moq_origin_close(origin), 0);
}
#[test]
fn unpublishable_catalog_containers_are_rejected() {
let origin = id(moq_origin_create());
let broadcast = publish_broadcast(origin, b"container-reject");
let name = "video";
let codec = "vp8";
let config = |container| moq_video_config {
name: name.as_ptr() as *const c_char,
name_len: name.len(),
label: std::ptr::null(),
label_len: 0,
codec: codec.as_ptr() as *const c_char,
codec_len: codec.len(),
description: std::ptr::null(),
description_len: 0,
coded_width: 0,
coded_height: 0,
container,
};
let unknown = config(moq_container {
kind: moq_container_kind::MOQ_CONTAINER_KIND_UNKNOWN as u32,
init: std::ptr::null(),
init_len: 0,
});
assert_eq!(
unsafe { moq_publish_video_config(broadcast, &unknown) },
-15,
"unknown container should return InvalidCode (-15)"
);
let garbage = config(moq_container {
kind: 12345,
init: std::ptr::null(),
init_len: 0,
});
assert_eq!(
unsafe { moq_publish_video_config(broadcast, &garbage) },
-15,
"out of range container should return InvalidCode (-15)"
);
let cmaf = config(moq_container {
kind: moq_container_kind::MOQ_CONTAINER_KIND_CMAF as u32,
init: std::ptr::null(),
init_len: 0,
});
assert_eq!(
unsafe { moq_publish_video_config(broadcast, &cmaf) },
-6,
"cmaf without an init segment should return InvalidPointer (-6)"
);
assert_eq!(moq_publish_finish(broadcast), 0);
assert_eq!(moq_origin_close(origin), 0);
}
#[test]
fn catalog_section_roundtrip() {
let origin = id(moq_origin_create());
let path = b"catalog-sections";
let broadcast = publish_broadcast(origin, path);
let name_a = b"viewers";
let json_a = br#"{"count":42}"#;
assert_eq!(
unsafe {
moq_publish_catalog_section(
broadcast,
name_a.as_ptr() as *const c_char,
name_a.len(),
json_a.as_ptr() as *const c_char,
json_a.len(),
)
},
0
);
let name_b = b"title";
let json_b = br#""hello world""#;
assert_eq!(
unsafe {
moq_publish_catalog_section(
broadcast,
name_b.as_ptr() as *const c_char,
name_b.len(),
json_b.as_ptr() as *const c_char,
json_b.len(),
)
},
0
);
let reserved = b"video";
let empty = b"{}";
assert!(
unsafe {
moq_publish_catalog_section(
broadcast,
reserved.as_ptr() as *const c_char,
reserved.len(),
empty.as_ptr() as *const c_char,
empty.len(),
)
} < 0,
"reserved section name should fail"
);
let bad = b"not json";
assert_eq!(
unsafe {
moq_publish_catalog_section(
broadcast,
name_a.as_ptr() as *const c_char,
name_a.len(),
bad.as_ptr() as *const c_char,
bad.len(),
)
},
-37,
"invalid JSON should return the Json error code"
);
let consume = request_broadcast(origin, path);
let catalog_cb = Callback::new();
let catalog_task = id(unsafe { moq_consume_catalog(consume, Some(channel_callback), catalog_cb.ptr) });
let catalog_id = id(catalog_cb.recv());
let count = moq_consume_catalog_section_count(catalog_id);
assert_eq!(count, 2, "expected two sections, got {count}");
let mut found_a = false;
let mut found_b = false;
for index in 0..count as u32 {
let mut section = moq_section {
name: std::ptr::null(),
name_len: 0,
json: std::ptr::null(),
json_len: 0,
};
assert_eq!(
unsafe { moq_consume_catalog_section_at(catalog_id, index, &mut section) },
0
);
let name = unsafe { std::slice::from_raw_parts(section.name.cast::<u8>(), section.name_len) };
let json = unsafe { std::slice::from_raw_parts(section.json.cast::<u8>(), section.json_len) };
match name {
n if n == name_a => {
found_a = true;
assert_eq!(json, json_a);
}
n if n == name_b => {
found_b = true;
assert_eq!(json, json_b);
}
other => panic!("unexpected section name: {:?}", std::str::from_utf8(other)),
}
}
assert!(found_a && found_b, "both sections should be present");
let mut value = moq_string {
data: std::ptr::null(),
len: 0,
};
assert_eq!(
unsafe { moq_consume_catalog_section(catalog_id, name_a.as_ptr() as *const c_char, name_a.len(), &mut value) },
0
);
let got = unsafe { std::slice::from_raw_parts(value.data.cast::<u8>(), value.len) };
assert_eq!(got, json_a);
let missing = b"nope";
assert!(
unsafe {
moq_consume_catalog_section(catalog_id, missing.as_ptr() as *const c_char, missing.len(), &mut value)
} < 0,
"missing section should fail"
);
assert_eq!(
unsafe { moq_publish_catalog_section_remove(broadcast, name_a.as_ptr() as *const c_char, name_a.len()) },
0
);
let catalog_id2 = id(catalog_cb.recv());
assert_eq!(
moq_consume_catalog_section_count(catalog_id2),
1,
"one section should remain after remove"
);
assert!(
unsafe { moq_consume_catalog_section(catalog_id2, name_a.as_ptr() as *const c_char, name_a.len(), &mut value) }
< 0,
"removed section should be gone"
);
assert_eq!(moq_consume_catalog_free(catalog_id), 0);
assert_eq!(moq_consume_catalog_free(catalog_id2), 0);
assert_eq!(moq_consume_catalog_cancel(catalog_task), 0);
assert_eq!(catalog_cb.recv_terminal(), 0, "catalog close delivers terminal 0");
assert_eq!(moq_consume_close(consume), 0);
assert_eq!(moq_publish_finish(broadcast), 0);
assert_eq!(moq_origin_close(origin), 0);
}
#[test]
fn publish_track_invalid_broadcast() {
let name = b"data";
assert!(unsafe { moq_publish_track(0, name.as_ptr() as *const c_char, name.len(), std::ptr::null()) } < 0);
let info = moq_track_info {
priority: 1,
max_age_us: 0,
max_age_present: false,
timescale: 0,
timescale_present: false,
};
assert!(unsafe { moq_publish_track(0, name.as_ptr() as *const c_char, name.len(), &info) } < 0);
assert!(moq_publish_track_group(9999) < 0);
assert!(unsafe { moq_publish_track_frame(9999, name.as_ptr(), name.len(), 0) } < 0);
assert!(unsafe { moq_publish_group_frame(9999, name.as_ptr(), name.len(), 0) } < 0);
assert!(moq_publish_track_finish(9999) < 0);
assert!(moq_publish_group_finish(9999) < 0);
let subscription = moq_subscription {
priority: 1,
max_age_us: 0,
group_start: 0,
group_start_present: false,
group_end: 0,
group_end_present: false,
};
assert!(unsafe { moq_consume_track_update(9999, &subscription) } < 0);
}
#[test]
fn publish_track_with_info_rejects_invalid_timescale() {
let origin = id(moq_origin_create());
let broadcast = publish_broadcast(origin, b"publish-track-with-info-rejects-invalid-timescale");
let name = b"data";
let info = moq_track_info {
priority: 0,
max_age_us: 0,
max_age_present: false,
timescale: 0,
timescale_present: true,
};
assert!(unsafe { moq_publish_track(broadcast, name.as_ptr() as *const c_char, name.len(), &info) } < 0);
assert_eq!(moq_publish_finish(broadcast), 0);
}
#[test]
fn raw_track_publish_consume() {
let origin = id(moq_origin_create());
let path = b"raw-track";
let broadcast = publish_broadcast(origin, path);
let track_name = b"data";
let track = id(unsafe {
moq_publish_track(
broadcast,
track_name.as_ptr() as *const c_char,
track_name.len(),
std::ptr::null(),
)
});
let consume = request_broadcast(origin, path);
let frame_cb = Callback::new();
let subscription = moq_subscription {
priority: 0,
max_age_us: 1_000_000,
group_start: 0,
group_start_present: false,
group_end: 0,
group_end_present: false,
};
let consumer = id(unsafe {
moq_consume_track(
consume,
track_name.as_ptr() as *const c_char,
track_name.len(),
&subscription,
Some(channel_callback),
frame_cb.ptr,
)
});
let payload = b"hello raw track";
let timestamp_us = 12_345;
assert_eq!(
unsafe { moq_publish_track_frame(track, payload.as_ptr(), payload.len(), timestamp_us) },
0
);
let frame_id = id(frame_cb.recv());
let mut frame = moq_frame {
payload: std::ptr::null(),
payload_size: 0,
timestamp_us: 0,
keyframe: true, };
assert_eq!(unsafe { moq_consume_track_frame(frame_id, &mut frame) }, 0);
let received = unsafe { std::slice::from_raw_parts(frame.payload, frame.payload_size) };
assert_eq!(received, payload);
assert_eq!(frame.timestamp_us, timestamp_us);
assert!(!frame.keyframe, "raw frames have no keyframe flag");
assert_eq!(moq_consume_track_frame_free(frame_id), 0);
let group = id(moq_publish_track_group(track));
let parts: [(&[u8], u64); 2] = [(b"part-0", 20_000), (b"part-1", 30_000)];
for (part, timestamp_us) in parts {
assert_eq!(
unsafe { moq_publish_group_frame(group, part.as_ptr(), part.len(), timestamp_us) },
0
);
}
assert_eq!(moq_publish_group_finish(group), 0);
for (expected, timestamp_us) in parts {
let frame_id = id(frame_cb.recv());
let mut frame = moq_frame {
payload: std::ptr::null(),
payload_size: 0,
timestamp_us: 0,
keyframe: false,
};
assert_eq!(unsafe { moq_consume_track_frame(frame_id, &mut frame) }, 0);
let received = unsafe { std::slice::from_raw_parts(frame.payload, frame.payload_size) };
assert_eq!(received, expected);
assert_eq!(frame.timestamp_us, timestamp_us);
assert_eq!(moq_consume_track_frame_free(frame_id), 0);
}
assert_eq!(moq_consume_track_cancel(consumer), 0);
assert_eq!(frame_cb.recv_terminal(), 0, "clean close delivers terminal 0");
assert!(moq_consume_track_cancel(consumer) < 0, "double-close should fail");
assert_eq!(moq_publish_track_finish(track), 0);
assert!(moq_publish_track_finish(track) < 0, "double-close should fail");
assert_eq!(moq_consume_close(consume), 0);
assert_eq!(moq_publish_finish(broadcast), 0);
assert_eq!(moq_origin_close(origin), 0);
}
#[test]
fn raw_track_datagram_publish_consume() {
let origin = id(moq_origin_create());
let path = b"raw-datagram";
let broadcast = publish_broadcast(origin, path);
let track_name = b"events";
let track = id(unsafe {
moq_publish_track(
broadcast,
track_name.as_ptr() as *const c_char,
track_name.len(),
std::ptr::null(),
)
});
let consume = request_broadcast(origin, path);
let dg_cb = Callback::new();
let consumer = id(unsafe {
moq_consume_datagrams(
consume,
track_name.as_ptr() as *const c_char,
track_name.len(),
Some(channel_callback),
dg_cb.ptr,
)
});
let payload = b"hello datagram";
let mut sequence: u64 = u64::MAX;
assert_eq!(
unsafe { moq_publish_track_datagram(track, payload.as_ptr(), payload.len(), 120_000, &mut sequence) },
0
);
let dg_id = id(dg_cb.recv());
let mut datagram = moq_datagram {
payload: std::ptr::null(),
payload_size: 0,
timestamp_us: 0,
sequence: 0,
};
assert_eq!(unsafe { moq_consume_datagram(dg_id, &mut datagram) }, 0);
let received = unsafe { std::slice::from_raw_parts(datagram.payload, datagram.payload_size) };
assert_eq!(received, payload);
assert_eq!(datagram.timestamp_us, 120_000);
assert_eq!(datagram.sequence, sequence);
assert_eq!(moq_consume_datagram_free(dg_id), 0);
assert_eq!(moq_consume_datagrams_cancel(consumer), 0);
assert_eq!(dg_cb.recv_terminal(), 0, "clean close delivers terminal 0");
assert!(moq_consume_datagrams_cancel(consumer) < 0, "double-close should fail");
assert_eq!(moq_publish_track_finish(track), 0);
assert_eq!(moq_consume_close(consume), 0);
assert_eq!(moq_publish_finish(broadcast), 0);
assert_eq!(moq_origin_close(origin), 0);
}
#[test]
fn raw_track_sparse_groups_and_known_end() {
let origin = id(moq_origin_create());
let broadcast = publish_broadcast(origin, b"raw-track-sparse-groups-and-known-end");
let name = b"sparse";
let track =
id(unsafe { moq_publish_track(broadcast, name.as_ptr() as *const c_char, name.len(), std::ptr::null()) });
let group = id(moq_publish_track_group_at(track, 2));
assert_eq!(moq_publish_group_finish(group), 0);
assert_eq!(moq_publish_track_finish_at(track, 5), 0);
let group = id(moq_publish_track_group_at(track, 4));
assert_eq!(moq_publish_group_finish(group), 0);
assert!(moq_publish_track_group_at(track, 5) < 0);
assert_eq!(moq_publish_track_finish(track), 0);
assert_eq!(moq_publish_finish(broadcast), 0);
}
#[test]
fn raw_track_and_group_abort_consume_their_handles() {
let origin = id(moq_origin_create());
let broadcast = publish_broadcast(origin, b"raw-track-and-group-abort-consume-their-handles");
let name = b"aborted";
let track =
id(unsafe { moq_publish_track(broadcast, name.as_ptr() as *const c_char, name.len(), std::ptr::null()) });
let group = id(moq_publish_track_group(track));
assert_eq!(moq_publish_group_abort(group, 409), 0);
assert!(moq_publish_group_finish(group) < 0);
assert_eq!(moq_publish_track_abort(track, 410), 0);
assert!(moq_publish_track_finish(track) < 0);
assert_eq!(moq_publish_finish(broadcast), 0);
}
#[test]
fn raw_track_subscription_options_and_update() {
let origin = id(moq_origin_create());
let path = b"raw-track-options";
let broadcast = publish_broadcast(origin, path);
let track_name = b"data";
let info = moq_track_info {
priority: 3,
max_age_us: 1_000_000,
max_age_present: true,
timescale: 1_000_000,
timescale_present: true,
};
let track =
id(unsafe { moq_publish_track(broadcast, track_name.as_ptr() as *const c_char, track_name.len(), &info) });
let payloads: [&[u8]; 3] = [b"zero", b"one", b"two"];
for (i, payload) in payloads.into_iter().enumerate() {
assert_eq!(
unsafe { moq_publish_track_frame(track, payload.as_ptr(), payload.len(), i as u64 * 20_000) },
0
);
}
let consume = request_broadcast(origin, path);
let frame_cb = Callback::new();
let subscription = moq_subscription {
priority: 5,
max_age_us: 25_000,
group_start: 1,
group_start_present: true,
group_end: 2,
group_end_present: true,
};
let consumer = id(unsafe {
moq_consume_track(
consume,
track_name.as_ptr() as *const c_char,
track_name.len(),
&subscription,
Some(channel_callback),
frame_cb.ptr,
)
});
let frame_id = id(frame_cb.recv());
let mut frame = moq_frame {
payload: std::ptr::null(),
payload_size: 0,
timestamp_us: 0,
keyframe: false,
};
assert_eq!(unsafe { moq_consume_track_frame(frame_id, &mut frame) }, 0);
let received = unsafe { std::slice::from_raw_parts(frame.payload, frame.payload_size) };
assert_eq!(received, b"one");
assert_eq!(frame.timestamp_us, 20_000);
assert_eq!(moq_consume_track_frame_free(frame_id), 0);
let update = moq_subscription {
group_end: 3,
..subscription
};
assert_eq!(unsafe { moq_consume_track_update(consumer, &update) }, 0);
let frame_id = id(frame_cb.recv());
let mut frame = moq_frame {
payload: std::ptr::null(),
payload_size: 0,
timestamp_us: 0,
keyframe: false,
};
assert_eq!(unsafe { moq_consume_track_frame(frame_id, &mut frame) }, 0);
let received = unsafe { std::slice::from_raw_parts(frame.payload, frame.payload_size) };
assert_eq!(received, b"two");
assert_eq!(frame.timestamp_us, 40_000);
assert_eq!(moq_consume_track_frame_free(frame_id), 0);
assert_eq!(moq_consume_track_cancel(consumer), 0);
assert_eq!(frame_cb.recv_terminal(), 0);
assert_eq!(moq_publish_track_finish(track), 0);
assert_eq!(moq_consume_close(consume), 0);
assert_eq!(moq_publish_finish(broadcast), 0);
assert_eq!(moq_origin_close(origin), 0);
}
#[test]
fn json_snapshot_publish_consume() {
let origin = id(moq_origin_create());
let path = b"json-snapshot";
let broadcast = publish_broadcast(origin, path);
let track_name = b"meta";
let config = moq_json_snapshot_config {
delta_ratio: 8,
compression: true,
};
let producer = id(unsafe {
moq_publish_json_snapshot(
broadcast,
track_name.as_ptr() as *const c_char,
track_name.len(),
&config,
)
});
let consume = request_broadcast(origin, path);
let value_cb = Callback::new();
let consumer = id(unsafe {
moq_consume_json_snapshot(
consume,
track_name.as_ptr() as *const c_char,
track_name.len(),
&config,
Some(channel_callback),
value_cb.ptr,
)
});
for expected in [r#"{"a":1}"#, r#"{"a":2}"#] {
assert_eq!(
unsafe { moq_publish_json_snapshot_update(producer, expected.as_ptr() as *const c_char, expected.len()) },
0
);
let value_id = id(value_cb.recv());
let mut value = moq_json_value {
json: std::ptr::null(),
json_len: 0,
};
assert_eq!(unsafe { moq_consume_json_value(value_id, &mut value) }, 0);
let received = unsafe { std::slice::from_raw_parts(value.json.cast::<u8>(), value.json_len) };
assert_eq!(
serde_json::from_slice::<serde_json::Value>(received).unwrap(),
serde_json::from_str::<serde_json::Value>(expected).unwrap()
);
assert_eq!(moq_consume_json_value_free(value_id), 0);
}
assert_eq!(moq_consume_json_cancel(consumer), 0);
assert_eq!(value_cb.recv_terminal(), 0, "clean close delivers terminal 0");
assert!(moq_consume_json_cancel(consumer) < 0, "double-close should fail");
assert_eq!(moq_publish_json_snapshot_finish(producer), 0);
assert!(
moq_publish_json_snapshot_finish(producer) < 0,
"double-close should fail"
);
assert_eq!(moq_consume_close(consume), 0);
assert_eq!(moq_publish_finish(broadcast), 0);
assert_eq!(moq_origin_close(origin), 0);
}
#[test]
fn json_stream_publish_consume() {
let origin = id(moq_origin_create());
let path = b"json-stream";
let broadcast = publish_broadcast(origin, path);
let track_name = b"events";
let config = moq_json_stream_config { compression: true };
let producer = id(unsafe {
moq_publish_json_stream(
broadcast,
track_name.as_ptr() as *const c_char,
track_name.len(),
&config,
)
});
let consume = request_broadcast(origin, path);
let value_cb = Callback::new();
let consumer = id(unsafe {
moq_consume_json_stream(
consume,
track_name.as_ptr() as *const c_char,
track_name.len(),
&config,
Some(channel_callback),
value_cb.ptr,
)
});
for expected in [r#"{"n":0}"#, r#"{"n":1}"#, r#"{"n":2}"#] {
assert_eq!(
unsafe { moq_publish_json_stream_append(producer, expected.as_ptr() as *const c_char, expected.len()) },
0
);
let value_id = id(value_cb.recv());
let mut value = moq_json_value {
json: std::ptr::null(),
json_len: 0,
};
assert_eq!(unsafe { moq_consume_json_value(value_id, &mut value) }, 0);
let received = unsafe { std::slice::from_raw_parts(value.json.cast::<u8>(), value.json_len) };
assert_eq!(
serde_json::from_slice::<serde_json::Value>(received).unwrap(),
serde_json::from_str::<serde_json::Value>(expected).unwrap()
);
assert_eq!(moq_consume_json_value_free(value_id), 0);
}
assert_eq!(moq_consume_json_cancel(consumer), 0);
assert_eq!(value_cb.recv_terminal(), 0, "clean close delivers terminal 0");
assert!(moq_consume_json_cancel(consumer) < 0, "double-close should fail");
assert_eq!(moq_publish_json_stream_finish(producer), 0);
assert!(moq_publish_json_stream_finish(producer) < 0, "double-close should fail");
assert_eq!(moq_consume_close(consume), 0);
assert_eq!(moq_publish_finish(broadcast), 0);
assert_eq!(moq_origin_close(origin), 0);
}
#[test]
fn close_invalid_or_zero_ids() {
assert!(moq_origin_close(9999) < 0);
assert!(moq_session_close(9999) < 0);
assert!(moq_publish_finish(9999) < 0);
assert!(moq_consume_close(9999) < 0);
assert!(moq_consume_frame_free(9999) < 0);
assert!(moq_origin_close(0) < 0);
assert!(moq_session_close(0) < 0);
assert!(moq_publish_finish(0) < 0);
}
#[test]
fn announced_free_lifecycle() {
let origin = id(moq_origin_create());
let path = b"announced-free";
let broadcast = publish_broadcast(origin, path);
let ann_cb = Callback::new();
let ann_task = id(unsafe {
moq_origin_announced(
origin,
std::ptr::null(),
0,
std::ptr::null(),
0,
Some(channel_callback),
ann_cb.ptr,
)
});
let announced = id(ann_cb.recv());
let mut info = moq_announce_update {
prefix: std::ptr::null(),
prefix_len: 0,
captures: std::ptr::null(),
captures_len: 0,
has_captures: false,
active: false,
};
assert_eq!(unsafe { moq_origin_announced_info(announced, &mut info) }, 0);
assert!(info.active, "broadcast should be active");
let got = unsafe { std::slice::from_raw_parts(info.prefix.cast::<u8>(), info.prefix_len) };
assert_eq!(got, path, "announced prefix should match");
assert_eq!(moq_origin_announced_free(announced), 0);
assert!(moq_origin_announced_free(announced) < 0, "double-free should fail");
assert!(
unsafe { moq_origin_announced_info(announced, &mut info) } < 0,
"info on a freed handle should fail"
);
assert_eq!(moq_origin_announced_cancel(ann_task), 0);
ann_cb.recv_terminal();
assert_eq!(moq_origin_close(origin), 0);
assert_eq!(moq_publish_finish(broadcast), 0);
}
#[test]
fn double_close_all_resource_types() {
let origin = id(moq_origin_create());
assert_eq!(moq_origin_close(origin), 0);
assert!(moq_origin_close(origin) < 0);
let origin = id(moq_origin_create());
let broadcast = publish_broadcast(origin, b"double-close-all-resource-types");
let init = opus_head();
let format = moq_audio_format::MOQ_AUDIO_FORMAT_OPUS;
let media = id(publish_audio(broadcast, format, &init, None));
assert_eq!(moq_publish_media_finish(media), 0);
assert!(moq_publish_media_finish(media) < 0);
assert_eq!(moq_publish_finish(broadcast), 0);
let origin = id(moq_origin_create());
let path = b"double-close-test";
let broadcast = publish_broadcast(origin, path);
let init = opus_head();
let media = id(publish_audio(broadcast, format, &init, None));
let consume = request_broadcast(origin, path);
let catalog_cb = Callback::new();
let catalog_task = id(unsafe { moq_consume_catalog(consume, Some(channel_callback), catalog_cb.ptr) });
let catalog_id = id(catalog_cb.recv());
let frame_cb = Callback::new();
let track = id(unsafe { moq_consume_audio(catalog_id, 0, 10_000_000, Some(channel_callback), frame_cb.ptr) });
let payload = b"test";
assert_eq!(
unsafe { moq_publish_media_frame(media, payload.as_ptr(), payload.len(), 1_000_000) },
0
);
let frame_id = id(frame_cb.recv());
assert_eq!(moq_consume_frame_free(frame_id), 0);
assert!(moq_consume_frame_free(frame_id) < 0);
assert_eq!(moq_consume_audio_cancel(track), 0);
assert_eq!(frame_cb.recv_terminal(), 0, "audio close delivers terminal 0");
assert!(moq_consume_audio_cancel(track) < 0);
assert_eq!(moq_consume_catalog_free(catalog_id), 0);
assert!(moq_consume_catalog_free(catalog_id) < 0);
assert_eq!(moq_consume_catalog_cancel(catalog_task), 0);
assert_eq!(catalog_cb.recv_terminal(), 0, "catalog close delivers terminal 0");
assert!(moq_consume_catalog_cancel(catalog_task) < 0);
assert_eq!(moq_consume_close(consume), 0);
assert_eq!(moq_publish_media_finish(media), 0);
assert_eq!(moq_publish_finish(broadcast), 0);
assert_eq!(moq_origin_close(origin), 0);
}
#[test]
fn media_cut_bounds_audio_groups() {
let origin = id(moq_origin_create());
let broadcast = publish_broadcast(origin, b"media-cut");
let init = opus_head();
let format = moq_audio_format::MOQ_AUDIO_FORMAT_OPUS;
let media = id(publish_audio(broadcast, format, &init, None));
let payload = b"test";
for i in 0..3u64 {
assert_eq!(
unsafe { moq_publish_media_frame(media, payload.as_ptr(), payload.len(), i * 20_000) },
0
);
assert_eq!(moq_publish_media_flush(media, i * 20_000), 0);
assert_eq!(moq_publish_media_cut(media), 0, "each packet is its own group");
}
assert_eq!(
unsafe { moq_publish_media_frame(media, payload.as_ptr(), payload.len(), 60_000) },
0
);
assert_eq!(moq_publish_media_seek(media, 42), 0);
assert!(moq_publish_media_flush(9999, 0) < 0);
assert!(moq_publish_media_flush(media, u64::MAX) < 0);
assert!(moq_publish_media_cut(9999) < 0);
assert!(moq_publish_media_seek(9999, 0) < 0);
assert_eq!(moq_publish_media_finish(media), 0);
assert_eq!(moq_publish_finish(broadcast), 0);
assert_eq!(moq_origin_close(origin), 0);
}
#[test]
fn unknown_format() {
let origin = id(moq_origin_create());
let broadcast = publish_broadcast(origin, b"unknown-format");
let _guard = Guard(Some(|| {
moq_publish_finish(broadcast);
}));
let config = moq_audio_init {
format: 9999,
init: std::ptr::null(),
init_len: 0,
label: std::ptr::null(),
label_len: 0,
};
let ret = unsafe { moq_publish_audio(broadcast, &config) };
assert!(ret < 0, "an out-of-range format code should fail");
}
#[test]
fn local_announce() {
let origin = id(moq_origin_create());
let cb = Callback::new();
let announced_task = id(unsafe {
moq_origin_announced(
origin,
std::ptr::null(),
0,
std::ptr::null(),
0,
Some(channel_callback),
cb.ptr,
)
});
let path = b"test/broadcast";
let broadcast = publish_broadcast(origin, path);
let announced_id = id(cb.recv());
let mut info = moq_announce_update {
prefix: std::ptr::null(),
prefix_len: 0,
captures: std::ptr::null(),
captures_len: 0,
has_captures: false,
active: false,
};
assert_eq!(unsafe { moq_origin_announced_info(announced_id, &mut info) }, 0);
assert!(info.active, "broadcast should be active");
let announced_prefix =
unsafe { std::str::from_utf8(std::slice::from_raw_parts(info.prefix.cast::<u8>(), info.prefix_len)).unwrap() };
assert_eq!(announced_prefix, "test/broadcast");
assert_eq!(moq_origin_announced_cancel(announced_task), 0);
assert_eq!(cb.recv_terminal(), 0, "announced close delivers terminal 0");
assert_eq!(moq_publish_finish(broadcast), 0);
assert_eq!(moq_origin_close(origin), 0);
}
#[test]
fn announced_filters_patterns_and_reports_captures() {
let origin = id(moq_origin_create());
let cb = Callback::new();
let prefix = b"room";
let filter = b"*/chat";
let announced_task = id(unsafe {
moq_origin_announced(
origin,
prefix.as_ptr().cast(),
prefix.len(),
filter.as_ptr().cast(),
filter.len(),
Some(channel_callback),
cb.ptr,
)
});
let audio = publish_broadcast(origin, b"room/alice/audio");
let chat = publish_broadcast(origin, b"room/alice/chat");
let announced_id = id(cb.recv());
let mut info = moq_announce_update {
prefix: std::ptr::null(),
prefix_len: 0,
captures: std::ptr::null(),
captures_len: 0,
has_captures: false,
active: false,
};
assert_eq!(unsafe { moq_origin_announced_info(announced_id, &mut info) }, 0);
let announced_prefix =
unsafe { std::str::from_utf8(std::slice::from_raw_parts(info.prefix.cast::<u8>(), info.prefix_len)).unwrap() };
assert_eq!(announced_prefix, "room/alice/chat");
assert!(info.has_captures);
assert_eq!(info.captures_len, 1);
let capture = unsafe { &*info.captures };
let capture = unsafe { std::str::from_utf8(std::slice::from_raw_parts(capture.data.cast(), capture.len)).unwrap() };
assert_eq!(capture, "alice");
assert_eq!(moq_origin_announced_free(announced_id), 0);
assert_eq!(moq_origin_announced_cancel(announced_task), 0);
assert_eq!(cb.recv_terminal(), 0);
assert_eq!(moq_publish_finish(audio), 0);
assert_eq!(moq_publish_finish(chat), 0);
assert_eq!(moq_origin_close(origin), 0);
}
#[test]
fn announced_deactivation() {
let origin = id(moq_origin_create());
let cb = Callback::new();
let announced_task = id(unsafe {
moq_origin_announced(
origin,
std::ptr::null(),
0,
std::ptr::null(),
0,
Some(channel_callback),
cb.ptr,
)
});
let path = b"deactivate/test";
let broadcast = publish_broadcast(origin, path);
let announced_id = id(cb.recv());
let mut info = moq_announce_update {
prefix: std::ptr::null(),
prefix_len: 0,
captures: std::ptr::null(),
captures_len: 0,
has_captures: false,
active: false,
};
assert_eq!(unsafe { moq_origin_announced_info(announced_id, &mut info) }, 0);
assert!(info.active);
assert_eq!(moq_publish_unannounce(broadcast), 0);
let deactivated_id = id(cb.recv());
assert_eq!(unsafe { moq_origin_announced_info(deactivated_id, &mut info) }, 0);
assert!(!info.active, "broadcast should be inactive after unannounce");
assert!(request_once(origin, path) < 0, "an unannounced broadcast is unroutable");
assert_eq!(unsafe { moq_publish_announce(broadcast, std::ptr::null()) }, 0);
let reannounced_id = id(cb.recv());
assert_eq!(unsafe { moq_origin_announced_info(reannounced_id, &mut info) }, 0);
assert!(info.active, "broadcast should be active again after announce");
let _ = request_broadcast(origin, path);
assert_eq!(moq_origin_announced_cancel(announced_task), 0);
assert_eq!(cb.recv_terminal(), 0, "announced close delivers terminal 0");
assert_eq!(moq_publish_finish(broadcast), 0);
assert_eq!(moq_origin_close(origin), 0);
}
#[test]
fn create_broadcast_is_unroutable_until_announced() {
let origin = id(moq_origin_create());
let cb = Callback::new();
let announced_task = id(unsafe {
moq_origin_announced(
origin,
std::ptr::null(),
0,
std::ptr::null(),
0,
Some(channel_callback),
cb.ptr,
)
});
let path = b"quiet";
let broadcast = id(unsafe { moq_origin_create_broadcast(origin, path.as_ptr() as *const c_char, path.len()) });
assert!(request_once(origin, path) < 0, "an unannounced broadcast is unroutable");
assert_eq!(unsafe { moq_publish_announce(broadcast, std::ptr::null()) }, 0);
let _ = request_broadcast(origin, path);
let announced_id = id(cb.recv());
let mut info = moq_announce_update {
prefix: std::ptr::null(),
prefix_len: 0,
captures: std::ptr::null(),
captures_len: 0,
has_captures: false,
active: false,
};
assert_eq!(unsafe { moq_origin_announced_info(announced_id, &mut info) }, 0);
assert!(info.active);
assert_eq!(moq_origin_announced_cancel(announced_task), 0);
assert_eq!(cb.recv_terminal(), 0);
assert_eq!(moq_publish_finish(broadcast), 0);
assert_eq!(moq_origin_close(origin), 0);
}
#[test]
fn announce_accepts_an_anonymous_hop() {
let origin = id(moq_origin_create());
let path = b"anon";
let broadcast = id(unsafe { moq_origin_create_broadcast(origin, path.as_ptr() as *const c_char, path.len()) });
let hops = [0u64];
let route = moq_route {
hops: hops.as_ptr(),
hops_len: hops.len(),
cost: 1,
cold: 0,
has_cold: false,
};
assert_eq!(unsafe { moq_publish_announce(broadcast, &route) }, 0);
assert_eq!(moq_publish_finish(broadcast), 0);
assert_eq!(moq_origin_close(origin), 0);
}
#[test]
fn dynamic_serves_a_request_under_a_prefix() {
let origin = id(moq_origin_create());
let cb = Callback::new();
let prefix = b"live";
let dynamic = id(unsafe {
moq_origin_dynamic(
origin,
prefix.as_ptr() as *const c_char,
prefix.len(),
std::ptr::null(),
Some(channel_callback),
cb.ptr,
)
});
let path = b"live/cam";
let req_cb = Callback::new();
let _task = id(unsafe {
moq_origin_request(
origin,
path.as_ptr() as *const c_char,
path.len(),
Some(channel_callback),
req_cb.ptr,
)
});
let request = id(cb.recv());
let mut info = moq_string {
data: std::ptr::null(),
len: 0,
};
assert_eq!(unsafe { moq_broadcast_request_path(request, &mut info) }, 0);
let got = unsafe { std::str::from_utf8(std::slice::from_raw_parts(info.data.cast::<u8>(), info.len)).unwrap() };
assert_eq!(got, "live/cam");
let served = id(unsafe { moq_origin_create_broadcast(origin, b"unused".as_ptr() as *const c_char, 6) });
assert_eq!(moq_broadcast_request_accept(request, served), 0);
assert!(req_cb.recv() > 0);
req_cb.recv_terminal();
assert_eq!(moq_origin_dynamic_cancel(dynamic), 0);
assert!(moq_origin_dynamic_cancel(dynamic) < 0, "double-cancel should fail");
assert_eq!(cb.recv_terminal(), 0);
assert_eq!(moq_publish_finish(served), 0);
assert_eq!(moq_origin_close(origin), 0);
}
#[test]
fn dynamic_refuses_a_missing_callback() {
let origin = id(moq_origin_create());
let prefix = b"live";
let code = unsafe {
moq_origin_dynamic(
origin,
prefix.as_ptr() as *const c_char,
prefix.len(),
std::ptr::null(),
None,
std::ptr::null_mut(),
)
};
assert!(code < 0, "a missing on_request must be refused, got {code}");
assert_eq!(moq_origin_close(origin), 0);
}
struct ProtocolCallback {
rx: mpsc::Receiver<(i32, Option<moq_protocol_error>)>,
ptr: *mut c_void,
}
impl ProtocolCallback {
fn new() -> Self {
let (tx, rx) = mpsc::channel();
let ptr = Box::into_raw(Box::new(tx)) as *mut c_void;
Self { rx, ptr }
}
fn recv(&self) -> (i32, Option<moq_protocol_error>) {
self.rx.recv_timeout(TIMEOUT).expect("callback timed out")
}
}
impl Drop for ProtocolCallback {
fn drop(&mut self) {
unsafe {
drop(Box::from_raw(
self.ptr as *mut mpsc::Sender<(i32, Option<moq_protocol_error>)>,
))
};
}
}
extern "C" fn protocol_callback(user_data: *mut c_void, code: i32) {
let tx = unsafe { &*(user_data as *const mpsc::Sender<(i32, Option<moq_protocol_error>)>) };
let mut out = moq_protocol_error {
scope: 0,
code: 0,
kind: 0,
};
let protocol = (code < 0 && unsafe { moq_error_protocol(&mut out) } == 0).then_some(out);
let _ = tx.send((code, protocol));
}
fn consume_local(broadcast: u32) -> u32 {
let mut state = State::lock();
let consumer = state
.publish
.producer(Id::try_from(broadcast).unwrap())
.unwrap()
.consume();
i32::from(state.consume.start(consumer, None).unwrap()) as u32
}
fn consume_track(consume: u32, name: &[u8], cb: &Callback) -> u32 {
id(unsafe {
moq_consume_track(
consume,
name.as_ptr() as *const c_char,
name.len(),
std::ptr::null(),
Some(channel_callback),
cb.ptr,
)
})
}
#[test]
fn track_demand_follows_subscribers() {
let origin = id(moq_origin_create());
let broadcast = publish_broadcast(origin, b"track-demand");
let name = b"data";
let track =
id(unsafe { moq_publish_track(broadcast, name.as_ptr() as *const c_char, name.len(), std::ptr::null()) });
let demand_cb = Callback::new();
let watcher = id(unsafe { moq_publish_track_demand(track, Some(channel_callback), demand_cb.ptr) });
assert_eq!(demand_cb.recv(), moq_demand::MOQ_DEMAND_UNUSED as i32);
let consume = consume_local(broadcast);
let frame_cb = Callback::new();
let consumer = consume_track(consume, name, &frame_cb);
assert_eq!(demand_cb.recv(), moq_demand::MOQ_DEMAND_USED as i32);
assert_eq!(moq_consume_track_cancel(consumer), 0);
assert_eq!(frame_cb.recv_terminal(), 0);
assert_eq!(demand_cb.recv(), moq_demand::MOQ_DEMAND_UNUSED as i32);
let late_cb = Callback::new();
let late = id(unsafe { moq_publish_track_demand(track, Some(channel_callback), late_cb.ptr) });
assert_eq!(late_cb.recv(), moq_demand::MOQ_DEMAND_UNUSED as i32);
assert_eq!(moq_publish_demand_cancel(late), 0);
assert_eq!(late_cb.recv_terminal(), 0);
assert!(moq_publish_demand_cancel(late) < 0, "double-close should fail");
assert_eq!(moq_publish_track_finish(track), 0);
assert_eq!(demand_cb.recv_terminal(), 0);
assert!(
moq_publish_demand_cancel(watcher) < 0,
"the watcher is gone after its terminal"
);
assert_eq!(moq_consume_close(consume), 0);
assert_eq!(moq_publish_finish(broadcast), 0);
assert_eq!(moq_origin_close(origin), 0);
}
#[test]
fn track_demand_reports_current_state_before_close() {
let origin = id(moq_origin_create());
let broadcast = publish_broadcast(origin, b"track-demand-close");
let name = b"data";
let track =
id(unsafe { moq_publish_track(broadcast, name.as_ptr() as *const c_char, name.len(), std::ptr::null()) });
let demand = State::lock()
.publish
.track_demand(Id::try_from(track).unwrap())
.unwrap();
let cb = Callback::new();
let on_demand = unsafe { crate::ffi::OnStatus::new(cb.ptr, Some(channel_callback)) }.unwrap();
let (close, closed) = tokio::sync::oneshot::channel();
drop(close);
crate::ffi::RUNTIME
.block_on(crate::publish::Publish::run_demand(on_demand, demand, closed))
.unwrap();
assert_eq!(cb.recv(), moq_demand::MOQ_DEMAND_UNUSED as i32);
assert!(
cb.rx.recv_timeout(Duration::from_millis(50)).is_err(),
"the terminal is the spawner's to deliver"
);
assert_eq!(moq_publish_track_finish(track), 0);
assert_eq!(moq_publish_finish(broadcast), 0);
assert_eq!(moq_origin_close(origin), 0);
}
#[test]
fn track_demand_reports_an_abort() {
let origin = id(moq_origin_create());
let broadcast = publish_broadcast(origin, b"track-demand-abort");
let name = b"data";
let track =
id(unsafe { moq_publish_track(broadcast, name.as_ptr() as *const c_char, name.len(), std::ptr::null()) });
let demand_cb = ProtocolCallback::new();
let _watcher = id(unsafe { moq_publish_track_demand(track, Some(protocol_callback), demand_cb.ptr) });
assert_eq!(demand_cb.recv().0, moq_demand::MOQ_DEMAND_UNUSED as i32);
assert_eq!(moq_publish_track_abort(track, 7), 0);
let (code, protocol) = demand_cb.recv();
assert!(code < 0, "an aborted track is a negative terminal, got {code}");
let protocol = protocol.expect("an app abort carries its protocol code");
assert_eq!(protocol.kind, moq_protocol_kind::MOQ_PROTOCOL_KIND_APP as u32);
assert_eq!(protocol.code, 64 + 7);
assert_eq!(moq_publish_finish(broadcast), 0);
assert_eq!(moq_origin_close(origin), 0);
}
#[test]
fn media_demand_refuses_a_container() {
let origin = id(moq_origin_create());
let broadcast = publish_broadcast(origin, b"media-demand");
let container = id(publish_container(
broadcast,
moq_container_format::MOQ_CONTAINER_FORMAT_FMP4,
&[],
));
let media = id(publish_audio(
broadcast,
moq_audio_format::MOQ_AUDIO_FORMAT_OPUS,
&opus_head(),
None,
));
let cb = Callback::new();
assert!(
unsafe { moq_publish_media_demand(container, Some(channel_callback), cb.ptr) } < 0,
"a container has no single demand"
);
assert!(
unsafe { moq_publish_media_demand(media, None, cb.ptr) } < 0,
"a missing on_demand must be refused"
);
let _watcher = id(unsafe { moq_publish_media_demand(media, Some(channel_callback), cb.ptr) });
assert_eq!(cb.recv(), moq_demand::MOQ_DEMAND_UNUSED as i32);
assert_eq!(moq_publish_media_finish(media), 0);
assert_eq!(cb.recv_terminal(), 0);
assert_eq!(moq_publish_container_finish(container), 0);
assert_eq!(moq_publish_finish(broadcast), 0);
assert_eq!(moq_origin_close(origin), 0);
}
#[test]
fn dynamic_serves_track_requests() {
let origin = id(moq_origin_create());
let broadcast = publish_broadcast(origin, b"dynamic-tracks");
assert!(
unsafe { moq_publish_dynamic(broadcast, None, std::ptr::null_mut()) } < 0,
"a missing on_request must be refused"
);
let request_cb = Callback::new();
let dynamic = id(unsafe { moq_publish_dynamic(broadcast, Some(channel_callback), request_cb.ptr) });
let consume = consume_local(broadcast);
let name = b"events";
let frame_cb = Callback::new();
let consumer = consume_track(consume, name, &frame_cb);
let request = id(request_cb.recv());
let mut got = moq_string {
data: std::ptr::null(),
len: 0,
};
assert_eq!(unsafe { moq_track_request_name(request, &mut got) }, 0);
assert_eq!(borrowed_string(got.data, got.len).as_deref(), Some("events"));
let track = id(unsafe { moq_track_request_accept(request, std::ptr::null()) });
assert!(
unsafe { moq_track_request_name(request, &mut got) } < 0,
"accept consumes the request"
);
let demand_cb = Callback::new();
let _watcher = id(unsafe { moq_publish_track_demand(track, Some(channel_callback), demand_cb.ptr) });
assert_eq!(demand_cb.recv(), moq_demand::MOQ_DEMAND_USED as i32);
let payload = b"hello dynamic track";
assert_eq!(
unsafe { moq_publish_track_frame(track, payload.as_ptr(), payload.len(), 42) },
0
);
let frame_id = id(frame_cb.recv());
let mut frame = moq_frame {
payload: std::ptr::null(),
payload_size: 0,
timestamp_us: 0,
keyframe: false,
};
assert_eq!(unsafe { moq_consume_track_frame(frame_id, &mut frame) }, 0);
assert_eq!(
unsafe { std::slice::from_raw_parts(frame.payload, frame.payload_size) },
payload
);
assert_eq!(frame.timestamp_us, 42);
assert_eq!(moq_consume_track_frame_free(frame_id), 0);
let denied_cb = ProtocolCallback::new();
let denied_name = b"denied";
let denied = id(unsafe {
moq_consume_track(
consume,
denied_name.as_ptr() as *const c_char,
denied_name.len(),
std::ptr::null(),
Some(protocol_callback),
denied_cb.ptr,
)
});
let request = id(request_cb.recv());
assert_eq!(moq_track_request_abort(request, 404), 0);
assert!(moq_track_request_free(request) < 0, "abort consumes the request");
let (code, protocol) = denied_cb.recv();
assert!(code < 0, "a rejected subscribe is a negative terminal, got {code}");
let protocol = protocol.expect("a rejection carries its protocol code");
assert_eq!(protocol.kind, moq_protocol_kind::MOQ_PROTOCOL_KIND_APP as u32);
assert_eq!(protocol.code, 64 + 404);
assert!(
moq_consume_track_cancel(denied) < 0,
"the subscriber is gone after its terminal"
);
let freed_cb = Callback::new();
let _freed = consume_track(consume, b"freed", &freed_cb);
let request = id(request_cb.recv());
assert_eq!(moq_track_request_free(request), 0);
assert!(freed_cb.recv_terminal() < 0, "a freed request fails its subscriber");
assert_eq!(moq_publish_dynamic_cancel(dynamic), 0);
assert_eq!(request_cb.recv_terminal(), 0);
assert!(moq_publish_dynamic_cancel(dynamic) < 0, "double-close should fail");
assert_eq!(moq_consume_track_cancel(consumer), 0);
assert_eq!(frame_cb.recv_terminal(), 0);
assert_eq!(demand_cb.recv(), moq_demand::MOQ_DEMAND_UNUSED as i32);
assert_eq!(moq_publish_track_finish(track), 0);
assert_eq!(demand_cb.recv_terminal(), 0);
assert_eq!(moq_consume_close(consume), 0);
assert_eq!(moq_publish_finish(broadcast), 0);
assert_eq!(moq_origin_close(origin), 0);
}
#[test]
fn dynamic_track_request_publishes_media() {
let origin = id(moq_origin_create());
let broadcast = publish_broadcast(origin, b"dynamic-media");
let request_cb = Callback::new();
let _dynamic = id(unsafe { moq_publish_dynamic(broadcast, Some(channel_callback), request_cb.ptr) });
let consume = consume_local(broadcast);
let frame_cb = Callback::new();
let consumer = consume_track(consume, b"audio", &frame_cb);
let request = id(request_cb.recv());
let init = opus_head();
let config = moq_audio_init {
format: moq_audio_format::MOQ_AUDIO_FORMAT_OPUS as u32,
init: init.as_ptr(),
init_len: init.len(),
label: std::ptr::null(),
label_len: 0,
};
let media = id(unsafe { moq_track_request_audio(request, &config) });
let demand_cb = Callback::new();
let _watcher = id(unsafe { moq_publish_media_demand(media, Some(channel_callback), demand_cb.ptr) });
assert_eq!(demand_cb.recv(), moq_demand::MOQ_DEMAND_USED as i32);
let payload = b"opus frame";
assert_eq!(
unsafe { moq_publish_media_frame(media, payload.as_ptr(), payload.len(), 1000) },
0
);
let frame_id = id(frame_cb.recv());
assert_eq!(moq_consume_track_frame_free(frame_id), 0);
assert_eq!(moq_consume_track_cancel(consumer), 0);
assert_eq!(frame_cb.recv_terminal(), 0);
assert_eq!(demand_cb.recv(), moq_demand::MOQ_DEMAND_UNUSED as i32);
assert_eq!(moq_publish_media_finish(media), 0);
assert_eq!(demand_cb.recv_terminal(), 0);
assert_eq!(moq_publish_finish(broadcast), 0);
assert_eq!(request_cb.recv_terminal(), 0);
assert_eq!(moq_consume_close(consume), 0);
assert_eq!(moq_origin_close(origin), 0);
}
fn fetch_group(
broadcast: u32,
name: &str,
sequence: u64,
priority: u8,
frame_start: u64,
) -> mpsc::Receiver<Result<(u64, Vec<u8>), moq_net::Error>> {
let consumer = State::lock()
.publish
.producer(Id::try_from(broadcast).unwrap())
.unwrap()
.consume();
let name = name.to_string();
let (tx, rx) = mpsc::channel();
crate::ffi::RUNTIME.spawn(async move {
let res = async {
let track = consumer.track(&name)?;
let options = moq_net::group::Fetch::default()
.with_priority(priority)
.with_frame_start(frame_start);
let mut group = track.fetch_group(sequence, options).await?;
let index = group.index();
let frame = group.read_frame().await?.expect("expected a fetched frame");
Ok((index, frame.payload.to_vec()))
}
.await;
let _ = tx.send(res);
});
rx
}
#[test]
fn track_dynamic_serves_a_fetch_miss() {
let origin = id(moq_origin_create());
let broadcast = publish_broadcast(origin, b"track-dynamic");
let name = b"data";
let track =
id(unsafe { moq_publish_track(broadcast, name.as_ptr() as *const c_char, name.len(), std::ptr::null()) });
let group_cb = Callback::new();
let dynamic = id(unsafe { moq_publish_track_dynamic(track, Some(channel_callback), group_cb.ptr) });
let fetch = fetch_group(broadcast, "data", 5, 11, 0);
let request = id(group_cb.recv());
let mut sequence = 0u64;
let mut priority = 0u8;
let mut frame_start = 1u64;
assert_eq!(unsafe { moq_group_request_sequence(request, &mut sequence) }, 0);
assert_eq!(unsafe { moq_group_request_priority(request, &mut priority) }, 0);
assert_eq!(unsafe { moq_group_request_frame_start(request, &mut frame_start) }, 0);
assert_eq!(sequence, 5);
assert_eq!(priority, 11);
assert_eq!(frame_start, 0);
let group = id(moq_group_request_accept(request));
assert!(moq_group_request_free(request) < 0, "accept consumes the request");
let payload = b"fetched";
assert_eq!(
unsafe { moq_publish_group_frame(group, payload.as_ptr(), payload.len(), 100_000) },
0
);
assert_eq!(moq_publish_group_finish(group), 0);
assert_eq!(fetch.recv_timeout(TIMEOUT).unwrap().unwrap(), (0, payload.to_vec()));
let fetch = fetch_group(broadcast, "data", 6, 0, 0);
let request = id(group_cb.recv());
assert_eq!(moq_group_request_abort(request, 9), 0);
let err = fetch
.recv_timeout(TIMEOUT)
.unwrap()
.expect_err("a rejected fetch fails");
assert!(matches!(err, moq_net::Error::App(9)), "got {err:?}");
assert_eq!(moq_publish_dynamic_cancel(dynamic), 0);
assert_eq!(group_cb.recv_terminal(), 0);
assert_eq!(moq_publish_track_finish(track), 0);
assert_eq!(moq_publish_finish(broadcast), 0);
assert_eq!(moq_origin_close(origin), 0);
}
#[test]
fn track_dynamic_serves_a_fetch_from_frame_start() {
let origin = id(moq_origin_create());
let broadcast = publish_broadcast(origin, b"track-dynamic-frame-start");
let name = b"data";
let track =
id(unsafe { moq_publish_track(broadcast, name.as_ptr() as *const c_char, name.len(), std::ptr::null()) });
let group_cb = Callback::new();
let dynamic = id(unsafe { moq_publish_track_dynamic(track, Some(channel_callback), group_cb.ptr) });
let fetch = fetch_group(broadcast, "data", 5, 0, 3);
let request = id(group_cb.recv());
let mut frame_start = 0u64;
assert_eq!(unsafe { moq_group_request_frame_start(request, &mut frame_start) }, 0);
assert_eq!(frame_start, 3);
let group = id(moq_group_request_accept(request));
let payload = b"tail";
assert_eq!(
unsafe { moq_publish_group_frame(group, payload.as_ptr(), payload.len(), 100_000) },
0
);
assert_eq!(moq_publish_group_finish(group), 0);
let (index, got) = fetch.recv_timeout(TIMEOUT).unwrap().unwrap();
assert_eq!(index, 3, "the consumer resumes at the requested frame");
assert_eq!(got, payload);
assert_eq!(moq_publish_dynamic_cancel(dynamic), 0);
assert_eq!(group_cb.recv_terminal(), 0);
assert_eq!(moq_publish_track_finish(track), 0);
assert_eq!(moq_publish_finish(broadcast), 0);
assert_eq!(moq_origin_close(origin), 0);
}
#[test]
fn track_request_dynamic_survives_accept() {
let origin = id(moq_origin_create());
let broadcast = publish_broadcast(origin, b"request-dynamic");
let request_cb = Callback::new();
let _dynamic = id(unsafe { moq_publish_dynamic(broadcast, Some(channel_callback), request_cb.ptr) });
let fetch = fetch_group(broadcast, "archive", 9, 0, 0);
let request = id(request_cb.recv());
let group_cb = Callback::new();
let track_dynamic = id(unsafe { moq_track_request_dynamic(request, Some(channel_callback), group_cb.ptr) });
let track = id(unsafe { moq_track_request_accept(request, std::ptr::null()) });
let group_request = id(group_cb.recv());
let mut sequence = 0u64;
assert_eq!(unsafe { moq_group_request_sequence(group_request, &mut sequence) }, 0);
assert_eq!(sequence, 9);
let group = id(moq_group_request_accept(group_request));
let payload = b"archive";
assert_eq!(
unsafe { moq_publish_group_frame(group, payload.as_ptr(), payload.len(), 180_000) },
0
);
assert_eq!(moq_publish_group_finish(group), 0);
assert_eq!(fetch.recv_timeout(TIMEOUT).unwrap().unwrap(), (0, payload.to_vec()));
assert_eq!(moq_publish_dynamic_cancel(track_dynamic), 0);
assert_eq!(group_cb.recv_terminal(), 0);
assert_eq!(moq_publish_track_finish(track), 0);
assert_eq!(moq_publish_finish(broadcast), 0);
assert_eq!(request_cb.recv_terminal(), 0);
assert_eq!(moq_origin_close(origin), 0);
}
#[test]
fn local_publish_consume() {
let origin = id(moq_origin_create());
let path = b"live";
let broadcast = publish_broadcast(origin, path);
let init = opus_head();
let format = moq_audio_format::MOQ_AUDIO_FORMAT_OPUS;
let media = id(publish_audio(broadcast, format, &init, None));
let consume = request_broadcast(origin, path);
let catalog_cb = Callback::new();
let catalog_task = id(unsafe { moq_consume_catalog(consume, Some(channel_callback), catalog_cb.ptr) });
let catalog_id = id(catalog_cb.recv());
let mut audio_cfg = moq_audio_config {
name: std::ptr::null(),
name_len: 0,
label: std::ptr::null(),
label_len: 0,
codec: std::ptr::null(),
codec_len: 0,
description: std::ptr::null(),
description_len: 0,
sample_rate: 0,
channel_count: 0,
container: moq_container::default(),
};
assert_eq!(unsafe { moq_consume_audio_config(catalog_id, 0, &mut audio_cfg) }, 0);
assert_eq!(audio_cfg.sample_rate, 48000);
assert_eq!(audio_cfg.channel_count, 2);
let codec = unsafe {
std::str::from_utf8(std::slice::from_raw_parts(
audio_cfg.codec.cast::<u8>(),
audio_cfg.codec_len,
))
}
.unwrap();
assert_eq!(codec, "opus");
let mut video_cfg = moq_video_config {
name: std::ptr::null(),
name_len: 0,
label: std::ptr::null(),
label_len: 0,
codec: std::ptr::null(),
codec_len: 0,
description: std::ptr::null(),
description_len: 0,
coded_width: 0,
coded_height: 0,
container: moq_container::default(),
};
assert!(
unsafe { moq_consume_video_config(catalog_id, 0, &mut video_cfg) } < 0,
"video config should fail (no video tracks)"
);
let frame_cb = Callback::new();
let track = id(unsafe { moq_consume_audio(catalog_id, 0, 10_000_000, Some(channel_callback), frame_cb.ptr) });
let payload = b"opus audio payload data";
let timestamp_us: u64 = 1_000_000;
assert_eq!(
unsafe { moq_publish_media_frame(media, payload.as_ptr(), payload.len(), timestamp_us) },
0
);
let frame_id = id(frame_cb.recv());
let mut frame = moq_frame {
payload: std::ptr::null(),
payload_size: 0,
timestamp_us: 0,
keyframe: false,
};
assert_eq!(unsafe { moq_consume_frame(frame_id, &mut frame) }, 0);
assert_eq!(frame.payload_size, payload.len());
assert_eq!(frame.timestamp_us, timestamp_us);
let received = unsafe { std::slice::from_raw_parts(frame.payload, frame.payload_size) };
assert_eq!(received, payload, "frame payload should match");
assert_eq!(moq_consume_frame_free(frame_id), 0);
assert_eq!(moq_consume_audio_cancel(track), 0);
assert_eq!(frame_cb.recv_terminal(), 0, "audio close delivers terminal 0");
assert_eq!(moq_consume_catalog_free(catalog_id), 0);
assert_eq!(moq_consume_catalog_cancel(catalog_task), 0);
assert_eq!(catalog_cb.recv_terminal(), 0, "catalog close delivers terminal 0");
assert_eq!(moq_consume_close(consume), 0);
assert_eq!(moq_publish_media_finish(media), 0);
assert_eq!(moq_publish_finish(broadcast), 0);
assert_eq!(moq_origin_close(origin), 0);
}
#[test]
fn consume_announced_local() {
let origin = id(moq_origin_create());
let cb = Callback::new();
let path = b"live";
let _task = id(unsafe {
moq_origin_announced_broadcast(
origin,
path.as_ptr() as *const c_char,
path.len(),
Some(channel_callback),
cb.ptr,
)
});
let broadcast = publish_broadcast(origin, path);
let init = opus_head();
let format = moq_audio_format::MOQ_AUDIO_FORMAT_OPUS;
let media = id(publish_audio(broadcast, format, &init, None));
let consume = id(cb.recv());
assert_eq!(cb.recv_terminal(), 0, "wait delivers terminal 0 after the handle");
let catalog_cb = Callback::new();
let catalog_task = id(unsafe { moq_consume_catalog(consume, Some(channel_callback), catalog_cb.ptr) });
let catalog_id = id(catalog_cb.recv());
let mut audio_cfg = moq_audio_config {
name: std::ptr::null(),
name_len: 0,
label: std::ptr::null(),
label_len: 0,
codec: std::ptr::null(),
codec_len: 0,
description: std::ptr::null(),
description_len: 0,
sample_rate: 0,
channel_count: 0,
container: moq_container::default(),
};
assert_eq!(unsafe { moq_consume_audio_config(catalog_id, 0, &mut audio_cfg) }, 0);
assert_eq!(audio_cfg.sample_rate, 48000);
assert_eq!(audio_cfg.channel_count, 2);
assert_eq!(moq_consume_catalog_free(catalog_id), 0);
assert_eq!(moq_consume_catalog_cancel(catalog_task), 0);
assert_eq!(catalog_cb.recv_terminal(), 0, "catalog close delivers terminal 0");
assert_eq!(moq_consume_close(consume), 0);
assert_eq!(moq_publish_media_finish(media), 0);
assert_eq!(moq_publish_finish(broadcast), 0);
assert_eq!(moq_origin_close(origin), 0);
}
#[test]
fn consume_audio_follows_a_sibling_broadcast_reference() {
let origin = id(moq_origin_create());
let source = publish_broadcast(origin, b"a/source");
let init = opus_head();
let media = id(publish_audio(
source,
moq_audio_format::MOQ_AUDIO_FORMAT_OPUS,
&init,
None,
));
let name = {
let mut state = State::lock();
let (_, catalog) = state.publish.pair_mut(Id::try_from(source).unwrap()).unwrap();
let catalog = catalog.modify().unwrap();
catalog
.audio
.renditions
.keys()
.next()
.expect("the importer publishes one audio rendition")
.clone()
};
let broadcast = publish_broadcast(origin, b"a/pub");
let codec = "opus";
let config = moq_audio_config {
name: name.as_ptr() as *const c_char,
name_len: name.len(),
label: std::ptr::null(),
label_len: 0,
codec: codec.as_ptr() as *const c_char,
codec_len: codec.len(),
description: std::ptr::null(),
description_len: 0,
sample_rate: 48000,
channel_count: 2,
container: moq_container::default(),
};
assert_eq!(unsafe { moq_publish_audio_config(broadcast, &config) }, 0);
{
let mut state = State::lock();
let (_, catalog) = state.publish.pair_mut(Id::try_from(broadcast).unwrap()).unwrap();
catalog
.modify()
.unwrap()
.audio
.renditions
.get_mut(&name)
.unwrap()
.broadcast = Some(moq_net::path::Relative::new("./source").into_owned());
}
let consume = request_broadcast(origin, b"a/pub");
let catalog_cb = Callback::new();
let catalog_task = id(unsafe { moq_consume_catalog(consume, Some(channel_callback), catalog_cb.ptr) });
let catalog_id = id(catalog_cb.recv());
let frame_cb = Callback::new();
let track = id(unsafe { moq_consume_audio(catalog_id, 0, 10_000_000, Some(channel_callback), frame_cb.ptr) });
let payload = b"opus audio payload data";
let timestamp_us: u64 = 1_000_000;
assert_eq!(
unsafe { moq_publish_media_frame(media, payload.as_ptr(), payload.len(), timestamp_us) },
0
);
let frame_id = id(frame_cb.recv());
let mut frame = moq_frame {
payload: std::ptr::null(),
payload_size: 0,
timestamp_us: 0,
keyframe: false,
};
assert_eq!(unsafe { moq_consume_frame(frame_id, &mut frame) }, 0);
let received = unsafe { std::slice::from_raw_parts(frame.payload, frame.payload_size) };
assert_eq!(received, payload, "the sibling broadcast's frame should arrive");
assert_eq!(frame.timestamp_us, timestamp_us);
assert_eq!(moq_consume_frame_free(frame_id), 0);
assert_eq!(moq_consume_audio_cancel(track), 0);
assert_eq!(frame_cb.recv_terminal(), 0, "audio close delivers terminal 0");
assert_eq!(moq_consume_catalog_free(catalog_id), 0);
assert_eq!(moq_consume_catalog_cancel(catalog_task), 0);
assert_eq!(catalog_cb.recv_terminal(), 0, "catalog close delivers terminal 0");
assert_eq!(moq_consume_close(consume), 0);
assert_eq!(moq_publish_media_finish(media), 0);
assert_eq!(moq_publish_finish(broadcast), 0);
assert_eq!(moq_publish_finish(source), 0);
assert_eq!(moq_origin_close(origin), 0);
}
#[test]
fn consume_announced_close_cancels() {
let origin = id(moq_origin_create());
let cb = Callback::new();
let path = b"never";
let task = id(unsafe {
moq_origin_announced_broadcast(
origin,
path.as_ptr() as *const c_char,
path.len(),
Some(channel_callback),
cb.ptr,
)
});
assert_eq!(moq_origin_announced_broadcast_cancel(task), 0);
assert_eq!(cb.recv_terminal(), 0, "close delivers terminal 0");
assert!(
moq_origin_announced_broadcast_cancel(task) < 0,
"double-close should fail"
);
assert_eq!(moq_origin_close(origin), 0);
}
#[test]
fn video_publish_consume() {
let origin = id(moq_origin_create());
let path = b"video-test";
let broadcast = publish_broadcast(origin, path);
let init = h264_init();
let format = moq_video_format::MOQ_VIDEO_FORMAT_AVC3;
let media = id(publish_video(broadcast, format, &init, None));
let consume = request_broadcast(origin, path);
let catalog_cb = Callback::new();
let catalog_task = id(unsafe { moq_consume_catalog(consume, Some(channel_callback), catalog_cb.ptr) });
let catalog_id = id(catalog_cb.recv());
let mut video_cfg = moq_video_config {
name: std::ptr::null(),
name_len: 0,
label: std::ptr::null(),
label_len: 0,
codec: std::ptr::null(),
codec_len: 0,
description: std::ptr::null(),
description_len: 0,
coded_width: 0,
coded_height: 0,
container: moq_container::default(),
};
assert_eq!(
unsafe { moq_consume_video_config(catalog_id, 0, &mut video_cfg) },
0,
"video config should succeed for avc3 H.264 track"
);
let codec = unsafe {
std::str::from_utf8(std::slice::from_raw_parts(
video_cfg.codec.cast::<u8>(),
video_cfg.codec_len,
))
}
.unwrap();
assert!(
codec.starts_with("avc1.") || codec.starts_with("avc3."),
"codec should be avc1/avc3, got {codec}"
);
assert_eq!(video_cfg.coded_width, 1280);
assert_eq!(video_cfg.coded_height, 720);
let mut audio_cfg = moq_audio_config {
name: std::ptr::null(),
name_len: 0,
label: std::ptr::null(),
label_len: 0,
codec: std::ptr::null(),
codec_len: 0,
description: std::ptr::null(),
description_len: 0,
sample_rate: 0,
channel_count: 0,
container: moq_container::default(),
};
assert!(
unsafe { moq_consume_audio_config(catalog_id, 0, &mut audio_cfg) } < 0,
"audio config should fail (no audio tracks)"
);
let frame_cb = Callback::new();
let track = id(unsafe { moq_consume_video(catalog_id, 0, 10_000_000, Some(channel_callback), frame_cb.ptr) });
let keyframe = [0x00, 0x00, 0x00, 0x01, 0x65, 0xAA, 0xBB, 0xCC];
assert_eq!(
unsafe { moq_publish_media_frame(media, keyframe.as_ptr(), keyframe.len(), 0) },
0
);
let frame_id = id(frame_cb.recv());
let mut frame = moq_frame {
payload: std::ptr::null(),
payload_size: 0,
timestamp_us: 0,
keyframe: false,
};
assert_eq!(unsafe { moq_consume_frame(frame_id, &mut frame) }, 0);
assert_eq!(frame.timestamp_us, 0);
assert!(frame.payload_size > 0, "frame should have payload data");
assert_eq!(moq_consume_frame_free(frame_id), 0);
assert_eq!(moq_consume_video_cancel(track), 0);
assert_eq!(frame_cb.recv_terminal(), 0, "video close delivers terminal 0");
assert_eq!(moq_consume_catalog_free(catalog_id), 0);
assert_eq!(moq_consume_catalog_cancel(catalog_task), 0);
assert_eq!(catalog_cb.recv_terminal(), 0, "catalog close delivers terminal 0");
assert_eq!(moq_consume_close(consume), 0);
assert_eq!(moq_publish_media_finish(media), 0);
assert_eq!(moq_publish_finish(broadcast), 0);
assert_eq!(moq_origin_close(origin), 0);
}
#[test]
fn audio_raw_publish() {
let origin = id(moq_origin_create());
let broadcast = publish_broadcast(origin, b"audio-raw-publish-test");
let name = b"audio";
let input = moq_audio_encoder_input {
format: moq_audio_sample_format::MOQ_AUDIO_SAMPLE_FORMAT_F32 as u32,
sample_rate: 48_000,
channels: 2,
};
let codec = b"opus";
let output = moq_audio_encoder_output {
codec: codec.as_ptr() as *const c_char,
codec_len: codec.len(),
sample_rate: 0,
channels: 0,
bitrate: 0,
frame_duration_us: 20_000,
};
let producer = id(unsafe {
moq_encode_audio(
broadcast,
name.as_ptr() as *const c_char,
name.len(),
&input,
&output,
0,
)
});
let samples = vec![0.0f32; 960 * 2];
let pcm = unsafe { std::slice::from_raw_parts(samples.as_ptr().cast::<u8>(), std::mem::size_of_val(&samples[..])) };
let frame = moq_audio_frame {
timestamp_us: 0,
data: pcm.as_ptr(),
data_size: pcm.len(),
};
assert_eq!(unsafe { moq_encode_audio_frame(producer, &frame) }, 0);
assert_eq!(moq_encode_audio_finish(producer), 0);
assert!(moq_encode_audio_finish(producer) < 0, "double-finish should fail");
assert!(
unsafe { moq_encode_audio_frame(producer, &frame) } < 0,
"a finished producer should take no more frames"
);
assert_eq!(moq_publish_finish(broadcast), 0);
assert_eq!(moq_origin_close(origin), 0);
}
#[test]
fn audio_raw_publish_frame_durations() {
let origin = id(moq_origin_create());
let broadcast = publish_broadcast(origin, b"audio-frame-duration-test");
let input = moq_audio_encoder_input {
format: moq_audio_sample_format::MOQ_AUDIO_SAMPLE_FORMAT_F32 as u32,
sample_rate: 48_000,
channels: 2,
};
let codec = b"opus";
let encode = |name: &[u8], frame_duration_us: u32| {
let output = moq_audio_encoder_output {
codec: codec.as_ptr() as *const c_char,
codec_len: codec.len(),
sample_rate: 0,
channels: 0,
bitrate: 0,
frame_duration_us,
};
unsafe {
moq_encode_audio(
broadcast,
name.as_ptr() as *const c_char,
name.len(),
&input,
&output,
0,
)
}
};
let producer = id(encode(b"fine", 2_500));
let samples = vec![0.0f32; 120 * 2];
let pcm = unsafe { std::slice::from_raw_parts(samples.as_ptr().cast::<u8>(), std::mem::size_of_val(&samples[..])) };
let frame = moq_audio_frame {
timestamp_us: 0,
data: pcm.as_ptr(),
data_size: pcm.len(),
};
assert_eq!(unsafe { moq_encode_audio_frame(producer, &frame) }, 0);
assert_eq!(moq_encode_audio_finish(producer), 0);
assert_eq!(moq_encode_audio_finish(id(encode(b"default", 0))), 0, "0 = 20 ms");
assert!(encode(b"rounded", 2_000) < 0, "2 ms is not an opus frame duration");
assert_eq!(moq_publish_finish(broadcast), 0);
assert_eq!(moq_origin_close(origin), 0);
}
fn gray_rgba(width: u32, height: u32) -> Vec<u8> {
vec![0x80u8; width as usize * height as usize * 4]
}
#[test]
fn video_raw_publish_consume() {
let origin = id(moq_origin_create());
let path = b"video-raw-publish-test";
let broadcast = publish_broadcast(origin, path);
let input = moq_video_encoder_input {
format: moq_video_pixel_format::MOQ_VIDEO_PIXEL_FORMAT_RGBA as u32,
width: 320,
height: 240,
framerate: 30,
};
let output = moq_video_encoder_output {
codec: moq_video_codec::MOQ_VIDEO_CODEC_H264 as u32,
bitrate: 0,
gop: 0,
kind: moq_video_encoder_kind::MOQ_VIDEO_ENCODER_KIND_SOFTWARE as u32,
encoder: std::ptr::null(),
encoder_len: 0,
};
let producer = id(unsafe { moq_encode_video(broadcast, &input, &output, 0) });
let rgba = gray_rgba(320, 240);
let publish = |index: u64| {
let frame = moq_video_encoder_frame {
timestamp_us: index * 33_333,
data: rgba.as_ptr(),
data_size: rgba.len(),
};
assert_eq!(unsafe { moq_encode_video_frame(producer, &frame) }, 0);
};
assert_eq!(moq_encode_video_cut(producer), 0);
for i in 0..5u64 {
publish(i);
}
let consume = request_broadcast(origin, path);
let catalog_cb = Callback::new();
let catalog_task = id(unsafe { moq_consume_catalog(consume, Some(channel_callback), catalog_cb.ptr) });
let catalog_id = id(catalog_cb.recv());
let decoder = moq_video_decoder_output {
max_age_us: 10_000_000,
format: moq_video_pixel_format::MOQ_VIDEO_PIXEL_FORMAT_I420 as u32,
width: 0,
height: 0,
};
let frame_cb = Callback::new();
let consumer = id(unsafe { moq_decode_video(catalog_id, 0, &decoder, Some(channel_callback), frame_cb.ptr) });
for i in 5..20u64 {
publish(i);
}
let frame_id = id(frame_cb.recv());
let mut frame = moq_video_frame {
timestamp_us: 0,
width: 0,
height: 0,
data: std::ptr::null(),
data_size: 0,
};
assert_eq!(unsafe { moq_decode_video_frame(frame_id, &mut frame) }, 0);
assert_eq!(frame.width, 320);
assert_eq!(frame.height, 240);
assert_eq!(frame.data_size, 320 * 240 * 3 / 2, "tightly-packed I420");
assert_eq!(moq_decode_video_frame_free(frame_id), 0);
assert_eq!(moq_decode_video_cancel(consumer), 0);
loop {
let code = frame_cb.recv();
if code > 0 {
assert_eq!(moq_decode_video_frame_free(id(code)), 0);
} else {
assert_eq!(code, 0, "raw video close delivers terminal 0");
break;
}
}
assert_eq!(moq_consume_catalog_free(catalog_id), 0);
assert_eq!(moq_consume_catalog_cancel(catalog_task), 0);
assert_eq!(catalog_cb.recv_catalog_terminal(), 0);
assert_eq!(moq_consume_close(consume), 0);
assert_eq!(moq_encode_video_finish(producer), 0);
assert_eq!(moq_publish_finish(broadcast), 0);
assert_eq!(moq_origin_close(origin), 0);
}
#[test]
fn video_raw_decode_output_rejected() {
use crate::ffi::ReturnCode;
let bad_format = moq_video_decoder_output {
max_age_us: 0,
format: 999,
width: 0,
height: 0,
};
assert_eq!(
unsafe { moq_decode_video(1, 0, &bad_format, Some(ignore_callback), std::ptr::null_mut()) },
Error::InvalidCode.code()
);
for (width, height) in [(320, 0), (0, 240), (321, 240), (320, 241), (u32::MAX - 1, u32::MAX - 1)] {
let bad_size = moq_video_decoder_output {
max_age_us: 0,
format: moq_video_pixel_format::MOQ_VIDEO_PIXEL_FORMAT_I420 as u32,
width,
height,
};
assert_eq!(
unsafe { moq_decode_video(1, 0, &bad_size, Some(ignore_callback), std::ptr::null_mut()) },
Error::InvalidConfig(String::new()).code(),
"size {width}x{height} must be refused"
);
}
assert_eq!(
unsafe { moq_decode_video(1, 0, std::ptr::null(), Some(ignore_callback), std::ptr::null_mut()) },
Error::InvalidPointer.code()
);
let valid = moq_video_decoder_output {
max_age_us: 0,
format: moq_video_pixel_format::MOQ_VIDEO_PIXEL_FORMAT_RGBA as u32,
width: 160,
height: 120,
};
assert_eq!(
unsafe { moq_decode_video(u32::MAX / 2, 0, &valid, Some(ignore_callback), std::ptr::null_mut(),) },
Error::CatalogNotFound.code(),
"a valid request fails on the catalog lookup, not on validation"
);
}
fn decode_first_frame(output: &moq_video_decoder_output) -> (u32, u32, usize) {
let origin = id(moq_origin_create());
let path = b"video-raw-output-test";
let broadcast = publish_broadcast(origin, path);
let input = moq_video_encoder_input {
format: moq_video_pixel_format::MOQ_VIDEO_PIXEL_FORMAT_RGBA as u32,
width: 320,
height: 240,
framerate: 30,
};
let encode_output = moq_video_encoder_output {
codec: moq_video_codec::MOQ_VIDEO_CODEC_H264 as u32,
bitrate: 0,
gop: 0,
kind: moq_video_encoder_kind::MOQ_VIDEO_ENCODER_KIND_SOFTWARE as u32,
encoder: std::ptr::null(),
encoder_len: 0,
};
let producer = id(unsafe { moq_encode_video(broadcast, &input, &encode_output, 0) });
let rgba = gray_rgba(320, 240);
let publish = |index: u64| {
let frame = moq_video_encoder_frame {
timestamp_us: index * 33_333,
data: rgba.as_ptr(),
data_size: rgba.len(),
};
assert_eq!(unsafe { moq_encode_video_frame(producer, &frame) }, 0);
};
assert_eq!(moq_encode_video_cut(producer), 0);
for i in 0..5u64 {
publish(i);
}
let consume = request_broadcast(origin, path);
let catalog_cb = Callback::new();
let catalog_task = id(unsafe { moq_consume_catalog(consume, Some(channel_callback), catalog_cb.ptr) });
let catalog_id = id(catalog_cb.recv());
let frame_cb = Callback::new();
let consumer = id(unsafe { moq_decode_video(catalog_id, 0, output, Some(channel_callback), frame_cb.ptr) });
for i in 5..20u64 {
publish(i);
}
let frame_id = id(frame_cb.recv());
let mut frame = moq_video_frame {
timestamp_us: 0,
width: 0,
height: 0,
data: std::ptr::null(),
data_size: 0,
};
assert_eq!(unsafe { moq_decode_video_frame(frame_id, &mut frame) }, 0);
let result = (frame.width, frame.height, frame.data_size);
assert_eq!(moq_decode_video_frame_free(frame_id), 0);
assert_eq!(moq_decode_video_cancel(consumer), 0);
loop {
let code = frame_cb.recv();
if code > 0 {
assert_eq!(moq_decode_video_frame_free(id(code)), 0);
} else {
assert_eq!(code, 0, "raw video close delivers terminal 0");
break;
}
}
assert_eq!(moq_consume_catalog_free(catalog_id), 0);
assert_eq!(moq_consume_catalog_cancel(catalog_task), 0);
assert_eq!(catalog_cb.recv_catalog_terminal(), 0);
assert_eq!(moq_consume_close(consume), 0);
assert_eq!(moq_encode_video_finish(producer), 0);
assert_eq!(moq_publish_finish(broadcast), 0);
assert_eq!(moq_origin_close(origin), 0);
result
}
#[test]
fn video_raw_decode_rgba() {
let output = moq_video_decoder_output {
max_age_us: 10_000_000,
format: moq_video_pixel_format::MOQ_VIDEO_PIXEL_FORMAT_RGBA as u32,
width: 0,
height: 0,
};
assert_eq!(decode_first_frame(&output), (320, 240, 320 * 240 * 4));
}
#[test]
fn video_raw_decode_resize() {
let output = moq_video_decoder_output {
max_age_us: 10_000_000,
format: moq_video_pixel_format::MOQ_VIDEO_PIXEL_FORMAT_I420 as u32,
width: 160,
height: 120,
};
assert_eq!(decode_first_frame(&output), (160, 120, 160 * 120 * 3 / 2));
}
#[test]
fn video_raw_publish_from_many_threads() {
let origin = id(moq_origin_create());
let broadcast = publish_broadcast(origin, b"video-raw-threads-test");
let input = moq_video_encoder_input {
format: moq_video_pixel_format::MOQ_VIDEO_PIXEL_FORMAT_RGBA as u32,
width: 320,
height: 240,
framerate: 30,
};
let output = moq_video_encoder_output {
codec: moq_video_codec::MOQ_VIDEO_CODEC_H264 as u32,
bitrate: 1_000_000,
gop: 0,
kind: moq_video_encoder_kind::MOQ_VIDEO_ENCODER_KIND_SOFTWARE as u32,
encoder: std::ptr::null(),
encoder_len: 0,
};
let producer = id(unsafe { moq_encode_video(broadcast, &input, &output, 0) });
let rgba = std::sync::Arc::new(gray_rgba(320, 240));
for i in 0..8u64 {
let rgba = rgba.clone();
std::thread::spawn(move || {
if i == 0 {
assert_eq!(moq_encode_video_cut(producer), 0);
}
let frame = moq_video_encoder_frame {
timestamp_us: i * 33_333,
data: rgba.as_ptr(),
data_size: rgba.len(),
};
assert_eq!(unsafe { moq_encode_video_frame(producer, &frame) }, 0);
assert_eq!(moq_encode_video_bitrate(producer, 900_000 - i), 0);
})
.join()
.unwrap();
}
std::thread::spawn(move || assert_eq!(moq_encode_video_finish(producer), 0))
.join()
.unwrap();
assert_eq!(moq_publish_finish(broadcast), 0);
assert_eq!(moq_origin_close(origin), 0);
}
fn publish_gray(producer: u32, rgba: &[u8]) -> i32 {
let frame = moq_video_encoder_frame {
timestamp_us: 0,
data: rgba.as_ptr(),
data_size: rgba.len(),
};
unsafe { moq_encode_video_frame(producer, &frame) }
}
#[test]
fn a_stalled_encode_does_not_block_unrelated_calls() {
let origin = id(moq_origin_create());
let broadcast = publish_broadcast(origin, b"video-raw-stall-test");
let input = moq_video_encoder_input {
format: moq_video_pixel_format::MOQ_VIDEO_PIXEL_FORMAT_RGBA as u32,
width: 320,
height: 240,
framerate: 30,
};
let output = moq_video_encoder_output {
codec: moq_video_codec::MOQ_VIDEO_CODEC_H264 as u32,
bitrate: 0,
gop: 0,
kind: moq_video_encoder_kind::MOQ_VIDEO_ENCODER_KIND_SOFTWARE as u32,
encoder: std::ptr::null(),
encoder_len: 0,
};
let stalled = id(unsafe { moq_encode_video(broadcast, &input, &output, 0) });
let other = id(unsafe { moq_encode_video(broadcast, &input, &output, 0) });
let handle = State::lock().video.producer(Id::try_from(stalled).unwrap()).unwrap();
let held = handle.lock();
let rgba = std::sync::Arc::new(gray_rgba(320, 240));
let stalling = {
let rgba = rgba.clone();
std::thread::spawn(move || publish_gray(stalled, &rgba))
};
let deadline = std::time::Instant::now() + TIMEOUT;
while handle.holders() < 3 {
assert!(
std::time::Instant::now() < deadline,
"the publish never reached the encoder"
);
std::thread::yield_now();
}
let (tx, rx) = mpsc::channel();
let unrelated = {
let rgba = rgba.clone();
std::thread::spawn(move || {
let _ = tx.send((moq_origin_create(), publish_gray(other, &rgba)));
})
};
let (created, published) = rx
.recv_timeout(TIMEOUT)
.expect("an unrelated call was waiting on the stalled encode");
assert!(created > 0, "creating an origin failed while a producer was stalled");
assert_eq!(
published, 0,
"a second producer could not encode while the first stalled"
);
unrelated.join().unwrap();
assert_eq!(handle.holders(), 3, "the stalled publish finished early");
drop(held);
assert_eq!(stalling.join().unwrap(), 0);
assert_eq!(moq_origin_close(id(created)), 0);
assert_eq!(moq_encode_video_finish(stalled), 0);
assert_eq!(moq_encode_video_finish(other), 0);
assert_eq!(moq_publish_finish(broadcast), 0);
assert_eq!(moq_origin_close(origin), 0);
}
#[test]
fn video_raw_publish_rejects_frame_size_mismatch() {
let origin = id(moq_origin_create());
let broadcast = publish_broadcast(origin, b"video-raw-mismatch-test");
let input = moq_video_encoder_input {
format: moq_video_pixel_format::MOQ_VIDEO_PIXEL_FORMAT_RGBA as u32,
width: 320,
height: 240,
framerate: 30,
};
let output = moq_video_encoder_output {
codec: moq_video_codec::MOQ_VIDEO_CODEC_H264 as u32,
bitrate: 0,
gop: 0,
kind: moq_video_encoder_kind::MOQ_VIDEO_ENCODER_KIND_SOFTWARE as u32,
encoder: std::ptr::null(),
encoder_len: 0,
};
let producer = id(unsafe { moq_encode_video(broadcast, &input, &output, 0) });
let rgba = gray_rgba(640, 480);
let frame = moq_video_encoder_frame {
timestamp_us: 0,
data: rgba.as_ptr(),
data_size: rgba.len(),
};
assert!(unsafe { moq_encode_video_frame(producer, &frame) } < 0);
assert_eq!(moq_encode_video_finish(producer), 0);
assert_eq!(moq_publish_finish(broadcast), 0);
assert_eq!(moq_origin_close(origin), 0);
}
#[test]
fn video_raw_publish_rejects_invalid_config() {
let origin = id(moq_origin_create());
let broadcast = publish_broadcast(origin, b"video-raw-invalid-test");
let valid_input = moq_video_encoder_input {
format: moq_video_pixel_format::MOQ_VIDEO_PIXEL_FORMAT_I420 as u32,
width: 320,
height: 240,
framerate: 30,
};
let valid_output = moq_video_encoder_output {
codec: moq_video_codec::MOQ_VIDEO_CODEC_H264 as u32,
bitrate: 0,
gop: 0,
kind: moq_video_encoder_kind::MOQ_VIDEO_ENCODER_KIND_SOFTWARE as u32,
encoder: std::ptr::null(),
encoder_len: 0,
};
assert!(unsafe { moq_encode_video(broadcast, std::ptr::null(), &valid_output, 0) } < 0);
assert!(unsafe { moq_encode_video(broadcast, &valid_input, std::ptr::null(), 0) } < 0);
let bad_format = moq_video_encoder_input {
format: 99,
..valid_input
};
assert!(unsafe { moq_encode_video(broadcast, &bad_format, &valid_output, 0) } < 0);
let zero_framerate = moq_video_encoder_input {
framerate: 0,
..valid_input
};
assert!(unsafe { moq_encode_video(broadcast, &zero_framerate, &valid_output, 0) } < 0);
let unrepresentable = moq_video_encoder_input {
width: u32::MAX - 1,
height: u32::MAX - 1,
..valid_input
};
assert!(unsafe { moq_encode_video(broadcast, &unrepresentable, &valid_output, 0) } < 0);
let merely_huge = moq_video_encoder_input {
width: 65534,
height: 65534,
..valid_input
};
let huge = unsafe { moq_encode_video(broadcast, &merely_huge, &valid_output, 0) };
if huge > 0 {
assert_eq!(moq_encode_video_finish(id(huge)), 0);
} else {
let reason = unsafe { std::ffi::CStr::from_ptr(moq_error()) }.to_str().unwrap();
assert!(
!reason.contains("too large to represent"),
"the representability check rejected a size it should have left to the backend: {reason}"
);
}
let bad_codec = moq_video_encoder_output {
codec: 99,
..valid_output
};
assert!(unsafe { moq_encode_video(broadcast, &valid_input, &bad_codec, 0) } < 0);
let bad_kind = moq_video_encoder_output {
kind: 99,
..valid_output
};
assert!(unsafe { moq_encode_video(broadcast, &valid_input, &bad_kind, 0) } < 0);
assert!(moq_encode_video_cut(0) < 0);
assert!(moq_encode_video_bitrate(0, 1_000_000) < 0);
assert!(moq_encode_video_finish(0) < 0);
assert_eq!(moq_publish_finish(broadcast), 0);
assert_eq!(moq_origin_close(origin), 0);
}
#[test]
fn video_raw_decode() {
let mut config = moq_video::encode::Config::new(320, 240, moq_video::Rate::new(30, 1).unwrap());
config.kind = moq_video::encode::Kind::Software;
let mut encoder = moq_video::encode::Encoder::new(&config).expect("openh264 encoder");
let gray = vec![0x80u8; 320 * 240 * 4];
let mut frames: Vec<moq_video::encode::Encoded> = Vec::new();
for i in 0..5u64 {
if i == 0 {
encoder.cut().unwrap();
}
let surface = moq_video::Surface::rgba(&gray, moq_video::Size::new(320, 240)).unwrap();
let frame = moq_video::Frame::new(surface, moq_net::Timestamp::from_micros(i * 33_333).unwrap());
frames.extend(encoder.encode(&frame).unwrap());
}
frames.extend(encoder.finish().unwrap());
assert!(!frames.is_empty(), "encoder produced no frames");
let origin = id(moq_origin_create());
let path = b"video-raw-test";
let broadcast = publish_broadcast(origin, path);
let init = h264_init();
let format = moq_video_format::MOQ_VIDEO_FORMAT_AVC3;
let media = id(publish_video(broadcast, format, &init, None));
let consume = request_broadcast(origin, path);
let catalog_cb = Callback::new();
let catalog_task = id(unsafe { moq_consume_catalog(consume, Some(channel_callback), catalog_cb.ptr) });
let catalog_id = id(catalog_cb.recv());
let output = moq_video_decoder_output {
max_age_us: 10_000_000,
format: moq_video_pixel_format::MOQ_VIDEO_PIXEL_FORMAT_I420 as u32,
width: 0,
height: 0,
};
let frame_cb = Callback::new();
let consumer = id(unsafe { moq_decode_video(catalog_id, 0, &output, Some(channel_callback), frame_cb.ptr) });
for (i, frame) in frames.iter().enumerate() {
assert_eq!(
unsafe { moq_publish_media_frame(media, frame.payload.as_ptr(), frame.payload.len(), (i as u64) * 33_000) },
0
);
}
let frame_id = id(frame_cb.recv());
let mut frame = moq_video_frame {
timestamp_us: 0,
width: 0,
height: 0,
data: std::ptr::null(),
data_size: 0,
};
assert_eq!(unsafe { moq_decode_video_frame(frame_id, &mut frame) }, 0);
assert_eq!(frame.width, 320);
assert_eq!(frame.height, 240);
assert_eq!(frame.data_size, 320 * 240 * 3 / 2, "tightly-packed I420");
assert!(!frame.data.is_null());
assert_eq!(moq_decode_video_frame_free(frame_id), 0);
assert_eq!(moq_decode_video_cancel(consumer), 0);
loop {
let code = frame_cb.recv();
if code > 0 {
assert_eq!(moq_decode_video_frame_free(id(code)), 0);
} else {
assert_eq!(code, 0, "raw video close delivers terminal 0");
break;
}
}
assert_eq!(moq_consume_catalog_free(catalog_id), 0);
assert_eq!(moq_consume_catalog_cancel(catalog_task), 0);
loop {
let code = catalog_cb.recv();
if code > 0 {
assert_eq!(moq_consume_catalog_free(id(code)), 0);
} else {
assert_eq!(code, 0, "catalog close delivers terminal 0");
break;
}
}
assert_eq!(moq_consume_close(consume), 0);
assert_eq!(moq_publish_media_finish(media), 0);
assert_eq!(moq_publish_finish(broadcast), 0);
assert_eq!(moq_origin_close(origin), 0);
}
#[test]
fn multiple_frames_ordering() {
let origin = id(moq_origin_create());
let path = b"ordering-test";
let broadcast = publish_broadcast(origin, path);
let init = opus_head();
let format = moq_audio_format::MOQ_AUDIO_FORMAT_OPUS;
let media = id(publish_audio(broadcast, format, &init, None));
let consume = request_broadcast(origin, path);
let catalog_cb = Callback::new();
let catalog_task = id(unsafe { moq_consume_catalog(consume, Some(channel_callback), catalog_cb.ptr) });
let catalog_id = id(catalog_cb.recv());
let frame_cb = Callback::new();
let track = id(unsafe { moq_consume_audio(catalog_id, 0, 10_000_000, Some(channel_callback), frame_cb.ptr) });
let timestamps: [u64; 5] = [0, 20_000, 40_000, 60_000, 80_000];
for (i, &ts) in timestamps.iter().enumerate() {
let payload = format!("frame-{i}");
assert_eq!(
unsafe { moq_publish_media_frame(media, payload.as_ptr(), payload.len(), ts) },
0
);
}
for (i, &expected_ts) in timestamps.iter().enumerate() {
let frame_id = id(frame_cb.recv());
let mut frame = moq_frame {
payload: std::ptr::null(),
payload_size: 0,
timestamp_us: 0,
keyframe: false,
};
assert_eq!(unsafe { moq_consume_frame(frame_id, &mut frame) }, 0);
assert_eq!(frame.timestamp_us, expected_ts, "frame {i} has wrong timestamp");
let received = unsafe { std::slice::from_raw_parts(frame.payload, frame.payload_size) };
let expected = format!("frame-{i}");
assert_eq!(received, expected.as_bytes(), "frame {i} has wrong payload");
assert_eq!(moq_consume_frame_free(frame_id), 0);
}
assert_eq!(moq_consume_audio_cancel(track), 0);
assert_eq!(frame_cb.recv_terminal(), 0, "audio close delivers terminal 0");
assert_eq!(moq_consume_catalog_free(catalog_id), 0);
assert_eq!(moq_consume_catalog_cancel(catalog_task), 0);
assert_eq!(
catalog_cb.recv_catalog_terminal(),
0,
"catalog close delivers terminal 0"
);
assert_eq!(moq_consume_close(consume), 0);
assert_eq!(moq_publish_media_finish(media), 0);
assert_eq!(moq_publish_finish(broadcast), 0);
assert_eq!(moq_origin_close(origin), 0);
}
#[test]
fn catalog_update_on_new_track() {
let origin = id(moq_origin_create());
let path = b"catalog-update";
let broadcast = publish_broadcast(origin, path);
let init = opus_head();
let format = moq_audio_format::MOQ_AUDIO_FORMAT_OPUS;
let media1 = id(publish_audio(broadcast, format, &init, None));
let consume = request_broadcast(origin, path);
let catalog_cb = Callback::new();
let catalog_task = id(unsafe { moq_consume_catalog(consume, Some(channel_callback), catalog_cb.ptr) });
let catalog_id1 = id(catalog_cb.recv());
let mut audio_cfg = moq_audio_config {
name: std::ptr::null(),
name_len: 0,
label: std::ptr::null(),
label_len: 0,
codec: std::ptr::null(),
codec_len: 0,
description: std::ptr::null(),
description_len: 0,
sample_rate: 0,
channel_count: 0,
container: moq_container::default(),
};
assert_eq!(unsafe { moq_consume_audio_config(catalog_id1, 0, &mut audio_cfg) }, 0);
assert!(unsafe { moq_consume_audio_config(catalog_id1, 1, &mut audio_cfg) } < 0);
let media2 = id(publish_audio(broadcast, format, &init, None));
let catalog_id2 = id(catalog_cb.recv());
assert_eq!(unsafe { moq_consume_audio_config(catalog_id2, 0, &mut audio_cfg) }, 0);
assert_eq!(unsafe { moq_consume_audio_config(catalog_id2, 1, &mut audio_cfg) }, 0);
assert_eq!(moq_consume_catalog_free(catalog_id1), 0);
assert_eq!(moq_consume_catalog_free(catalog_id2), 0);
assert_eq!(moq_consume_catalog_cancel(catalog_task), 0);
assert_eq!(catalog_cb.recv_terminal(), 0, "catalog close delivers terminal 0");
assert_eq!(moq_consume_close(consume), 0);
assert_eq!(moq_publish_media_finish(media1), 0);
assert_eq!(moq_publish_media_finish(media2), 0);
assert_eq!(moq_publish_finish(broadcast), 0);
assert_eq!(moq_origin_close(origin), 0);
}
#[test]
fn null_pointer_handling() {
assert_eq!(
unsafe { moq_consume_frame(9999, std::ptr::null_mut()) },
-6,
"null dst should return InvalidPointer (-6)"
);
assert_eq!(
unsafe { moq_consume_video_config(9999, 0, std::ptr::null_mut()) },
-6,
"null dst should return InvalidPointer (-6)"
);
assert_eq!(
unsafe { moq_consume_audio_config(9999, 0, std::ptr::null_mut()) },
-6,
"null dst should return InvalidPointer (-6)"
);
assert_eq!(
unsafe { moq_origin_announced_info(9999, std::ptr::null_mut()) },
-6,
"null dst should return InvalidPointer (-6)"
);
}
#[test]
fn session_connect_invalid_url() {
let url = b"not a valid url!!!";
let ret = unsafe {
moq_session_connect(
url.as_ptr() as *const c_char,
url.len(),
std::ptr::null(),
0,
0,
None,
std::ptr::null_mut(),
)
};
assert!(ret < 0, "connecting with an invalid URL should fail immediately");
}
#[test]
fn session_connect_and_close() {
let cb = Callback::new();
let url = b"moqt://localhost:1";
let session = id(unsafe {
moq_session_connect(
url.as_ptr() as *const c_char,
url.len(),
std::ptr::null(),
0,
0,
Some(channel_callback),
cb.ptr,
)
});
assert_eq!(moq_session_close(session), 0);
assert!(cb.recv() <= 0, "session close delivers a terminal code");
}
fn moq_str(s: &str) -> moq_string {
moq_string {
data: s.as_ptr() as *const c_char,
len: s.len(),
}
}
fn client_config() -> moq_client_config {
unsafe { std::mem::zeroed() }
}
fn dial(config: Option<&moq_client_config>) -> i32 {
let url = b"moqt://localhost:1";
unsafe {
moq_session_connect(
url.as_ptr() as *const c_char,
url.len(),
config.map_or(std::ptr::null(), |c| c as *const _),
0,
0,
Some(ignore_callback),
std::ptr::null_mut(),
)
}
}
fn parsed(config: &moq_client_config) -> crate::client::Config {
unsafe { crate::parse_client(Some(config)) }.expect("config should parse")
}
#[test]
fn a_null_config_dials_with_the_defaults() {
let defaults = crate::client::Config::default();
let parsed = unsafe { crate::parse_client(None) }.expect("NULL is the defaults");
assert_eq!(parsed.connect.backoff.initial, defaults.connect.backoff.initial);
assert_eq!(
parsed.connect.websocket.resolve().enabled,
defaults.connect.websocket.resolve().enabled
);
}
#[test]
fn a_zeroed_config_is_the_defaults() {
let defaults = crate::client::Config::default();
let parsed = parsed(&client_config());
assert_eq!(parsed.connect.backoff.initial, defaults.connect.backoff.initial);
assert_eq!(parsed.connect.backoff.multiplier, defaults.connect.backoff.multiplier);
assert_eq!(parsed.connect.backoff.max, defaults.connect.backoff.max);
assert_eq!(parsed.connect.backoff.timeout, defaults.connect.backoff.timeout);
assert_eq!(
parsed.connect.websocket.resolve().enabled,
defaults.connect.websocket.resolve().enabled
);
assert_eq!(
parsed.connect.websocket.resolve().delay,
defaults.connect.websocket.resolve().delay
);
assert_eq!(parsed.connect.version, defaults.connect.version);
assert_eq!(parsed.connect.bind, defaults.connect.bind);
assert!(parsed.connect.tls.fingerprint.is_empty());
assert!(parsed.connect.tls.root.is_empty());
assert!(parsed.quic.gso.is_none());
assert!(parsed.quic.mtu_discovery.is_none());
}
#[test]
fn defaults_report_what_a_zeroed_config_dials() {
let config = moq_client_defaults();
let expected = crate::client::Config::default();
let quic = moq_tokio::quic::Resolved::default();
let connect = expected.connect.resolve();
assert!(config.has_connect_timeout);
assert_eq!(config.connect_timeout_us, connect.timeout.as_micros() as u64);
assert!(config.has_failover_delay);
assert_eq!(config.failover_delay_us, connect.race.as_micros() as u64);
assert!(config.has_resolution_delay);
assert_eq!(config.resolution_delay_us, connect.resolution_delay.as_micros() as u64);
assert!(config.has_backoff_initial);
assert_eq!(
config.backoff_initial_us,
expected.connect.backoff.initial.as_micros() as u64
);
assert!(config.has_backoff_multiplier);
assert_eq!(config.backoff_multiplier, expected.connect.backoff.multiplier);
assert!(config.has_backoff_max);
assert_eq!(config.backoff_max_us, expected.connect.backoff.max.as_micros() as u64);
assert!(config.has_backoff_timeout);
assert_eq!(
config.backoff_timeout_us,
expected.connect.backoff.timeout.as_micros() as u64
);
let websocket = expected.connect.websocket.resolve();
assert!(config.has_websocket_enabled);
assert_eq!(config.websocket_enabled, websocket.enabled);
assert!(config.has_websocket_delay);
assert_eq!(config.websocket_delay_us, websocket.delay.as_micros() as u64);
assert!(config.has_quic_max_streams);
assert_eq!(config.quic_max_streams, quic.max_streams);
assert!(config.has_quic_idle_timeout);
assert_eq!(config.quic_idle_timeout_us, quic.idle_timeout.as_micros() as u64);
assert_eq!(
config.has_quic_keep_alive.then_some(config.quic_keep_alive_us),
quic.keep_alive.map(|d| d.as_micros() as u64)
);
assert!(!config.has_quic_gso);
assert!(!config.has_quic_mtu_discovery);
assert!(!config.has_tls_system_roots);
assert!(config.quic_congestion_control.is_null());
let expected = crate::client::Config::default();
let reparsed = parsed(&config);
assert_eq!(reparsed.connect.backoff.initial, expected.connect.backoff.initial);
assert_eq!(
reparsed.connect.websocket.resolve().enabled,
expected.connect.websocket.resolve().enabled
);
}
#[test]
fn zero_with_a_flag_set_is_a_real_value() {
let mut config = client_config();
config.backoff_timeout_us = 0;
config.has_backoff_timeout = true;
config.quic_keep_alive_us = 0;
config.has_quic_keep_alive = true;
let explicit = parsed(&config);
assert_eq!(explicit.connect.backoff.timeout, std::time::Duration::ZERO);
assert_eq!(explicit.quic.keep_alive, std::time::Duration::ZERO);
let defaults = parsed(&client_config());
assert_ne!(defaults.connect.backoff.timeout, std::time::Duration::ZERO);
assert_eq!(defaults.quic.keep_alive, std::time::Duration::from_secs(5));
}
#[test]
fn config_versions_round_trip() {
let versions = [moq_str("moq-lite-05"), moq_str("moq-transport-19")];
let mut config = client_config();
config.versions = versions.as_ptr();
config.versions_len = versions.len();
assert_eq!(parsed(&config).connect.version.len(), 2);
let mut config = client_config();
config.versions_len = 0;
assert_eq!(
parsed(&config).connect.version,
crate::client::Config::default().connect.version
);
}
#[test]
fn config_rejects_an_unknown_version() {
let versions = [moq_str("moq-lite-05"), moq_str("moq-carrier-pigeon-01")];
let mut config = client_config();
config.versions = versions.as_ptr();
config.versions_len = versions.len();
assert_eq!(dial(Some(&config)), Error::InvalidConfig(String::new()).code());
}
#[test]
fn config_rejects_a_bad_bind_address() {
let good = "127.0.0.1:0";
let mut config = client_config();
config.bind = good.as_ptr() as *const c_char;
config.bind_len = good.len();
assert_eq!(parsed(&config).connect.bind, Some(good.parse().unwrap()));
let bad = "not-an-address";
let mut config = client_config();
config.bind = bad.as_ptr() as *const c_char;
config.bind_len = bad.len();
assert_eq!(dial(Some(&config)), Error::InvalidConfig(String::new()).code());
}
#[test]
fn config_optional_strings_are_unset_when_null_or_empty() {
let name = "relay.example.com";
let mut config = client_config();
config.tls_host_name = name.as_ptr() as *const c_char;
config.tls_host_name_len = name.len();
assert_eq!(parsed(&config).connect.tls.host_name.as_deref(), Some(name));
let config = client_config();
assert_eq!(parsed(&config).connect.tls.host_name, None);
let mut config = client_config();
config.tls_host_name = "".as_ptr() as *const c_char;
config.tls_host_name_len = 0;
assert_eq!(parsed(&config).connect.tls.host_name, None);
}
#[test]
fn config_rejects_malformed_tls_fingerprints() {
for invalid in ["not-hex", "abcd"] {
let fingerprints = [moq_str(invalid)];
let mut config = client_config();
config.tls_fingerprints = fingerprints.as_ptr();
config.tls_fingerprints_len = fingerprints.len();
assert_eq!(dial(Some(&config)), Error::InvalidConfig(String::new()).code());
}
let valid_value = "ab".repeat(32);
let valid = [moq_str(&valid_value)];
let mut config = client_config();
config.tls_fingerprints = valid.as_ptr();
config.tls_fingerprints_len = valid.len();
assert_eq!(parsed(&config).connect.tls.fingerprint, vec![valid_value]);
}
#[test]
fn config_rejects_unknown_congestion_control() {
let bogus = "sideways";
let mut config = client_config();
config.quic_congestion_control = bogus.as_ptr() as *const c_char;
config.quic_congestion_control_len = bogus.len();
assert_eq!(dial(Some(&config)), Error::InvalidConfig(String::new()).code());
let delay = "delay";
let mut config = client_config();
config.quic_congestion_control = delay.as_ptr() as *const c_char;
config.quic_congestion_control_len = delay.len();
assert!(parsed(&config).quic.congestion_control.is_some());
assert!(parsed(&client_config()).quic.congestion_control.is_none());
}
#[test]
fn config_quic_and_backoff_knobs_apply() {
let mut config = client_config();
config.backoff_initial_us = 500_000;
config.has_backoff_initial = true;
config.backoff_multiplier = 3;
config.has_backoff_multiplier = true;
config.backoff_max_us = 10_000_000;
config.has_backoff_max = true;
config.quic_max_streams = 4096;
config.has_quic_max_streams = true;
config.quic_idle_timeout_us = 15_000_000;
config.has_quic_idle_timeout = true;
config.quic_gso = false;
config.has_quic_gso = true;
config.quic_mtu_discovery = true;
config.has_quic_mtu_discovery = true;
let dir = "/tmp/qlog";
config.quic_qlog = dir.as_ptr() as *const c_char;
config.quic_qlog_len = dir.len();
let parsed = parsed(&config);
assert_eq!(
parsed.connect.backoff.initial,
std::time::Duration::from_micros(500_000)
);
assert_eq!(parsed.connect.backoff.multiplier, 3);
assert_eq!(parsed.connect.backoff.max, std::time::Duration::from_micros(10_000_000));
assert_eq!(parsed.quic.max_streams, Some(4096));
assert_eq!(parsed.quic.idle_timeout, std::time::Duration::from_millis(15_000));
assert_eq!(parsed.quic.gso, Some(false));
assert_eq!(parsed.quic.mtu_discovery, Some(true));
assert_eq!(parsed.quic.qlog.as_deref(), Some(std::path::Path::new(dir)));
}
#[test]
fn qlog_support_matches_what_a_dial_accepts() {
let dir = std::env::temp_dir().join(format!("moq-qlog-test-{}", std::process::id()));
std::fs::create_dir_all(&dir).expect("create the qlog directory");
let path = dir.to_str().expect("temp dir is UTF-8").to_string();
let mut config = client_config();
config.quic_qlog = path.as_ptr() as *const c_char;
config.quic_qlog_len = path.len();
assert!(unsafe { crate::parse_client(Some(&config)) }.is_ok());
let ret = dial(Some(&config));
match moq_qlog_supported() {
true => {
assert!(ret > 0, "qlog is supported, so the dial must start: {ret}");
moq_session_close(id(ret));
}
false => assert!(ret < 0, "qlog is unsupported, so the dial must be refused"),
}
std::fs::remove_dir_all(&dir).ok();
}
#[test]
fn versions_lists_the_offered_set() {
let count = unsafe { moq_versions(std::ptr::null_mut(), 0) };
assert!(count > 0, "expected at least one offered version, got {count}");
let mut names = vec![
moq_string {
data: std::ptr::null(),
len: 0
};
count as usize
];
assert_eq!(unsafe { moq_versions(names.as_mut_ptr(), names.len()) }, count);
for name in &names {
let name = unsafe { ffi::parse_str(name.data, name.len) }.expect("version name is UTF-8");
let one = [moq_str(name)];
let mut config = client_config();
config.versions = one.as_ptr();
config.versions_len = one.len();
assert!(
unsafe { crate::parse_client(Some(&config)) }.is_ok(),
"listed version {name} must be settable"
);
}
}
#[test]
fn dial_applies_the_config() {
let versions = [moq_str("moq-lite-05")];
let mut config = client_config();
config.versions = versions.as_ptr();
config.versions_len = versions.len();
config.connect_timeout_us = 100_000;
config.has_connect_timeout = true;
let cb = Callback::new();
let url = b"moqt://localhost:1";
let session = id(unsafe {
moq_session_connect(
url.as_ptr() as *const c_char,
url.len(),
&config,
0,
0,
Some(channel_callback),
cb.ptr,
)
});
assert_eq!(moq_session_close(session), 0);
assert!(cb.recv() <= 0, "session close delivers a terminal code");
}
fn test_allocator() -> (u32, moq_net::bandwidth::Producer) {
let estimate = moq_net::bandwidth::Producer::new();
let allocator = moq_net::bandwidth::Allocator::new(estimate.consume());
let bandwidth = {
let mut state = State::lock();
id(i32::from(state.bandwidth.insert(allocator).unwrap()))
};
(bandwidth, estimate)
}
fn grant(reservation: u32) -> Option<u64> {
let mut bps = 0;
let mut present = false;
assert_eq!(unsafe { moq_reservation_grant(reservation, &mut bps, &mut present) }, 0);
present.then_some(bps)
}
fn wait_grant(reservation: u32) -> u64 {
let deadline = std::time::Instant::now() + TIMEOUT;
loop {
if let Some(bps) = grant(reservation) {
return bps;
}
assert!(
std::time::Instant::now() < deadline,
"timed out waiting for a reservation grant"
);
std::thread::sleep(Duration::from_millis(10));
}
}
fn publish_named_track(broadcast: u32, name: &[u8]) -> u32 {
id(unsafe { moq_publish_track(broadcast, name.as_ptr() as *const c_char, name.len(), std::ptr::null()) })
}
fn subscribe_named_track(consume: u32, name: &[u8]) -> (u32, Callback) {
let cb = Callback::new();
let consumer = id(unsafe {
moq_consume_track(
consume,
name.as_ptr() as *const c_char,
name.len(),
std::ptr::null(),
Some(channel_callback),
cb.ptr,
)
});
(consumer, cb)
}
#[test]
fn bandwidth_reservations_split_the_estimate() {
let (bandwidth, estimate) = test_allocator();
let origin = id(moq_origin_create());
let path = b"bandwidth-split";
let broadcast = publish_broadcast(origin, path);
let consume = request_broadcast(origin, path);
let first_track = publish_named_track(broadcast, b"a");
let second_track = publish_named_track(broadcast, b"b");
let (first_sub, first_cb) = subscribe_named_track(consume, b"a");
let (second_sub, second_cb) = subscribe_named_track(consume, b"b");
let first = id(moq_bandwidth_reserve(bandwidth, first_track, 4_000_000));
let second = id(moq_bandwidth_reserve(bandwidth, second_track, 2_000_000));
estimate
.set(Some(moq_net::bandwidth::Rate::from_bps(3_000_000)))
.unwrap();
wait_grant(first);
wait_grant(second);
let a = grant(first).expect("first grant");
let b = grant(second).expect("second grant");
assert!(a + b <= 3_000_000, "{a} + {b} oversubscribed");
assert_eq!(a, 1_500_000);
assert_eq!(b, 1_500_000);
assert_eq!(moq_reservation_close(first), 0);
assert_eq!(grant(second), Some(2_000_000));
assert_eq!(moq_reservation_close(second), 0);
assert_eq!(moq_bandwidth_close(bandwidth), 0);
assert_eq!(moq_consume_track_cancel(first_sub), 0);
assert_eq!(moq_consume_track_cancel(second_sub), 0);
let _ = first_cb.recv_terminal();
let _ = second_cb.recv_terminal();
assert_eq!(moq_consume_close(consume), 0);
assert_eq!(moq_publish_finish(broadcast), 0);
assert_eq!(moq_origin_close(origin), 0);
}
#[test]
fn bandwidth_handles_share_the_registry() {
let estimate = moq_net::bandwidth::Producer::new();
estimate
.set(Some(moq_net::bandwidth::Rate::from_bps(3_000_000)))
.unwrap();
let allocator = moq_net::bandwidth::Allocator::new(estimate.consume());
let (first, second) = {
let mut state = State::lock();
let first = id(i32::from(state.bandwidth.insert(allocator.clone()).unwrap()));
let second = id(i32::from(state.bandwidth.insert(allocator).unwrap()));
(first, second)
};
let origin = id(moq_origin_create());
let path = b"bandwidth-shared";
let broadcast = publish_broadcast(origin, path);
let consume = request_broadcast(origin, path);
let first_track = publish_named_track(broadcast, b"a");
let second_track = publish_named_track(broadcast, b"b");
let (first_sub, first_cb) = subscribe_named_track(consume, b"a");
let (second_sub, second_cb) = subscribe_named_track(consume, b"b");
let reserved = id(moq_bandwidth_reserve(first, first_track, 4_000_000));
assert_eq!(wait_grant(reserved), 3_000_000);
let other = id(moq_bandwidth_reserve(second, second_track, 4_000_000));
assert_eq!(wait_grant(other), 1_500_000);
assert_eq!(grant(reserved), Some(1_500_000));
assert_eq!(moq_reservation_close(reserved), 0);
assert_eq!(moq_reservation_close(other), 0);
assert_eq!(moq_bandwidth_close(first), 0);
assert_eq!(moq_bandwidth_close(second), 0);
assert_eq!(moq_consume_track_cancel(first_sub), 0);
assert_eq!(moq_consume_track_cancel(second_sub), 0);
let _ = first_cb.recv_terminal();
let _ = second_cb.recv_terminal();
assert_eq!(moq_consume_close(consume), 0);
assert_eq!(moq_publish_finish(broadcast), 0);
assert_eq!(moq_origin_close(origin), 0);
}
#[test]
fn encode_video_bitrate_caps_the_reservation() {
let (bandwidth, estimate) = test_allocator();
let origin = id(moq_origin_create());
let path = b"bitrate-cap";
let broadcast = publish_broadcast(origin, path);
let consume = request_broadcast(origin, path);
let input = moq_video_encoder_input {
format: moq_video_pixel_format::MOQ_VIDEO_PIXEL_FORMAT_RGBA as u32,
width: 320,
height: 240,
framerate: 30,
};
let output = moq_video_encoder_output {
codec: moq_video_codec::MOQ_VIDEO_CODEC_H264 as u32,
bitrate: 4_000_000,
gop: 0,
kind: moq_video_encoder_kind::MOQ_VIDEO_ENCODER_KIND_SOFTWARE as u32,
encoder: std::ptr::null(),
encoder_len: 0,
};
let producer = id(unsafe { moq_encode_video(broadcast, &input, &output, bandwidth) });
let reservation = id(moq_encode_video_reservation(producer));
let (sub, cb) = subscribe_named_track(consume, b"0.avc3");
estimate
.set(Some(moq_net::bandwidth::Rate::from_bps(4_000_000)))
.unwrap();
assert_eq!(wait_grant(reservation), 4_000_000);
assert_eq!(moq_encode_video_bitrate(producer, 1_000_000), 0);
assert_eq!(grant(reservation), Some(1_000_000));
estimate
.set(Some(moq_net::bandwidth::Rate::from_bps(3_000_000)))
.unwrap();
assert_eq!(grant(reservation), Some(1_000_000));
assert_eq!(moq_reservation_close(reservation), 0);
assert_eq!(moq_consume_track_cancel(sub), 0);
let _ = cb.recv_terminal();
assert_eq!(moq_encode_video_finish(producer), 0);
assert_eq!(moq_bandwidth_close(bandwidth), 0);
assert_eq!(moq_consume_close(consume), 0);
assert_eq!(moq_publish_finish(broadcast), 0);
assert_eq!(moq_origin_close(origin), 0);
}
fn listen(cb: &Callback) -> (u32, String, String) {
let bind = "127.0.0.1:0";
let generate = [moq_str("localhost")];
let mut config: moq_server_config = unsafe { std::mem::zeroed() };
config.bind = bind.as_ptr() as *const c_char;
config.bind_len = bind.len();
config.tls_generate = generate.as_ptr();
config.tls_generate_len = generate.len();
let server = id(unsafe { moq_server_listen(&config, Some(channel_callback), cb.ptr) });
let mut addr = moq_str("");
assert_eq!(unsafe { moq_server_addr(server, &mut addr) }, 0);
let addr = borrowed_string(addr.data, addr.len).unwrap();
let count = unsafe { moq_server_fingerprints(server, std::ptr::null_mut(), 0) };
assert!(count > 0, "a generated certificate has a fingerprint, got {count}");
let mut fingerprints = vec![moq_str(""); count as usize];
assert_eq!(
unsafe { moq_server_fingerprints(server, fingerprints.as_mut_ptr(), fingerprints.len()) },
count
);
let fingerprint = borrowed_string(fingerprints[0].data, fingerprints[0].len).unwrap();
(server, addr, fingerprint)
}
fn dial_pinned(
url: &str,
fingerprint: &str,
consume: u32,
on_status: extern "C" fn(*mut c_void, i32),
user_data: *mut c_void,
) -> u32 {
let fingerprints = [moq_str(fingerprint)];
let mut config = client_config();
config.tls_fingerprints = fingerprints.as_ptr();
config.tls_fingerprints_len = fingerprints.len();
config.tls_host_name = "localhost".as_ptr() as *const c_char;
config.tls_host_name_len = "localhost".len();
id(unsafe {
moq_session_connect(
url.as_ptr() as *const c_char,
url.len(),
&config,
0,
consume,
Some(on_status),
user_data,
)
})
}
#[test]
fn server_accepts_a_session() {
let served = id(moq_origin_create());
let broadcast = publish_broadcast(served, b"demo");
let server_cb = Callback::new();
let (server, addr, fingerprint) = listen(&server_cb);
let received = id(moq_origin_create());
let client_cb = Callback::new();
let client = dial_pinned(
&format!("https://{addr}/room?jwt=abc"),
&fingerprint,
received,
channel_callback,
client_cb.ptr,
);
let request = id(server_cb.recv());
let mut path = moq_str("");
assert_eq!(unsafe { moq_session_request_path(request, &mut path) }, 0);
assert_eq!(borrowed_string(path.data, path.len).as_deref(), Some("/room"));
let mut query = moq_str("");
assert_eq!(unsafe { moq_session_request_query(request, &mut query) }, 0);
assert_eq!(borrowed_string(query.data, query.len).as_deref(), Some("jwt=abc"));
let session_cb = Callback::new();
let session = id(unsafe { moq_session_request_accept(request, served, 0, Some(channel_callback), session_cb.ptr) });
assert!(
unsafe { moq_session_request_path(request, &mut path) } < 0,
"accept consumes the request"
);
assert_eq!(session_cb.recv(), 1);
assert_eq!(client_cb.recv(), 1);
let mut stats: moq_connection_stats = unsafe { std::mem::zeroed() };
assert_eq!(unsafe { moq_session_stats(session, &mut stats) }, 0);
let bandwidth = id(moq_session_bandwidth(session));
assert_eq!(moq_bandwidth_close(bandwidth), 0);
let consumed = request_broadcast(received, b"demo");
assert_eq!(moq_consume_close(consumed), 0);
assert_eq!(moq_session_close(session), 0);
assert_eq!(session_cb.recv_terminal(), 0);
assert_eq!(moq_session_close(client), 0);
client_cb.recv_terminal();
assert_eq!(moq_server_close(server), 0);
assert_eq!(server_cb.recv_terminal(), 0);
assert!(
moq_server_close(server) < 0,
"the server handle is gone after its terminal callback"
);
assert_eq!(moq_publish_finish(broadcast), 0);
assert_eq!(moq_origin_close(served), 0);
assert_eq!(moq_origin_close(received), 0);
}
#[test]
fn server_rejects_a_session() {
let server_cb = Callback::new();
let (server, addr, fingerprint) = listen(&server_cb);
let client_cb = ProtocolCallback::new();
let client = dial_pinned(
&format!("https://{addr}/"),
&fingerprint,
0,
protocol_callback,
client_cb.ptr,
);
let request = id(server_cb.recv());
let mut path = moq_str("x");
assert_eq!(unsafe { moq_session_request_path(request, &mut path) }, 0);
assert_eq!(path.len, 0, "the root is an empty path");
let mut query = moq_str("x");
assert_eq!(unsafe { moq_session_request_query(request, &mut query) }, 0);
assert!(query.data.is_null(), "no query is NULL");
assert_eq!(moq_session_request_reject(request, 403), 0);
assert!(moq_session_request_free(request) < 0, "reject consumes the request");
let mut status = client_cb.recv();
if status.0 == 1 {
status = client_cb.recv();
}
let (code, protocol) = status;
assert_eq!(code, MOQ_ERROR_MOQ);
let protocol = protocol.expect("the rejection carries a protocol close");
assert_eq!(protocol.scope, moq_error_scope::MOQ_ERROR_SCOPE_SESSION as u32);
assert_eq!(protocol.kind, moq_protocol_kind::MOQ_PROTOCOL_KIND_UNAUTHORIZED as u32);
assert!(moq_session_close(client) < 0, "the client session already ended");
assert_eq!(moq_server_close(server), 0);
assert_eq!(server_cb.recv_terminal(), 0);
}
#[test]
fn server_listen_refuses_bad_config() {
assert_eq!(
unsafe { moq_server_listen(std::ptr::null(), Some(ignore_callback), std::ptr::null_mut()) },
MOQ_ERROR_INVALID_POINTER
);
let mut config: moq_server_config = unsafe { std::mem::zeroed() };
let bind = "not an address";
config.bind = bind.as_ptr() as *const c_char;
config.bind_len = bind.len();
assert_eq!(
unsafe { moq_server_listen(&config, Some(ignore_callback), std::ptr::null_mut()) },
MOQ_ERROR_INVALID_CONFIG
);
let bind = "127.0.0.1:0";
config.bind = bind.as_ptr() as *const c_char;
config.bind_len = bind.len();
assert_eq!(
unsafe { moq_server_listen(&config, Some(ignore_callback), std::ptr::null_mut()) },
MOQ_ERROR_INVALID_CONFIG
);
}
fn read_raw_frames(consume: u32, name: &[u8], count: usize) -> Vec<Vec<u8>> {
let frame_cb = Callback::new();
let subscription = moq_subscription {
priority: 0,
max_age_us: 1_000_000,
group_start: 0,
group_start_present: false,
group_end: 0,
group_end_present: false,
};
let track = id(unsafe {
moq_consume_track(
consume,
name.as_ptr() as *const c_char,
name.len(),
&subscription,
Some(channel_callback),
frame_cb.ptr,
)
});
let mut payloads = Vec::with_capacity(count);
for _ in 0..count {
let frame_id = id(frame_cb.recv());
let mut frame = moq_frame {
payload: std::ptr::null(),
payload_size: 0,
timestamp_us: 0,
keyframe: false,
};
assert_eq!(unsafe { moq_consume_track_frame(frame_id, &mut frame) }, 0);
payloads.push(unsafe { std::slice::from_raw_parts(frame.payload, frame.payload_size) }.to_vec());
assert_eq!(moq_consume_track_frame_free(frame_id), 0);
}
assert_eq!(moq_consume_track_cancel(track), 0);
assert_eq!(frame_cb.recv_terminal(), 0, "clean cancel delivers terminal 0");
payloads
}
fn published_catalog(broadcast: u32) -> moq_mux::catalog::hang::Catalog<moq_mux::catalog::hang::Extra> {
let id = crate::Id::try_from(broadcast).expect("valid broadcast id");
crate::State::lock()
.publish
.catalog_snapshot(id)
.expect("broadcast exists")
}
#[test]
fn json_tracks_are_advertised_in_the_catalog() {
let origin = id(moq_origin_create());
let broadcast = publish_broadcast(origin, b"json-catalog");
let status = b"status";
let snapshot = id(unsafe {
moq_publish_json_snapshot(
broadcast,
status.as_ptr() as *const c_char,
status.len(),
&moq_json_snapshot_config {
delta_ratio: 4,
compression: true,
},
)
});
let events = b"events";
let stream = id(unsafe {
moq_publish_json_stream(
broadcast,
events.as_ptr() as *const c_char,
events.len(),
&moq_json_stream_config { compression: false },
)
});
let catalog = published_catalog(broadcast);
let entry = catalog.json.tracks.get("status").expect("snapshot track advertised");
assert_eq!(entry.mode, hang::catalog::Mode::Snapshot);
assert_eq!(entry.compression, Some(hang::catalog::Compression::Deflate));
let entry = catalog.json.tracks.get("events").expect("stream track advertised");
assert_eq!(entry.mode, hang::catalog::Mode::Stream);
assert_eq!(entry.compression, None);
assert_eq!(moq_publish_json_snapshot_finish(snapshot), 0);
let catalog = published_catalog(broadcast);
assert!(
!catalog.json.tracks.contains_key("status"),
"finished track still advertised"
);
assert!(catalog.json.tracks.contains_key("events"));
assert_eq!(moq_publish_json_stream_finish(stream), 0);
assert!(published_catalog(broadcast).json.tracks.is_empty());
assert_eq!(moq_publish_finish(broadcast), 0);
assert_eq!(moq_origin_close(origin), 0);
}
#[test]
fn binary_snapshot_is_advertised_and_delivered() {
let origin = id(moq_origin_create());
let path = b"binary-snapshot";
let broadcast = publish_broadcast(origin, path);
let name = b"thumbnail";
let mime = b"image/jpeg";
let producer = id(unsafe {
moq_publish_binary_snapshot(
broadcast,
name.as_ptr() as *const c_char,
name.len(),
&moq_binary_config {
compression: false,
mime: mime.as_ptr() as *const c_char,
mime_len: mime.len(),
},
)
});
let catalog = published_catalog(broadcast);
let entry = catalog.binary.tracks.get("thumbnail").expect("binary track advertised");
assert_eq!(entry.mode, hang::catalog::Mode::Snapshot);
assert_eq!(entry.mime.as_deref(), Some("image/jpeg"));
let payload = [0xff_u8, 0xd8, 0xff, 0xe0, 1, 2, 3];
assert_eq!(
unsafe { moq_publish_binary_snapshot_update(producer, payload.as_ptr(), payload.len()) },
0
);
let consume = request_broadcast(origin, path);
let frames = read_raw_frames(consume, name, 1);
assert_eq!(frames, vec![payload.to_vec()]);
assert_eq!(moq_publish_binary_snapshot_finish(producer), 0);
assert!(
moq_publish_binary_snapshot_finish(producer) < 0,
"double-finish should fail"
);
assert!(published_catalog(broadcast).binary.tracks.is_empty());
assert_eq!(moq_consume_close(consume), 0);
assert_eq!(moq_publish_finish(broadcast), 0);
assert_eq!(moq_origin_close(origin), 0);
}
#[test]
fn binary_stream_is_advertised_and_delivered() {
let origin = id(moq_origin_create());
let path = b"binary-stream";
let broadcast = publish_broadcast(origin, path);
let name = b"blobs";
let producer = id(unsafe {
moq_publish_binary_stream(
broadcast,
name.as_ptr() as *const c_char,
name.len(),
&moq_binary_config {
compression: false,
mime: std::ptr::null(),
mime_len: 0,
},
)
});
let catalog = published_catalog(broadcast);
let entry = catalog.binary.tracks.get("blobs").expect("binary stream advertised");
assert_eq!(entry.mode, hang::catalog::Mode::Stream);
assert_eq!(entry.mime, None);
let payloads: [&[u8]; 2] = [b"first", b"second"];
for payload in payloads {
assert_eq!(
unsafe { moq_publish_binary_stream_append(producer, payload.as_ptr(), payload.len()) },
0
);
}
let consume = request_broadcast(origin, path);
let frames = read_raw_frames(consume, name, payloads.len());
assert_eq!(frames, payloads.map(<[u8]>::to_vec));
assert_eq!(moq_publish_binary_stream_finish(producer), 0);
assert!(published_catalog(broadcast).binary.tracks.is_empty());
assert_eq!(moq_consume_close(consume), 0);
assert_eq!(moq_publish_finish(broadcast), 0);
assert_eq!(moq_origin_close(origin), 0);
}
#[test]
fn data_track_names_cannot_collide() {
let origin = id(moq_origin_create());
let broadcast = publish_broadcast(origin, b"data-collide");
let name = b"state";
let first = id(unsafe {
moq_publish_json_snapshot(
broadcast,
name.as_ptr() as *const c_char,
name.len(),
&moq_json_snapshot_config {
delta_ratio: 0,
compression: false,
},
)
});
assert!(
unsafe {
moq_publish_binary_stream(
broadcast,
name.as_ptr() as *const c_char,
name.len(),
&moq_binary_config {
compression: false,
mime: std::ptr::null(),
mime_len: 0,
},
)
} < 0,
"a duplicate data track name should fail"
);
assert_eq!(
published_catalog(broadcast)
.json
.tracks
.get("state")
.map(|e| e.mode.clone()),
Some(hang::catalog::Mode::Snapshot),
"the refused duplicate must leave the first entry in place"
);
let other = b"other";
assert!(
unsafe {
moq_publish_binary_stream(
broadcast,
other.as_ptr() as *const c_char,
other.len(),
std::ptr::null(),
)
} < 0
);
assert_eq!(moq_publish_json_snapshot_finish(first), 0);
assert_eq!(moq_publish_finish(broadcast), 0);
assert_eq!(moq_origin_close(origin), 0);
}