mirror of
https://github.com/block/buzz.git
synced 2026-08-18 06:50:31 +02:00
fix(db): bound sessions idle inside an open transaction
statement_timeout only counts time spent executing, so a session that opens a transaction and then stops issuing statements holds its pooled connection -- and every lock it already took -- with nothing for the statement timeout to cancel. That is the same connection-starvation failure the other two limits close, reached a different way. Set idle_in_transaction_session_timeout alongside them, defaulting to 60s: the budget covers only the gaps between a transaction's statements, so continuous work is never at risk. The three limits move into a RuntimeTimeouts struct rather than growing apply_runtime_connection_timeouts to three same-typed arguments, where a swapped pair would be a silent misconfiguration. Migration exemption assertions now iterate every GUC the applier sets, so a fourth limit added without a matching exemption fails the test. The armed-pool assertions compare milliseconds read from pg_settings instead of SHOW text: Postgres re-spells settings on the way out (60s reads back as 1min), which coupled the test to server formatting. Signed-off-by: Eli Foster <efoster@squareup.com>
This commit is contained in:
+11
-5
@@ -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)
|
||||
|
||||
@@ -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<Db> {
|
||||
// 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?;
|
||||
|
||||
+149
-62
@@ -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<u128> {
|
||||
}
|
||||
}
|
||||
|
||||
/// 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<PgPool> {
|
||||
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"));
|
||||
|
||||
@@ -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::<crate::PgTimeout>()
|
||||
.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::<crate::PgTimeout>()
|
||||
.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<i64> {
|
||||
@@ -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,
|
||||
|
||||
@@ -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<u32>,
|
||||
/// 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::<u32>().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]
|
||||
|
||||
@@ -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| {
|
||||
|
||||
Reference in New Issue
Block a user