use clap::Parser;
use openrtc::media::{
EncodedMediaChunk, MediaCodec, MediaKind, MediaProtocolError, MediaPublicationConfig,
MediaQueue, PortableMediaSession, PublicationId,
};
use serde_json::json;
use std::collections::{HashMap, VecDeque};
use std::time::{Duration, Instant};
const TICK_MS: u64 = 20;
const MAX_LATENCY_SAMPLES: usize = 100_000;
const KEYFRAME_INTERVAL: u64 = 30;
const MAX_QUEUE_ITEMS: usize = 8;
const MAX_QUEUE_BYTES: usize = 8 * 1024;
const MAX_QUEUE_AGE_MS: u64 = 1_000;
const VIDEO_QUEUE_HOLD_TICKS: u64 = 12;
const MAX_VIDEO_RECOVERY_US: u64 = 1_000_000;
#[derive(Debug, Parser)]
#[command(name = "media-endurance-harness")]
struct Options {
#[arg(long, default_value_t = 180)]
duration_seconds: u64,
#[arg(long, default_value = "stable")]
fault_profile: String,
#[arg(long, default_value_t = 10)]
checkpoint_seconds: u64,
#[arg(long, default_value_t = false)]
accelerate_faults: bool,
}
#[derive(Debug, Clone, Copy)]
struct Publication {
id: PublicationId,
kind: MediaKind,
codec: MediaCodec,
}
impl Publication {
fn config(self, media_generation: u32) -> MediaPublicationConfig {
let video = self.kind == MediaKind::Video;
MediaPublicationConfig {
publication_id: self.id,
media_generation,
kind: self.kind,
codec: self.codec,
clock_rate: if video { 90_000 } else { 48_000 },
coded_width: video.then_some(640),
coded_height: video.then_some(360),
channels: (!video).then_some(2),
}
}
}
#[derive(Debug, Default)]
struct PlaybackState {
decoded: u64,
rendered: u64,
last_rendered_sequence: Option<u64>,
last_audio_sequence: Option<u64>,
audio_gaps: u64,
awaiting_keyframe_since_us: Option<u64>,
max_recovery_us: u64,
recovery_events: u64,
recovery_bound_violations: u64,
}
fn oracle_payload(
publication: Publication,
index: usize,
media_generation: u32,
sequence: u64,
timestamp_us: u64,
) -> Vec<u8> {
let mut payload = Vec::with_capacity(43);
payload.extend_from_slice(&publication.id.0);
payload.extend_from_slice(&media_generation.to_be_bytes());
payload.extend_from_slice(&sequence.to_be_bytes());
payload.extend_from_slice(×tamp_us.to_be_bytes());
payload.push(index as u8);
payload.push(if publication.kind == MediaKind::Audio {
1
} else {
2
});
payload.push(match publication.codec {
MediaCodec::Opus => 1,
MediaCodec::Vp8 => 3,
_ => 0,
});
while payload.len() < 900 {
payload.push(
(sequence as usize)
.wrapping_add(payload.len())
.wrapping_add(index) as u8,
);
}
payload
}
fn mark_video_discontinuity(state: &mut PlaybackState, timestamp_us: u64) {
state.awaiting_keyframe_since_us.get_or_insert(timestamp_us);
}
fn render_chunk(
publication: Publication,
chunk: EncodedMediaChunk,
now_us: u64,
state: &mut PlaybackState,
) {
if publication.kind == MediaKind::Audio {
if let Some(previous) = state.last_audio_sequence {
if chunk.sequence != previous.saturating_add(1) {
state.audio_gaps = state.audio_gaps.saturating_add(1);
}
}
state.last_audio_sequence = Some(chunk.sequence);
}
if publication.kind == MediaKind::Video {
if let Some(started) = state.awaiting_keyframe_since_us {
if !chunk.keyframe {
return;
}
state.awaiting_keyframe_since_us = None;
state.max_recovery_us = state.max_recovery_us.max(now_us.saturating_sub(started));
state.recovery_events = state.recovery_events.saturating_add(1);
if now_us.saturating_sub(started) > MAX_VIDEO_RECOVERY_US {
state.recovery_bound_violations = state.recovery_bound_violations.saturating_add(1);
}
}
}
state.decoded = state.decoded.saturating_add(1);
state.rendered = state.rendered.saturating_add(1);
state.last_rendered_sequence = Some(chunk.sequence);
}
#[tokio::main]
async fn main() -> anyhow::Result<()> {
let options = Options::parse();
if options.duration_seconds == 0 || options.checkpoint_seconds == 0 {
anyhow::bail!("duration and checkpoint cadence must be positive");
}
if options.fault_profile != "stable" && options.fault_profile != "churn" {
anyhow::bail!("fault profile must be stable or churn");
}
let publications = [
Publication {
id: PublicationId([1; 16]),
kind: MediaKind::Audio,
codec: MediaCodec::Opus,
},
Publication {
id: PublicationId([2; 16]),
kind: MediaKind::Video,
codec: MediaCodec::Vp8,
},
Publication {
id: PublicationId([3; 16]),
kind: MediaKind::Audio,
codec: MediaCodec::Opus,
},
Publication {
id: PublicationId([4; 16]),
kind: MediaKind::Video,
codec: MediaCodec::Vp8,
},
];
let mut generations: HashMap<PublicationId, u32> = publications
.iter()
.map(|publication| (publication.id, 1))
.collect();
let mut session = PortableMediaSession::default();
for publication in publications {
let (config, control) = session.begin_publication(publication.config(1))?;
session.decode_control(&control)?;
generations.insert(publication.id, config.media_generation);
}
let mut queues: HashMap<PublicationId, MediaQueue> = publications
.iter()
.map(|publication| {
(
publication.id,
MediaQueue::new(MAX_QUEUE_ITEMS, MAX_QUEUE_BYTES, MAX_QUEUE_AGE_MS),
)
})
.collect();
let mut playback: HashMap<PublicationId, PlaybackState> = publications
.iter()
.map(|publication| (publication.id, PlaybackState::default()))
.collect();
let mut last_sequences: HashMap<PublicationId, u64> = HashMap::new();
let mut latency_samples_us = VecDeque::with_capacity(MAX_LATENCY_SAMPLES);
let mut accepted = 0_u64;
let mut decoded = 0_u64;
let mut rendered = 0_u64;
let mut payload_mismatches = 0_u64;
let mut audio_gaps = 0_u64;
let mut video_recovery_bound_violations = 0_u64;
let mut injected_loss = 0_u64;
let mut rejected_duplicates = 0_u64;
let mut rejected_corruption = 0_u64;
let mut route_switches = 0_u64;
let mut physical_replacements = 0_u64;
let started = Instant::now();
let deadline = started + Duration::from_secs(options.duration_seconds);
let checkpoint = Duration::from_secs(options.checkpoint_seconds);
let mut next_checkpoint = started + checkpoint;
let mut ticker = tokio::time::interval(Duration::from_millis(TICK_MS));
ticker.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
let mut sequence = 0_u64;
let (
route_interval,
replacement_interval,
loss_interval,
duplicate_interval,
corruption_interval,
) = if options.accelerate_faults {
(10_u64, 20, 7, 11, 13)
} else {
(900_u64, 3_600, 997, 887, 4_999)
};
let video_hold_ticks = VIDEO_QUEUE_HOLD_TICKS.min(route_interval.saturating_sub(1));
while Instant::now() < deadline {
ticker.tick().await;
let churn = options.fault_profile == "churn";
if churn && sequence > 0 && sequence % route_interval == 0 {
route_switches += 1;
}
if churn && sequence > 0 && sequence % replacement_interval == 0 {
physical_replacements += 1;
for publication in publications {
let (config, control) =
session.begin_publication(publication.config(generations[&publication.id]))?;
session.decode_control(&control)?;
generations.insert(publication.id, config.media_generation);
if publication.kind == MediaKind::Video {
mark_video_discontinuity(
playback
.get_mut(&publication.id)
.expect("publication playback"),
sequence.saturating_mul(TICK_MS * 1_000),
);
}
}
}
for (index, publication) in publications.iter().enumerate() {
if churn
&& sequence > 0
&& sequence % loss_interval == 0
&& publication.kind == MediaKind::Video
{
injected_loss += 1;
mark_video_discontinuity(
playback
.get_mut(&publication.id)
.expect("publication playback"),
sequence.saturating_mul(TICK_MS * 1_000),
);
continue;
}
let processing_started = Instant::now();
let media_generation = generations[&publication.id];
let timestamp_us = sequence.saturating_mul(TICK_MS * 1_000);
let chunk = EncodedMediaChunk {
publication_id: publication.id,
media_generation,
sequence,
timestamp_us,
duration_us: (TICK_MS * 1_000) as u32,
kind: publication.kind,
codec: publication.codec,
keyframe: publication.kind == MediaKind::Video && sequence % KEYFRAME_INTERVAL == 0,
discardable: publication.kind == MediaKind::Video
&& sequence % KEYFRAME_INTERVAL != 0,
payload: oracle_payload(
*publication,
index,
media_generation,
sequence,
timestamp_us,
),
};
let mut encoded = chunk.encode()?;
if churn && sequence > 0 && sequence % corruption_interval == 0 && index == 1 {
*encoded.last_mut().expect("encoded media payload") ^= 1;
if EncodedMediaChunk::decode(&encoded) != Err(MediaProtocolError::IntegrityFailure)
{
anyhow::bail!("injected corruption reached the replay/decoder boundary");
}
rejected_corruption += 1;
mark_video_discontinuity(
playback
.get_mut(&publication.id)
.expect("publication playback"),
timestamp_us,
);
continue;
}
let received = session.decode_chunk(&encoded)?;
if received != chunk {
payload_mismatches += 1;
anyhow::bail!(
"payload identity mismatch for publication {} sequence {} (mismatches={})",
publication.id,
sequence,
payload_mismatches
);
}
accepted += 1;
last_sequences.insert(publication.id, sequence);
if churn && sequence > 0 && sequence % duplicate_interval == 0 && index == 0 {
match session.decode_chunk(&encoded) {
Err(MediaProtocolError::Replay) => rejected_duplicates += 1,
other => anyhow::bail!("duplicate media was not rejected: {other:?}"),
}
}
queues
.get_mut(&publication.id)
.expect("publication queue")
.push(received, timestamp_us / 1_000);
if latency_samples_us.len() == MAX_LATENCY_SAMPLES {
latency_samples_us.pop_front();
}
latency_samples_us.push_back(processing_started.elapsed().as_micros() as u64);
}
sequence = sequence.saturating_add(1);
let now_us = sequence.saturating_mul(TICK_MS * 1_000);
drain_queues(
churn && sequence.saturating_sub(1) % route_interval < video_hold_ticks,
now_us,
&publications,
&mut queues,
&mut playback,
);
refresh_totals(
&playback,
&mut decoded,
&mut rendered,
&mut audio_gaps,
&mut video_recovery_bound_violations,
);
if Instant::now() >= next_checkpoint {
print_checkpoint(
"checkpoint",
started,
accepted,
decoded,
rendered,
payload_mismatches,
audio_gaps,
video_recovery_bound_violations,
injected_loss,
rejected_duplicates,
rejected_corruption,
route_switches,
physical_replacements,
&publications,
&last_sequences,
&playback,
&queues,
&latency_samples_us,
);
next_checkpoint += checkpoint;
}
}
let recovery_deadline_sequence = sequence.saturating_add(KEYFRAME_INTERVAL);
while publications.iter().any(|publication| {
publication.kind == MediaKind::Video
&& playback
.get(&publication.id)
.is_some_and(|state| state.awaiting_keyframe_since_us.is_some())
}) && sequence <= recovery_deadline_sequence
{
ticker.tick().await;
for (index, publication) in publications.iter().enumerate() {
if publication.kind != MediaKind::Video {
continue;
}
let media_generation = generations[&publication.id];
let timestamp_us = sequence.saturating_mul(TICK_MS * 1_000);
let chunk = EncodedMediaChunk {
publication_id: publication.id,
media_generation,
sequence,
timestamp_us,
duration_us: (TICK_MS * 1_000) as u32,
kind: publication.kind,
codec: publication.codec,
keyframe: sequence % KEYFRAME_INTERVAL == 0,
discardable: sequence % KEYFRAME_INTERVAL != 0,
payload: oracle_payload(
*publication,
index,
media_generation,
sequence,
timestamp_us,
),
};
let received = session.decode_chunk(&chunk.encode()?)?;
if received != chunk {
anyhow::bail!(
"payload identity mismatch during recovery for publication {} sequence {}",
publication.id,
sequence
);
}
accepted = accepted.saturating_add(1);
last_sequences.insert(publication.id, sequence);
queues
.get_mut(&publication.id)
.expect("publication queue")
.push(received, timestamp_us / 1_000);
}
sequence = sequence.saturating_add(1);
drain_queues(
false,
sequence.saturating_mul(TICK_MS * 1_000),
&publications,
&mut queues,
&mut playback,
);
refresh_totals(
&playback,
&mut decoded,
&mut rendered,
&mut audio_gaps,
&mut video_recovery_bound_violations,
);
}
drain_queues(
false,
sequence.saturating_mul(TICK_MS * 1_000),
&publications,
&mut queues,
&mut playback,
);
refresh_totals(
&playback,
&mut decoded,
&mut rendered,
&mut audio_gaps,
&mut video_recovery_bound_violations,
);
if last_sequences.len() != publications.len() {
anyhow::bail!("not all four bidirectional audio/video publications advanced");
}
if options.fault_profile == "churn"
&& (injected_loss == 0
|| rejected_duplicates == 0
|| rejected_corruption == 0
|| route_switches == 0
|| physical_replacements == 0)
{
anyhow::bail!("churn run ended before every fault class was exercised");
}
if payload_mismatches > 0
|| audio_gaps > 0
|| video_recovery_bound_violations > 0
|| decoded == 0
|| rendered == 0
|| decoded != rendered
{
anyhow::bail!("media playback oracle invariants failed");
}
for publication in publications {
let state = playback.get(&publication.id).expect("publication playback");
let stats = queues
.get(&publication.id)
.expect("publication queue")
.stats();
if stats.queued != 0
|| stats.queued_bytes != 0
|| stats.max_queued > MAX_QUEUE_ITEMS as u64
|| stats.max_queued_bytes > MAX_QUEUE_BYTES
|| (options.fault_profile == "churn"
&& publication.kind == MediaKind::Video
&& (stats.max_queued != MAX_QUEUE_ITEMS as u64
|| stats.max_queued_bytes < MAX_QUEUE_BYTES * 3 / 4))
{
anyhow::bail!("media queue bound or plateau invariant failed");
}
if options.fault_profile == "churn"
&& publication.kind == MediaKind::Video
&& (state.awaiting_keyframe_since_us.is_some() || state.recovery_events == 0)
{
anyhow::bail!("video did not recover to a keyframe");
}
}
print_checkpoint(
"complete",
started,
accepted,
decoded,
rendered,
payload_mismatches,
audio_gaps,
video_recovery_bound_violations,
injected_loss,
rejected_duplicates,
rejected_corruption,
route_switches,
physical_replacements,
&publications,
&last_sequences,
&playback,
&queues,
&latency_samples_us,
);
Ok(())
}
fn drain_queues(
hold_video: bool,
now_us: u64,
publications: &[Publication; 4],
queues: &mut HashMap<PublicationId, MediaQueue>,
playback: &mut HashMap<PublicationId, PlaybackState>,
) {
for publication in publications {
if hold_video && publication.kind == MediaKind::Video {
continue;
}
let queue = queues.get_mut(&publication.id).expect("publication queue");
while let Some(chunk) = queue.pop(now_us / 1_000) {
render_chunk(
*publication,
chunk,
now_us,
playback
.get_mut(&publication.id)
.expect("publication playback"),
);
}
}
}
fn refresh_totals(
playback: &HashMap<PublicationId, PlaybackState>,
decoded: &mut u64,
rendered: &mut u64,
audio_gaps: &mut u64,
recovery_bound_violations: &mut u64,
) {
*decoded = playback.values().map(|state| state.decoded).sum();
*rendered = playback.values().map(|state| state.rendered).sum();
*audio_gaps = playback.values().map(|state| state.audio_gaps).sum();
*recovery_bound_violations = playback
.values()
.map(|state| state.recovery_bound_violations)
.sum();
}
#[allow(clippy::too_many_arguments)]
fn print_checkpoint(
event: &str,
started: Instant,
accepted: u64,
decoded: u64,
rendered: u64,
payload_mismatches: u64,
audio_gaps: u64,
video_recovery_bound_violations: u64,
injected_loss: u64,
rejected_duplicates: u64,
rejected_corruption: u64,
route_switches: u64,
physical_replacements: u64,
publications: &[Publication; 4],
last_sequences: &HashMap<PublicationId, u64>,
playback: &HashMap<PublicationId, PlaybackState>,
queues: &HashMap<PublicationId, MediaQueue>,
latency_samples_us: &VecDeque<u64>,
) {
let mut latencies: Vec<u64> = latency_samples_us.iter().copied().collect();
latencies.sort_unstable();
let p95 = latencies
.get(latencies.len().saturating_mul(95) / 100)
.copied()
.unwrap_or(0);
let publication_progress: Vec<_> = publications
.iter()
.map(|publication| {
let state = playback.get(&publication.id).expect("publication playback");
let stats = queues
.get(&publication.id)
.expect("publication queue")
.stats();
json!({
"publicationId": publication.id.to_string(),
"sequence": last_sequences.get(&publication.id).copied().unwrap_or(0),
"decoded": state.decoded,
"rendered": state.rendered,
"audioGaps": state.audio_gaps,
"videoRecoveryEvents": state.recovery_events,
"maxRecoveryUs": state.max_recovery_us,
"queueCurrentItems": stats.queued,
"queueMaxItems": stats.max_queued,
"queueCurrentBytes": stats.queued_bytes,
"queueMaxBytes": stats.max_queued_bytes,
"queuePlateau": stats.max_queued == MAX_QUEUE_ITEMS as u64,
"queueByteUtilization": stats.max_queued_bytes as f64 / MAX_QUEUE_BYTES as f64,
})
})
.collect();
println!(
"{}",
json!({
"event": event,
"elapsedMs": started.elapsed().as_millis(),
"accepted": accepted,
"decoded": decoded,
"rendered": rendered,
"payloadMismatches": payload_mismatches,
"audioGaps": audio_gaps,
"videoRecoveryBoundViolations": video_recovery_bound_violations,
"injectedLoss": injected_loss,
"rejectedDuplicates": rejected_duplicates,
"rejectedCorruption": rejected_corruption,
"routeSwitches": route_switches,
"physicalReplacements": physical_replacements,
"localProcessingLatencyP95Us": p95,
"latencyKind": "measured-local-protocol-processing-not-glass-to-glass",
"publications": publication_progress,
})
);
}