Skip to main content

linkprobe_core/backends/
librespeed.rs

1use std::time::Instant;
2
3use reqwest::blocking::{Body, Client};
4use url::Url;
5
6use crate::measurement::{Measurement, Throughput};
7use crate::server::Server;
8use crate::{Error, MeasurementEngine};
9
10const PING_SAMPLES: usize = 10;
11const DOWNLOAD_CHUNK_SIZE: u32 = 4; // 4 MiB from garbage.php
12const UPLOAD_BYTES: usize = 2 * 1024 * 1024;
13const HTTP_ATTEMPTS: usize = 3;
14
15/// LibreSpeed-compatible measurement over blocking HTTPS.
16///
17/// Each HTTP phase (ping, download, upload) is retried up to three times on timeout,
18/// connection failure, or truncated bodies.
19#[derive(Debug, Clone)]
20pub struct LibreSpeedEngine {
21    client: Client,
22}
23
24impl LibreSpeedEngine {
25    /// Build an engine with a default reqwest client (60 s timeout, linkprobe user agent).
26    pub fn new() -> Result<Self, Error> {
27        let client = Client::builder()
28            .user_agent(concat!("linkprobe/", env!("CARGO_PKG_VERSION")))
29            .timeout(std::time::Duration::from_secs(60))
30            .build()?;
31        Ok(Self { client })
32    }
33
34    fn join(base: &str, path: &str) -> Result<Url, Error> {
35        let base = if base.ends_with('/') {
36            base.to_string()
37        } else {
38            format!("{base}/")
39        };
40        Ok(Url::parse(&base)?.join(path)?)
41    }
42
43    fn is_retryable(err: &reqwest::Error) -> bool {
44        err.is_timeout() || err.is_connect() || err.is_decode() || err.is_body()
45    }
46
47    fn with_retries<T>(
48        phase: &'static str,
49        mut op: impl FnMut() -> Result<T, reqwest::Error>,
50    ) -> Result<T, Error> {
51        let mut last = None;
52        for _ in 0..HTTP_ATTEMPTS {
53            match op() {
54                Ok(v) => return Ok(v),
55                Err(e) if Self::is_retryable(&e) => last = Some(e),
56                Err(e) => return Err(Error::from_reqwest(phase, e)),
57            }
58        }
59        Err(Error::from_reqwest(phase, last.expect("retry loop")))
60    }
61
62    fn ping_jitter(&self, server: &Server) -> Result<(f64, f64), Error> {
63        let url = Self::join(&server.base_url, &server.ping_path)?;
64        let mut samples_ms = Vec::with_capacity(PING_SAMPLES);
65
66        for _ in 0..PING_SAMPLES {
67            let start = Instant::now();
68            Self::with_retries("ping", || {
69                let resp = self.client.get(url.clone()).send()?.error_for_status()?;
70                let _ = resp.bytes()?;
71                Ok(())
72            })?;
73            samples_ms.push(start.elapsed().as_secs_f64() * 1000.0);
74        }
75
76        let latency = samples_ms.iter().sum::<f64>() / samples_ms.len() as f64;
77        let mut jitter_acc = 0.0;
78        for w in samples_ms.windows(2) {
79            jitter_acc += (w[1] - w[0]).abs();
80        }
81        let jitter = if samples_ms.len() > 1 {
82            jitter_acc / (samples_ms.len() - 1) as f64
83        } else {
84            0.0
85        };
86        Ok((latency, jitter))
87    }
88
89    fn download(&self, server: &Server) -> Result<Throughput, Error> {
90        let mut url = Self::join(&server.base_url, &server.dl_path)?;
91        url.query_pairs_mut()
92            .append_pair("ckSize", &DOWNLOAD_CHUNK_SIZE.to_string());
93
94        let (n, secs) = Self::with_retries("download", || {
95            let start = Instant::now();
96            let resp = self.client.get(url.clone()).send()?.error_for_status()?;
97            let n = resp.bytes()?.len();
98            Ok((n, start.elapsed().as_secs_f64().max(1e-6)))
99        })?;
100        Ok(Throughput::from_bps((n as f64) * 8.0 / secs))
101    }
102
103    fn upload(&self, server: &Server) -> Result<Throughput, Error> {
104        let url = Self::join(&server.base_url, &server.ul_path)?;
105        let len = UPLOAD_BYTES as f64;
106
107        let secs = Self::with_retries("upload", || {
108            let payload = vec![0_u8; UPLOAD_BYTES];
109            let start = Instant::now();
110            let resp = self
111                .client
112                .post(url.clone())
113                .header(reqwest::header::CONTENT_TYPE, "application/octet-stream")
114                .body(Body::from(payload))
115                .send()?
116                .error_for_status()?;
117            let _ = resp.bytes()?.len();
118            Ok(start.elapsed().as_secs_f64().max(1e-6))
119        })?;
120        Ok(Throughput::from_bps(len * 8.0 / secs))
121    }
122
123    pub fn measure_with_failover(
124        &self,
125        candidates: &[Server],
126    ) -> Result<(Server, Measurement), Error> {
127        if candidates.is_empty() {
128            return Err(Error::Message("no LibreSpeed candidates".into()));
129        }
130        let mut last_err: Option<Error> = None;
131        for (i, server) in candidates.iter().enumerate() {
132            match self.measure(server) {
133                Ok(m) => return Ok((server.clone(), m)),
134                Err(e) => {
135                    if i + 1 < candidates.len() {
136                        eprintln!("linkprobe: {} failed, trying next server", server.name);
137                    }
138                    last_err = Some(e);
139                }
140            }
141        }
142        Err(last_err.expect("non-empty candidates"))
143    }
144}
145
146impl Default for LibreSpeedEngine {
147    fn default() -> Self {
148        Self::new().expect("failed to build reqwest client")
149    }
150}
151
152impl MeasurementEngine for LibreSpeedEngine {
153    fn measure(&self, server: &Server) -> Result<Measurement, Error> {
154        let (latency_ms, jitter_ms) = self.ping_jitter(server)?;
155        let download = self.download(server)?;
156        let upload = self.upload(server)?;
157        Ok(Measurement {
158            latency_ms: Some(latency_ms),
159            jitter_ms: Some(jitter_ms),
160            download: Some(download),
161            upload: Some(upload),
162            packet_loss: None,
163        })
164    }
165}
166
167#[cfg(test)]
168mod tests {
169    use super::*;
170    use mockito::Server as MockServer;
171
172    #[test]
173    fn measures_against_mock_librespeed() {
174        let mut server = MockServer::new();
175
176        let ping = server
177            .mock("GET", "/backend/empty.php")
178            .with_status(200)
179            .with_body("")
180            .expect_at_least(PING_SAMPLES)
181            .create();
182
183        let dl_body = vec![1_u8; DOWNLOAD_CHUNK_SIZE as usize * 1024 * 1024];
184        let download = server
185            .mock("GET", "/backend/garbage.php")
186            .match_query(mockito::Matcher::UrlEncoded(
187                "ckSize".into(),
188                DOWNLOAD_CHUNK_SIZE.to_string(),
189            ))
190            .with_status(200)
191            .with_body(dl_body)
192            .create();
193
194        let upload = server
195            .mock("POST", "/backend/empty.php")
196            .with_status(200)
197            .with_body("")
198            .create();
199
200        let engine = LibreSpeedEngine::new().unwrap();
201        let target = Server::librespeed(server.url());
202        let m = engine.measure(&target).unwrap();
203
204        assert!(m.latency_ms.unwrap() >= 0.0);
205        assert!(m.jitter_ms.unwrap() >= 0.0);
206        assert!(m.download.unwrap().bps > 0.0);
207        assert!(m.upload.unwrap().bps > 0.0);
208
209        ping.assert();
210        download.assert();
211        upload.assert();
212    }
213
214    #[test]
215    fn failovers_to_second_server() {
216        let mut bad = MockServer::new();
217        let mut good = MockServer::new();
218
219        let _ping_bad = bad
220            .mock("GET", "/backend/empty.php")
221            .with_status(200)
222            .with_body("")
223            .expect_at_least(PING_SAMPLES)
224            .create();
225        let _dl_bad = bad
226            .mock("GET", "/backend/garbage.php")
227            .match_query(mockito::Matcher::UrlEncoded(
228                "ckSize".into(),
229                DOWNLOAD_CHUNK_SIZE.to_string(),
230            ))
231            .with_status(500)
232            .with_body("nope")
233            .create();
234
235        let _ping_good = good
236            .mock("GET", "/backend/empty.php")
237            .with_status(200)
238            .with_body("")
239            .expect_at_least(PING_SAMPLES)
240            .create();
241        let dl_body = vec![1_u8; DOWNLOAD_CHUNK_SIZE as usize * 1024 * 1024];
242        let _dl_good = good
243            .mock("GET", "/backend/garbage.php")
244            .match_query(mockito::Matcher::UrlEncoded(
245                "ckSize".into(),
246                DOWNLOAD_CHUNK_SIZE.to_string(),
247            ))
248            .with_status(200)
249            .with_body(dl_body)
250            .create();
251        let _ul_good = good
252            .mock("POST", "/backend/empty.php")
253            .with_status(200)
254            .with_body("")
255            .create();
256
257        let engine = LibreSpeedEngine::new().unwrap();
258        let a = Server::librespeed(bad.url());
259        let b = Server::librespeed(good.url());
260        let (used, m) = engine.measure_with_failover(&[a, b.clone()]).unwrap();
261        assert_eq!(used.base_url, b.base_url);
262        assert!(m.download.unwrap().bps > 0.0);
263    }
264}