Skip to main content

systemprompt_database/lifecycle/
validation.rs

1//! Pre-flight validation helpers used by the boot path and tests.
2//!
3//! Copyright (c) systemprompt.io — Business Source License 1.1.
4//! See <https://systemprompt.io> for licensing details.
5
6use crate::error::{DatabaseResult, RepositoryError};
7use crate::services::{Database, DatabaseProvider};
8
9pub async fn validate_database_connection(db: &dyn DatabaseProvider) -> DatabaseResult<()> {
10    db.test_connection()
11        .await
12        .map_err(|e| RepositoryError::Connection(Box::new(e)))
13}
14
15pub async fn validate_write_pool_is_primary(db: &Database) -> DatabaseResult<()> {
16    let result = db
17        .write()
18        .query_raw(&"SELECT pg_is_in_recovery() as in_recovery")
19        .await?;
20
21    let in_recovery = result
22        .first()
23        .and_then(|row| row.get("in_recovery"))
24        .and_then(serde_json::Value::as_bool)
25        .ok_or_else(|| {
26            RepositoryError::invalid_data("in_recovery", "pg_is_in_recovery() returned no boolean")
27        })?;
28
29    if !in_recovery {
30        return Ok(());
31    }
32
33    Err(if db.has_write_pool() {
34        RepositoryError::invalid_argument(
35            "database_write_url",
36            "points at a read-only standby. Writes, migrations and LISTEN/NOTIFY all require \
37             the primary — point it at the primary and restart",
38        )
39    } else {
40        RepositoryError::invalid_argument(
41            "database_url",
42            "points at a read-only standby and no `database_write_url` is set, so the write \
43             pool falls back to it. Set `database_write_url` (or `DATABASE_WRITE_URL` with the \
44             env secrets source) to the primary and restart",
45        )
46    })
47}
48
49#[derive(Debug, Clone, Copy, PartialEq)]
50pub struct ReplicaStatus {
51    pub in_recovery: bool,
52    pub replay_lag_secs: Option<f64>,
53}
54
55pub async fn replica_status(db: &dyn DatabaseProvider) -> DatabaseResult<ReplicaStatus> {
56    let result = db
57        .query_raw(
58            &"SELECT pg_is_in_recovery() AS in_recovery, CASE WHEN pg_is_in_recovery() THEN \
59              EXTRACT(EPOCH FROM (now() - pg_last_xact_replay_timestamp()))::double precision \
60              ELSE NULL END AS lag_secs",
61        )
62        .await?;
63    let row = result.first().ok_or_else(|| {
64        RepositoryError::invalid_data("in_recovery", "replica status probe returned no row")
65    })?;
66    let in_recovery = row
67        .get("in_recovery")
68        .and_then(serde_json::Value::as_bool)
69        .ok_or_else(|| {
70            RepositoryError::invalid_data("in_recovery", "replica status probe returned no boolean")
71        })?;
72    let replay_lag_secs = row.get("lag_secs").and_then(serde_json::Value::as_f64);
73    Ok(ReplicaStatus {
74        in_recovery,
75        replay_lag_secs,
76    })
77}
78
79pub async fn validate_table_exists(
80    db: &dyn DatabaseProvider,
81    table_name: &str,
82) -> DatabaseResult<bool> {
83    let result = db
84        .query_raw_with(
85            &"SELECT EXISTS (SELECT 1 FROM information_schema.tables WHERE table_schema = \
86              'public' AND table_name = $1) as exists",
87            &[&table_name],
88        )
89        .await?;
90
91    result
92        .first()
93        .and_then(|row| row.get("exists"))
94        .and_then(serde_json::Value::as_bool)
95        .ok_or_else(|| {
96            RepositoryError::invalid_data(
97                "exists",
98                format!("table existence probe for '{table_name}' returned no boolean"),
99            )
100        })
101}
102
103pub async fn validate_column_exists(
104    db: &dyn DatabaseProvider,
105    table_name: &str,
106    column_name: &str,
107) -> DatabaseResult<bool> {
108    let result = db
109        .query_raw_with(
110            &"SELECT EXISTS (SELECT 1 FROM information_schema.columns WHERE table_schema = \
111              'public' AND table_name = $1 AND column_name = $2) as exists",
112            &[&table_name, &column_name],
113        )
114        .await?;
115
116    result
117        .first()
118        .and_then(|row| row.get("exists"))
119        .and_then(serde_json::Value::as_bool)
120        .ok_or_else(|| {
121            RepositoryError::invalid_data(
122                "exists",
123                format!(
124                    "column existence probe for '{table_name}.{column_name}' returned no boolean"
125                ),
126            )
127        })
128}