systemprompt_database/repository/service/
claim.rs1use sqlx::{Connection, PgConnection};
15use systemprompt_identifiers::InstanceId;
16use systemprompt_traits::RepositoryError;
17use tracing::warn;
18
19use super::repo::ServiceRepository;
20
21pub const INSTANCE_CLAIM_CLASS: i32 = 0x5350_0001;
25
26#[derive(Debug, thiserror::Error)]
28pub enum InstanceClaimError {
29 #[error(
30 "instance id '{instance_id}' is already claimed by a live process; give each replica its \
31 own server.instance_id or leave it unset so HOSTNAME is used"
32 )]
33 Claimed { instance_id: InstanceId },
34
35 #[error("could not claim instance id: {0}")]
36 Repository(#[from] RepositoryError),
37}
38
39impl From<sqlx::Error> for InstanceClaimError {
40 fn from(err: sqlx::Error) -> Self {
41 Self::Repository(RepositoryError::from(err))
42 }
43}
44
45pub struct InstanceClaim {
47 session: Option<PgConnection>,
48 instance_id: InstanceId,
49}
50
51impl InstanceClaim {
52 pub const fn instance_id(&self) -> &InstanceId {
53 &self.instance_id
54 }
55
56 pub async fn release(mut self) {
57 if let Some(mut session) = self.session.take() {
58 if let Err(e) = sqlx::query_scalar!(
59 "SELECT pg_advisory_unlock($1, hashtext($2))",
60 INSTANCE_CLAIM_CLASS,
61 self.instance_id.as_str()
62 )
63 .fetch_one(&mut session)
64 .await
65 {
66 warn!(
67 instance_id = %self.instance_id,
68 error = %e,
69 "Failed to release instance claim; closing its session"
70 );
71 }
72 if let Err(e) = session.close().await {
73 warn!(
74 instance_id = %self.instance_id,
75 error = %e,
76 "Failed to close instance claim session"
77 );
78 }
79 }
80 }
81}
82
83impl std::fmt::Debug for InstanceClaim {
84 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
85 f.debug_struct("InstanceClaim")
86 .field("instance_id", &self.instance_id)
87 .field("held", &self.session.is_some())
88 .finish()
89 }
90}
91
92impl ServiceRepository {
93 pub async fn claim_instance(&self) -> Result<InstanceClaim, InstanceClaimError> {
94 let mut session = self.write_pool.acquire().await?.detach();
95 let acquired = sqlx::query_scalar!(
96 r#"SELECT pg_try_advisory_lock($1, hashtext($2)) AS "acquired!""#,
97 INSTANCE_CLAIM_CLASS,
98 self.instance_id.as_str()
99 )
100 .fetch_one(&mut session)
101 .await?;
102
103 if !acquired {
104 if let Err(e) = session.close().await {
105 warn!(
106 instance_id = %self.instance_id,
107 error = %e,
108 "Failed to close refused claim session"
109 );
110 }
111 return Err(InstanceClaimError::Claimed {
112 instance_id: self.instance_id.clone(),
113 });
114 }
115
116 Ok(InstanceClaim {
117 session: Some(session),
118 instance_id: self.instance_id.clone(),
119 })
120 }
121}