use super::claim::INSTANCE_CLAIM_CLASS;
use super::repo::ServiceRepository;
use crate::error::DatabaseResult;
impl ServiceRepository {
pub async fn cleanup_stale_entries(&self) -> DatabaseResult<u64> {
let result = sqlx::query!(
r#"
DELETE FROM services
WHERE instance_id = $1
AND (status = 'error'
OR (status = 'running' AND pid IS NULL))
"#,
self.instance_id.as_str()
)
.execute(&*self.write_pool)
.await?;
Ok(result.rows_affected())
}
pub async fn touch_heartbeat(&self) -> DatabaseResult<u64> {
let result = sqlx::query!(
r#"UPDATE services SET heartbeat_at = CURRENT_TIMESTAMP WHERE instance_id = $1"#,
self.instance_id.as_str()
)
.execute(&*self.write_pool)
.await?;
Ok(result.rows_affected())
}
pub async fn delete_dead_instances(&self, older_than_secs: i64) -> DatabaseResult<u64> {
let result = sqlx::query!(
r#"
DELETE FROM services
WHERE heartbeat_at < CURRENT_TIMESTAMP - make_interval(secs => $1::double precision)
AND NOT EXISTS (
SELECT 1 FROM pg_locks l
WHERE l.locktype = 'advisory'
AND l.granted
AND l.classid = ($2::int)::oid
AND l.objid = (hashtext(services.instance_id))::oid
AND l.objsubid = 2
)
"#,
older_than_secs as f64,
INSTANCE_CLAIM_CLASS
)
.execute(&*self.write_pool)
.await?;
Ok(result.rows_affected())
}
}