Skip to content
File

Blob: firmware/vendor/str0m/tests/stats.rs

rust227 lines
1use std::net::Ipv4Addr;
2use std::time::{Duration, Instant};
3 
4use str0m::format::Codec;
5use str0m::media::{Direction, MediaKind};
6use str0m::rtp::{RtpWrite, Ssrc};
7use str0m::stats::{MediaEgressStats, PeerStats};
8use str0m::{Event, RtcConfig, RtcError};
9use tracing::info_span;
10 
11mod common;
12use common::{TestRtc, connect_l_r_with_rtc, init_crypto_default, init_log, progress};
13 
14#[test]
15pub fn stats() -> Result<(), RtcError> {
16 init_log();
17 init_crypto_default();
18 
19 let l_config = RtcConfig::new().set_stats_interval(Some(Duration::from_secs(10)));
20 let r_config = RtcConfig::new().set_stats_interval(Some(Duration::from_secs(10)));
21 
22 let now = Instant::now();
23 let mut l = TestRtc::new_with_rtc(info_span!("L"), l_config.build(now));
24 let mut r = TestRtc::new_with_rtc(info_span!("R"), r_config.build(now));
25 
26 l.add_host_candidate((Ipv4Addr::new(1, 1, 1, 1), 1000).into());
27 r.add_host_candidate((Ipv4Addr::new(2, 2, 2, 2), 2000).into());
28 
29 let mut change = l.sdp_api();
30 let mid = change.add_media(MediaKind::Audio, Direction::SendRecv, None, None, None);
31 let (offer, pending) = change.apply().unwrap();
32 
33 let answer = r.rtc.sdp_api().accept_offer(offer)?;
34 l.rtc.sdp_api().accept_answer(pending, answer)?;
35 
36 loop {
37 if l.is_connected() || r.is_connected() {
38 break;
39 }
40 progress(&mut l, &mut r)?;
41 }
42 
43 let max = l.last.max(r.last);
44 l.last = max;
45 r.last = max;
46 
47 let params = l.params_opus();
48 assert_eq!(params.spec().codec, Codec::Opus);
49 let pt = params.pt();
50 
51 let data_a = vec![1_u8; 80];
52 let data_b = vec![2_u8; 80];
53 
54 l.set_forced_time_advance(Duration::from_millis(1));
55 r.set_forced_time_advance(Duration::from_millis(1));
56 
57 loop {
58 {
59 let wallclock = l.start + l.duration();
60 let time = l.duration().into();
61 l.writer(mid)
62 .unwrap()
63 .write(pt, wallclock, time, data_a.clone())?;
64 }
65 
66 {
67 let wallclock = r.start + r.duration();
68 let time = l.duration().into();
69 r.writer(mid)
70 .unwrap()
71 .write(pt, wallclock, time, data_b.clone())?;
72 }
73 
74 progress(&mut l, &mut r)?;
75 
76 if l.duration() > Duration::from_secs(25) {
77 break;
78 }
79 }
80 
81 let media_count_r = r
82 .events
83 .iter()
84 .filter(|(_, e)| matches!(e, Event::MediaData(_)))
85 .count();
86 
87 assert!(
88 media_count_r > 170,
89 "Not enough MediaData at R: {}",
90 media_count_r
91 );
92 
93 let media_count_l = l
94 .events
95 .iter()
96 .filter(|(_, e)| matches!(e, Event::MediaData(_)))
97 .count();
98 
99 let egress_stats_l: Vec<MediaEgressStats> = l
100 .events
101 .iter()
102 .filter(|(_, e)| matches!(e, Event::MediaEgressStats(_)))
103 .map(|(_, e)| {
104 if let Event::MediaEgressStats(stats) = e {
105 stats.clone()
106 } else {
107 panic!("Unexpected event type!")
108 }
109 })
110 .collect();
111 
112 egress_stats_l
113 .iter()
114 .filter_map(|egress_stat_l| egress_stat_l.rtt)
115 // rtt should be under 100ms in this scenario
116 .for_each(|rtt| assert!(rtt < Duration::from_millis(100)));
117 
118 let egress_stats_r: Vec<MediaEgressStats> = l
119 .events
120 .iter()
121 .filter(|(_, e)| matches!(e, Event::MediaEgressStats(_)))
122 .map(|(_, e)| {
123 if let Event::MediaEgressStats(stats) = e {
124 stats.clone()
125 } else {
126 panic!("Unexpected event type!")
127 }
128 })
129 .collect();
130 
131 egress_stats_r
132 .iter()
133 .filter_map(|egress_stat_l| egress_stat_l.rtt)
134 // rtt should be under 100ms in this scenario
135 .for_each(|rtt| assert!(rtt < Duration::from_millis(100)));
136 assert!(
137 media_count_l > 1100,
138 "Not enough MediaData at L: {}",
139 media_count_l
140 );
141 
142 Ok(())
143}
144 
145#[test]
146pub fn peer_media_stats_do_not_drop_when_streams_are_removed() -> Result<(), RtcError> {
147 init_log();
148 init_crypto_default();
149 
150 let now = Instant::now();
151 let config_l = RtcConfig::new()
152 .set_rtp_mode(true)
153 .set_stats_interval(Some(Duration::from_secs(1)));
154 let config_r = RtcConfig::new()
155 .set_rtp_mode(true)
156 .set_stats_interval(Some(Duration::from_secs(1)));
157 let (mut l, mut r) = connect_l_r_with_rtc(config_l.build(now), config_r.build(now));
158 
159 let mid = "aud".into();
160 let ssrc: Ssrc = 1.into();
161 
162 l.direct_api().declare_media(mid, MediaKind::Audio);
163 l.direct_api().declare_stream_tx(ssrc, None, mid, None);
164 r.direct_api().declare_media(mid, MediaKind::Audio);
165 r.direct_api().expect_stream_rx(ssrc, None, mid, None);
166 
167 let max = l.last.max(r.last);
168 l.last = max;
169 r.last = max;
170 
171 let params = l.params_opus();
172 assert_eq!(params.spec().codec, Codec::Opus);
173 let pt = params.pt();
174 
175 for i in 0..3 {
176 let wallclock = l.start + l.duration();
177 let time = (48_000 * i) as u32;
178 let seq_no = (1000 + i).into();
179 let payload = vec![i as u8; 80];
180 
181 l.direct_api()
182 .stream_tx(&ssrc)
183 .unwrap()
184 .write_rtp(RtpWrite::new(pt, seq_no, time, wallclock, payload));
185 
186 progress(&mut l, &mut r)?;
187 }
188 
189 let before_l = wait_for_peer_stats(&mut l, &mut r, true, |s| s.bytes_tx > 0)?;
190 let before_r = wait_for_peer_stats(&mut l, &mut r, false, |s| s.bytes_rx > 0)?;
191 
192 l.direct_api().remove_media(mid);
193 r.direct_api().remove_media(mid);
194 
195 let after_l = wait_for_peer_stats(&mut l, &mut r, true, |s| s.timestamp > before_l.timestamp)?;
196 let after_r = wait_for_peer_stats(&mut l, &mut r, false, |s| s.timestamp > before_r.timestamp)?;
197 
198 assert_eq!(after_l.bytes_tx, before_l.bytes_tx);
199 assert_eq!(after_r.bytes_rx, before_r.bytes_rx);
200 
201 Ok(())
202}
203 
204fn wait_for_peer_stats(
205 l: &mut TestRtc,
206 r: &mut TestRtc,
207 left: bool,
208 predicate: impl Fn(&PeerStats) -> bool,
209) -> Result<PeerStats, RtcError> {
210 for _ in 0..1000 {
211 progress(l, r)?;
212 
213 let events = if left { &l.events } else { &r.events };
214 if let Some(stats) = events.iter().rev().find_map(|(_, e)| {
215 if let Event::PeerStats(stats) = e {
216 predicate(stats).then(|| stats.clone())
217 } else {
218 None
219 }
220 }) {
221 return Ok(stats);
222 }
223 }
224 
225 panic!("timed out waiting for PeerStats");
226}