@@ -14,15 +14,23 @@ use serde::Serialize;
1414use std:: {
1515 fmt:: Display ,
1616 io:: Write ,
17- sync:: atomic:: { AtomicBool , Ordering } ,
17+ sync:: {
18+ atomic:: { AtomicBool , Ordering } ,
19+ LazyLock ,
20+ } ,
1821 thread,
1922 time:: { Duration , Instant } ,
2023} ;
2124
2225const BASE_URL : & str = "https://speed.cloudflare.com" ;
2326const DOWNLOAD_URL : & str = "__down?bytes=" ;
2427const UPLOAD_URL : & str = "__up" ;
25- static WARNED_NEGATIVE_LATENCY : AtomicBool = AtomicBool :: new ( false ) ;
28+ static RE_CF_REQUEST_DURATION : LazyLock < Regex > =
29+ LazyLock :: new ( || Regex :: new ( r"cfRequestDuration;dur=([\d.]+)" ) . unwrap ( ) ) ;
30+ static RE_CFL4_RTT : LazyLock < Regex > =
31+ LazyLock :: new ( || Regex :: new ( r"[?&]rtt=(\d+)" ) . unwrap ( ) ) ;
32+ static WARNED_NO_HEADER : AtomicBool = AtomicBool :: new ( false ) ;
33+ static WARNED_UNKNOWN_HEADER : AtomicBool = AtomicBool :: new ( false ) ;
2634const TIME_THRESHOLD : Duration = Duration :: from_secs ( 5 ) ;
2735const MAX_ATTEMPT_FACTOR : u32 = 4 ;
2836const RETRY_BASE_BACKOFF : Duration = Duration :: from_millis ( 250 ) ;
@@ -181,58 +189,94 @@ pub fn run_latency_test(
181189 if output_format == OutputFormat :: StdOut {
182190 print_progress ( "latency test" , i + 1 , nr_latency_tests) ;
183191 }
184- let latency = test_latency ( client) ;
185- measurements. push ( latency) ;
192+ if let Some ( latency) = try_test_latency ( client) {
193+ measurements. push ( latency) ;
194+ }
186195 }
187- let avg_latency = measurements. iter ( ) . sum :: < f64 > ( ) / measurements. len ( ) as f64 ;
196+ let avg_latency = if measurements. is_empty ( ) {
197+ 0.0
198+ } else {
199+ measurements. iter ( ) . sum :: < f64 > ( ) / measurements. len ( ) as f64
200+ } ;
188201
189202 if output_format == OutputFormat :: StdOut {
190- println ! (
191- "\n Avg GET request latency {avg_latency:.2} ms (RTT excluding server processing time)\n "
192- ) ;
203+ println ! ( "\n Avg GET request latency {avg_latency:.2} ms\n " ) ;
193204 }
194205 ( measurements, avg_latency)
195206}
196207
197- pub fn test_latency ( client : & Client ) -> f64 {
208+ // Parse latency from a Server-Timing header value. Supports the legacy
209+ // cfRequestDuration format and the newer cfL4 format. Returns None if
210+ // the header doesn't match either.
211+ fn parse_latency_from_server_timing ( header : & str , total_ms : f64 ) -> Option < f64 > {
212+ // Legacy: cfRequestDuration;dur=<milliseconds>
213+ if let Some ( caps) = RE_CF_REQUEST_DURATION . captures ( header) {
214+ if let Some ( dur_match) = caps. get ( 1 ) {
215+ if let Ok ( server_duration) = dur_match. as_str ( ) . parse :: < f64 > ( ) {
216+ let latency = total_ms - server_duration;
217+ return Some ( if latency < 0.0 { 0.0 } else { latency } ) ;
218+ }
219+ }
220+ }
221+
222+ // Current: cfL4;desc="?...&rtt=<microseconds>&..."
223+ // [?&] anchor prevents matching min_rtt= or rtt_var=
224+ if header. contains ( "cfL4" ) {
225+ if let Some ( caps) = RE_CFL4_RTT . captures ( header) {
226+ if let Some ( rtt_match) = caps. get ( 1 ) {
227+ if let Ok ( rtt_us) = rtt_match. as_str ( ) . parse :: < f64 > ( ) {
228+ return Some ( rtt_us / 1_000.0 ) ;
229+ }
230+ }
231+ }
232+ }
233+
234+ None
235+ }
236+
237+ fn try_test_latency ( client : & Client ) -> Option < f64 > {
198238 let url = & format ! ( "{}/{}{}" , BASE_URL , DOWNLOAD_URL , 0 ) ;
199239 let req_builder = client. get ( url) ;
200240
201241 let start = Instant :: now ( ) ;
202- let mut response = req_builder. send ( ) . expect ( "failed to get response" ) ;
242+ let mut response = match req_builder. send ( ) {
243+ Ok ( resp) => resp,
244+ Err ( e) => {
245+ log:: debug!( "Latency test request failed: {e}" ) ;
246+ return None ;
247+ }
248+ } ;
203249 let _status_code = response. status ( ) ;
204- // Drain body to complete the request; ignore errors.
205250 let _ = std:: io:: copy ( & mut response, & mut std:: io:: sink ( ) ) ;
206251 let total_ms = start. elapsed ( ) . as_secs_f64 ( ) * 1_000.0 ;
207252
208- let re = Regex :: new ( r"cfRequestDuration;dur=([\d.]+)" ) . unwrap ( ) ;
209253 let server_timing = response
210254 . headers ( )
211255 . get ( "Server-Timing" )
212- . expect ( "No Server-Timing in response header" )
213- . to_str ( )
214- . unwrap ( ) ;
215- let cf_req_duration: f64 = re
216- . captures ( server_timing)
217- . unwrap ( )
218- . get ( 1 )
219- . unwrap ( )
220- . as_str ( )
221- . parse ( )
222- . unwrap ( ) ;
223- let mut req_latency = total_ms - cf_req_duration;
224- log:: debug!(
225- "latency debug: total_ms={total_ms:.3} cf_req_duration_ms={cf_req_duration:.3} req_latency_total={req_latency:.3} server_timing={server_timing}"
226- ) ;
227- if req_latency < 0.0 {
228- if !WARNED_NEGATIVE_LATENCY . swap ( true , Ordering :: Relaxed ) {
229- log:: warn!(
230- "negative latency after server timing subtraction; clamping to 0.0 (total_ms={total_ms:.3} cf_req_duration_ms={cf_req_duration:.3})"
231- ) ;
256+ . and_then ( |v| v. to_str ( ) . ok ( ) ) ;
257+
258+ if let Some ( header) = server_timing {
259+ if let Some ( latency) = parse_latency_from_server_timing ( header, total_ms) {
260+ log:: debug!( "latency: total_ms={total_ms:.3} parsed={latency:.3}" ) ;
261+ return Some ( latency) ;
262+ }
263+ if !WARNED_UNKNOWN_HEADER . swap ( true , Ordering :: Relaxed ) {
264+ log:: warn!( "Server-Timing header format not recognized, falling back to raw RTT" ) ;
265+ }
266+ } else {
267+ if !WARNED_NO_HEADER . swap ( true , Ordering :: Relaxed ) {
268+ log:: warn!( "No Server-Timing header in response, falling back to raw RTT" ) ;
232269 }
233- req_latency = 0.0
234270 }
235- req_latency
271+ log:: debug!( "latency fallback: total_ms={total_ms:.3}" ) ;
272+ Some ( total_ms)
273+ }
274+
275+ pub fn test_latency ( client : & Client ) -> f64 {
276+ try_test_latency ( client) . unwrap_or_else ( || {
277+ log:: debug!( "Latency measurement failed, returning 0.0" ) ;
278+ 0.0
279+ } )
236280}
237281
238282#[ derive( Debug ) ]
@@ -1457,4 +1501,104 @@ mod tests {
14571501 metadata. ip, metadata. colo, metadata. country
14581502 ) ;
14591503 }
1504+
1505+ #[ test]
1506+ fn test_parse_latency_from_legacy_header ( ) {
1507+ // Old cfRequestDuration format
1508+ let header = "cfRequestDuration;dur=3.456" ;
1509+ let total_ms = 50.0 ;
1510+ let result = parse_latency_from_server_timing ( header, total_ms) ;
1511+ assert ! ( result. is_some( ) , "Should parse legacy cfRequestDuration header" ) ;
1512+ let latency = result. unwrap ( ) ;
1513+ // latency = total_ms - server_duration = 50.0 - 3.456 = 46.544
1514+ assert ! ( ( latency - 46.544 ) . abs( ) < 0.001 ) ;
1515+ }
1516+
1517+ #[ test]
1518+ fn test_parse_latency_from_cfl4_header ( ) {
1519+ // New cfL4 format - rtt is in microseconds
1520+ let header = r#"cfL4;desc="?proto=TCP&rtt=5003&min_rtt=4257&rtt_var=2477&sent=6&recv=6&lost=0""# ;
1521+ let total_ms = 50.0 ;
1522+ let result = parse_latency_from_server_timing ( header, total_ms) ;
1523+ assert ! ( result. is_some( ) , "Should parse cfL4 rtt header" ) ;
1524+ let latency = result. unwrap ( ) ;
1525+ // rtt=5003 microseconds = 5.003 milliseconds
1526+ assert ! ( ( latency - 5.003 ) . abs( ) < 0.001 ) ;
1527+ }
1528+
1529+ #[ test]
1530+ fn test_parse_latency_missing_header ( ) {
1531+ let header = "some-unrelated;value=123" ;
1532+ let total_ms = 50.0 ;
1533+ let result = parse_latency_from_server_timing ( header, total_ms) ;
1534+ assert ! ( result. is_none( ) , "Should return None for unrecognized header" ) ;
1535+ }
1536+
1537+ #[ test]
1538+ fn test_parse_latency_prefers_legacy_over_cfl4 ( ) {
1539+ // If both are present (unlikely but defensive), prefer legacy
1540+ let header = "cfRequestDuration;dur=3.456, cfL4;desc=\" ?proto=TCP&rtt=5003\" " ;
1541+ let total_ms = 50.0 ;
1542+ let result = parse_latency_from_server_timing ( header, total_ms) ;
1543+ assert ! ( result. is_some( ) ) ;
1544+ let latency = result. unwrap ( ) ;
1545+ assert ! ( ( latency - 46.544 ) . abs( ) < 0.001 ) ;
1546+ }
1547+
1548+ #[ test]
1549+ fn test_parse_latency_cfl4_zero_rtt ( ) {
1550+ let header = r#"cfL4;desc="?proto=TCP&rtt=0&min_rtt=0""# ;
1551+ let total_ms = 50.0 ;
1552+ let result = parse_latency_from_server_timing ( header, total_ms) ;
1553+ assert ! ( result. is_some( ) ) ;
1554+ assert ! ( ( result. unwrap( ) - 0.0 ) . abs( ) < 0.001 ) ;
1555+ }
1556+
1557+ #[ test]
1558+ fn test_parse_latency_negative_clamp ( ) {
1559+ // Legacy header where server processing > total RTT (clock skew)
1560+ let header = "cfRequestDuration;dur=100.0" ;
1561+ let total_ms = 50.0 ;
1562+ let result = parse_latency_from_server_timing ( header, total_ms) ;
1563+ assert ! ( result. is_some( ) ) ;
1564+ // Should clamp to 0, not return negative
1565+ assert ! ( ( result. unwrap( ) - 0.0 ) . abs( ) < 0.001 ) ;
1566+ }
1567+
1568+ #[ test]
1569+ fn test_parse_latency_cfl4_does_not_match_min_rtt ( ) {
1570+ // Regression: rtt= regex must not match min_rtt= or rtt_var=
1571+ // Header with min_rtt before rtt - if regex is naive, it grabs min_rtt's value
1572+ let header = r#"cfL4;desc="?proto=TCP&min_rtt=4257&rtt_var=2477&rtt=5003&sent=6""# ;
1573+ let total_ms = 50.0 ;
1574+ let result = parse_latency_from_server_timing ( header, total_ms) ;
1575+ assert ! ( result. is_some( ) ) ;
1576+ let latency = result. unwrap ( ) ;
1577+ // Must match rtt=5003, NOT min_rtt=4257
1578+ assert ! ( ( latency - 5.003 ) . abs( ) < 0.001 ) ;
1579+ }
1580+
1581+ #[ test]
1582+ fn test_test_latency_no_panic_on_request_failure ( ) {
1583+ // Proxy pointed at a port with nothing listening - instant connection refused.
1584+ let client = reqwest:: blocking:: Client :: builder ( )
1585+ . proxy ( reqwest:: Proxy :: all ( "http://127.0.0.1:1" ) . unwrap ( ) )
1586+ . build ( )
1587+ . unwrap ( ) ;
1588+ // Must not panic. Returns 0.0 as the fallback for request failure.
1589+ let result = test_latency ( & client) ;
1590+ assert ! ( ( result - 0.0 ) . abs( ) < 0.001 ) ;
1591+ }
1592+
1593+ #[ test]
1594+ fn test_run_latency_test_all_failures_returns_zero_avg ( ) {
1595+ // Every request fails immediately - failed samples are skipped
1596+ let client = reqwest:: blocking:: Client :: builder ( )
1597+ . proxy ( reqwest:: Proxy :: all ( "http://127.0.0.1:1" ) . unwrap ( ) )
1598+ . build ( )
1599+ . unwrap ( ) ;
1600+ let ( measurements, avg) = run_latency_test ( & client, 3 , OutputFormat :: Json ) ;
1601+ assert ! ( measurements. is_empty( ) , "Failed requests should be skipped" ) ;
1602+ assert ! ( ( avg - 0.0 ) . abs( ) < 0.001 , "Average should be 0.0" ) ;
1603+ }
14601604}
0 commit comments