systemprompt_database/lifecycle/
validation.rs1use 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}