Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 4 additions & 0 deletions locales/en.json
Original file line number Diff line number Diff line change
Expand Up @@ -532,6 +532,10 @@
"advancedAudioL4sAndFrameRateView": "Advanced · audio, L4S and frame-rate view",
"videoFormatUnavailable": "Video format unavailable",
"streamFpsLabel": "STREAM FPS",
"videoBitrateLabel": "VIDEO BITRATE",
"streamUdpReceiveLabel": "STREAM UDP RECEIVE",
"knownSessionPeerUdpDatagramBytes": "known session peer · UDP datagram bytes",
"udpRx": "UDP RX",
"localOutputFps": "LOCAL OUTPUT FPS",
"frameGeneration": "FRAME GENERATION",
"warmingUp": "Warming up",
Expand Down
94 changes: 93 additions & 1 deletion native/opennow-streamer/crates/opennow-streamer-core/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -100,6 +100,9 @@ trait NvstSessionResources {
fn network_metrics(&self) -> Option<(f64, f64)> {
None
}
fn socket_receive_bytes(&self) -> Option<u64> {
None
}
fn frame_stage_timings(&self) -> Option<FrameStageTimings> {
None
}
Expand Down Expand Up @@ -139,6 +142,9 @@ impl NvstSessionResources for ActiveNvstResources {
fn network_metrics(&self) -> Option<(f64, f64)> {
self.feedback.recent_network_metrics(Instant::now())
}
fn socket_receive_bytes(&self) -> Option<u64> {
Some(self.feedback.socket_receive_bytes())
}
fn frame_stage_timings(&self) -> Option<FrameStageTimings> {
let timings = self.feedback.frame_stage_timings();
(!timings.is_empty()).then_some(timings)
Expand Down Expand Up @@ -1744,6 +1750,7 @@ fn forward_nvst_session_events<R: NvstSessionResources>(
transport: resources,
} = event_resources;
let mut feedback_state = NvstMediaFeedbackState::new(false);
feedback_state.previous_socket_receive_bytes = resources.socket_receive_bytes().unwrap_or(0);
feedback_state.start_id = start_id.clone();
if let Some(queue) = captured_input.as_ref() {
queue.set_text_ready(generation, false);
Expand Down Expand Up @@ -1957,6 +1964,14 @@ fn forward_nvst_session_events<R: NvstSessionResources>(
}
}
}
if feedback_state.telemetry_started
|| resources
.socket_receive_bytes()
.is_some_and(|bytes| bytes > feedback_state.previous_socket_receive_bytes)
{
feedback_state.telemetry_started = true;
flush_nvst_telemetry(output, &resources, &mut feedback_state);
}
}
if let Some(queue) = captured_input.as_ref() {
queue.set_text_ready(generation, false);
Expand Down Expand Up @@ -2331,6 +2346,8 @@ struct NvstMediaFeedbackState {
telemetry_window_started: Instant,
telemetry_frames: u64,
telemetry_bytes: u64,
telemetry_started: bool,
previous_socket_receive_bytes: u64,
peak_bitrate_mbps: f64,
decode_timings: Option<DecodeTimingsReport>,
decode_progress: DecodeProgressWatchdog,
Expand All @@ -2348,6 +2365,8 @@ impl NvstMediaFeedbackState {
telemetry_window_started: Instant::now(),
telemetry_frames: 0,
telemetry_bytes: 0,
telemetry_started: false,
previous_socket_receive_bytes: 0,
peak_bitrate_mbps: 0.0,
decode_timings: None,
decode_progress: DecodeProgressWatchdog::default(),
Expand All @@ -2369,13 +2388,23 @@ fn flush_nvst_telemetry<R: NvstSessionResources>(
let elapsed_seconds = elapsed.as_secs_f64();
let frames_per_second = state.telemetry_frames as f64 / elapsed_seconds;
let bitrate_mbps = state.telemetry_bytes as f64 * 8.0 / elapsed_seconds / 1_000_000.0;
let socket_bytes = resources.socket_receive_bytes();
let receive_bitrate_mbps = socket_bytes.and_then(|bytes| {
bytes
.checked_sub(state.previous_socket_receive_bytes)
.map(|delta| delta as f64 * 8.0 / elapsed_seconds / 1_000_000.0)
});
if let Some(bytes) = socket_bytes {
state.previous_socket_receive_bytes = bytes;
}
state.peak_bitrate_mbps = state.peak_bitrate_mbps.max(bitrate_mbps);
let network = resources.network_metrics();
let _ = output.send(event(
"telemetry",
json!({
"framesPerSecond": frames_per_second,
"bitrateMbps": bitrate_mbps,
"receiveBitrateMbps": receive_bitrate_mbps,
"peakBitrateMbps": state.peak_bitrate_mbps,
"pingMs": resources.ping_ms(),
"jitterMs": network.map(|metrics| metrics.0),
Expand Down Expand Up @@ -2427,6 +2456,7 @@ fn forward_nvst_media_feedback<R: NvstSessionResources>(
}
state.telemetry_frames = state.telemetry_frames.saturating_add(1);
state.telemetry_bytes = state.telemetry_bytes.saturating_add(u64::from(bytes));
state.telemetry_started = true;
flush_nvst_telemetry(output, resources, state);
}
MediaFeedback::PlaybackStarted { backend } => {
Expand Down Expand Up @@ -2977,7 +3007,7 @@ fn consume_encoded_media(
mod tests {
use super::*;
use std::net::UdpSocket;
use std::sync::atomic::{AtomicUsize, Ordering};
use std::sync::atomic::{AtomicU64, AtomicUsize, Ordering};
use std::time::Instant;

fn command(value: Value) -> Command {
Expand Down Expand Up @@ -3241,6 +3271,7 @@ mod tests {
struct TestNvstResources {
rumble: Arc<Mutex<[Option<NvstControllerRumble>; 4]>>,
ping_ms: Option<f64>,
socket_bytes: Option<Arc<AtomicU64>>,
frame_stage_timings: Option<FrameStageTimings>,
decode_progress_policy: Option<DecodeProgressPolicy>,
keyframe_requests: Arc<AtomicUsize>,
Expand All @@ -3259,6 +3290,12 @@ mod tests {
self.ping_ms
}

fn socket_receive_bytes(&self) -> Option<u64> {
self.socket_bytes
.as_ref()
.map(|bytes| bytes.load(Ordering::Relaxed))
}

fn frame_stage_timings(&self) -> Option<FrameStageTimings> {
self.frame_stage_timings
}
Expand Down Expand Up @@ -3573,6 +3610,61 @@ mod tests {
assert_eq!(telemetry["peakBitrateMbps"], telemetry["bitrateMbps"]);
}

#[test]
fn socket_receive_rate_tracks_cumulative_bytes_through_idle_and_reset() {
let (sender, receiver) = std::sync::mpsc::channel();
let sender = EventSender::unbounded(sender);
let bytes = Arc::new(AtomicU64::new(50_000));
let resources = TestNvstResources {
socket_bytes: Some(Arc::clone(&bytes)),
..Default::default()
};
let mut state = NvstMediaFeedbackState::new(true);
state.previous_socket_receive_bytes = resources.socket_receive_bytes().unwrap();
bytes.store(300_000, Ordering::Relaxed);
state.telemetry_window_started = Instant::now() - Duration::from_secs(2);
flush_nvst_telemetry(&sender, &resources, &mut state);
let first = receiver.recv().unwrap();
assert!((first["receiveBitrateMbps"].as_f64().unwrap() - 1.0).abs() < 0.01);
assert_eq!(first["bitrateMbps"], json!(0.0));

bytes.store(550_000, Ordering::Relaxed);
state.telemetry_window_started = Instant::now() - Duration::from_secs(1);
flush_nvst_telemetry(&sender, &resources, &mut state);
let second = receiver.recv().unwrap();
assert!((second["receiveBitrateMbps"].as_f64().unwrap() - 2.0).abs() < 0.02);

state.telemetry_window_started = Instant::now() - Duration::from_secs(1);
flush_nvst_telemetry(&sender, &resources, &mut state);
assert_eq!(receiver.recv().unwrap()["receiveBitrateMbps"], json!(0.0));

bytes.store(100, Ordering::Relaxed);
state.telemetry_window_started = Instant::now() - Duration::from_secs(1);
flush_nvst_telemetry(&sender, &resources, &mut state);
assert_eq!(receiver.recv().unwrap()["receiveBitrateMbps"], Value::Null);

bytes.store(125_100, Ordering::Relaxed);
state.telemetry_window_started = Instant::now() - Duration::from_secs(1);
flush_nvst_telemetry(&sender, &resources, &mut state);
assert!(
(receiver.recv().unwrap()["receiveBitrateMbps"]
.as_f64()
.unwrap()
- 1.0)
.abs()
< 0.01
);

let mut next_session = NvstMediaFeedbackState::new(true);
let new_resources = TestNvstResources {
socket_bytes: Some(Arc::new(AtomicU64::new(0))),
..Default::default()
};
next_session.telemetry_window_started = Instant::now() - Duration::from_secs(1);
flush_nvst_telemetry(&sender, &new_resources, &mut next_session);
assert_eq!(receiver.recv().unwrap()["receiveBitrateMbps"], json!(0.0));
}

#[test]
fn accepted_video_telemetry_preserves_measured_and_unavailable_ping() {
for ping_ms in [None, Some(0.0), Some(25.5)] {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -536,6 +536,12 @@ struct RecentNetworkMetrics {
previous: Option<(Instant, NetworkCounters)>,
}

#[derive(Debug, Clone, Copy)]
enum StreamSocket {
Bundle,
Video,
}

impl RecentNetworkMetrics {
fn sample(&mut self, now: Instant, counters: NetworkCounters) -> Option<f64> {
if counters.ssrc == 0 || counters.base == u32::MAX || counters.highest < counters.base {
Expand Down Expand Up @@ -584,6 +590,8 @@ pub struct NvstFeedbackState {
received_packets: AtomicU32,
report_prior: Mutex<(u32, u32)>,
recent_network_metrics: Mutex<RecentNetworkMetrics>,
bundle_receive_bytes: AtomicU64,
video_receive_bytes: AtomicU64,
reception_timing: Mutex<ReceptionTiming>,
ice_ping: Mutex<Option<(Instant, Duration)>>,
video_ping: Mutex<Option<(Instant, Duration)>>,
Expand All @@ -607,6 +615,8 @@ impl Default for NvstFeedbackState {
received_packets: AtomicU32::new(0),
report_prior: Mutex::new((0, 0)),
recent_network_metrics: Mutex::new(RecentNetworkMetrics::default()),
bundle_receive_bytes: AtomicU64::new(0),
video_receive_bytes: AtomicU64::new(0),
reception_timing: Mutex::new(ReceptionTiming::default()),
ice_ping: Mutex::new(None),
video_ping: Mutex::new(None),
Expand All @@ -622,6 +632,31 @@ impl Default for NvstFeedbackState {
}

impl NvstFeedbackState {
fn record_socket_receive(
&self,
socket: StreamSocket,
peer_matches: bool,
authenticated: bool,
length: usize,
) {
if !peer_matches || !authenticated {
return;
}
let counter = match socket {
StreamSocket::Bundle => &self.bundle_receive_bytes,
StreamSocket::Video => &self.video_receive_bytes,
};
let _ = counter.fetch_update(Ordering::Relaxed, Ordering::Relaxed, |bytes| {
Some(bytes.saturating_add(length as u64))
});
}

pub fn socket_receive_bytes(&self) -> u64 {
self.bundle_receive_bytes
.load(Ordering::Relaxed)
.saturating_add(self.video_receive_bytes.load(Ordering::Relaxed))
}

fn update_ice_ping(
&self,
pair: Option<&CandidatePairStats>,
Expand Down Expand Up @@ -7050,6 +7085,7 @@ fn run_nvst_webrtc_bundle(
if source != bundle_peer {
continue;
}
feedback.record_socket_receive(StreamSocket::Bundle, true, dtls_ready, length);
if let Some(credentials) = stun_credentials.as_ref() {
let received_at = Instant::now();
if let Some(elapsed) = ping_tracker.receive(
Expand Down Expand Up @@ -7294,6 +7330,13 @@ fn run_nvst_udp_receiver(
);
}
let expected_source = receiver.config.accepts_video_source(source);
let authenticated_before = receiver.last_authenticated_packet.is_some();
feedback.record_socket_receive(
StreamSocket::Video,
expected_source,
authenticated_before,
length,
);
if !expected_source {
wrong_source += 1;
}
Expand Down Expand Up @@ -7330,6 +7373,14 @@ fn run_nvst_udp_receiver(
}
let received_at = Instant::now();
let events = receiver.process_datagram(source, &datagram[..length], received_at);
if !authenticated_before {
feedback.record_socket_receive(
StreamSocket::Video,
expected_source,
receiver.last_authenticated_packet.is_some(),
length,
);
}
for event in events {
if !forward_receive_event(
&media_consumer,
Expand Down Expand Up @@ -7529,6 +7580,30 @@ fn forward_receive_event(
mod tests {
static PREFERRED_NVST_PORTS: std::sync::Mutex<()> = std::sync::Mutex::new(());

#[test]
fn stream_socket_receive_counts_peer_datagram_lengths_once_across_sockets() {
use super::{NvstFeedbackState, StreamSocket};

let feedback = NvstFeedbackState::default();
feedback.record_socket_receive(StreamSocket::Bundle, false, true, 99_000);
feedback.record_socket_receive(StreamSocket::Bundle, true, false, 88_000);
assert_eq!(feedback.socket_receive_bytes(), 0);

let bundle_control = 1_300;
let bundle_audio = 800;
let video_data = 1_400;
let video_fec = 300;
let video_retransmission = 1_400;
feedback.record_socket_receive(StreamSocket::Bundle, true, true, bundle_control);
feedback.record_socket_receive(StreamSocket::Bundle, true, true, bundle_audio);
feedback.record_socket_receive(StreamSocket::Video, true, true, video_data);
feedback.record_socket_receive(StreamSocket::Video, true, true, video_fec);
feedback.record_socket_receive(StreamSocket::Video, true, true, video_retransmission);
feedback.record_socket_receive(StreamSocket::Video, false, true, 77_000);
assert_eq!(feedback.socket_receive_bytes(), 5_200);
assert_eq!(NvstFeedbackState::default().socket_receive_bytes(), 0);
}

#[test]
fn typed_text_queue_preserves_order_reservation_and_readiness() {
use opennow_streamer_protocol::text_input::{TextInputError, TextInputSlot};
Expand Down Expand Up @@ -11564,6 +11639,7 @@ mod tests {
let mut config = NvstVideoConfig::from_legacy_handoff(&handoff, None).unwrap();
config.stun_credentials = Some(stun_credentials());
config.ping_payload = b"setup-ping".to_vec();
let feedback = config.feedback();
let packet = protect_for_test(
&test_srtp(&config),
build_plaintext_rtp(
Expand Down Expand Up @@ -11621,6 +11697,7 @@ mod tests {
assert_eq!(frame.frame_index, Some(42));
assert!(frame.keyframe);
session.stop();
assert_eq!(feedback.socket_receive_bytes(), packet.len() as u64);
}
}
}
Expand Down
Loading
Loading