pub mod models;
use std::env;
use std::process::Stdio;
use std::sync::Mutex;
use std::time::{Duration, Instant};
use actix_web::{get, HttpResponse, Responder};
use async_trait::async_trait;
use once_cell::sync::Lazy;
use tokio::process::Command;
use tokio::time;
pub use models::*;
pub static LAST_RESULT: Lazy<Mutex<Option<(SpeedTestResult, Instant)>>> = Lazy::new(|| Mutex::new(None));
pub fn get_last_result() -> Option<SpeedTestResult> {
let cache = LAST_RESULT.lock().unwrap();
cache.as_ref().map(|(result, _)| result.clone())
}
pub fn set_last_result_for_test(result: SpeedTestResult) {
let mut cache = LAST_RESULT.lock().unwrap();
*cache = Some((result, Instant::now()));
}
pub fn clear_last_result_for_test() {
let mut cache = LAST_RESULT.lock().unwrap();
*cache = None;
}
#[get("/speed")]
pub async fn speedtest() -> impl Responder {
let cache = LAST_RESULT.lock().unwrap();
if let Some((cached_result, _timestamp)) = &*cache {
HttpResponse::Ok().json(cached_result)
} else {
HttpResponse::ServiceUnavailable().body("Speedtest result not available yet.")
}
}
pub fn min_frequency_duration() -> Duration {
let minutes = env::var("INTERVAL_MINUTES")
.ok()
.and_then(|s| s.parse::<u64>().ok())
.unwrap_or(10); Duration::from_secs(minutes * 60)
}
#[async_trait]
pub trait SpeedtestRunner: Send + Sync {
async fn run_speedtest(&self) -> Result<String, String>;
}
pub struct RealSpeedtestRunner;
#[async_trait]
impl SpeedtestRunner for RealSpeedtestRunner {
async fn run_speedtest(&self) -> Result<String, String> {
let output = Command::new("speedtest-cli")
.arg("--json")
.stdout(Stdio::piped())
.stderr(Stdio::piped())
.output()
.await
.map_err(|e| format!("Failed to run speedtest-cli: {}", e))?;
if output.status.success() {
Ok(String::from_utf8_lossy(&output.stdout).to_string())
} else {
let stderr = String::from_utf8_lossy(&output.stderr);
Err(format!("speedtest-cli failed: {}", stderr))
}
}
}
pub async fn run_speedtest_and_cache_with_runner(runner: &dyn SpeedtestRunner) {
match runner.run_speedtest().await {
Ok(stdout) => match serde_json::from_str::<SpeedTestResponse>(&stdout) {
Ok(data) => {
let result = SpeedTestResult {
bytes_received: data.bytes_received,
bytes_sent: data.bytes_sent,
download_bps: data.download,
upload_bps: data.upload,
download_mbps: data.download / 1_000_000.0,
upload_mbps: data.upload / 1_000_000.0,
ping_ms: data.ping,
client: data.client,
server: data.server,
share: data.share,
timestamp: data.timestamp,
};
let mut cache = LAST_RESULT.lock().unwrap();
*cache = Some((result.clone(), Instant::now()));
println!("Speedtest updated at {}", result.timestamp);
}
Err(e) => eprintln!("Failed to parse speedtest-cli JSON: {}", e),
},
Err(e) => eprintln!("{}", e),
}
}
pub async fn spawn_speedtest_scheduler() {
let interval = min_frequency_duration();
let runner = RealSpeedtestRunner;
run_speedtest_and_cache_with_runner(&runner).await;
let mut ticker = time::interval(interval);
loop {
ticker.tick().await;
run_speedtest_and_cache_with_runner(&runner).await;
}
}
pub async fn get_cached_speedtest_result() -> Result<SpeedTestResult, String> {
let cache = LAST_RESULT.lock().unwrap();
if let Some((cached_result, _)) = &*cache {
Ok(cached_result.clone())
} else {
Err("Speedtest result not available yet.".to_string())
}
}