Skip to main content

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}