turnframe_store_postgres/store.rs
1//! [`PgStores`]: one PostgreSQL pool behind every persistence trait.
2
3use std::fmt;
4use std::sync::Arc;
5
6use sqlx::pool::PoolConnection;
7use sqlx::postgres::{PgConnectOptions, PgPool, PgPoolOptions};
8use sqlx::{Executor, Postgres, Transaction};
9use turnframe_store::error::StoreError;
10use turnframe_store::stores::{Stores, StoresBuilderError};
11
12use crate::config::PgStoreConfig;
13use crate::error::{commit_failed, store_error};
14
15/// The migrations of this crate, embedded in the binary at compile time.
16///
17/// `sqlx::migrate!` reads `migrations/` while the crate is being compiled and
18/// bakes the files in, so a deployed binary needs no directory beside it and no
19/// database while it is being built. Applying them is
20/// [`PgStores::migrate`]; running them from your own pool is
21/// `MIGRATOR.run(&pool)`.
22pub static MIGRATOR: sqlx::migrate::Migrator = sqlx::migrate!("./migrations");
23
24/// Every Turnframe store over one PostgreSQL pool.
25///
26/// Implements [`ConversationStore`](turnframe_store::conversation::ConversationStore),
27/// [`InteractionStore`](turnframe_store::interaction::InteractionStore),
28/// [`CommandJournal`](turnframe_store::journal::CommandJournal),
29/// [`EventJournal`](turnframe_store::events::EventJournal),
30/// [`OutboxStore`](turnframe_store::outbox::OutboxStore),
31/// [`ReplayStore`](turnframe_store::replay::ReplayStore) and
32/// [`CommitStore`](turnframe_store::commit::CommitStore), so a card written
33/// through the interaction trait is the card a commit bundle settles.
34///
35/// Cloning shares the pool; it never opens new connections.
36#[derive(Clone)]
37pub struct PgStores {
38 pool: PgPool,
39 schema: Option<String>,
40}
41
42impl PgStores {
43 /// Opens a pool on `url` with the default [`PgStoreConfig`].
44 ///
45 /// The URL is a standard PostgreSQL connection string
46 /// (`postgres://user:password@host:port/database`). It is never stored in a
47 /// field this type can print.
48 ///
49 /// # Errors
50 /// * `Other` when the URL cannot be parsed, `Unavailable` when the first
51 /// connection cannot be opened.
52 pub async fn connect(url: &str) -> Result<Self, StoreError> {
53 Self::connect_with(url, &PgStoreConfig::new()).await
54 }
55
56 /// Opens a pool on `url` with `config`.
57 ///
58 /// # Errors
59 /// * `Other` when the URL cannot be parsed, `Unavailable` when the first
60 /// connection cannot be opened.
61 pub async fn connect_with(url: &str, config: &PgStoreConfig) -> Result<Self, StoreError> {
62 let options: PgConnectOptions = url.parse().map_err(|error| store_error(&error))?;
63 let pool = config
64 .apply_pool(PgPoolOptions::new())
65 .connect_with(config.apply_connection(options))
66 .await
67 .map_err(|error| store_error(&error))?;
68 Ok(Self {
69 pool,
70 schema: config.schema_name().map(ToOwned::to_owned),
71 })
72 }
73
74 /// Uses a pool the caller already has.
75 ///
76 /// Nothing is configured on it: the pool's own settings, its `search_path`
77 /// and its statement timeout are whatever the caller set. Use this to share
78 /// one pool with the rest of an application, so the store and the domain
79 /// tables are reached through the same connections and can be enlisted in
80 /// the same transaction (see [`Self::commit_in`]).
81 #[must_use]
82 pub fn from_pool(pool: PgPool) -> Self {
83 Self { pool, schema: None }
84 }
85
86 /// Puts the tables of a caller-supplied pool in `schema`.
87 ///
88 /// Only tells [`Self::migrate`] which schema to create; the pool must
89 /// already resolve to it, normally through a `search_path` in its connect
90 /// options.
91 ///
92 /// # Errors
93 /// * [`ConfigError::InvalidSchemaName`](crate::ConfigError::InvalidSchemaName)
94 /// when the name is not a plain unquoted identifier.
95 pub fn with_schema(mut self, schema: impl Into<String>) -> Result<Self, crate::ConfigError> {
96 let config = PgStoreConfig::new().schema(schema)?;
97 self.schema = config.schema_name().map(ToOwned::to_owned);
98 Ok(self)
99 }
100
101 /// Applies every migration this crate carries, creating the configured
102 /// schema first when there is one.
103 ///
104 /// It is safe to call on every start-up: `sqlx` records applied versions and
105 /// skips them, and the statements are written to be applied twice anyway.
106 /// Two processes racing it is safe too — the migrator takes a database-wide
107 /// advisory lock — though a deployment normally runs it once, before the
108 /// instances that will use the schema come up.
109 ///
110 /// # Errors
111 /// * [`MigrateError`](sqlx::migrate::MigrateError) when the schema cannot be
112 /// created, a migration fails, or an already-applied migration no longer
113 /// matches its recorded checksum.
114 pub async fn migrate(&self) -> Result<(), sqlx::migrate::MigrateError> {
115 if let Some(schema) = &self.schema {
116 // The name is a validated plain identifier, which is why it can be
117 // interpolated: `CREATE SCHEMA` takes no bind parameter.
118 let statement = format!("CREATE SCHEMA IF NOT EXISTS {schema}");
119 match self.pool.execute(statement.as_str()).await {
120 Ok(_) => {}
121 // `IF NOT EXISTS` checks and creates in two steps, so two
122 // instances starting at the same moment can both pass the check
123 // and one of them loses. The schema exists either way, which is
124 // all this call was asking for.
125 Err(error) if created_concurrently(&error) => {}
126 Err(error) => return Err(error.into()),
127 }
128 }
129 tracing::debug!(
130 migrations = MIGRATOR.iter().len(),
131 "applying turnframe store migrations"
132 );
133 MIGRATOR.run(&self.pool).await
134 }
135
136 /// The pool underneath, for health checks, metrics, or statements this
137 /// adapter does not make.
138 #[must_use]
139 pub fn pool(&self) -> &PgPool {
140 &self.pool
141 }
142
143 /// The schema the tables live in, when one was configured.
144 #[must_use]
145 pub fn schema(&self) -> Option<&str> {
146 self.schema.as_deref()
147 }
148
149 /// Closes the pool and waits for its connections to be released.
150 pub async fn close(&self) {
151 self.pool.close().await;
152 }
153
154 /// Bundles this store behind all seven traits.
155 ///
156 /// # Errors
157 /// * [`StoresBuilderError`] — never in practice. Every role is supplied
158 /// here, and the `Result` exists only because `turnframe-store` offers no
159 /// infallible "one backend for every role" constructor for a backend
160 /// other than its own in-memory one.
161 pub fn stores(&self) -> Result<Stores, StoresBuilderError> {
162 let backend = Arc::new(self.clone());
163 Stores::builder()
164 .conversations(backend.clone())
165 .interactions(backend.clone())
166 .journal(backend.clone())
167 .events(backend.clone())
168 .outbox(backend.clone())
169 .replay(backend.clone())
170 .commit(backend)
171 .build()
172 }
173
174 /// Takes a connection from the pool for a single read.
175 pub(crate) async fn connection(&self) -> Result<PoolConnection<Postgres>, StoreError> {
176 self.pool
177 .acquire()
178 .await
179 .map_err(|error| store_error(&error))
180 }
181
182 /// Begins the transaction every write of this adapter runs inside.
183 pub(crate) async fn transaction(&self) -> Result<Transaction<'_, Postgres>, StoreError> {
184 self.pool.begin().await.map_err(|error| store_error(&error))
185 }
186}
187
188/// Returns `true` when a `CREATE SCHEMA IF NOT EXISTS` lost a race with another
189/// process running the same statement.
190fn created_concurrently(error: &sqlx::Error) -> bool {
191 /// `duplicate_schema`.
192 const DUPLICATE_SCHEMA: &str = "42P06";
193 /// The unique violation on `pg_namespace` a tighter race raises instead.
194 const UNIQUE_VIOLATION: &str = "23505";
195
196 error
197 .as_database_error()
198 .and_then(|database| database.code())
199 .is_some_and(|code| code == DUPLICATE_SCHEMA || code == UNIQUE_VIOLATION)
200}
201
202/// Commits a transaction, classifying a commit that does not answer as
203/// [`StoreError::Timeout`].
204pub(crate) async fn commit(transaction: Transaction<'_, Postgres>) -> Result<(), StoreError> {
205 transaction
206 .commit()
207 .await
208 .map_err(|error| commit_failed(&error))
209}
210
211impl fmt::Debug for PgStores {
212 /// Names the schema and the pool's size, never the connection string: this
213 /// output is safe to log.
214 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
215 f.debug_struct("PgStores")
216 .field("schema", &self.schema.as_deref().unwrap_or("<default>"))
217 .field("connections", &self.pool.size())
218 .field("closed", &self.pool.is_closed())
219 .finish()
220 }
221}
222
223#[cfg(test)]
224mod tests {
225 use super::*;
226
227 #[test]
228 fn the_migrations_are_embedded() {
229 assert_eq!(MIGRATOR.iter().len(), 2, "both migrations are embedded");
230 assert!(
231 MIGRATOR
232 .iter()
233 .all(|migration| !migration.sql.trim().is_empty())
234 );
235 }
236
237 #[tokio::test]
238 async fn a_schema_name_is_validated_before_it_reaches_ddl() {
239 let store = PgStores::from_pool(PgPool::connect_lazy("postgres://x/y").unwrap());
240 assert!(store.clone().with_schema("tf_test").is_ok());
241 assert!(store.with_schema("tf\";DROP SCHEMA public").is_err());
242 }
243
244 #[tokio::test]
245 async fn debug_output_carries_no_connection_string() {
246 let store =
247 PgStores::from_pool(PgPool::connect_lazy("postgres://secret:hunter2@host/db").unwrap());
248 let rendered = format!("{store:?}");
249 assert!(!rendered.contains("hunter2"), "{rendered}");
250 assert!(rendered.contains("<default>"), "{rendered}");
251 }
252}