Skip to main content

turnframe_store_postgres/
config.rs

1//! Pool settings, the optional dedicated schema, and the statement-timeout
2//! advisory.
3
4use std::fmt;
5use std::time::Duration;
6
7/// How a [`PgStores`](crate::PgStores) opens and shapes its connection pool.
8///
9/// The defaults are meant for a service that handles turns: a pool small enough
10/// that PostgreSQL is not the thing that falls over first, an acquire timeout
11/// short enough that a saturated pool surfaces as
12/// [`Unavailable`](turnframe_store::error::StoreError::Unavailable) instead of a
13/// hung request, and connections recycled often enough that a rolling database
14/// upgrade drains cleanly.
15///
16/// ```rust
17/// use std::time::Duration;
18///
19/// use turnframe_store_postgres::PgStoreConfig;
20///
21/// # fn main() -> Result<(), turnframe_store_postgres::ConfigError> {
22/// let config = PgStoreConfig::new()
23///     .max_connections(16)
24///     .statement_timeout(Duration::from_secs(5))
25///     .schema("turnframe")?;
26/// assert_eq!(config.schema_name(), Some("turnframe"));
27/// # Ok(())
28/// # }
29/// ```
30#[derive(Debug, Clone, PartialEq, Eq)]
31pub struct PgStoreConfig {
32    max_connections: u32,
33    min_connections: u32,
34    acquire_timeout: Duration,
35    idle_timeout: Option<Duration>,
36    max_lifetime: Option<Duration>,
37    statement_timeout: Option<Duration>,
38    schema: Option<String>,
39}
40
41impl PgStoreConfig {
42    /// Upper bound on pooled connections. PostgreSQL serves every connection
43    /// with a backend process, so this is a budget shared with every other
44    /// service on the same cluster.
45    pub const DEFAULT_MAX_CONNECTIONS: u32 = 10;
46    /// Connections kept open while idle, so a quiet period does not make the
47    /// next turn pay for a handshake.
48    pub const DEFAULT_MIN_CONNECTIONS: u32 = 1;
49    /// How long a caller waits for a connection before the store reports
50    /// `Unavailable`.
51    pub const DEFAULT_ACQUIRE_TIMEOUT: Duration = Duration::from_secs(5);
52    /// How long an unused connection is kept.
53    pub const DEFAULT_IDLE_TIMEOUT: Duration = Duration::from_secs(600);
54    /// How long any connection is kept, however busy. Recycling bounds the
55    /// damage of a leaked server-side state and lets a rolling upgrade drain.
56    pub const DEFAULT_MAX_LIFETIME: Duration = Duration::from_secs(1800);
57
58    /// The defaults described above, on the `public` schema, with no statement
59    /// timeout of the adapter's own.
60    #[must_use]
61    pub fn new() -> Self {
62        Self {
63            max_connections: Self::DEFAULT_MAX_CONNECTIONS,
64            min_connections: Self::DEFAULT_MIN_CONNECTIONS,
65            acquire_timeout: Self::DEFAULT_ACQUIRE_TIMEOUT,
66            idle_timeout: Some(Self::DEFAULT_IDLE_TIMEOUT),
67            max_lifetime: Some(Self::DEFAULT_MAX_LIFETIME),
68            statement_timeout: None,
69            schema: None,
70        }
71    }
72
73    /// Sets the maximum number of pooled connections.
74    #[must_use]
75    pub fn max_connections(mut self, connections: u32) -> Self {
76        self.max_connections = connections;
77        self
78    }
79
80    /// Sets the number of connections kept open while idle.
81    #[must_use]
82    pub fn min_connections(mut self, connections: u32) -> Self {
83        self.min_connections = connections;
84        self
85    }
86
87    /// Sets how long a caller waits for a connection from the pool.
88    ///
89    /// This is *not* a statement timeout: it bounds the wait for a connection,
90    /// not the work done once one is held. See [`Self::statement_timeout`].
91    #[must_use]
92    pub fn acquire_timeout(mut self, timeout: Duration) -> Self {
93        self.acquire_timeout = timeout;
94        self
95    }
96
97    /// Sets how long an unused connection is kept; `None` keeps it forever.
98    #[must_use]
99    pub fn idle_timeout(mut self, timeout: Option<Duration>) -> Self {
100        self.idle_timeout = timeout;
101        self
102    }
103
104    /// Sets how long any connection is kept; `None` keeps it forever.
105    #[must_use]
106    pub fn max_lifetime(mut self, lifetime: Option<Duration>) -> Self {
107        self.max_lifetime = lifetime;
108        self
109    }
110
111    /// Sets `statement_timeout` on every connection this pool opens.
112    ///
113    /// # The advisory
114    ///
115    /// A statement timeout is the only thing that bounds a query the database
116    /// decides to run slowly, and without one a single lock wait can hold a
117    /// pooled connection until the pool is empty and every turn fails. Set one.
118    ///
119    /// Set it with these three consequences in mind.
120    ///
121    /// *It cancels statements, not transactions.* PostgreSQL raises
122    /// `query_canceled` on the statement that ran too long; this adapter reports
123    /// [`Timeout`](turnframe_store::error::StoreError::Timeout) and abandons the
124    /// transaction, so a cancelled statement inside
125    /// [`CommitStore::commit`](turnframe_store::commit::CommitStore::commit)
126    /// discards the whole bundle. That is the correct outcome — a partially
127    /// written bundle is an invariant violation — but it means the timeout must
128    /// be generous enough for the largest bundle a turn produces, not for the
129    /// median statement.
130    ///
131    /// *`Timeout` is not a signal to retry.* The contract reads it as "the write
132    /// may have landed": the caller re-reads and resumes by idempotency key. A
133    /// timeout set so tight that healthy commits trip it turns every one of them
134    /// into a recovery.
135    ///
136    /// *The outbox claim is the exception to keep short.*
137    /// [`claim_due`](turnframe_store::outbox::OutboxWriter::claim_due) takes row
138    /// locks with `SKIP LOCKED`, so it never waits on another dispatcher; if it
139    /// is slow, something else is wrong and cutting it off is right.
140    ///
141    /// A few seconds suits a service handling turns. Leave it `None` to inherit
142    /// whatever the role or the server sets, which is the better choice when the
143    /// database is administered separately.
144    #[must_use]
145    pub fn statement_timeout(mut self, timeout: Duration) -> Self {
146        self.statement_timeout = Some(timeout);
147        self
148    }
149
150    /// Clears the adapter's statement timeout, inheriting the server's.
151    #[must_use]
152    pub fn inherit_statement_timeout(mut self) -> Self {
153        self.statement_timeout = None;
154        self
155    }
156
157    /// Puts every table in a dedicated schema instead of the connection's
158    /// default `search_path`.
159    ///
160    /// [`PgStores::migrate`](crate::PgStores::migrate) creates the schema when
161    /// it is missing, so pointing a fresh deployment at an empty database is one
162    /// call. Tests use it to give each run a private namespace inside one
163    /// database.
164    ///
165    /// # Errors
166    /// * [`ConfigError::InvalidSchemaName`] when the name is not a plain
167    ///   unquoted PostgreSQL identifier. The name is interpolated into `CREATE
168    ///   SCHEMA` and into `search_path`, neither of which can take a bind
169    ///   parameter, so it is validated here rather than escaped later.
170    pub fn schema(mut self, schema: impl Into<String>) -> Result<Self, ConfigError> {
171        let schema = schema.into();
172        if !is_plain_identifier(&schema) {
173            return Err(ConfigError::InvalidSchemaName);
174        }
175        self.schema = Some(schema);
176        Ok(self)
177    }
178
179    /// The configured schema, if any.
180    #[must_use]
181    pub fn schema_name(&self) -> Option<&str> {
182        self.schema.as_deref()
183    }
184
185    /// The configured statement timeout, if any.
186    #[must_use]
187    pub fn statement_timeout_value(&self) -> Option<Duration> {
188        self.statement_timeout
189    }
190
191    /// The configured pool bounds, as `(min, max)`.
192    #[must_use]
193    pub fn connection_bounds(&self) -> (u32, u32) {
194        (self.min_connections, self.max_connections)
195    }
196
197    /// Applies the pool settings to a `sqlx` pool builder.
198    pub(crate) fn apply_pool(
199        &self,
200        options: sqlx::postgres::PgPoolOptions,
201    ) -> sqlx::postgres::PgPoolOptions {
202        options
203            .max_connections(self.max_connections)
204            .min_connections(self.min_connections)
205            .acquire_timeout(self.acquire_timeout)
206            .idle_timeout(self.idle_timeout)
207            .max_lifetime(self.max_lifetime)
208    }
209
210    /// Applies the per-connection settings to `sqlx` connect options.
211    pub(crate) fn apply_connection(
212        &self,
213        mut options: sqlx::postgres::PgConnectOptions,
214    ) -> sqlx::postgres::PgConnectOptions {
215        if let Some(schema) = &self.schema {
216            options = options.options([("search_path", schema.as_str())]);
217        }
218        if let Some(timeout) = self.statement_timeout {
219            let millis = timeout.as_millis().to_string();
220            options = options.options([("statement_timeout", millis.as_str())]);
221        }
222        options
223    }
224}
225
226impl Default for PgStoreConfig {
227    fn default() -> Self {
228        Self::new()
229    }
230}
231
232/// Why a [`PgStoreConfig`] was refused.
233#[derive(Debug, Clone, Copy, PartialEq, Eq, thiserror::Error)]
234#[non_exhaustive]
235pub enum ConfigError {
236    /// The schema name is not a plain unquoted PostgreSQL identifier: it must
237    /// start with a letter or an underscore, continue with letters, digits or
238    /// underscores, and fit in 63 bytes.
239    #[error("schema name is not a plain unquoted PostgreSQL identifier")]
240    InvalidSchemaName,
241}
242
243/// Returns `true` for a name that needs neither quoting nor escaping.
244fn is_plain_identifier(name: &str) -> bool {
245    /// PostgreSQL truncates identifiers past `NAMEDATALEN - 1`.
246    const MAX_IDENTIFIER_BYTES: usize = 63;
247
248    let mut characters = name.chars();
249    let Some(first) = characters.next() else {
250        return false;
251    };
252    name.len() <= MAX_IDENTIFIER_BYTES
253        && (first.is_ascii_lowercase() || first.is_ascii_uppercase() || first == '_')
254        && characters.all(|character| character.is_ascii_alphanumeric() || character == '_')
255}
256
257impl fmt::Display for PgStoreConfig {
258    /// Names the settings, never a connection string: this value is safe to log.
259    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
260        write!(
261            f,
262            "pool {}..{}, schema {}",
263            self.min_connections,
264            self.max_connections,
265            self.schema.as_deref().unwrap_or("<default>")
266        )
267    }
268}
269
270#[cfg(test)]
271mod tests {
272    use super::*;
273
274    #[test]
275    fn defaults_are_the_documented_ones() {
276        let config = PgStoreConfig::new();
277        assert_eq!(config.connection_bounds(), (1, 10));
278        assert_eq!(config.schema_name(), None);
279        assert_eq!(config.statement_timeout_value(), None);
280        assert_eq!(config, PgStoreConfig::default());
281    }
282
283    #[test]
284    fn a_schema_name_must_be_a_plain_identifier() {
285        assert!(PgStoreConfig::new().schema("turnframe").is_ok());
286        assert!(PgStoreConfig::new().schema("_tf_test_1").is_ok());
287        for rejected in [
288            "",
289            "1leading_digit",
290            "has space",
291            "quote\"injection",
292            "semicolon;drop",
293            "dotted.name",
294            "unicodé",
295        ] {
296            assert_eq!(
297                PgStoreConfig::new().schema(rejected).unwrap_err(),
298                ConfigError::InvalidSchemaName,
299                "{rejected:?} must be refused"
300            );
301        }
302        let too_long = "a".repeat(64);
303        assert_eq!(
304            PgStoreConfig::new().schema(too_long).unwrap_err(),
305            ConfigError::InvalidSchemaName
306        );
307    }
308
309    #[test]
310    fn display_never_carries_a_connection_string() {
311        let config = PgStoreConfig::new().schema("turnframe").unwrap();
312        assert_eq!(config.to_string(), "pool 1..10, schema turnframe");
313        assert_eq!(
314            PgStoreConfig::new().to_string(),
315            "pool 1..10, schema <default>"
316        );
317    }
318
319    #[test]
320    fn timeouts_are_settable_and_clearable() {
321        let config = PgStoreConfig::new()
322            .statement_timeout(Duration::from_secs(3))
323            .acquire_timeout(Duration::from_secs(1))
324            .idle_timeout(None)
325            .max_lifetime(None)
326            .min_connections(2)
327            .max_connections(4);
328        assert_eq!(
329            config.statement_timeout_value(),
330            Some(Duration::from_secs(3))
331        );
332        assert_eq!(config.connection_bounds(), (2, 4));
333        assert_eq!(
334            config.inherit_statement_timeout().statement_timeout_value(),
335            None
336        );
337    }
338}