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}