Skip to main content

systemprompt_database/repository/service/
maintenance.rs

1//! Liveness maintenance on the `services` registry: heartbeats, stale-row
2//! cleanup for this instance, and the cross-instance reap of dead replicas.
3//!
4//! Copyright (c) systemprompt.io — Business Source License 1.1.
5//! See <https://systemprompt.io> for licensing details.
6
7use super::claim::INSTANCE_CLAIM_CLASS;
8use super::repo::ServiceRepository;
9use crate::error::DatabaseResult;
10
11impl ServiceRepository {
12    pub async fn cleanup_stale_entries(&self) -> DatabaseResult<u64> {
13        let result = sqlx::query!(
14            r#"
15            DELETE FROM services
16            WHERE instance_id = $1
17              AND (status = 'error'
18                   OR (status = 'running' AND pid IS NULL))
19            "#,
20            self.instance_id.as_str()
21        )
22        .execute(&*self.write_pool)
23        .await?;
24        Ok(result.rows_affected())
25    }
26
27    pub async fn touch_heartbeat(&self) -> DatabaseResult<u64> {
28        let result = sqlx::query!(
29            r#"UPDATE services SET heartbeat_at = CURRENT_TIMESTAMP WHERE instance_id = $1"#,
30            self.instance_id.as_str()
31        )
32        .execute(&*self.write_pool)
33        .await?;
34        Ok(result.rows_affected())
35    }
36
37    pub async fn delete_dead_instances(&self, older_than_secs: i64) -> DatabaseResult<u64> {
38        let result = sqlx::query!(
39            r#"
40            DELETE FROM services
41            WHERE heartbeat_at < CURRENT_TIMESTAMP - make_interval(secs => $1::double precision)
42              AND NOT EXISTS (
43                  SELECT 1 FROM pg_locks l
44                  WHERE l.locktype = 'advisory'
45                    AND l.granted
46                    AND l.classid = ($2::int)::oid
47                    AND l.objid = (hashtext(services.instance_id))::oid
48                    AND l.objsubid = 2
49              )
50            "#,
51            older_than_secs as f64,
52            INSTANCE_CLAIM_CLASS
53        )
54        .execute(&*self.write_pool)
55        .await?;
56        Ok(result.rows_affected())
57    }
58}