linkprobe_core/backends/
iperf3.rs1use std::path::PathBuf;
2use std::process::Command;
3
4use serde::Deserialize;
5
6use crate::measurement::{Measurement, Throughput};
7use crate::server::Server;
8use crate::{Error, MeasurementEngine};
9
10const DEFAULT_DURATION_SECS: u64 = 5;
11const DEFAULT_PORT: u16 = 5201;
12
13#[derive(Debug, Clone)]
18pub struct Iperf3Engine {
19 binary: PathBuf,
20 duration_secs: u64,
21 udp: bool,
22 bandwidth: String,
23}
24
25impl Iperf3Engine {
26 pub fn new() -> Self {
28 Self {
29 binary: PathBuf::from("iperf3"),
30 duration_secs: DEFAULT_DURATION_SECS,
31 udp: false,
32 bandwidth: "10M".into(),
33 }
34 }
35
36 pub fn with_binary(mut self, binary: impl Into<PathBuf>) -> Self {
38 self.binary = binary.into();
39 self
40 }
41
42 pub fn with_duration_secs(mut self, secs: u64) -> Self {
44 self.duration_secs = secs.max(1);
45 self
46 }
47
48 pub fn with_udp(mut self, udp: bool) -> Self {
50 self.udp = udp;
51 self
52 }
53
54 pub fn with_bandwidth(mut self, bandwidth: impl Into<String>) -> Self {
56 let b = bandwidth.into();
57 self.bandwidth = if b.is_empty() { "10M".into() } else { b };
58 self
59 }
60
61 fn port(server: &Server) -> u16 {
62 server.port.unwrap_or(DEFAULT_PORT)
63 }
64
65 fn run_json(&self, server: &Server, reverse: bool) -> Result<Iperf3Json, Error> {
66 let mut cmd = Command::new(&self.binary);
67 cmd.arg("-c")
68 .arg(&server.base_url)
69 .arg("-p")
70 .arg(Self::port(server).to_string())
71 .arg("-t")
72 .arg(self.duration_secs.to_string())
73 .arg("-J");
74 if self.udp {
75 cmd.arg("-u").arg("-b").arg(&self.bandwidth);
76 }
77 if reverse {
78 cmd.arg("-R");
79 }
80
81 let output = cmd.output().map_err(|e| {
82 if e.kind() == std::io::ErrorKind::NotFound {
83 Error::Iperf3Missing
84 } else {
85 Error::probe("iperf3", e)
86 }
87 })?;
88
89 if !output.status.success() {
90 let stderr = String::from_utf8_lossy(&output.stderr);
91 return Err(Error::probe(
92 "iperf3",
93 Error::Message(format!("exited {}: {stderr}", output.status)),
94 ));
95 }
96
97 Ok(serde_json::from_slice(&output.stdout)?)
98 }
99
100 pub fn bps_from_json(json: &Iperf3Json, reverse: bool) -> Result<f64, Error> {
102 let end = json
103 .end
104 .as_ref()
105 .ok_or_else(|| Error::Message("iperf3 JSON missing end block".into()))?;
106
107 let bps = if reverse {
108 end.sum_received
109 .as_ref()
110 .or(end.sum_sent.as_ref())
111 .map(|s| s.bits_per_second)
112 } else {
113 end.sum_sent
114 .as_ref()
115 .or(end.sum_received.as_ref())
116 .map(|s| s.bits_per_second)
117 };
118
119 bps.ok_or_else(|| Error::Message("iperf3 JSON missing bps".into()))
120 }
121
122 fn receiver_sum(end: &IperfEnd) -> Option<&IperfSum> {
123 end.sum_received
124 .as_ref()
125 .or(end.sum.as_ref())
126 .or(end.sum_sent.as_ref())
127 }
128
129 pub fn udp_stats(json: &Iperf3Json) -> (Option<f64>, Option<f64>) {
130 let Some(end) = json.end.as_ref() else {
131 return (None, None);
132 };
133 let Some(sum) = Self::receiver_sum(end) else {
134 return (None, None);
135 };
136 let jitter = sum.jitter_ms;
137 let loss = sum.lost_percent.map(|p| (p / 100.0).clamp(0.0, 1.0));
138 (jitter, loss)
139 }
140}
141
142impl Default for Iperf3Engine {
143 fn default() -> Self {
144 Self::new()
145 }
146}
147
148impl MeasurementEngine for Iperf3Engine {
149 fn measure(&self, server: &Server) -> Result<Measurement, Error> {
150 let dl = self.run_json(server, true)?;
151 let ul = self.run_json(server, false)?;
152 let (jitter_ms, packet_loss) = if self.udp {
153 let (j1, l1) = Self::udp_stats(&dl);
154 let (j2, l2) = Self::udp_stats(&ul);
155 let jitter = match (j1, j2) {
156 (Some(a), Some(b)) => Some((a + b) / 2.0),
157 (a, b) => a.or(b),
158 };
159 let loss = match (l1, l2) {
160 (Some(a), Some(b)) => Some((a + b) / 2.0),
161 (a, b) => a.or(b),
162 };
163 (jitter, loss)
164 } else {
165 (None, None)
166 };
167 Ok(Measurement {
168 latency_ms: None,
169 jitter_ms,
170 download: Some(Throughput::from_bps(Self::bps_from_json(&dl, true)?)),
171 upload: Some(Throughput::from_bps(Self::bps_from_json(&ul, false)?)),
172 packet_loss,
173 })
174 }
175}
176
177#[derive(Debug, Deserialize)]
178pub struct Iperf3Json {
179 pub end: Option<IperfEnd>,
180}
181
182#[derive(Debug, Deserialize)]
183pub struct IperfEnd {
184 pub sum: Option<IperfSum>,
185 pub sum_sent: Option<IperfSum>,
186 pub sum_received: Option<IperfSum>,
187}
188
189#[derive(Debug, Deserialize)]
190pub struct IperfSum {
191 pub bits_per_second: f64,
192 pub jitter_ms: Option<f64>,
193 pub lost_percent: Option<f64>,
194}
195
196#[cfg(test)]
197mod tests {
198 use super::*;
199
200 #[test]
201 fn parses_forward_and_reverse_fixtures() {
202 let forward = include_str!("../../tests/fixtures/iperf3/forward.json");
203 let reverse = include_str!("../../tests/fixtures/iperf3/reverse.json");
204 let fwd: Iperf3Json = serde_json::from_str(forward).unwrap();
205 let rev: Iperf3Json = serde_json::from_str(reverse).unwrap();
206
207 let up = Iperf3Engine::bps_from_json(&fwd, false).unwrap();
208 let down = Iperf3Engine::bps_from_json(&rev, true).unwrap();
209 assert!(up > 0.0);
210 assert!(down > 0.0);
211 }
212
213 #[test]
214 fn parses_udp_jitter_and_loss() {
215 let json = include_str!("../../tests/fixtures/iperf3/udp.json");
216 let parsed: Iperf3Json = serde_json::from_str(json).unwrap();
217 let (jitter, loss) = Iperf3Engine::udp_stats(&parsed);
218 assert_eq!(jitter, Some(1.5));
219 assert!((loss.unwrap() - 0.02).abs() < 1e-9);
220 assert!(Iperf3Engine::bps_from_json(&parsed, true).unwrap() > 0.0);
221 }
222}