linkprobe_core/backends/
librespeed.rs1use 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; const UPLOAD_BYTES: usize = 2 * 1024 * 1024;
13const HTTP_ATTEMPTS: usize = 3;
14
15#[derive(Debug, Clone)]
20pub struct LibreSpeedEngine {
21 client: Client,
22}
23
24impl LibreSpeedEngine {
25 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}