use std::sync::Arc;
use std::time::Duration;
use systemprompt_database::ServiceRepository;
use systemprompt_runtime::AppContext;
pub(in crate::services::server) const HEARTBEAT_INTERVAL: Duration = Duration::from_secs(15);
const GC_EVERY_BEATS: u32 = 4;
const DEAD_AFTER_SECS: i64 = 90;
pub(in crate::services::server) fn start_registry_heartbeat(ctx: &AppContext) {
let repository = Arc::clone(ctx.service_repository());
ctx.background_tasks()
.spawn_cancellable("registry_heartbeat", |cancel| async move {
let mut tick = tokio::time::interval(HEARTBEAT_INTERVAL);
tick.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
let mut beats: u32 = 0;
loop {
tokio::select! {
() = cancel.cancelled() => break,
_ = tick.tick() => {},
}
beat(&repository, &mut beats).await;
}
});
}
async fn beat(repository: &ServiceRepository, beats: &mut u32) {
if let Err(error) = repository.touch_heartbeat().await {
tracing::warn!(
instance_id = %repository.instance_id(),
error = %error,
"service registry heartbeat failed"
);
}
*beats = beats.wrapping_add(1);
if !beats.is_multiple_of(GC_EVERY_BEATS) {
return;
}
match repository.delete_dead_instances(DEAD_AFTER_SECS).await {
Ok(0) => {},
Ok(reaped) => tracing::info!(reaped, "service registry reaped dead instances"),
Err(error) => {
tracing::warn!(error = %error, "service registry reap failed");
},
}
}