use crate::balancer::LoadBalancer;
use crate::endpoint::RpcEndpoint;
use crate::metrics::COOLDOWN_SECONDS_GAUGE;
use std::{
sync::Arc,
time::{Duration, Instant},
};
use tokio::sync::watch;
use tokio::time::{interval, MissedTickBehavior};
use tracing::{info, warn};
pub fn trigger_cooldown(ep: &mut RpcEndpoint, base_cooldown_secs: u64, max_cooldown_secs: u64) {
ep.cooldown_attempts = ep.cooldown_attempts.saturating_add(1);
let exponential_cooldown = base_cooldown_secs.saturating_mul(2u64.pow(ep.cooldown_attempts));
let cooldown_duration = exponential_cooldown.min(max_cooldown_secs);
ep.cooldown_until = Some(Instant::now() + Duration::from_secs(cooldown_duration));
ep.last_check = Instant::now();
warn!(
endpoint = %ep.name,
cooldown_secs = cooldown_duration,
attempts = ep.cooldown_attempts,
"Endpoint has been put into cooldown."
);
COOLDOWN_SECONDS_GAUGE.with_label_values(&[&ep.name]).set(cooldown_duration as i64);
}
pub async fn cooldown_gauge_updater(
balancer: Arc<LoadBalancer>,
mut shutdown_rx: watch::Receiver<()>,
) {
let mut ticker = interval(Duration::from_secs(1));
ticker.set_missed_tick_behavior(MissedTickBehavior::Delay);
loop {
tokio::select! {
biased; _ = shutdown_rx.changed() => {
info!("Cooldown gauge updater received shutdown signal, exiting.");
return;
}
_ = ticker.tick() => {
{
let endpoints = balancer.endpoints.read();
for ep in endpoints.iter() {
let remaining = ep.cooldown_remaining_secs();
COOLDOWN_SECONDS_GAUGE
.with_label_values(&[&ep.name])
.set(remaining);
}
}
balancer.update_healthy_count().await;
}
}
}
}