systemprompt_database/repository/service/
maintenance.rs1use 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}