use std::sync::atomic::Ordering;
use std::time::Duration;
use crate::server::state::AppState;
pub fn start_cluster_tasks(state: &AppState) -> Vec<tokio::task::JoinHandle<()>> {
if !state.cluster.enabled {
return Vec::new();
}
vec![start_epoch_watcher(state.clone())]
}
fn start_epoch_watcher(state: AppState) -> tokio::task::JoinHandle<()> {
tokio::spawn(async move {
let mut interval = tokio::time::interval(Duration::from_millis(
state.config.cluster.epoch_poll_interval_ms,
));
interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
interval.tick().await;
tracing::info!(
poll_interval_ms = state.config.cluster.epoch_poll_interval_ms,
"Epoch watcher started"
);
loop {
interval.tick().await;
let row = match state.cluster.repo.get_epoch().await {
Ok(row) => row,
Err(e) => {
crate::metrics::record_error("epoch_watcher");
tracing::warn!(error = %e, "Epoch watcher: failed to read config epoch");
continue;
}
};
let tick_ok = advance_config_epoch(&state.cluster.last_seen_epoch, row.epoch, || {
crate::engine::resync_from_db(&state)
})
.await;
let last_breaker = state
.cluster
.last_seen_breaker_epoch
.load(Ordering::Acquire);
if row.breaker_epoch > last_breaker {
if !row.breaker_key.is_empty() {
let found = state
.connector_registry
.reset_circuit_breaker(&row.breaker_key)
.await;
tracing::info!(
key = %row.breaker_key,
found,
"Applied cluster-wide circuit breaker reset"
);
}
state
.cluster
.last_seen_breaker_epoch
.fetch_max(row.breaker_epoch, Ordering::AcqRel);
}
if tick_ok {
crate::metrics::record_job_success("epoch_watcher");
}
}
})
}
async fn advance_config_epoch<F, Fut>(
last_seen: &std::sync::atomic::AtomicI64,
epoch: i64,
resync: F,
) -> bool
where
F: FnOnce() -> Fut,
Fut: std::future::Future<Output = Result<(), crate::errors::OrionError>>,
{
let last = last_seen.load(Ordering::Acquire);
if epoch <= last {
return true;
}
tracing::info!(
from = last,
to = epoch,
"Config epoch advanced on another node; resyncing from DB"
);
match resync().await {
Ok(()) => {
last_seen.fetch_max(epoch, Ordering::AcqRel);
true
}
Err(e) => {
crate::metrics::record_error("epoch_watcher");
tracing::warn!(error = %e, "Epoch watcher: resync failed; will retry");
false
}
}
}
#[cfg(test)]
mod tests {
use super::advance_config_epoch;
use std::sync::atomic::{AtomicBool, AtomicI64, Ordering};
#[tokio::test]
async fn a_failed_resync_does_not_advance_the_watermark() {
let last_seen = AtomicI64::new(5);
let ok = advance_config_epoch(&last_seen, 7, || async {
Err(crate::errors::OrionError::internal("resync exploded"))
})
.await;
assert!(!ok, "a failed resync is a failed tick");
assert_eq!(
last_seen.load(Ordering::Acquire),
5,
"the watermark must not record a config that was never applied"
);
let ok = advance_config_epoch(&last_seen, 7, || async { Ok(()) }).await;
assert!(ok);
assert_eq!(last_seen.load(Ordering::Acquire), 7);
}
#[tokio::test]
async fn an_unchanged_epoch_never_resyncs() {
let last_seen = AtomicI64::new(7);
let ran = AtomicBool::new(false);
let ok = advance_config_epoch(&last_seen, 7, || {
ran.store(true, Ordering::SeqCst);
async { Ok(()) }
})
.await;
assert!(ok, "nothing to do is a successful tick");
assert!(!ran.load(Ordering::SeqCst), "resync must not have run");
assert_eq!(last_seen.load(Ordering::Acquire), 7);
}
}