@@ -14,15 +14,22 @@ 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 > = LazyLock :: new ( || Regex :: new ( r"[?&]rtt=(\d+)" ) . unwrap ( ) ) ;
31+ static WARNED_NO_HEADER : AtomicBool = AtomicBool :: new ( false ) ;
32+ static WARNED_UNKNOWN_HEADER : AtomicBool = AtomicBool :: new ( false ) ;
2633const TIME_THRESHOLD : Duration = Duration :: from_secs ( 5 ) ;
2734const MAX_ATTEMPT_FACTOR : u32 = 4 ;
2835const RETRY_BASE_BACKOFF : Duration = Duration :: from_millis ( 250 ) ;
@@ -181,58 +188,94 @@ pub fn run_latency_test(
181188 if output_format == OutputFormat :: StdOut {
182189 print_progress ( "latency test" , i + 1 , nr_latency_tests) ;
183190 }
184- let latency = test_latency ( client) ;
185- measurements. push ( latency) ;
191+ if let Some ( latency) = try_test_latency ( client) {
192+ measurements. push ( latency) ;
193+ }
186194 }
187- let avg_latency = measurements. iter ( ) . sum :: < f64 > ( ) / measurements. len ( ) as f64 ;
195+ let avg_latency = if measurements. is_empty ( ) {
196+ 0.0
197+ } else {
198+ measurements. iter ( ) . sum :: < f64 > ( ) / measurements. len ( ) as f64
199+ } ;
188200
189201 if output_format == OutputFormat :: StdOut {
190- println ! (
191- "\n Avg GET request latency {avg_latency:.2} ms (RTT excluding server processing time)\n "
192- ) ;
202+ println ! ( "\n Avg GET request latency {avg_latency:.2} ms\n " ) ;
193203 }
194204 ( measurements, avg_latency)
195205}
196206
197- pub fn test_latency ( client : & Client ) -> f64 {
207+ // Parse latency from a Server-Timing header value. Supports the legacy
208+ // cfRequestDuration format and the newer cfL4 format. Returns None if
209+ // the header doesn't match either.
210+ fn parse_latency_from_server_timing ( header : & str , total_ms : f64 ) -> Option < f64 > {
211+ // Legacy: cfRequestDuration;dur=<milliseconds>
212+ if let Some ( caps) = RE_CF_REQUEST_DURATION . captures ( header) {
213+ if let Some ( dur_match) = caps. get ( 1 ) {
214+ if let Ok ( server_duration) = dur_match. as_str ( ) . parse :: < f64 > ( ) {
215+ let latency = total_ms - server_duration;
216+ return Some ( if latency < 0.0 { 0.0 } else { latency } ) ;
217+ }
218+ }
219+ }
220+
221+ // Current: cfL4;desc="?...&rtt=<microseconds>&..."
222+ // [?&] anchor prevents matching min_rtt= or rtt_var=
223+ if header. contains ( "cfL4" ) {
224+ if let Some ( caps) = RE_CFL4_RTT . captures ( header) {
225+ if let Some ( rtt_match) = caps. get ( 1 ) {
226+ if let Ok ( rtt_us) = rtt_match. as_str ( ) . parse :: < f64 > ( ) {
227+ return Some ( rtt_us / 1_000.0 ) ;
228+ }
229+ }
230+ }
231+ }
232+
233+ None
234+ }
235+
236+ fn try_test_latency ( client : & Client ) -> Option < f64 > {
198237 let url = & format ! ( "{}/{}{}" , BASE_URL , DOWNLOAD_URL , 0 ) ;
199238 let req_builder = client. get ( url) ;
200239
201240 let start = Instant :: now ( ) ;
202- let mut response = req_builder. send ( ) . expect ( "failed to get response" ) ;
241+ let mut response = match req_builder. send ( ) {
242+ Ok ( resp) => resp,
243+ Err ( e) => {
244+ log:: debug!( "Latency test request failed: {e}" ) ;
245+ return None ;
246+ }
247+ } ;
203248 let _status_code = response. status ( ) ;
204- // Drain body to complete the request; ignore errors.
205249 let _ = std:: io:: copy ( & mut response, & mut std:: io:: sink ( ) ) ;
206250 let total_ms = start. elapsed ( ) . as_secs_f64 ( ) * 1_000.0 ;
207251
208- let re = Regex :: new ( r"cfRequestDuration;dur=([\d.]+)" ) . unwrap ( ) ;
209252 let server_timing = response
210253 . headers ( )
211254 . 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- ) ;
255+ . and_then ( |v| v. to_str ( ) . ok ( ) ) ;
256+
257+ if let Some ( header) = server_timing {
258+ if let Some ( latency) = parse_latency_from_server_timing ( header, total_ms) {
259+ log:: debug!( "latency: total_ms={total_ms:.3} parsed={latency:.3}" ) ;
260+ return Some ( latency) ;
261+ }
262+ if !WARNED_UNKNOWN_HEADER . swap ( true , Ordering :: Relaxed ) {
263+ log:: warn!( "Server-Timing header format not recognized, falling back to raw RTT" ) ;
264+ }
265+ } else {
266+ if !WARNED_NO_HEADER . swap ( true , Ordering :: Relaxed ) {
267+ log:: warn!( "No Server-Timing header in response, falling back to raw RTT" ) ;
232268 }
233- req_latency = 0.0
234269 }
235- req_latency
270+ log:: debug!( "latency fallback: total_ms={total_ms:.3}" ) ;
271+ Some ( total_ms)
272+ }
273+
274+ pub fn test_latency ( client : & Client ) -> f64 {
275+ try_test_latency ( client) . unwrap_or_else ( || {
276+ log:: debug!( "Latency measurement failed, returning 0.0" ) ;
277+ 0.0
278+ } )
236279}
237280
238281#[ derive( Debug ) ]
@@ -1457,4 +1500,111 @@ mod tests {
14571500 metadata. ip, metadata. colo, metadata. country
14581501 ) ;
14591502 }
1503+
1504+ #[ test]
1505+ fn test_parse_latency_from_legacy_header ( ) {
1506+ // Old cfRequestDuration format
1507+ let header = "cfRequestDuration;dur=3.456" ;
1508+ let total_ms = 50.0 ;
1509+ let result = parse_latency_from_server_timing ( header, total_ms) ;
1510+ assert ! (
1511+ result. is_some( ) ,
1512+ "Should parse legacy cfRequestDuration header"
1513+ ) ;
1514+ let latency = result. unwrap ( ) ;
1515+ // latency = total_ms - server_duration = 50.0 - 3.456 = 46.544
1516+ assert ! ( ( latency - 46.544 ) . abs( ) < 0.001 ) ;
1517+ }
1518+
1519+ #[ test]
1520+ fn test_parse_latency_from_cfl4_header ( ) {
1521+ // New cfL4 format - rtt is in microseconds
1522+ let header =
1523+ r#"cfL4;desc="?proto=TCP&rtt=5003&min_rtt=4257&rtt_var=2477&sent=6&recv=6&lost=0""# ;
1524+ let total_ms = 50.0 ;
1525+ let result = parse_latency_from_server_timing ( header, total_ms) ;
1526+ assert ! ( result. is_some( ) , "Should parse cfL4 rtt header" ) ;
1527+ let latency = result. unwrap ( ) ;
1528+ // rtt=5003 microseconds = 5.003 milliseconds
1529+ assert ! ( ( latency - 5.003 ) . abs( ) < 0.001 ) ;
1530+ }
1531+
1532+ #[ test]
1533+ fn test_parse_latency_missing_header ( ) {
1534+ let header = "some-unrelated;value=123" ;
1535+ let total_ms = 50.0 ;
1536+ let result = parse_latency_from_server_timing ( header, total_ms) ;
1537+ assert ! (
1538+ result. is_none( ) ,
1539+ "Should return None for unrecognized header"
1540+ ) ;
1541+ }
1542+
1543+ #[ test]
1544+ fn test_parse_latency_prefers_legacy_over_cfl4 ( ) {
1545+ // If both are present (unlikely but defensive), prefer legacy
1546+ let header = "cfRequestDuration;dur=3.456, cfL4;desc=\" ?proto=TCP&rtt=5003\" " ;
1547+ let total_ms = 50.0 ;
1548+ let result = parse_latency_from_server_timing ( header, total_ms) ;
1549+ assert ! ( result. is_some( ) ) ;
1550+ let latency = result. unwrap ( ) ;
1551+ assert ! ( ( latency - 46.544 ) . abs( ) < 0.001 ) ;
1552+ }
1553+
1554+ #[ test]
1555+ fn test_parse_latency_cfl4_zero_rtt ( ) {
1556+ let header = r#"cfL4;desc="?proto=TCP&rtt=0&min_rtt=0""# ;
1557+ let total_ms = 50.0 ;
1558+ let result = parse_latency_from_server_timing ( header, total_ms) ;
1559+ assert ! ( result. is_some( ) ) ;
1560+ assert ! ( ( result. unwrap( ) - 0.0 ) . abs( ) < 0.001 ) ;
1561+ }
1562+
1563+ #[ test]
1564+ fn test_parse_latency_negative_clamp ( ) {
1565+ // Legacy header where server processing > total RTT (clock skew)
1566+ let header = "cfRequestDuration;dur=100.0" ;
1567+ let total_ms = 50.0 ;
1568+ let result = parse_latency_from_server_timing ( header, total_ms) ;
1569+ assert ! ( result. is_some( ) ) ;
1570+ // Should clamp to 0, not return negative
1571+ assert ! ( ( result. unwrap( ) - 0.0 ) . abs( ) < 0.001 ) ;
1572+ }
1573+
1574+ #[ test]
1575+ fn test_parse_latency_cfl4_does_not_match_min_rtt ( ) {
1576+ // Regression: rtt= regex must not match min_rtt= or rtt_var=
1577+ // Header with min_rtt before rtt - if regex is naive, it grabs min_rtt's value
1578+ let header = r#"cfL4;desc="?proto=TCP&min_rtt=4257&rtt_var=2477&rtt=5003&sent=6""# ;
1579+ let total_ms = 50.0 ;
1580+ let result = parse_latency_from_server_timing ( header, total_ms) ;
1581+ assert ! ( result. is_some( ) ) ;
1582+ let latency = result. unwrap ( ) ;
1583+ // Must match rtt=5003, NOT min_rtt=4257
1584+ assert ! ( ( latency - 5.003 ) . abs( ) < 0.001 ) ;
1585+ }
1586+
1587+ #[ test]
1588+ fn test_test_latency_no_panic_on_request_failure ( ) {
1589+ // Proxy pointed at a port with nothing listening - instant connection refused.
1590+ let client = reqwest:: blocking:: Client :: builder ( )
1591+ . proxy ( reqwest:: Proxy :: all ( "http://127.0.0.1:1" ) . unwrap ( ) )
1592+ . build ( )
1593+ . unwrap ( ) ;
1594+ // Must not panic. Returns 0.0 as the fallback for request failure.
1595+ let result = test_latency ( & client) ;
1596+ assert ! ( ( result - 0.0 ) . abs( ) < 0.001 ) ;
1597+ }
1598+
1599+ #[ test]
1600+ fn test_run_latency_test_all_failures_returns_zero_avg ( ) {
1601+ // Every request fails immediately - failed samples are skipped
1602+ let client = reqwest:: blocking:: Client :: builder ( )
1603+ . proxy ( reqwest:: Proxy :: all ( "http://127.0.0.1:1" ) . unwrap ( ) )
1604+ . build ( )
1605+ . unwrap ( ) ;
1606+ let ( measurements, avg) = run_latency_test ( & client, 3 , OutputFormat :: Json ) ;
1607+ assert ! ( measurements. is_empty( ) , "Failed requests should be skipped" ) ;
1608+ assert ! ( ( avg - 0.0 ) . abs( ) < 0.001 , "Average should be 0.0" ) ;
1609+ }
14601610}
0 commit comments