use crate::route::SharedTable;
use std::time::Duration;
fn health_url(uri: &http::Uri, path: &str) -> String {
let base = uri.to_string();
format!("{}{path}", base.trim_end_matches('/'))
}
fn should_probe(up: &crate::route::Upstream) -> bool {
!up.health_disabled()
}
pub async fn health_loop(table: SharedTable, interval: Duration) {
let client = match reqwest::Client::builder()
.timeout(Duration::from_secs(2))
.build()
{
Ok(c) => c,
Err(e) => {
tracing::error!(?e, "failed to build health-check reqwest client");
return;
}
};
let mut ticker = tokio::time::interval(interval);
ticker.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
loop {
ticker.tick().await;
let current = table.load_full();
for (_, up) in ¤t.rules {
if !should_probe(up) {
continue;
}
let url = health_url(&up.uri, up.health_path());
let host = up
.uri
.authority()
.map_or_else(String::new, |a| a.to_string());
match client.get(&url).send().await {
Ok(r) if r.status().is_success() => {
up.mark_success();
metrics::gauge!("ferryman_upstream_alive", "upstream" => host).set(1.0);
}
_ => {
up.mark_failed();
metrics::gauge!("ferryman_upstream_alive", "upstream" => host).set(0.0);
}
}
}
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn no_double_slash_for_bare_authority() {
let uri: http::Uri = "http://localhost:8001".parse().unwrap();
assert_eq!(health_url(&uri, "/health"), "http://localhost:8001/health");
}
#[test]
fn preserves_non_root_path() {
let uri: http::Uri = "http://localhost:8001/base".parse().unwrap();
assert_eq!(
health_url(&uri, "/health"),
"http://localhost:8001/base/health"
);
}
#[test]
fn custom_path() {
let uri: http::Uri = "http://localhost:8001/base/".parse().unwrap();
assert_eq!(
health_url(&uri, "/ready"),
"http://localhost:8001/base/ready"
);
}
#[tokio::test]
async fn loop_skips_disabled_upstream() {
use crate::route::{RouteTable, Upstream};
let uri: http::Uri = "http://127.0.0.1:1".parse().unwrap();
let up = Upstream::new(uri, 30).with_health(None, true);
let table: SharedTable = std::sync::Arc::new(arc_swap::ArcSwap::from_pointee(
RouteTable::new(vec![("/".into(), up)]),
));
let _ = tokio::time::timeout(
Duration::from_millis(500),
health_loop(table.clone(), Duration::from_secs(60)),
)
.await;
assert!(table.load().rules[0].1.is_routable());
}
#[test]
fn disabled_upstream_is_not_probed() {
let uri: http::Uri = "http://localhost:8001".parse().unwrap();
let up = crate::route::Upstream::new(uri, 30);
assert!(should_probe(&up));
assert!(!should_probe(&up.with_health(None, true)));
}
}