#![cfg(feature = "rtmp")]
use bytes::Bytes;
use ez_ffmpeg::rtmp::embed_rtmp_server::EmbedRtmpServer;
use ez_ffmpeg::{FfmpegContext, Input};
use rml_rtmp::handshake::{Handshake, HandshakeProcessResult, PeerType};
use rml_rtmp::sessions::{
ClientSession, ClientSessionConfig, ClientSessionEvent, ClientSessionResult,
};
use std::io::{ErrorKind, Read, Write};
use std::net::{SocketAddr, TcpStream};
use std::time::{Duration, Instant};
const WATCHDOG: Duration = Duration::from_secs(30);
#[derive(Debug)]
enum WatcherEvent {
Video(Bytes),
Audio,
Metadata,
Status(String),
}
struct Watcher {
stream: TcpStream,
session: ClientSession,
events: Vec<WatcherEvent>,
connected: bool,
playing: bool,
eof: bool,
}
impl Watcher {
fn connect(addr: SocketAddr, app: &str, stream_key: &str, watchdog: Duration) -> Watcher {
let deadline = Instant::now() + watchdog;
let mut stream = TcpStream::connect(addr).expect("watcher connect");
stream.set_nodelay(true).ok();
stream
.set_read_timeout(Some(Duration::from_millis(100)))
.expect("set read timeout");
let mut handshake = Handshake::new(PeerType::Client);
let c0c1 = handshake
.generate_outbound_p0_and_p1()
.expect("handshake c0+c1");
stream.write_all(&c0c1).expect("send c0+c1");
let mut buf = [0u8; 8192];
let leftover = loop {
let n = read_some(&mut stream, &mut buf, deadline);
assert!(n > 0, "server closed during handshake");
match handshake.process_bytes(&buf[..n]).expect("handshake") {
HandshakeProcessResult::InProgress { response_bytes } => {
if !response_bytes.is_empty() {
stream.write_all(&response_bytes).expect("handshake send");
}
}
HandshakeProcessResult::Completed {
response_bytes,
remaining_bytes,
} => {
if !response_bytes.is_empty() {
stream.write_all(&response_bytes).expect("handshake send");
}
break remaining_bytes;
}
}
};
let (session, initial_results) =
ClientSession::new(ClientSessionConfig::new()).expect("client session");
let mut watcher = Watcher {
stream,
session,
events: Vec::new(),
connected: false,
playing: false,
eof: false,
};
watcher.apply(initial_results);
watcher.feed(&leftover);
let result = watcher
.session
.request_connection(app.to_string())
.expect("request_connection");
watcher.apply(vec![result]);
watcher.pump_until(watchdog, |w| w.connected);
assert!(watcher.connected, "server never accepted the connection");
let result = watcher
.session
.request_playback(stream_key.to_string())
.expect("request_playback");
watcher.apply(vec![result]);
watcher.pump_until(watchdog, |w| w.playing);
assert!(watcher.playing, "server never accepted the play request");
watcher
}
fn apply(&mut self, results: Vec<ClientSessionResult>) {
for result in results {
match result {
ClientSessionResult::OutboundResponse(packet) => {
self.stream
.write_all(&packet.bytes)
.expect("watcher send to server");
}
ClientSessionResult::RaisedEvent(event) => match event {
ClientSessionEvent::ConnectionRequestAccepted => self.connected = true,
ClientSessionEvent::PlaybackRequestAccepted => self.playing = true,
ClientSessionEvent::VideoDataReceived { data, .. } => {
self.events.push(WatcherEvent::Video(data))
}
ClientSessionEvent::AudioDataReceived { .. } => {
self.events.push(WatcherEvent::Audio)
}
ClientSessionEvent::StreamMetadataReceived { .. } => {
self.events.push(WatcherEvent::Metadata)
}
ClientSessionEvent::UnhandleableOnStatusCode { code } => {
self.events.push(WatcherEvent::Status(code))
}
_ => {}
},
ClientSessionResult::UnhandleableMessageReceived(_) => {}
}
}
}
fn feed(&mut self, bytes: &[u8]) {
if bytes.is_empty() {
return;
}
let results = self.session.handle_input(bytes).expect("handle_input");
self.apply(results);
}
fn pump_until(&mut self, watchdog: Duration, done: impl Fn(&Watcher) -> bool) {
let deadline = Instant::now() + watchdog;
let mut buf = [0u8; 8192];
while !done(self) && !self.eof {
assert!(
Instant::now() < deadline,
"watchdog expired; events so far: {:?}",
self.events
);
match self.stream.read(&mut buf) {
Ok(0) => self.eof = true,
Ok(n) => {
let bytes = buf[..n].to_vec();
self.feed(&bytes);
}
Err(ref e)
if e.kind() == ErrorKind::WouldBlock || e.kind() == ErrorKind::TimedOut => {}
Err(e) => panic!("watcher socket error: {e:?}"),
}
}
}
fn video_payloads(&self) -> Vec<&Bytes> {
self.events
.iter()
.filter_map(|e| match e {
WatcherEvent::Video(data) => Some(data),
_ => None,
})
.collect()
}
fn idr_count(&self) -> usize {
self.video_payloads()
.iter()
.filter(|d| is_h264_idr(d))
.count()
}
}
fn is_h264_idr(tag: &[u8]) -> bool {
if tag.len() < 5 || tag[0] != 0x17 || tag[1] != 0x01 {
return false;
}
let mut i = 5; while i + 4 <= tag.len() {
let len = u32::from_be_bytes([tag[i], tag[i + 1], tag[i + 2], tag[i + 3]]) as usize;
i += 4;
match i.checked_add(len) {
Some(end) if len != 0 && end <= tag.len() => {
if tag[i] & 0x1f == 5 {
return true;
}
i = end;
}
_ => break,
}
}
false
}
#[test]
fn is_h264_idr_rejects_malformed_avcc_without_panicking() {
assert!(!is_h264_idr(b""));
assert!(!is_h264_idr(&[0x27, 0x01, 0, 0, 0]));
assert!(!is_h264_idr(&[
0x17, 0x01, 0, 0, 0, 0xFF, 0xFF, 0xFF, 0xFF, 0x65
]));
assert!(!is_h264_idr(&[
0x17, 0x01, 0, 0, 0, 0x00, 0x00, 0x00, 0x64, 0x65
]));
assert!(is_h264_idr(&[
0x17, 0x01, 0, 0, 0, 0x00, 0x00, 0x00, 0x01, 0x65
]));
assert!(!is_h264_idr(&[
0x17, 0x01, 0, 0, 0, 0x00, 0x00, 0x00, 0x01, 0x41
]));
}
fn read_some(stream: &mut TcpStream, buf: &mut [u8], deadline: Instant) -> usize {
loop {
assert!(Instant::now() < deadline, "watchdog expired in read");
match stream.read(buf) {
Ok(n) => return n,
Err(ref e) if e.kind() == ErrorKind::WouldBlock || e.kind() == ErrorKind::TimedOut => {}
Err(e) => panic!("watcher socket error: {e:?}"),
}
}
}
#[test]
fn late_joiner_gets_headers_then_idr_first() {
let server = EmbedRtmpServer::new_with_gop_limit("127.0.0.1:0", 2)
.start()
.expect("server start");
let addr = server.local_addr().expect("bound address");
let output = server.create_rtmp_input("app", "live").expect("rtmp input");
let scheduler = FfmpegContext::builder()
.input(Input::from("test.mp4").set_stream_loop(-1))
.output(output)
.build()
.expect("context")
.start()
.expect("ffmpeg start");
let mut early = Watcher::connect(addr, "app", "live", WATCHDOG);
early.pump_until(WATCHDOG, |w| w.idr_count() >= 2);
assert!(early.idr_count() >= 2, "publisher never produced two IDRs");
let mut late = Watcher::connect(addr, "app", "live", WATCHDOG);
late.pump_until(WATCHDOG, |w| w.video_payloads().len() >= 2);
let videos = late.video_payloads();
assert!(videos.len() >= 2, "late joiner received too little video");
assert!(
videos[0].len() >= 2 && videos[0][0] == 0x17 && videos[0][1] == 0x00,
"first video tag must be the AVC sequence header, got {:02x?}",
&videos[0][..videos[0].len().min(2)]
);
let first_nalu = videos
.iter()
.find(|d| !(d.len() >= 2 && d[0] == 0x17 && d[1] == 0x00))
.expect("a video NALU after the sequence header");
assert!(
is_h264_idr(first_nalu),
"the first NALU after the sequence header must be a real IDR (AVCC NAL type 5): {:02x?}",
&first_nalu[..first_nalu.len().min(8)]
);
scheduler.abort();
server.stop();
}
#[test]
fn publisher_finish_delivers_stream_eof_to_watcher() {
let server = EmbedRtmpServer::new_with_gop_limit("127.0.0.1:0", 2)
.start()
.expect("server start");
let addr = server.local_addr().expect("bound address");
let output = server.create_rtmp_input("app", "live").expect("rtmp input");
let mut watcher = Watcher::connect(addr, "app", "live", WATCHDOG);
let scheduler = FfmpegContext::builder()
.input(Input::from("test.mp4"))
.output(output)
.build()
.expect("context")
.start()
.expect("ffmpeg start");
scheduler.wait().expect("publish completes");
watcher.pump_until(WATCHDOG, |w| w.eof);
assert!(watcher.eof, "watcher never reached EOF");
assert!(
watcher
.events
.iter()
.any(|e| matches!(e, WatcherEvent::Status(code) if code == "NetStream.Play.Complete")),
"the play-complete status must arrive before the close; events: {:?}",
watcher
.events
.iter()
.map(|e| match e {
WatcherEvent::Video(_) => "video",
WatcherEvent::Audio => "audio",
WatcherEvent::Metadata => "metadata",
WatcherEvent::Status(code) => code.as_str(),
})
.collect::<Vec<_>>()
);
assert!(
!watcher.video_payloads().is_empty(),
"media must have flowed before the finish"
);
server.stop();
}
#[test]
fn stream_builder_session_releases_port() {
let handle = EmbedRtmpServer::stream_builder()
.address("127.0.0.1:0")
.app_name("app")
.stream_key("live")
.input_file("test.mp4")
.readrate(64.0)
.gop_limit(2)
.start()
.expect("stream builder start");
let addr = handle.local_addr().expect("server bound address");
handle.wait().expect("stream completes");
let deadline = Instant::now() + WATCHDOG;
loop {
match std::net::TcpListener::bind(addr) {
Ok(_) => break,
Err(e) => {
assert!(
Instant::now() < deadline,
"port not released after StreamHandle wait+drop: {e:?}"
);
std::thread::sleep(Duration::from_millis(20));
}
}
}
}
use ez_ffmpeg::rtmp::embed_rtmp_server::RtmpStreamSender;
use rml_rtmp::chunk_io::ChunkSerializer;
use rml_rtmp::messages::RtmpMessage;
use rml_rtmp::time::RtmpTimestamp;
use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
use std::sync::{Arc, OnceLock};
fn epoch() -> Instant {
static EPOCH: OnceLock<Instant> = OnceLock::new();
*EPOCH.get_or_init(Instant::now)
}
fn nanos_since_epoch() -> u64 {
epoch().elapsed().as_nanos() as u64
}
const LOAD_VIDEO_FPS: u64 = 30;
const LOAD_KEY_INTERVAL: u64 = 30; const LOAD_KEY_BYTES: usize = 64 * 1024;
const LOAD_DELTA_BYTES: usize = 16 * 1024;
const LOAD_AUDIO_FPS: u64 = 43;
const LOAD_AUDIO_BYTES: usize = 512;
fn video_tag(keyframe: bool, seq: u32, total_len: usize) -> Vec<u8> {
let mut v = vec![0u8; total_len.max(22)];
v[0] = if keyframe { 0x17 } else { 0x27 };
v[1] = 0x01;
let nal_len = (v.len() - 9) as u32;
v[5..9].copy_from_slice(&nal_len.to_be_bytes());
v[9] = if keyframe { 0x65 } else { 0x41 }; v[10..14].copy_from_slice(&seq.to_be_bytes());
v[14..22].copy_from_slice(&nanos_since_epoch().to_be_bytes());
v
}
fn audio_tag(seq: u32, total_len: usize) -> Vec<u8> {
let mut v = vec![0u8; total_len.max(14)];
v[0] = 0xAF;
v[1] = 0x01;
v[2..6].copy_from_slice(&seq.to_be_bytes());
v[6..14].copy_from_slice(&nanos_since_epoch().to_be_bytes());
v
}
fn video_sequence_header_tag() -> Vec<u8> {
let mut v = vec![0u8; 46];
v[0] = 0x17;
v[1] = 0x00;
v[5] = 0x01; v
}
fn audio_sequence_header_tag() -> Vec<u8> {
vec![0xAF, 0x00, 0x12, 0x10] }
fn parse_video_stamp(data: &[u8]) -> Option<(u32, u64)> {
if data.len() < 22 || data[1] != 0x01 {
return None; }
let seq = u32::from_be_bytes(data[10..14].try_into().ok()?);
let stamp = u64::from_be_bytes(data[14..22].try_into().ok()?);
Some((seq, stamp))
}
fn parse_audio_stamp(data: &[u8]) -> Option<(u32, u64)> {
if data.len() < 14 || data[1] != 0x01 {
return None;
}
let seq = u32::from_be_bytes(data[2..6].try_into().ok()?);
let stamp = u64::from_be_bytes(data[6..14].try_into().ok()?);
Some((seq, stamp))
}
struct PublisherReport {
video_sent: u64,
audio_sent: u64,
stall: Duration,
}
fn publish_tag(
sender: &RtmpStreamSender,
serializer: &mut ChunkSerializer,
is_video: bool,
body: Vec<u8>,
ts_ms: u32,
) -> Duration {
let data = Bytes::from(body);
let msg = if is_video {
RtmpMessage::VideoData { data }
} else {
RtmpMessage::AudioData { data }
};
let payload = msg
.into_message_payload(RtmpTimestamp::new(ts_ms), 1)
.expect("payload conversion");
let packet = serializer
.serialize(&payload, false, true)
.expect("publisher serialize");
let t0 = Instant::now();
sender.send(packet.bytes).expect("publisher send");
t0.elapsed()
}
fn run_publisher(sender: RtmpStreamSender, stop: Arc<AtomicBool>) -> PublisherReport {
let mut serializer = ChunkSerializer::new();
let announce = serializer
.set_max_chunk_size(4096, RtmpTimestamp::new(0))
.expect("publisher chunk size");
sender.send(announce.bytes).expect("send SetChunkSize");
for (is_video, body) in [
(true, video_sequence_header_tag()),
(false, audio_sequence_header_tag()),
] {
let data = Bytes::from(body);
let msg = if is_video {
RtmpMessage::VideoData { data }
} else {
RtmpMessage::AudioData { data }
};
let payload = msg
.into_message_payload(RtmpTimestamp::new(0), 1)
.expect("payload conversion");
let packet = serializer
.serialize(&payload, false, false)
.expect("publisher serialize");
sender.send(packet.bytes).expect("send sequence header");
}
let start = Instant::now();
let video_interval = Duration::from_nanos(1_000_000_000 / LOAD_VIDEO_FPS);
let audio_interval = Duration::from_nanos(1_000_000_000 / LOAD_AUDIO_FPS);
let mut next_video = start;
let mut next_audio = start;
let mut video_seq = 0u32;
let mut audio_seq = 0u32;
let mut stall = Duration::ZERO;
while !stop.load(Ordering::Relaxed) {
let now = Instant::now();
if now >= next_video {
let keyframe = (video_seq as u64).is_multiple_of(LOAD_KEY_INTERVAL);
let size = if keyframe {
LOAD_KEY_BYTES
} else {
LOAD_DELTA_BYTES
};
let ts_ms = start.elapsed().as_millis() as u32;
stall += publish_tag(
&sender,
&mut serializer,
true,
video_tag(keyframe, video_seq, size),
ts_ms,
);
video_seq = video_seq.wrapping_add(1);
next_video += video_interval;
continue;
}
if now >= next_audio {
let ts_ms = start.elapsed().as_millis() as u32;
stall += publish_tag(
&sender,
&mut serializer,
false,
audio_tag(audio_seq, LOAD_AUDIO_BYTES),
ts_ms,
);
audio_seq = audio_seq.wrapping_add(1);
next_audio += audio_interval;
continue;
}
let next = next_video.min(next_audio);
std::thread::sleep(next.saturating_duration_since(now).min(Duration::from_millis(5)));
}
PublisherReport {
video_sent: video_seq as u64,
audio_sent: audio_seq as u64,
stall,
}
}
#[derive(Clone, Copy, PartialEq)]
enum ReaderKind {
Fast,
Slow,
}
const LOAD_WIRE_BYTES_PER_SEC: u64 = 560_000;
fn slow_reader_stall() -> Duration {
#[cfg(target_os = "linux")]
fn max_socket_buffers() -> u64 {
fn third_field(path: &str) -> Option<u64> {
std::fs::read_to_string(path)
.ok()?
.split_whitespace()
.nth(2)?
.parse()
.ok()
}
let wmem = third_field("/proc/sys/net/ipv4/tcp_wmem").unwrap_or(4 * 1024 * 1024);
let rmem = third_field("/proc/sys/net/ipv4/tcp_rmem").unwrap_or(6 * 1024 * 1024);
wmem + rmem
}
#[cfg(not(target_os = "linux"))]
fn max_socket_buffers() -> u64 {
16 * 1024 * 1024
}
let stall_bytes = max_socket_buffers() + 2 * 1024 * 1024;
let secs = stall_bytes as f64 / LOAD_WIRE_BYTES_PER_SEC as f64 + 2.0;
Duration::from_secs_f64(secs.clamp(8.0, 30.0))
}
#[derive(Default)]
struct ClassWindow {
received: u64,
first_seq: Option<u32>,
last_seq: Option<u32>,
}
impl ClassWindow {
fn record(&mut self, seq: u32) {
self.received += 1;
if self.first_seq.is_none() {
self.first_seq = Some(seq);
}
self.last_seq = Some(seq);
}
fn gap_drops(&self) -> u64 {
match (self.first_seq, self.last_seq) {
(Some(first), Some(last)) => {
let span = last.wrapping_sub(first) as u64 + 1;
span.saturating_sub(self.received)
}
_ => 0,
}
}
}
struct SubscriberReport {
kind: ReaderKind,
video: ClassWindow,
audio: ClassWindow,
latencies: Vec<u64>,
peak_lag_ns: u64,
unexpected_eof: bool,
}
fn run_subscriber(
addr: SocketAddr,
kind: ReaderKind,
slow_stall: Duration,
stop: Arc<AtomicBool>,
recording: Arc<AtomicBool>,
connected: Arc<AtomicUsize>,
watchdog: Duration,
) -> SubscriberReport {
let mut watcher = Watcher::connect(addr, "app", "live", watchdog);
connected.fetch_add(1, Ordering::SeqCst);
watcher.events.clear();
let mut report = SubscriberReport {
kind,
video: ClassWindow::default(),
audio: ClassWindow::default(),
latencies: Vec::with_capacity(4096),
peak_lag_ns: 0,
unexpected_eof: false,
};
if kind == ReaderKind::Slow {
let stall_deadline = Instant::now() + slow_stall;
while Instant::now() < stall_deadline && !stop.load(Ordering::Relaxed) {
std::thread::sleep(Duration::from_millis(250));
}
}
let mut buf = vec![0u8; 64 * 1024];
while !stop.load(Ordering::Relaxed) {
let n = match watcher.stream.read(&mut buf) {
Ok(0) => {
report.unexpected_eof = !stop.load(Ordering::Relaxed);
break;
}
Ok(n) => n,
Err(ref e) if e.kind() == ErrorKind::WouldBlock || e.kind() == ErrorKind::TimedOut => {
continue;
}
Err(e) => panic!("subscriber socket error: {e:?}"),
};
let results = watcher.session.handle_input(&buf[..n]).expect("handle_input");
let now_ns = nanos_since_epoch();
let is_recording = recording.load(Ordering::Relaxed);
for result in results {
match result {
ClientSessionResult::OutboundResponse(packet) => {
watcher
.stream
.write_all(&packet.bytes)
.expect("subscriber send to server");
}
ClientSessionResult::RaisedEvent(event) => match event {
ClientSessionEvent::VideoDataReceived { data, .. } => {
if !is_recording {
continue;
}
if let Some((seq, stamp)) = parse_video_stamp(&data) {
report.video.record(seq);
let age = now_ns.saturating_sub(stamp);
match kind {
ReaderKind::Fast => report.latencies.push(age),
ReaderKind::Slow => {
report.peak_lag_ns = report.peak_lag_ns.max(age)
}
}
}
}
ClientSessionEvent::AudioDataReceived { data, .. } => {
if !is_recording {
continue;
}
if let Some((seq, stamp)) = parse_audio_stamp(&data) {
report.audio.record(seq);
let age = now_ns.saturating_sub(stamp);
match kind {
ReaderKind::Fast => report.latencies.push(age),
ReaderKind::Slow => {
report.peak_lag_ns = report.peak_lag_ns.max(age)
}
}
}
}
_ => {}
},
ClientSessionResult::UnhandleableMessageReceived(_) => {}
}
}
}
report
}
#[derive(Clone, Copy, Default)]
struct ReactorSnapshot {
runtime_ns: u64,
voluntary_switches: u64,
involuntary_switches: u64,
}
#[cfg(target_os = "linux")]
fn reactor_tid() -> Option<u64> {
let tasks = std::fs::read_dir("/proc/self/task").ok()?;
for task in tasks.flatten() {
let Ok(tid) = task.file_name().to_string_lossy().parse::<u64>() else {
continue;
};
let comm = std::fs::read_to_string(task.path().join("comm")).unwrap_or_default();
if comm.trim() == "rtmp-server-wor" || comm.trim() == "rtmp-server-worker" {
return Some(tid);
}
}
None
}
#[cfg(not(target_os = "linux"))]
fn reactor_tid() -> Option<u64> {
None
}
#[cfg(target_os = "linux")]
fn reactor_snapshot(tid: u64) -> ReactorSnapshot {
let mut snapshot = ReactorSnapshot::default();
let base = format!("/proc/self/task/{tid}");
if let Ok(sched) = std::fs::read_to_string(format!("{base}/schedstat")) {
let mut fields = sched.split_whitespace();
snapshot.runtime_ns = fields.next().and_then(|f| f.parse().ok()).unwrap_or(0);
}
if let Ok(status) = std::fs::read_to_string(format!("{base}/status")) {
for line in status.lines() {
let mut kv = line.split_whitespace();
match kv.next() {
Some("voluntary_ctxt_switches:") => {
snapshot.voluntary_switches =
kv.next().and_then(|v| v.parse().ok()).unwrap_or(0);
}
Some("nonvoluntary_ctxt_switches:") => {
snapshot.involuntary_switches =
kv.next().and_then(|v| v.parse().ok()).unwrap_or(0);
}
_ => {}
}
}
}
snapshot
}
#[cfg(not(target_os = "linux"))]
fn reactor_snapshot(_tid: u64) -> ReactorSnapshot {
ReactorSnapshot::default()
}
struct LoadScenario {
name: &'static str,
watchers: usize,
slow_watchers: usize,
warmup: Duration,
window: Duration,
slow_stall: Duration,
}
fn percentile(sorted: &[u64], p: f64) -> u64 {
if sorted.is_empty() {
return 0;
}
let rank = ((sorted.len() as f64 - 1.0) * p).round() as usize;
sorted[rank.min(sorted.len() - 1)]
}
fn run_load_scenario(scenario: LoadScenario) {
let watchdog = Duration::from_secs(120);
let server = EmbedRtmpServer::new_with_gop_limit("127.0.0.1:0", 2)
.start()
.expect("server start");
let addr = server.local_addr().expect("bound address");
let sender = server
.create_stream_sender("app", "live")
.expect("stream sender");
let tid = reactor_tid();
println!("load_report,{},reactor_tid,{:?}", scenario.name, tid);
if let Some(tid) = tid {
println!(
"load_report,{},perf_hint,perf stat -t {tid} \
-e cycles:u,cycles:k,instructions:u,instructions:k,task-clock \
-e context-switches,cpu-migrations,syscalls:sys_enter_writev \
-- sleep {}",
scenario.name,
scenario.window.as_secs()
);
}
let stop = Arc::new(AtomicBool::new(false));
let recording = Arc::new(AtomicBool::new(false));
let connected = Arc::new(AtomicUsize::new(0));
let publisher = {
let stop = stop.clone();
std::thread::Builder::new()
.name("load-publisher".into())
.spawn(move || run_publisher(sender, stop))
.expect("spawn publisher")
};
let mut subscribers = Vec::with_capacity(scenario.watchers);
for i in 0..scenario.watchers {
let kind = if i < scenario.slow_watchers {
ReaderKind::Slow
} else {
ReaderKind::Fast
};
let stop = stop.clone();
let recording = recording.clone();
let connected = connected.clone();
let slow_stall = scenario.slow_stall;
let stagger = Duration::from_millis((i as u64) * if scenario.watchers >= 500 { 2 } else { 1 });
subscribers.push(
std::thread::Builder::new()
.name(format!("load-sub-{i}"))
.spawn(move || {
std::thread::sleep(stagger);
run_subscriber(addr, kind, slow_stall, stop, recording, connected, watchdog)
})
.expect("spawn subscriber"),
);
}
let connect_deadline = Instant::now() + watchdog;
while connected.load(Ordering::SeqCst) < scenario.watchers {
assert!(
Instant::now() < connect_deadline,
"subscribers stuck connecting: {}/{}",
connected.load(Ordering::SeqCst),
scenario.watchers
);
std::thread::sleep(Duration::from_millis(20));
}
std::thread::sleep(scenario.warmup);
let cpu_before = tid.map(reactor_snapshot);
recording.store(true, Ordering::SeqCst);
let window_start = Instant::now();
std::thread::sleep(scenario.window);
recording.store(false, Ordering::SeqCst);
let wall = window_start.elapsed();
let cpu_after = tid.map(reactor_snapshot);
stop.store(true, Ordering::SeqCst);
let publisher_report = publisher.join().expect("publisher thread");
let mut reports = Vec::with_capacity(scenario.watchers);
for handle in subscribers {
reports.push(handle.join().expect("subscriber thread"));
}
server.stop();
let name = scenario.name;
println!(
"load_report,{name},population,watchers={} slow={} window_secs={:.1}",
scenario.watchers,
scenario.slow_watchers,
wall.as_secs_f64()
);
println!(
"load_report,{name},publisher,video_sent={} audio_sent={} stall_ms={:.1}",
publisher_report.video_sent,
publisher_report.audio_sent,
publisher_report.stall.as_secs_f64() * 1e3
);
if let (Some(before), Some(after)) = (cpu_before, cpu_after) {
let cpu_ns = after.runtime_ns.saturating_sub(before.runtime_ns);
println!(
"load_report,{name},reactor_cpu,busy_ms={:.1} pct_core={:.1} vol_switch={} invol_switch={}",
cpu_ns as f64 / 1e6,
cpu_ns as f64 / wall.as_nanos() as f64 * 100.0,
after.voluntary_switches.saturating_sub(before.voluntary_switches),
after
.involuntary_switches
.saturating_sub(before.involuntary_switches),
);
}
let mut all_latencies: Vec<u64> = Vec::new();
let mut worst_fast_p99 = 0u64;
let mut fast_video_rx = 0u64;
let mut fast_audio_rx = 0u64;
let mut fast_video_drops = 0u64;
let mut fast_audio_drops = 0u64;
let mut slow_video_rx = 0u64;
let mut slow_audio_rx = 0u64;
let mut slow_video_drops = 0u64;
let mut slow_audio_drops = 0u64;
let mut slow_peak_lag_ns = 0u64;
let mut eofs = 0u64;
for report in &mut reports {
if report.unexpected_eof {
eofs += 1;
}
match report.kind {
ReaderKind::Fast => {
fast_video_rx += report.video.received;
fast_audio_rx += report.audio.received;
fast_video_drops += report.video.gap_drops();
fast_audio_drops += report.audio.gap_drops();
report.latencies.sort_unstable();
worst_fast_p99 = worst_fast_p99.max(percentile(&report.latencies, 0.99));
all_latencies.append(&mut report.latencies);
}
ReaderKind::Slow => {
slow_video_rx += report.video.received;
slow_audio_rx += report.audio.received;
slow_video_drops += report.video.gap_drops();
slow_audio_drops += report.audio.gap_drops();
slow_peak_lag_ns = slow_peak_lag_ns.max(report.peak_lag_ns);
}
}
}
all_latencies.sort_unstable();
println!(
"load_report,{name},fast_readers,video_rx={fast_video_rx} audio_rx={fast_audio_rx} video_gap_drops={fast_video_drops} audio_gap_drops={fast_audio_drops}"
);
println!(
"load_report,{name},slow_readers,video_rx={slow_video_rx} audio_rx={slow_audio_rx} video_gap_drops={slow_video_drops} audio_gap_drops={slow_audio_drops} peak_lag_ms={:.0}",
slow_peak_lag_ns as f64 / 1e6
);
println!(
"load_report,{name},glass_to_glass_ms,p50={:.2} p99={:.2} max={:.2} worst_watcher_p99={:.2} samples={}",
percentile(&all_latencies, 0.50) as f64 / 1e6,
percentile(&all_latencies, 0.99) as f64 / 1e6,
all_latencies.last().copied().unwrap_or(0) as f64 / 1e6,
worst_fast_p99 as f64 / 1e6,
all_latencies.len()
);
println!("load_report,{name},unexpected_eof,{eofs}");
if scenario.watchers > scenario.slow_watchers {
assert!(
fast_video_rx > 0,
"no video reached the fast readers — the harness itself is broken"
);
assert!(
fast_audio_rx > 0,
"no audio reached the fast readers — the harness itself is broken"
);
assert_eq!(
(fast_video_drops, fast_audio_drops),
(0, 0),
"fast readers must not be shed (video/audio gap drops)"
);
}
if scenario.slow_watchers > 0 {
assert!(
slow_video_rx > 0,
"no video reached the slow readers — the harness itself is broken"
);
assert!(
slow_video_drops > 0,
"stalled readers observed zero gaps — shedding never engaged"
);
}
assert_eq!(eofs, 0, "a subscriber was disconnected mid-run");
}
#[test]
#[ignore]
fn bench_rtmp_load_fast_w10() {
run_load_scenario(LoadScenario {
name: "fast_w10",
watchers: 10,
slow_watchers: 0,
warmup: Duration::from_secs(2),
window: Duration::from_secs(10),
slow_stall: Duration::ZERO,
});
}
#[test]
#[ignore]
fn bench_rtmp_load_fast_w100() {
run_load_scenario(LoadScenario {
name: "fast_w100",
watchers: 100,
slow_watchers: 0,
warmup: Duration::from_secs(2),
window: Duration::from_secs(10),
slow_stall: Duration::ZERO,
});
}
#[test]
#[ignore]
fn bench_rtmp_load_slow_w10() {
let stall = slow_reader_stall();
run_load_scenario(LoadScenario {
name: "slow_w10",
watchers: 10,
slow_watchers: 10,
warmup: Duration::from_secs(2),
window: stall + Duration::from_secs(8),
slow_stall: stall,
});
}
#[test]
#[ignore]
fn bench_rtmp_load_mixed_w100() {
let stall = slow_reader_stall();
run_load_scenario(LoadScenario {
name: "mixed_w100",
watchers: 100,
slow_watchers: 20,
warmup: Duration::from_secs(2),
window: stall + Duration::from_secs(8),
slow_stall: stall,
});
}
#[test]
#[ignore]
fn bench_rtmp_load_fast_w1000() {
run_load_scenario(LoadScenario {
name: "fast_w1000",
watchers: 1000,
slow_watchers: 0,
warmup: Duration::from_secs(4),
window: Duration::from_secs(10),
slow_stall: Duration::ZERO,
});
}