From 6e50ba25a917390631e12fb7e738311b48e903da Mon Sep 17 00:00:00 2001 From: John Coffey Date: Thu, 24 Sep 2026 13:11:17 -0700 Subject: [PATCH] SQL pools time out; task locks are a renewed five-minute lease A 3-node rehearsal (PostgreSQL + NATS + Garage) found two ways a crash leaves work stuck: Pool hangs. The PostgreSQL pool (deadpool) was built with no timeouts, so a request waited for a free connection, and for one to be opened or recycled, for as long as it took: forever when the server stopped answering. MySQL's pool (mysql_async) has no wait timeout at all. - PostgreSQL: wait 30 s (or the store's timeout if longer), create the store's timeout or 15 s (it bounds the whole handshake, where tokio-postgres's connect_timeout covers only the TCP connect), recycle 10 s. The pool config is now always set, not only with poolMaxConnections. - MySQL: every connection is taken through MysqlStore::conn(), which gives up after 30 s. - Both: TCP keepalive after 60 s idle, so a server that vanished without closing the connection is noticed in minutes rather than the two-hour system default. The DataStore schema has no pool timeout settings, so these are fixed defaults; the store's own timeout bounds connecting on PostgreSQL. Task locks. A task lock lasted an hour, so after a hard crash the dead node's tasks waited up to an hour and five minutes. The lock is now a five-minute lease: while this node runs a task, the task manager renews its lock every third of the lifetime (InMemoryStore::renew_lock, a compare-and-set on the store backends and SET XX EX on Redis, which leaves a lock that already expired alone). A killed node's tasks run elsewhere within about five minutes plus the claim recheck. A task this node holds isn't handed to a worker again by the scan. store::pool_timeout (new): a local listener that accepts connections and never answers plays a hung server; a PostgreSQL store with a 2 s timeout returns an error in about 4 s, and a MySQL store in 30 s. Without the timeouts both wait for good. store::task_locks gains a task held for 1.5 lock lifetimes: its lease is still held, and released when the task ends. --- crates/common/src/ipc.rs | 19 +++- crates/services/src/task_manager/lock.rs | 39 +++++++++ crates/services/src/task_manager/manager.rs | 30 ++++++- crates/store/src/backend/mysql/blob.rs | 8 +- crates/store/src/backend/mysql/lookup.rs | 4 +- crates/store/src/backend/mysql/main.rs | 7 +- crates/store/src/backend/mysql/mod.rs | 27 ++++++ crates/store/src/backend/mysql/read.rs | 10 ++- crates/store/src/backend/mysql/search.rs | 6 +- crates/store/src/backend/mysql/write.rs | 8 +- crates/store/src/backend/postgres/main.rs | 42 ++++++++- crates/store/src/backend/redis/lookup.rs | 42 +++++++++ crates/store/src/dispatch/lookup.rs | 51 +++++++++++ tests/src/store/mod.rs | 2 + tests/src/store/pool_timeout.rs | 96 +++++++++++++++++++++ tests/src/store/task_locks.rs | 28 +++++- 16 files changed, 395 insertions(+), 24 deletions(-) create mode 100644 tests/src/store/pool_timeout.rs diff --git a/crates/common/src/ipc.rs b/crates/common/src/ipc.rs index a7fab1c..75bd155 100644 --- a/crates/common/src/ipc.rs +++ b/crates/common/src/ipc.rs @@ -345,8 +345,13 @@ pub struct TaskLocks { } impl TaskLocks { - /// How long a task lock lasts, in seconds, unless it is released first. - pub const DEFAULT_EXPIRY: u64 = 60 * 60; + /// How long a task lock lasts, in seconds, unless it is released first + /// or renewed. inbuxa: upstream held a lock for an hour, so a killed + /// node's tasks waited that long; the lock is now a five-minute lease + /// that the task manager renews every third of it while the task runs + /// (renew_task_locks), so a dead node's tasks run elsewhere within + /// minutes. + pub const DEFAULT_EXPIRY: u64 = 5 * 60; pub fn is_stopping(&self) -> bool { self.stopping.load(Ordering::Acquire) @@ -370,6 +375,16 @@ impl TaskLocks { self.held.lock().len() } + /// inbuxa: the tasks this node holds, to renew their locks. + pub fn held_ids(&self) -> Vec { + self.held.lock().iter().copied().collect() + } + + /// inbuxa: whether this node holds (and is running) the task. + pub fn is_held(&self, id: u64) -> bool { + self.held.lock().contains(&id) + } + pub fn expiry(&self) -> u64 { self.expiry.load(Ordering::Relaxed) } diff --git a/crates/services/src/task_manager/lock.rs b/crates/services/src/task_manager/lock.rs index 861d5b9..a1098ed 100644 --- a/crates/services/src/task_manager/lock.rs +++ b/crates/services/src/task_manager/lock.rs @@ -82,3 +82,42 @@ pub async fn release_task_locks(server: &Server) -> usize { } ids.len() } + +/// inbuxa: renews the lease on every task this node is running, so it stays +/// claimed for as long as it runs while a node that dies loses its claims +/// within one lock lifetime. Returns how many leases were renewed and how +/// many were found lost (expired, perhaps taken by another node). +pub async fn renew_task_locks(server: &Server) -> (usize, usize) { + let locks = &server.inner.ipc.task_locks; + let expiry = locks.expiry(); + let (mut renewed, mut lost) = (0, 0); + for id in locks.held_ids() { + match server + .in_memory_store() + .renew_lock(KV_LOCK_TASK, &id.to_be_bytes(), expiry) + .await + { + Ok(true) => renewed += 1, + Ok(false) => { + // Still held here as far as this node knows; the task + // finishes and its lock is removed as usual + if locks.is_held(id) { + lost += 1; + trc::event!( + TaskManager(TaskManagerEvent::TaskLocked), + Id = id, + Details = "Task lock expired while the task was running", + ); + } + } + Err(err) => { + trc::error!( + err.details("Failed to renew task lock") + .ctx(trc::Key::Id, id) + .caused_by(trc::location!()) + ); + } + } + } + (renewed, lost) +} diff --git a/crates/services/src/task_manager/manager.rs b/crates/services/src/task_manager/manager.rs index 1906e5b..e7d0d1d 100644 --- a/crates/services/src/task_manager/manager.rs +++ b/crates/services/src/task_manager/manager.rs @@ -13,7 +13,7 @@ use crate::task_manager::dkim::DkimManagementTask; use crate::task_manager::dns::DnsManagementTask; use crate::task_manager::imip::SendImipTask; use crate::task_manager::index::SearchIndexTask; -use crate::task_manager::lock::TaskLockManager; +use crate::task_manager::lock::{TaskLockManager, renew_task_locks}; use crate::task_manager::maintenance::MaintenanceTask; use crate::task_manager::merge_threads::MergeThreadsTask; use crate::task_manager::report::{self, SubmitReportTask}; @@ -75,6 +75,28 @@ pub fn spawn_task_manager(inner: Arc) { trc::event!(TaskManager(TaskManagerEvent::ManagerStarted)); + // inbuxa: keep the leases of running tasks alive, every third of a lock + // lifetime, until the node stops + { + let inner = inner.clone(); + tokio::spawn(async move { + let mut renewed_at = Instant::now(); + loop { + tokio::time::sleep(Duration::from_secs(1)).await; + let locks = &inner.ipc.task_locks; + if locks.is_stopping() { + break; + } + if renewed_at.elapsed() >= Duration::from_secs((locks.expiry() / 3).max(1)) { + renewed_at = Instant::now(); + if locks.held() > 0 { + renew_task_locks(&inner.build_server()).await; + } + } + } + }); + } + // Create dummy server instance for alarms let server_instance = Arc::new(ServerInstance { id: "_local".to_string(), @@ -292,6 +314,12 @@ impl TaskQueueManager for Server { .caused_by(trc::location!()) .ctx(trc::Key::Value, value) })?; + // inbuxa: running here under a lease this node + // renews; don't hand it to a worker again + if task_locks.is_held(task_id) { + return Ok(true); + } + let enabled = task_enabled(roles, task_type); if !enabled { diff --git a/crates/store/src/backend/mysql/blob.rs b/crates/store/src/backend/mysql/blob.rs index 983a358..4866229 100644 --- a/crates/store/src/backend/mysql/blob.rs +++ b/crates/store/src/backend/mysql/blob.rs @@ -2,6 +2,8 @@ * SPDX-FileCopyrightText: 2020 Stalwart Labs LLC * * SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-SEL + * + * Modified by Coffey Labs in 2026 for INBUXA. */ use std::ops::Range; @@ -16,7 +18,7 @@ impl MysqlStore { key: &[u8], range: Range, ) -> trc::Result>> { - let mut conn = self.conn_pool.get_conn().await.map_err(into_error)?; + let mut conn = self.conn().await?; let s = conn .prep("SELECT v FROM t WHERE k = ?") .await @@ -39,7 +41,7 @@ impl MysqlStore { } pub(crate) async fn put_blob(&self, key: &[u8], data: &[u8]) -> trc::Result<()> { - let mut conn = self.conn_pool.get_conn().await.map_err(into_error)?; + let mut conn = self.conn().await?; let s = conn .prep("INSERT INTO t (k, v) VALUES (?, ?) ON DUPLICATE KEY UPDATE v = VALUES(v)") .await @@ -51,7 +53,7 @@ impl MysqlStore { } pub(crate) async fn delete_blob(&self, key: &[u8]) -> trc::Result { - let mut conn = self.conn_pool.get_conn().await.map_err(into_error)?; + let mut conn = self.conn().await?; let s = conn .prep("DELETE FROM t WHERE k = ?") .await diff --git a/crates/store/src/backend/mysql/lookup.rs b/crates/store/src/backend/mysql/lookup.rs index 321bc1e..a9bedff 100644 --- a/crates/store/src/backend/mysql/lookup.rs +++ b/crates/store/src/backend/mysql/lookup.rs @@ -2,6 +2,8 @@ * SPDX-FileCopyrightText: 2020 Stalwart Labs LLC * * SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-SEL + * + * Modified by Coffey Labs in 2026 for INBUXA. */ use mysql_async::{Params, Row, prelude::Queryable}; @@ -16,7 +18,7 @@ impl MysqlStore { query: &str, params: &[Value<'_>], ) -> trc::Result { - let mut conn = self.conn_pool.get_conn().await.map_err(into_error)?; + let mut conn = self.conn().await?; let s = conn.prep(query).await.map_err(into_error)?; let params = Params::Positional(params.iter().map(Into::into).collect()); diff --git a/crates/store/src/backend/mysql/main.rs b/crates/store/src/backend/mysql/main.rs index 47a6ce6..0e8a6dd 100644 --- a/crates/store/src/backend/mysql/main.rs +++ b/crates/store/src/backend/mysql/main.rs @@ -32,6 +32,9 @@ impl MysqlStore { .max_allowed_packet(config.max_allowed_packet.map(|v| v as usize)) .wait_timeout(config.timeout.map(|t| t.as_secs() as usize)) .client_found_rows(true) + // inbuxa: notice a server that went away without closing the + // connection in minutes, not the system default of two hours + .tcp_keepalive(Some(super::POOL_KEEPALIVE_IDLE)) .tcp_port(config.port as u16); if config.use_tls { @@ -95,7 +98,7 @@ impl MysqlStore { } pub(crate) async fn create_storage_tables(&self) -> trc::Result<()> { - let mut conn = self.conn_pool.get_conn().await.map_err(into_error)?; + let mut conn = self.conn().await?; for table in [ SUBSPACE_ACL, @@ -169,7 +172,7 @@ impl MysqlStore { } pub(crate) async fn create_search_tables(&self) -> trc::Result<()> { - let mut conn = self.conn_pool.get_conn().await.map_err(into_error)?; + let mut conn = self.conn().await?; create_search_tables::(&mut conn).await?; create_search_tables::(&mut conn).await?; diff --git a/crates/store/src/backend/mysql/mod.rs b/crates/store/src/backend/mysql/mod.rs index 88ae06e..023d63c 100644 --- a/crates/store/src/backend/mysql/mod.rs +++ b/crates/store/src/backend/mysql/mod.rs @@ -27,6 +27,33 @@ pub struct MysqlStore { pub(crate) conn_pool: Pool, } +/// inbuxa: how long a request waits for a pooled connection (including +/// opening one). mysql_async's pool has no wait timeout, so upstream waited +/// forever when the server stopped answering. +pub(crate) const POOL_WAIT_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(30); +/// inbuxa: idle time before TCP keepalive probes start. +pub(crate) const POOL_KEEPALIVE_IDLE: std::time::Duration = std::time::Duration::from_secs(60); + +impl MysqlStore { + /// inbuxa: a pooled connection, or an error once POOL_WAIT_TIMEOUT has + /// passed without one. + pub(crate) async fn conn(&self) -> trc::Result { + pool_conn(&self.conn_pool, POOL_WAIT_TIMEOUT).await + } +} + +pub(crate) async fn pool_conn( + pool: &Pool, + wait: std::time::Duration, +) -> trc::Result { + match tokio::time::timeout(wait, pool.get_conn()).await { + Ok(result) => result.map_err(into_error), + Err(_) => Err(trc::StoreEvent::MysqlError + .reason("Timed out waiting for a database connection") + .details(format!("No connection within {} s", wait.as_secs()))), + } +} + #[inline(always)] pub(crate) fn into_error(err: impl Display) -> trc::Error { trc::StoreEvent::MysqlError.reason(err) diff --git a/crates/store/src/backend/mysql/read.rs b/crates/store/src/backend/mysql/read.rs index 3fc387e..74f03d4 100644 --- a/crates/store/src/backend/mysql/read.rs +++ b/crates/store/src/backend/mysql/read.rs @@ -2,6 +2,8 @@ * SPDX-FileCopyrightText: 2020 Stalwart Labs LLC * * SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-SEL + * + * Modified by Coffey Labs in 2026 for INBUXA. */ use super::{MysqlStore, into_error, is_timeout_error}; @@ -14,7 +16,7 @@ impl MysqlStore { where U: Deserialize + 'static, { - let mut conn = self.conn_pool.get_conn().await.map_err(into_error)?; + let mut conn = self.conn().await?; let s = conn .prep(format!( "SELECT v FROM {} WHERE k = ?", @@ -36,7 +38,7 @@ impl MysqlStore { } pub(crate) async fn key_exists(&self, key: impl Key) -> trc::Result { - let mut conn = self.conn_pool.get_conn().await.map_err(into_error)?; + let mut conn = self.conn().await?; let s = conn .prep(format!( "SELECT 1 FROM {} WHERE k = ?", @@ -56,7 +58,7 @@ impl MysqlStore { params: IterateParams, mut cb: impl for<'x> FnMut(&'x [u8], &'x [u8]) -> trc::Result + Sync + Send, ) -> trc::Result<()> { - let mut conn = self.conn_pool.get_conn().await.map_err(into_error)?; + let mut conn = self.conn().await?; let table = char::from(params.begin.subspace()); let begin = params.begin.serialize(0); let end = params.end.serialize(0); @@ -155,7 +157,7 @@ impl MysqlStore { let key = key.into(); let table = char::from(key.subspace()); let key = key.serialize(0); - let mut conn = self.conn_pool.get_conn().await.map_err(into_error)?; + let mut conn = self.conn().await?; let s = conn .prep(format!("SELECT v FROM {table} WHERE k = ?")) .await diff --git a/crates/store/src/backend/mysql/search.rs b/crates/store/src/backend/mysql/search.rs index 7df89af..f8c1a67 100644 --- a/crates/store/src/backend/mysql/search.rs +++ b/crates/store/src/backend/mysql/search.rs @@ -26,7 +26,7 @@ use std::fmt::Write; impl MysqlStore { pub async fn index(&self, documents: Vec) -> trc::Result<()> { - let mut conn = self.conn_pool.get_conn().await.map_err(into_error)?; + let mut conn = self.conn().await?; let mut tx_opts = TxOpts::default(); tx_opts .with_consistent_snapshot(false) @@ -96,7 +96,7 @@ impl MysqlStore { build_sort(&mut query, sort); } - let mut conn = self.conn_pool.get_conn().await.map_err(into_error)?; + let mut conn = self.conn().await?; let s = conn.prep(query).await.map_err(into_error)?; conn.exec::(s, params) @@ -110,7 +110,7 @@ impl MysqlStore { let mut query = format!("DELETE FROM {table} "); let params = build_filter(&mut query, &filter.filters); - let mut conn = self.conn_pool.get_conn().await.map_err(into_error)?; + let mut conn = self.conn().await?; let s = conn.prep(&query).await.map_err(into_error)?; match conn.exec_drop(s, params.clone()).await { diff --git a/crates/store/src/backend/mysql/write.rs b/crates/store/src/backend/mysql/write.rs index d5b183a..150ad3b 100644 --- a/crates/store/src/backend/mysql/write.rs +++ b/crates/store/src/backend/mysql/write.rs @@ -2,6 +2,8 @@ * SPDX-FileCopyrightText: 2020 Stalwart Labs LLC * * SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-SEL + * + * Modified by Coffey Labs in 2026 for INBUXA. */ use super::{DELETE_CHUNK_SIZE, MIN_DELETE_CHUNK_SIZE, MysqlStore, into_error, is_timeout_error}; @@ -29,7 +31,7 @@ impl MysqlStore { pub(crate) async fn write(&self, mut batch: Batch<'_>) -> trc::Result { let start = Instant::now(); let mut retry_count = 0; - let mut conn = self.conn_pool.get_conn().await.map_err(into_error)?; + let mut conn = self.conn().await?; loop { let err = match self.write_trx(&mut conn, &mut batch).await { @@ -382,7 +384,7 @@ impl MysqlStore { } pub(crate) async fn purge_store(&self) -> trc::Result<()> { - let mut conn = self.conn_pool.get_conn().await.map_err(into_error)?; + let mut conn = self.conn().await?; for subspace in [SUBSPACE_QUOTA, SUBSPACE_COUNTER, SUBSPACE_IN_MEMORY_COUNTER] { purge_table(&mut conn, char::from(subspace)).await?; } @@ -391,7 +393,7 @@ impl MysqlStore { } pub(crate) async fn delete_range(&self, from: impl Key, to: impl Key) -> trc::Result<()> { - let mut conn = self.conn_pool.get_conn().await.map_err(into_error)?; + let mut conn = self.conn().await?; let table = char::from(from.subspace()); let mut from = from.serialize(0); let to = to.serialize(0); diff --git a/crates/store/src/backend/postgres/main.rs b/crates/store/src/backend/postgres/main.rs index 60f5d2b..6dbd853 100644 --- a/crates/store/src/backend/postgres/main.rs +++ b/crates/store/src/backend/postgres/main.rs @@ -22,11 +22,34 @@ use crate::{ use ::registry::schema::{enums::PostgreSqlRecyclingMethod, structs}; use ahash::AHashSet; use deadpool_postgres::{ - Config, ManagerConfig, Object, Pool, PoolConfig, RecyclingMethod, Runtime, + Config, ManagerConfig, Object, Pool, PoolConfig, RecyclingMethod, Runtime, Timeouts, }; +use std::time::Duration; use tokio_postgres::NoTls; use utils::tls::rustls_client_config; +/// inbuxa: how long a request waits for a pooled connection. +pub(crate) const POOL_WAIT_TIMEOUT: Duration = Duration::from_secs(30); +/// inbuxa: how long opening a connection may take when the store sets no +/// timeout of its own. +pub(crate) const POOL_CREATE_TIMEOUT: Duration = Duration::from_secs(15); +/// inbuxa: how long checking a pooled connection before reuse may take. +pub(crate) const POOL_RECYCLE_TIMEOUT: Duration = Duration::from_secs(10); +/// inbuxa: idle time before TCP keepalive probes start. +pub(crate) const POOL_KEEPALIVE_IDLE: Duration = Duration::from_secs(60); + +/// inbuxa: the pool's timeouts. Opening a connection is bounded by the +/// store's own timeout when it has one; waiting for one covers at least that +/// long, so a slow connect isn't cut short by the wait. +pub(crate) fn pool_timeouts(connect_timeout: Option) -> Timeouts { + let create = connect_timeout.unwrap_or(POOL_CREATE_TIMEOUT); + Timeouts { + wait: POOL_WAIT_TIMEOUT.max(create).into(), + create: create.into(), + recycle: POOL_RECYCLE_TIMEOUT.into(), + } +} + impl PostgresStore { pub async fn open(config: structs::PostgreSqlStore) -> Result { // inbuxa: ST-15: where the primary is, to tell a replica from it @@ -46,9 +69,20 @@ impl PostgresStore { PostgreSqlRecyclingMethod::Clean => RecyclingMethod::Clean, }, }); - if let Some(max_conn) = config.pool_max_connections { - cfg.pool = PoolConfig::new(max_conn as usize).into(); - } + // inbuxa: upstream set no pool timeouts, so a request waited for a + // free connection, or for one to be made or recycled, for as long as + // it took: forever when the server stopped answering. A worker now + // gets an error instead and the task or request is retried. + let mut pool = config + .pool_max_connections + .map(|max_conn| PoolConfig::new(max_conn as usize)) + .unwrap_or_default(); + pool.timeouts = pool_timeouts(cfg.connect_timeout); + cfg.pool = pool.into(); + // Notice a server that went away without closing the connection in + // minutes rather than the system default of two hours + cfg.keepalives = true.into(); + cfg.keepalives_idle = POOL_KEEPALIVE_IDLE.into(); let primary_pool = if config.use_tls { cfg.create_pool( diff --git a/crates/store/src/backend/redis/lookup.rs b/crates/store/src/backend/redis/lookup.rs index 300624b..f1c9cc4 100644 --- a/crates/store/src/backend/redis/lookup.rs +++ b/crates/store/src/backend/redis/lookup.rs @@ -2,6 +2,8 @@ * SPDX-FileCopyrightText: 2020 Stalwart Labs LLC * * SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-SEL + * + * Modified by Coffey Labs in 2026 for INBUXA. */ use super::{RedisPool, RedisStore, into_error}; @@ -79,6 +81,30 @@ impl RedisStore { } } + // inbuxa: see InMemoryStore::renew_lock + pub async fn renew_lock(&self, key: &[u8], expires: u64) -> trc::Result { + match &self.pool { + RedisPool::Single(pool) => { + with_conn(pool, async |conn| { + Self::renew_lock_(conn, key, expires).await + }) + .await + } + RedisPool::Cluster(pool) => { + with_conn(pool, async |conn| { + Self::renew_lock_(conn, key, expires).await + }) + .await + } + RedisPool::Sentinel(pool) => { + with_conn(pool, async |conn| { + Self::renew_lock_(conn, key, expires).await + }) + .await + } + } + } + pub async fn key_delete(&self, key: &[u8]) -> trc::Result<()> { match &self.pool { RedisPool::Single(pool) => { @@ -226,6 +252,22 @@ impl RedisStore { .map(|reply| reply.is_some()) } + async fn renew_lock_( + conn: &mut impl AsyncCommands, + key: &[u8], + expires: u64, + ) -> RedisResult { + redis::cmd("SET") + .arg(key) + .arg(now() + expires) + .arg("XX") + .arg("EX") + .arg(expires as i64) + .query_async::>(conn) + .await + .map(|reply| reply.is_some()) + } + async fn key_delete_(conn: &mut impl AsyncCommands, key: &[u8]) -> RedisResult<()> { conn.del(key).await } diff --git a/crates/store/src/dispatch/lookup.rs b/crates/store/src/dispatch/lookup.rs index b433260..804b87b 100644 --- a/crates/store/src/dispatch/lookup.rs +++ b/crates/store/src/dispatch/lookup.rs @@ -401,6 +401,57 @@ impl InMemoryStore { } } + /// inbuxa: extends a lock this node holds to `duration` seconds from now. + /// Returns false when the lock is gone or has expired: it may have been + /// taken by someone else since, so it is left alone. + pub async fn renew_lock(&self, prefix: u8, key: &[u8], duration: u64) -> trc::Result { + match self { + InMemoryStore::Store(store) => { + let key = KeyValue::<()>::build_key(prefix, key); + let key = ValueClass::InMemory(InMemoryClass::Key(key)); + let Some(lock_expiry) = store + .get_value::(ValueKey::from(key.clone())) + .await + .caused_by(trc::location!())? + else { + return Ok(false); + }; + let now = now(); + if lock_expiry <= now { + return Ok(false); + } + + let mut batch = BatchBuilder::new(); + batch.assert_value(key.clone(), AssertValue::U64(lock_expiry)); + batch.set(key, (now + duration).serialize()); + match store.write(batch.build_all()).await { + Ok(_) => Ok(true), + Err(err) if err.is_assertion_failure() => Ok(false), + Err(err) => Err(err + .details("Failed to renew lock.") + .caused_by(trc::location!())), + } + } + InMemoryStore::Sharded(store) => { + Box::pin( + store + .member(&KeyValue::<()>::build_key(prefix, key)) + .renew_lock(prefix, key, duration), + ) + .await + } + #[cfg(feature = "redis")] + InMemoryStore::Redis(store) => { + store + .renew_lock(&KeyValue::<()>::build_key(prefix, key), duration) + .await + } + InMemoryStore::Static(_) | InMemoryStore::Http(_) => { + Err(trc::StoreEvent::NotSupported.into_err()) + } + } + } + pub async fn remove_lock(&self, prefix: u8, key: &[u8]) -> trc::Result<()> { self.key_delete(KeyValue::<()>::build_key(prefix, key)) .await diff --git a/tests/src/store/mod.rs b/tests/src/store/mod.rs index e863401..490bb4f 100644 --- a/tests/src/store/mod.rs +++ b/tests/src/store/mod.rs @@ -10,6 +10,8 @@ pub mod blob; pub mod import_export; pub mod lookup; pub mod ops; +#[cfg(any(feature = "postgres", feature = "mysql"))] +pub mod pool_timeout; // inbuxa: SQL pools give up instead of hanging pub mod query; pub mod registry; #[cfg(feature = "postgres")] diff --git a/tests/src/store/pool_timeout.rs b/tests/src/store/pool_timeout.rs new file mode 100644 index 0000000..b9b829a --- /dev/null +++ b/tests/src/store/pool_timeout.rs @@ -0,0 +1,96 @@ +/* + * SPDX-FileCopyrightText: 2026 Coffey Labs + * + * SPDX-License-Identifier: AGPL-3.0-only + */ + +//! A database that accepts connections and then says nothing (a hung or +//! half-dead server, a black-holed failover) gives a worker an error within +//! the pool's timeouts. Upstream's pools had none, so the worker waited for +//! good. No database is needed: a local listener that never answers plays +//! the server. + +use registry::schema::structs::DataStore; +use std::time::{Duration, Instant}; +use store::{Store, ValueKey, write::ValueClass}; +use tokio::net::TcpListener; + +/// Accepts connections on a local port and never sends a byte. +async fn silent_server() -> u16 { + let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); + let port = listener.local_addr().unwrap().port(); + tokio::spawn(async move { + let mut held = Vec::new(); + while let Ok((socket, _)) = listener.accept().await { + held.push(socket); + } + }); + port +} + +/// Builds the store and reads a key; both must end, with an error for the +/// read, well within `limit`. +async fn assert_times_out(data_store: DataStore, limit: Duration) { + let started = Instant::now(); + let result = tokio::time::timeout(limit, async { + match Store::build(data_store).await { + Ok(store) => store + .get_value::(ValueKey::from(ValueClass::Property(0))) + .await + .map(|_| ()) + .map_err(|err| err.to_string()), + Err(err) => Err(err.to_string()), + } + }) + .await; + let elapsed = started.elapsed(); + match result { + Ok(Err(err)) => println!("Got {err} after {elapsed:?}"), + Ok(Ok(())) => panic!("a silent server answered?"), + Err(_) => panic!("still waiting for a connection after {elapsed:?}"), + } +} + +#[cfg(feature = "postgres")] +#[tokio::test(flavor = "multi_thread")] +pub async fn postgres_pool_timeout() { + use registry::schema::structs::PostgreSqlStore; + + let port = silent_server().await; + println!("Running PostgreSQL pool timeout test..."); + // The store's own timeout bounds opening a connection, handshake + // included (tokio-postgres's connect_timeout covers only the TCP connect) + assert_times_out( + DataStore::PostgreSql(PostgreSqlStore { + host: "127.0.0.1".into(), + port: port as u64, + database: "none".into(), + timeout: Some(Duration::from_secs(2).into()), + use_tls: false, + ..Default::default() + }), + Duration::from_secs(20), + ) + .await; +} + +#[cfg(feature = "mysql")] +#[tokio::test(flavor = "multi_thread")] +pub async fn mysql_pool_timeout() { + use registry::schema::structs::MySqlStore; + + let port = silent_server().await; + println!("Running MySQL pool timeout test..."); + // mysql_async has no pool timeout; the store waits 30 s for a connection + assert_times_out( + DataStore::MySql(MySqlStore { + host: "127.0.0.1".into(), + port: port as u64, + database: "none".into(), + use_tls: false, + ..Default::default() + }), + Duration::from_secs(60), + ) + .await; +} diff --git a/tests/src/store/task_locks.rs b/tests/src/store/task_locks.rs index bcfacb6..d742785 100644 --- a/tests/src/store/task_locks.rs +++ b/tests/src/store/task_locks.rs @@ -79,7 +79,33 @@ pub async fn task_lock_tests() { "ran before the other node's locks expired: {elapsed:?}" ); - // 3. A graceful stop releases the locks this node holds: another node + // 3. A task that runs longer than a lock lifetime keeps its claim: the + // task manager renews the lease while this node holds it, and the claim + // ends when the task does. (Before, a lock simply lasted an hour.) + let [id] = new_task_ids(1)[..] else { + unreachable!() + }; + assert!(server.try_lock_task(id).await, "claim {id}"); + tokio::time::sleep(Duration::from_secs(LOCK_EXPIRY + LOCK_EXPIRY / 2)).await; + assert!( + !foreign_lock(&server, id, LOCK_EXPIRY).await, + "lease lapsed while the task ran" + ); + server.remove_index_lock(id).await; + assert!( + foreign_lock(&server, id, LOCK_EXPIRY).await, + "released when the task ended" + ); + let _ = server + .in_memory_store() + .remove_lock(KV_LOCK_TASK, &id.to_be_bytes()) + .await; + assert!( + common::ipc::TaskLocks::DEFAULT_EXPIRY <= 5 * 60, + "a dead node's tasks wait no more than a few minutes" + ); + + // 4. A graceful stop releases the locks this node holds: another node // can claim those tasks at once, and this one claims nothing more let ids = new_task_ids(3); for id in &ids { -- 2.54.0