Skip to main content

systemprompt_database/lifecycle/installation/extension/
lock.rs

1//! Session-pinned Postgres advisory lock serialising concurrent bootstraps.
2//!
3//! Copyright (c) systemprompt.io — Business Source License 1.1.
4//! See <https://systemprompt.io> for licensing details.
5
6use sqlx::Postgres;
7use sqlx::pool::PoolConnection;
8use systemprompt_extension::LoaderError;
9use systemprompt_identifiers::ExtensionId;
10use tracing::{debug, warn};
11
12use crate::services::DatabaseProvider;
13
14pub const BOOTSTRAP_ADVISORY_LOCK_KEY: i64 = 0x73_70_72_6F_6D_70_74_01;
15
16/// Session-pinned advisory lock serialising concurrent bootstraps.
17///
18/// Only the acquiring Postgres session can release its advisory lock, so the
19/// guard pins that connection. Dropping the guard without `release` closes
20/// the session instead of returning it to the pool, so the lock never
21/// outlives a cancelled or panicking holder.
22pub struct BootstrapLockGuard {
23    conn: Option<PoolConnection<Postgres>>,
24}
25
26impl std::fmt::Debug for BootstrapLockGuard {
27    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
28        f.debug_struct("BootstrapLockGuard")
29            .field("key", &BOOTSTRAP_ADVISORY_LOCK_KEY)
30            .field("held", &self.conn.is_some())
31            .finish()
32    }
33}
34
35impl BootstrapLockGuard {
36    pub async fn acquire(db: &dyn DatabaseProvider) -> Result<Self, LoaderError> {
37        let mut conn = db.get_postgres_pool().acquire().await.map_err(|e| {
38            LoaderError::SchemaInstallationStepFailed {
39                extension: ExtensionId::new("database"),
40                context: "Failed to acquire bootstrap lock connection".to_owned(),
41                source: Box::new(e),
42            }
43        })?;
44
45        sqlx::query!("SELECT pg_advisory_lock($1)", BOOTSTRAP_ADVISORY_LOCK_KEY)
46            .execute(conn.as_mut())
47            .await
48            .map_err(|e| LoaderError::SchemaInstallationStepFailed {
49                extension: ExtensionId::new("database"),
50                context: "Failed to acquire bootstrap advisory lock".to_owned(),
51                source: Box::new(e),
52            })?;
53
54        debug!(
55            key = BOOTSTRAP_ADVISORY_LOCK_KEY,
56            "Acquired bootstrap advisory lock"
57        );
58
59        Ok(Self { conn: Some(conn) })
60    }
61
62    pub async fn release(mut self) {
63        let Some(mut conn) = self.conn.take() else {
64            return;
65        };
66        match sqlx::query_scalar!("SELECT pg_advisory_unlock($1)", BOOTSTRAP_ADVISORY_LOCK_KEY)
67            .fetch_one(conn.as_mut())
68            .await
69        {
70            Ok(Some(true)) => drop(conn),
71            Ok(released) => {
72                warn!(
73                    key = BOOTSTRAP_ADVISORY_LOCK_KEY,
74                    ?released,
75                    "Bootstrap advisory lock was not held by this session at release"
76                );
77                drop(conn);
78            },
79            Err(e) => {
80                warn!(
81                    error = %e,
82                    "Failed to release bootstrap advisory lock; closing its session instead of pooling it"
83                );
84                let session = conn.detach();
85                drop(session);
86            },
87        }
88    }
89}
90
91impl Drop for BootstrapLockGuard {
92    fn drop(&mut self) {
93        // Why: a connection returned to the pool keeps its session, and with it
94        // the advisory lock; detaching closes the session so a cancelled or
95        // panicking install cannot leave every other replica blocked.
96        if let Some(conn) = self.conn.take() {
97            warn!(
98                key = BOOTSTRAP_ADVISORY_LOCK_KEY,
99                "BootstrapLockGuard dropped without explicit release; closing its session"
100            );
101            let session = conn.detach();
102            drop(session);
103        }
104    }
105}