diff --git a/.env.example b/.env.example index f46d2f389..205e6610f 100644 --- a/.env.example +++ b/.env.example @@ -38,13 +38,19 @@ REDIS_URL=redis://localhost:6379 # READ_DATABASE_URL is set, reader (default 50). # BUZZ_DB_POOL_SIZE=50 -# Postgres statement_timeout and lock_timeout applied to every runtime -# connection. Accepts an integer (milliseconds) with an optional us/ms/s/min/h/d -# unit; `0` disables the limit. Postgres stores both as int milliseconds, so the -# ceiling is 2147483647ms (~24 days); anything above it, or otherwise malformed, -# falls back to the default. Schema migrations always run with both lifted. +# Per-session Postgres limits applied to every runtime connection, so no single +# caller can hold a pooled connection indefinitely. Each accepts an integer +# (milliseconds) with an optional us/ms/s/min/h/d unit; `0` disables that limit. +# Postgres stores all three as int milliseconds, so the ceiling is 2147483647ms +# (~24 days); anything above it, or otherwise malformed, falls back to the +# default. Schema migrations always run with all three lifted. +# +# IDLE_IN_TRANSACTION covers only the gaps between a transaction's statements — +# it reclaims a session that opened a transaction and stopped working, which +# STATEMENT_TIMEOUT cannot see because no statement is running. # BUZZ_DB_STATEMENT_TIMEOUT=30s # BUZZ_DB_LOCK_TIMEOUT=5s +# BUZZ_DB_IDLE_IN_TRANSACTION_TIMEOUT=60s # ----------------------------------------------------------------------------- # Typesense (search) diff --git a/crates/buzz-admin/src/main.rs b/crates/buzz-admin/src/main.rs index 07695c4d5..381cf2252 100644 --- a/crates/buzz-admin/src/main.rs +++ b/crates/buzz-admin/src/main.rs @@ -25,7 +25,7 @@ use std::sync::Arc; use anyhow::Result; use buzz_core::kind::KIND_NIP43_MEMBERSHIP_LIST; use buzz_core::tenant::{relay_url_authority, TenantContext}; -use buzz_db::{Db, DbConfig, PgTimeout}; +use buzz_db::{Db, DbConfig, RuntimeTimeouts}; use buzz_pubsub::{EventTopic, PubSubManager}; use clap::{Parser, Subcommand}; use nostr::{EventBuilder, Keys, Kind, Tag}; @@ -426,11 +426,7 @@ async fn connect_db() -> Result { // env vars still apply for an operator who wants a bound here. let db = Db::new(&DbConfig { database_url: db_url, - statement_timeout: PgTimeout::from_env_or( - "BUZZ_DB_STATEMENT_TIMEOUT", - PgTimeout::disabled(), - ), - lock_timeout: PgTimeout::from_env_or("BUZZ_DB_LOCK_TIMEOUT", PgTimeout::disabled()), + timeouts: RuntimeTimeouts::from_env_or(RuntimeTimeouts::disabled()), ..DbConfig::default() }) .await?; diff --git a/crates/buzz-db/src/lib.rs b/crates/buzz-db/src/lib.rs index 9eaa21eda..4c89ade24 100644 --- a/crates/buzz-db/src/lib.rs +++ b/crates/buzz-db/src/lib.rs @@ -69,12 +69,18 @@ use buzz_core::{CommunityId, StoredEvent}; pub const RUNTIME_STATEMENT_TIMEOUT: &str = "30s"; /// Default maximum time a runtime query may wait to acquire a lock. pub const RUNTIME_LOCK_TIMEOUT: &str = "5s"; +/// Default maximum time a runtime session may sit idle inside an open +/// transaction. Deliberately looser than [`RUNTIME_STATEMENT_TIMEOUT`]: this +/// budget covers only the gaps *between* a transaction's statements, so a +/// transaction doing continuous work is never at risk, and a wide bound still +/// reclaims a session that has stopped working entirely. +pub const RUNTIME_IDLE_IN_TRANSACTION_TIMEOUT: &str = "60s"; /// Postgres spelling of "no limit", used for schema migrations. pub const TIMEOUT_DISABLED: &str = "0"; -/// `statement_timeout` and `lock_timeout` are `int` GUCs measured in -/// milliseconds, so Postgres refuses anything larger regardless of the unit it -/// is spelled with. Private on purpose: [`PgTimeout`] is the only way to build a -/// timeout, so no caller needs to range-check by hand. +/// The runtime timeouts are `int` GUCs measured in milliseconds, so Postgres +/// refuses anything larger regardless of the unit it is spelled with. Private on +/// purpose: [`PgTimeout`] is the only way to build a timeout, so no caller needs +/// to range-check by hand. /// `pg_timeout_max_millis_matches_postgres` pins it to the live server. const PG_TIMEOUT_MAX_MILLIS: u128 = i32::MAX as u128; @@ -129,6 +135,12 @@ impl PgTimeout { Self(RUNTIME_LOCK_TIMEOUT.to_string()) } + /// The default idle-in-transaction timeout + /// ([`RUNTIME_IDLE_IN_TRANSACTION_TIMEOUT`]). + pub fn idle_in_transaction_default() -> Self { + Self(RUNTIME_IDLE_IN_TRANSACTION_TIMEOUT.to_string()) + } + /// No limit at all — what schema migrations and one-shot operator tools /// run with. pub fn disabled() -> Self { @@ -222,19 +234,75 @@ fn pg_timeout_millis(magnitude: &str, unit: &str) -> Option { } } +/// The per-session limits that stop one caller from holding a pooled connection +/// indefinitely. +/// +/// Grouped rather than passed as three same-typed arguments, which would make a +/// swapped pair a silent misconfiguration. +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct RuntimeTimeouts { + /// Bounds a single statement's execution time. + pub statement: PgTimeout, + /// Bounds heavyweight and row lock waits. Advisory-lock waits are bounded + /// by [`Self::statement`] instead. + pub lock: PgTimeout, + /// Bounds a session sitting idle inside an open transaction. + /// + /// [`Self::statement`] only counts time spent *executing*, so a session + /// that opens a transaction and then stops issuing statements holds its + /// connection — and every lock it already took — with no statement running + /// for the statement timeout to cancel. + pub idle_in_transaction: PgTimeout, +} + +impl Default for RuntimeTimeouts { + fn default() -> Self { + Self { + statement: PgTimeout::statement_default(), + lock: PgTimeout::lock_default(), + idle_in_transaction: PgTimeout::idle_in_transaction_default(), + } + } +} + +impl RuntimeTimeouts { + /// Every limit lifted — schema migrations and one-shot operator tools. + pub fn disabled() -> Self { + Self { + statement: PgTimeout::disabled(), + lock: PgTimeout::disabled(), + idle_in_transaction: PgTimeout::disabled(), + } + } + + /// Read each limit from its environment variable, falling back to the + /// corresponding field of `defaults` with a warning. + pub fn from_env_or(defaults: Self) -> Self { + Self { + statement: PgTimeout::from_env_or("BUZZ_DB_STATEMENT_TIMEOUT", defaults.statement), + lock: PgTimeout::from_env_or("BUZZ_DB_LOCK_TIMEOUT", defaults.lock), + idle_in_transaction: PgTimeout::from_env_or( + "BUZZ_DB_IDLE_IN_TRANSACTION_TIMEOUT", + defaults.idle_in_transaction, + ), + } + } +} + /// Apply the runtime safety limits shared by writer, reader, audit, and search -/// pools. [`PgTimeout::disabled`] lifts a limit entirely. +/// pools. [`RuntimeTimeouts::disabled`] lifts them entirely. pub async fn apply_runtime_connection_timeouts( connection: &mut PgConnection, - statement_timeout: &PgTimeout, - lock_timeout: &PgTimeout, + timeouts: &RuntimeTimeouts, ) -> std::result::Result<(), sqlx::Error> { sqlx::query( "SELECT set_config('statement_timeout', $1, false), \ - set_config('lock_timeout', $2, false)", + set_config('lock_timeout', $2, false), \ + set_config('idle_in_transaction_session_timeout', $3, false)", ) - .bind(statement_timeout.as_str()) - .bind(lock_timeout.as_str()) + .bind(timeouts.statement.as_str()) + .bind(timeouts.lock.as_str()) + .bind(timeouts.idle_in_transaction.as_str()) .execute(connection) .await?; Ok(()) @@ -702,14 +770,10 @@ pub struct DbConfig { /// than the staleness gate never routes anyway, so a larger budget /// would only misrepresent the config. pub replica_read_max_age_ms: u64, - /// Postgres `statement_timeout` applied to every runtime connection. An - /// operator running a backfill or working an incident can widen this without - /// a code change; [`PgTimeout::disabled`] removes the cap. - pub statement_timeout: PgTimeout, - /// Postgres `lock_timeout` applied to every runtime connection. Bounds - /// heavyweight and row lock waits only — advisory-lock waits are bounded by - /// [`Self::statement_timeout`] instead. - pub lock_timeout: PgTimeout, + /// Per-session limits applied to every runtime connection. An operator + /// running a backfill or working an incident can widen these without a code + /// change; [`RuntimeTimeouts::disabled`] removes them. + pub timeouts: RuntimeTimeouts, } impl Default for DbConfig { @@ -727,8 +791,7 @@ impl Default for DbConfig { max_lifetime_secs: 1800, idle_timeout_secs: 600, replica_read_max_age_ms: 0, - statement_timeout: PgTimeout::statement_default(), - lock_timeout: PgTimeout::lock_default(), + timeouts: RuntimeTimeouts::default(), } } } @@ -867,13 +930,11 @@ impl Db { .acquire_timeout(Duration::from_secs(config.acquire_timeout_secs)) .max_lifetime(Duration::from_secs(config.max_lifetime_secs)) .idle_timeout(Duration::from_secs(config.idle_timeout_secs)); - let statement_timeout = config.statement_timeout.clone(); - let lock_timeout = config.lock_timeout.clone(); + let timeouts = config.timeouts.clone(); options = options.after_connect(move |conn, _meta| { - let statement_timeout = statement_timeout.clone(); - let lock_timeout = lock_timeout.clone(); + let timeouts = timeouts.clone(); Box::pin(async move { - apply_runtime_connection_timeouts(conn, &statement_timeout, &lock_timeout).await?; + apply_runtime_connection_timeouts(conn, &timeouts).await?; if arm_floor_guard { // `SET` cannot take bind parameters; `set_config` can. sqlx::query("SELECT set_config('buzz.created_at_floor', $1, false)") @@ -911,8 +972,7 @@ impl Db { /// No floor guard: replica sessions are read-only, the trigger never /// fires there (see [`Db::connect_pool`]). fn connect_read_pool(config: &DbConfig, url: &str, max_connections: u32) -> Result { - let statement_timeout = config.statement_timeout.clone(); - let lock_timeout = config.lock_timeout.clone(); + let timeouts = config.timeouts.clone(); Ok(PgPoolOptions::new() .max_connections(max_connections) .min_connections(0) @@ -920,12 +980,10 @@ impl Db { .max_lifetime(Duration::from_secs(config.max_lifetime_secs)) .idle_timeout(Duration::from_secs(config.idle_timeout_secs)) .after_connect(move |connection, _meta| { - let statement_timeout = statement_timeout.clone(); - let lock_timeout = lock_timeout.clone(); - Box::pin(async move { - apply_runtime_connection_timeouts(connection, &statement_timeout, &lock_timeout) - .await - }) + let timeouts = timeouts.clone(); + Box::pin( + async move { apply_runtime_connection_timeouts(connection, &timeouts).await }, + ) }) .connect_lazy(url)?) } @@ -6737,8 +6795,11 @@ mod tests { database_url: scratch_url, max_connections: 2, min_connections: 2, - statement_timeout: TIGHT.parse().expect("test timeout literal"), - lock_timeout: TIGHT.parse().expect("test timeout literal"), + timeouts: RuntimeTimeouts { + statement: TIGHT.parse().expect("test timeout literal"), + lock: TIGHT.parse().expect("test timeout literal"), + idle_in_transaction: TIGHT.parse().expect("test timeout literal"), + }, ..DbConfig::default() }) .await @@ -6761,8 +6822,14 @@ mod tests { .fetch_one(&mut *connection) .await .expect("SHOW lock_timeout"); + let idle_timeout: String = + sqlx::query_scalar("SHOW idle_in_transaction_session_timeout") + .fetch_one(&mut *connection) + .await + .expect("SHOW idle_in_transaction_session_timeout"); assert_eq!(statement_timeout, TIGHT); assert_eq!(lock_timeout, TIGHT); + assert_eq!(idle_timeout, TIGHT); held.push(connection); } drop(held); @@ -8586,28 +8653,36 @@ mod tests { let cid = CommunityId::from_uuid(community); // Assert the effective session values, not only pool-builder intent. - let statement_timeout: String = sqlx::query_scalar("SHOW statement_timeout") - .fetch_one(&db.pool) - .await - .expect("SHOW statement_timeout"); - let lock_timeout: String = sqlx::query_scalar("SHOW lock_timeout") - .fetch_one(&db.pool) - .await - .expect("SHOW lock_timeout"); - assert_eq!(statement_timeout, RUNTIME_STATEMENT_TIMEOUT); - assert_eq!(lock_timeout, RUNTIME_LOCK_TIMEOUT); - - let read_pool = db.read_pool.as_ref().expect("read pool configured"); - let reader_statement_timeout: String = sqlx::query_scalar("SHOW statement_timeout") - .fetch_one(read_pool) - .await - .expect("SHOW reader statement_timeout"); - let reader_lock_timeout: String = sqlx::query_scalar("SHOW lock_timeout") - .fetch_one(read_pool) - .await - .expect("SHOW reader lock_timeout"); - assert_eq!(reader_statement_timeout, RUNTIME_STATEMENT_TIMEOUT); - assert_eq!(reader_lock_timeout, RUNTIME_LOCK_TIMEOUT); + // Compare milliseconds, not spellings: Postgres re-spells a setting on + // the way out (`60s` reads back as `1min`), so asserting on `SHOW` + // text couples the test to the server's formatting rather than to the + // limit actually in force. + for (pool, label) in [ + (&db.pool, "writer"), + ( + db.read_pool.as_ref().expect("read pool configured"), + "reader", + ), + ] { + for (setting, expected_millis) in [ + ("statement_timeout", 30_000), + ("lock_timeout", 5_000), + ("idle_in_transaction_session_timeout", 60_000), + ] { + let millis: i64 = + sqlx::query_scalar("SELECT setting::bigint FROM pg_settings WHERE name = $1") + .bind(setting) + .fetch_one(pool) + .await + .unwrap_or_else(|error| { + panic!("read {setting} on the {label} pool: {error}") + }); + assert_eq!( + millis, expected_millis, + "{label} pool must carry the default {setting}" + ); + } + } let effective: String = sqlx::query_scalar("SHOW buzz.created_at_floor") .fetch_one(&db.pool) @@ -8671,6 +8746,16 @@ mod tests { db.pool.close().await; } + /// Every limit set to the same value, for tests that care about one + /// spelling rather than about the individual limits. + fn uniform_timeouts(timeout: &PgTimeout) -> RuntimeTimeouts { + RuntimeTimeouts { + statement: timeout.clone(), + lock: timeout.clone(), + idle_in_transaction: timeout.clone(), + } + } + /// The built-in defaults bypass parsing, so prove they would survive it — /// otherwise a typo in a literal ships a value Postgres refuses. #[test] @@ -8791,7 +8876,7 @@ mod tests { ] { let parsed = PgTimeout::parse_or(Some(raw), PgTimeout::statement_default()); assert_eq!(parsed.as_str(), canonical); - apply_runtime_connection_timeouts(&mut connection, &parsed, &parsed) + apply_runtime_connection_timeouts(&mut connection, &uniform_timeouts(&parsed)) .await .unwrap_or_else(|error| panic!("{raw:?} normalized to {parsed}: {error}")); } @@ -8811,7 +8896,7 @@ mod tests { .expect("connect"); let boundary = PgTimeout::unchecked(PG_TIMEOUT_MAX_MILLIS.to_string()); - apply_runtime_connection_timeouts(&mut conn, &boundary, &boundary) + apply_runtime_connection_timeouts(&mut conn, &uniform_timeouts(&boundary)) .await .expect("PG_TIMEOUT_MAX_MILLIS must be settable"); @@ -8819,7 +8904,7 @@ mod tests { // range check performs must land inside the range too. let in_seconds = (PG_TIMEOUT_MAX_MILLIS / 1_000).to_string(); let boundary_seconds = PgTimeout::unchecked(format!("{in_seconds}s")); - apply_runtime_connection_timeouts(&mut conn, &boundary_seconds, &boundary_seconds) + apply_runtime_connection_timeouts(&mut conn, &uniform_timeouts(&boundary_seconds)) .await .expect("the boundary in seconds must be settable"); @@ -8831,8 +8916,10 @@ mod tests { ] { let err = apply_runtime_connection_timeouts( &mut conn, - &PgTimeout::unchecked(over.clone()), - &PgTimeout::lock_default(), + &RuntimeTimeouts { + statement: PgTimeout::unchecked(over.clone()), + ..RuntimeTimeouts::default() + }, ) .await .expect_err(&format!("{over} must be rejected by Postgres")); diff --git a/crates/buzz-db/src/migration.rs b/crates/buzz-db/src/migration.rs index 283c2af53..dd6a556df 100644 --- a/crates/buzz-db/src/migration.rs +++ b/crates/buzz-db/src/migration.rs @@ -71,12 +71,8 @@ async fn migrate_on_exempt_connection(connection: &mut sqlx::PgConnection) -> Re /// Remove both runtime limits from one connection's session. async fn lift_runtime_timeouts(connection: &mut sqlx::PgConnection) -> Result<()> { - crate::apply_runtime_connection_timeouts( - connection, - &crate::PgTimeout::disabled(), - &crate::PgTimeout::disabled(), - ) - .await?; + crate::apply_runtime_connection_timeouts(connection, &crate::RuntimeTimeouts::disabled()) + .await?; Ok(()) } @@ -1174,11 +1170,18 @@ mod tests { .expect("create public schema"); } - /// Parse a literal the tests control, so a typo in one fails loudly here - /// rather than silently falling back to the runtime default. - fn tight_timeout(raw: &str) -> crate::PgTimeout { - raw.parse::() - .expect("test timeout literal must be valid") + /// Every limit set to the same tight value. Parsed rather than constructed + /// so a typo in a literal fails loudly here instead of silently falling + /// back to the runtime default. + fn tight_timeouts(raw: &str) -> crate::RuntimeTimeouts { + let timeout = raw + .parse::() + .expect("test timeout literal must be valid"); + crate::RuntimeTimeouts { + statement: timeout.clone(), + lock: timeout.clone(), + idle_in_transaction: timeout, + } } async fn applied_versions(pool: &PgPool) -> Vec { @@ -1215,9 +1218,9 @@ mod tests { let pool = sqlx::postgres::PgPoolOptions::new() .max_connections(1) .after_connect(|connection, _meta| { - let tight = tight_timeout(TIGHT); + let tight = tight_timeouts(TIGHT); Box::pin(async move { - crate::apply_runtime_connection_timeouts(connection, &tight, &tight).await + crate::apply_runtime_connection_timeouts(connection, &tight).await }) }) .connect(&database_url) @@ -1231,7 +1234,7 @@ mod tests { let mut connection = exempt_migration_connection(&pool) .await .expect("acquire exempt migration connection"); - for setting in ["statement_timeout", "lock_timeout"] { + for setting in RUNTIME_TIMEOUT_SETTINGS { assert_eq!( show_timeout(&mut connection, setting).await, "0", @@ -1244,7 +1247,7 @@ mod tests { drop(connection); let mut fresh = pool.acquire().await.expect("re-acquire"); - for setting in ["statement_timeout", "lock_timeout"] { + for setting in RUNTIME_TIMEOUT_SETTINGS { assert_eq!( show_timeout(&mut fresh, setting).await, TIGHT, @@ -1255,6 +1258,14 @@ mod tests { pool.close().await; } + /// Every GUC `apply_runtime_connection_timeouts` sets, so a limit added + /// there without a matching exemption fails these assertions. + const RUNTIME_TIMEOUT_SETTINGS: [&str; 3] = [ + "statement_timeout", + "lock_timeout", + "idle_in_transaction_session_timeout", + ]; + async fn show_timeout(connection: &mut sqlx::PgConnection, setting: &str) -> String { sqlx::query_scalar(sqlx::AssertSqlSafe(format!("SHOW {setting}"))) .fetch_one(connection) @@ -1344,9 +1355,9 @@ mod tests { sqlx::postgres::PgPoolOptions::new() .max_connections(2) .after_connect(move |connection, _meta| { - let timeout = tight_timeout(timeout); + let timeouts = tight_timeouts(timeout); Box::pin(async move { - crate::apply_runtime_connection_timeouts(connection, &timeout, &timeout).await + crate::apply_runtime_connection_timeouts(connection, &timeouts).await }) }) .connect(&database_url) @@ -1453,7 +1464,7 @@ mod tests { // And the caps are still in force for traffic afterwards. let mut runtime = capped.acquire().await.expect("acquire after migration"); - for setting in ["statement_timeout", "lock_timeout"] { + for setting in RUNTIME_TIMEOUT_SETTINGS { assert_eq!( show_timeout(&mut runtime, setting).await, TIGHT, diff --git a/crates/buzz-relay/src/config.rs b/crates/buzz-relay/src/config.rs index 567d6fcb6..c89d8ff1c 100644 --- a/crates/buzz-relay/src/config.rs +++ b/crates/buzz-relay/src/config.rs @@ -3,7 +3,7 @@ use std::net::SocketAddr; use std::time::Duration; -use buzz_db::PgTimeout; +use buzz_db::RuntimeTimeouts; use sha2::{Digest, Sha256}; use thiserror::Error; use tracing::warn; @@ -100,14 +100,13 @@ pub struct Config { /// independently so reader capacity can be tuned against the replica's /// headroom without touching the writer pool. pub db_read_pool_size: Option, - /// Postgres `statement_timeout` for every runtime connection - /// (`BUZZ_DB_STATEMENT_TIMEOUT`, e.g. `45s`, `500ms`, `0` to disable). - /// Tunable so a backfill or an incident does not need a code change. - pub db_statement_timeout: PgTimeout, - /// Postgres `lock_timeout` for every runtime connection - /// (`BUZZ_DB_LOCK_TIMEOUT`). Schema migrations always run with both limits - /// lifted — see `buzz_db::migration::run_migrations`. - pub db_lock_timeout: PgTimeout, + /// Per-session Postgres limits for every runtime connection + /// (`BUZZ_DB_STATEMENT_TIMEOUT`, `BUZZ_DB_LOCK_TIMEOUT`, + /// `BUZZ_DB_IDLE_IN_TRANSACTION_TIMEOUT`; e.g. `45s`, `500ms`, `0` to + /// disable). Tunable so a backfill or an incident does not need a code + /// change. Schema migrations always run with all of them lifted — see + /// `buzz_db::migration::run_migrations`. + pub db_timeouts: RuntimeTimeouts, /// Public WebSocket URL of this relay, advertised in NIP-11. pub relay_url: String, /// Public WebSocket URL of the dedicated device-pairing relay, when configured. @@ -540,10 +539,7 @@ impl Config { .and_then(|v| v.parse::().ok()) .filter(|&v| v > 0); - let db_statement_timeout = - PgTimeout::from_env_or("BUZZ_DB_STATEMENT_TIMEOUT", PgTimeout::statement_default()); - let db_lock_timeout = - PgTimeout::from_env_or("BUZZ_DB_LOCK_TIMEOUT", PgTimeout::lock_default()); + let db_timeouts = RuntimeTimeouts::from_env_or(RuntimeTimeouts::default()); let relay_url = std::env::var("RELAY_URL").unwrap_or_else(|_| "ws://localhost:3000".to_string()); @@ -1010,8 +1006,7 @@ impl Config { redis_pool_size, db_pool_size, db_read_pool_size, - db_statement_timeout, - db_lock_timeout, + db_timeouts, relay_url, pairing_relay_url, max_connections, @@ -1285,35 +1280,43 @@ mod tests { let _guard = ENV_MUTEX.lock().unwrap(); let previous_statement = std::env::var_os("BUZZ_DB_STATEMENT_TIMEOUT"); let previous_lock = std::env::var_os("BUZZ_DB_LOCK_TIMEOUT"); + let previous_idle = std::env::var_os("BUZZ_DB_IDLE_IN_TRANSACTION_TIMEOUT"); std::env::remove_var("BUZZ_DB_STATEMENT_TIMEOUT"); std::env::remove_var("BUZZ_DB_LOCK_TIMEOUT"); + std::env::remove_var("BUZZ_DB_IDLE_IN_TRANSACTION_TIMEOUT"); let defaults = Config::from_env().expect("config"); assert_eq!( - defaults.db_statement_timeout.as_str(), + defaults.db_timeouts.statement.as_str(), buzz_db::RUNTIME_STATEMENT_TIMEOUT ); assert_eq!( - defaults.db_lock_timeout.as_str(), + defaults.db_timeouts.lock.as_str(), buzz_db::RUNTIME_LOCK_TIMEOUT ); + assert_eq!( + defaults.db_timeouts.idle_in_transaction.as_str(), + buzz_db::RUNTIME_IDLE_IN_TRANSACTION_TIMEOUT + ); std::env::set_var("BUZZ_DB_STATEMENT_TIMEOUT", "90s"); std::env::set_var("BUZZ_DB_LOCK_TIMEOUT", "250ms"); + std::env::set_var("BUZZ_DB_IDLE_IN_TRANSACTION_TIMEOUT", "2min"); let overridden = Config::from_env().expect("config"); - assert_eq!(overridden.db_statement_timeout.as_str(), "90s"); - assert_eq!(overridden.db_lock_timeout.as_str(), "250ms"); + assert_eq!(overridden.db_timeouts.statement.as_str(), "90s"); + assert_eq!(overridden.db_timeouts.lock.as_str(), "250ms"); + assert_eq!(overridden.db_timeouts.idle_in_transaction.as_str(), "2min"); std::env::set_var("BUZZ_DB_STATEMENT_TIMEOUT", "45S"); std::env::set_var("BUZZ_DB_LOCK_TIMEOUT", " 2 MIN "); let canonicalized = Config::from_env().expect("config"); - assert_eq!(canonicalized.db_statement_timeout.as_str(), "45s"); - assert_eq!(canonicalized.db_lock_timeout.as_str(), "2min"); + assert_eq!(canonicalized.db_timeouts.statement.as_str(), "45s"); + assert_eq!(canonicalized.db_timeouts.lock.as_str(), "2min"); std::env::set_var("BUZZ_DB_STATEMENT_TIMEOUT", "half a minute"); let junk = Config::from_env().expect("config"); assert_eq!( - junk.db_statement_timeout.as_str(), + junk.db_timeouts.statement.as_str(), buzz_db::RUNTIME_STATEMENT_TIMEOUT, "a malformed value must not be handed to Postgres" ); @@ -1325,16 +1328,21 @@ mod tests { "999999999999999999999999999999999999999999d", ); std::env::set_var("BUZZ_DB_LOCK_TIMEOUT", "25d"); + std::env::set_var("BUZZ_DB_IDLE_IN_TRANSACTION_TIMEOUT", "25d"); let out_of_range = Config::from_env().expect("config"); assert_eq!( - out_of_range.db_statement_timeout.as_str(), + out_of_range.db_timeouts.statement.as_str(), buzz_db::RUNTIME_STATEMENT_TIMEOUT, "an out-of-range value must not be handed to Postgres" ); assert_eq!( - out_of_range.db_lock_timeout.as_str(), + out_of_range.db_timeouts.lock.as_str(), buzz_db::RUNTIME_LOCK_TIMEOUT ); + assert_eq!( + out_of_range.db_timeouts.idle_in_transaction.as_str(), + buzz_db::RUNTIME_IDLE_IN_TRANSACTION_TIMEOUT + ); match previous_statement { Some(value) => std::env::set_var("BUZZ_DB_STATEMENT_TIMEOUT", value), @@ -1344,6 +1352,10 @@ mod tests { Some(value) => std::env::set_var("BUZZ_DB_LOCK_TIMEOUT", value), None => std::env::remove_var("BUZZ_DB_LOCK_TIMEOUT"), } + match previous_idle { + Some(value) => std::env::set_var("BUZZ_DB_IDLE_IN_TRANSACTION_TIMEOUT", value), + None => std::env::remove_var("BUZZ_DB_IDLE_IN_TRANSACTION_TIMEOUT"), + } } #[test] diff --git a/crates/buzz-relay/src/main.rs b/crates/buzz-relay/src/main.rs index c962b549e..507b91d29 100644 --- a/crates/buzz-relay/src/main.rs +++ b/crates/buzz-relay/src/main.rs @@ -48,19 +48,12 @@ fn runtime_timeout_hook( + Send + Sync + 'static { - let statement_timeout = db_config.statement_timeout.clone(); - let lock_timeout = db_config.lock_timeout.clone(); + let timeouts = db_config.timeouts.clone(); move |connection, _meta| { - let statement_timeout = statement_timeout.clone(); - let lock_timeout = lock_timeout.clone(); - Box::pin(async move { - buzz_db::apply_runtime_connection_timeouts( - connection, - &statement_timeout, - &lock_timeout, - ) - .await - }) + let timeouts = timeouts.clone(); + Box::pin( + async move { buzz_db::apply_runtime_connection_timeouts(connection, &timeouts).await }, + ) } } @@ -198,8 +191,7 @@ async fn main() -> anyhow::Result<()> { replica_read_max_age_ms: config.replica_read_max_age_ms, max_connections: config.db_pool_size, read_max_connections: config.db_read_pool_size, - statement_timeout: config.db_statement_timeout.clone(), - lock_timeout: config.db_lock_timeout.clone(), + timeouts: config.db_timeouts.clone(), ..DbConfig::default() }; let db = Db::new(&db_config).await.map_err(|e| {