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.
This commit is contained in:
@@ -345,8 +345,13 @@ pub struct TaskLocks {
|
|||||||
}
|
}
|
||||||
|
|
||||||
impl TaskLocks {
|
impl TaskLocks {
|
||||||
/// How long a task lock lasts, in seconds, unless it is released first.
|
/// How long a task lock lasts, in seconds, unless it is released first
|
||||||
pub const DEFAULT_EXPIRY: u64 = 60 * 60;
|
/// 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 {
|
pub fn is_stopping(&self) -> bool {
|
||||||
self.stopping.load(Ordering::Acquire)
|
self.stopping.load(Ordering::Acquire)
|
||||||
@@ -370,6 +375,16 @@ impl TaskLocks {
|
|||||||
self.held.lock().len()
|
self.held.lock().len()
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// inbuxa: the tasks this node holds, to renew their locks.
|
||||||
|
pub fn held_ids(&self) -> Vec<u64> {
|
||||||
|
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 {
|
pub fn expiry(&self) -> u64 {
|
||||||
self.expiry.load(Ordering::Relaxed)
|
self.expiry.load(Ordering::Relaxed)
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -82,3 +82,42 @@ pub async fn release_task_locks(server: &Server) -> usize {
|
|||||||
}
|
}
|
||||||
ids.len()
|
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)
|
||||||
|
}
|
||||||
|
|||||||
@@ -13,7 +13,7 @@ use crate::task_manager::dkim::DkimManagementTask;
|
|||||||
use crate::task_manager::dns::DnsManagementTask;
|
use crate::task_manager::dns::DnsManagementTask;
|
||||||
use crate::task_manager::imip::SendImipTask;
|
use crate::task_manager::imip::SendImipTask;
|
||||||
use crate::task_manager::index::SearchIndexTask;
|
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::maintenance::MaintenanceTask;
|
||||||
use crate::task_manager::merge_threads::MergeThreadsTask;
|
use crate::task_manager::merge_threads::MergeThreadsTask;
|
||||||
use crate::task_manager::report::{self, SubmitReportTask};
|
use crate::task_manager::report::{self, SubmitReportTask};
|
||||||
@@ -75,6 +75,28 @@ pub fn spawn_task_manager(inner: Arc<Inner>) {
|
|||||||
|
|
||||||
trc::event!(TaskManager(TaskManagerEvent::ManagerStarted));
|
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
|
// Create dummy server instance for alarms
|
||||||
let server_instance = Arc::new(ServerInstance {
|
let server_instance = Arc::new(ServerInstance {
|
||||||
id: "_local".to_string(),
|
id: "_local".to_string(),
|
||||||
@@ -292,6 +314,12 @@ impl TaskQueueManager for Server {
|
|||||||
.caused_by(trc::location!())
|
.caused_by(trc::location!())
|
||||||
.ctx(trc::Key::Value, value)
|
.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);
|
let enabled = task_enabled(roles, task_type);
|
||||||
|
|
||||||
if !enabled {
|
if !enabled {
|
||||||
|
|||||||
@@ -2,6 +2,8 @@
|
|||||||
* SPDX-FileCopyrightText: 2020 Stalwart Labs LLC <[email protected]>
|
* SPDX-FileCopyrightText: 2020 Stalwart Labs LLC <[email protected]>
|
||||||
*
|
*
|
||||||
* SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-SEL
|
* SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-SEL
|
||||||
|
*
|
||||||
|
* Modified by Coffey Labs in 2026 for INBUXA.
|
||||||
*/
|
*/
|
||||||
|
|
||||||
use std::ops::Range;
|
use std::ops::Range;
|
||||||
@@ -16,7 +18,7 @@ impl MysqlStore {
|
|||||||
key: &[u8],
|
key: &[u8],
|
||||||
range: Range<usize>,
|
range: Range<usize>,
|
||||||
) -> trc::Result<Option<Vec<u8>>> {
|
) -> trc::Result<Option<Vec<u8>>> {
|
||||||
let mut conn = self.conn_pool.get_conn().await.map_err(into_error)?;
|
let mut conn = self.conn().await?;
|
||||||
let s = conn
|
let s = conn
|
||||||
.prep("SELECT v FROM t WHERE k = ?")
|
.prep("SELECT v FROM t WHERE k = ?")
|
||||||
.await
|
.await
|
||||||
@@ -39,7 +41,7 @@ impl MysqlStore {
|
|||||||
}
|
}
|
||||||
|
|
||||||
pub(crate) async fn put_blob(&self, key: &[u8], data: &[u8]) -> trc::Result<()> {
|
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
|
let s = conn
|
||||||
.prep("INSERT INTO t (k, v) VALUES (?, ?) ON DUPLICATE KEY UPDATE v = VALUES(v)")
|
.prep("INSERT INTO t (k, v) VALUES (?, ?) ON DUPLICATE KEY UPDATE v = VALUES(v)")
|
||||||
.await
|
.await
|
||||||
@@ -51,7 +53,7 @@ impl MysqlStore {
|
|||||||
}
|
}
|
||||||
|
|
||||||
pub(crate) async fn delete_blob(&self, key: &[u8]) -> trc::Result<bool> {
|
pub(crate) async fn delete_blob(&self, key: &[u8]) -> trc::Result<bool> {
|
||||||
let mut conn = self.conn_pool.get_conn().await.map_err(into_error)?;
|
let mut conn = self.conn().await?;
|
||||||
let s = conn
|
let s = conn
|
||||||
.prep("DELETE FROM t WHERE k = ?")
|
.prep("DELETE FROM t WHERE k = ?")
|
||||||
.await
|
.await
|
||||||
|
|||||||
@@ -2,6 +2,8 @@
|
|||||||
* SPDX-FileCopyrightText: 2020 Stalwart Labs LLC <[email protected]>
|
* SPDX-FileCopyrightText: 2020 Stalwart Labs LLC <[email protected]>
|
||||||
*
|
*
|
||||||
* SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-SEL
|
* 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};
|
use mysql_async::{Params, Row, prelude::Queryable};
|
||||||
@@ -16,7 +18,7 @@ impl MysqlStore {
|
|||||||
query: &str,
|
query: &str,
|
||||||
params: &[Value<'_>],
|
params: &[Value<'_>],
|
||||||
) -> trc::Result<T> {
|
) -> trc::Result<T> {
|
||||||
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 s = conn.prep(query).await.map_err(into_error)?;
|
||||||
let params = Params::Positional(params.iter().map(Into::into).collect());
|
let params = Params::Positional(params.iter().map(Into::into).collect());
|
||||||
|
|
||||||
|
|||||||
@@ -32,6 +32,9 @@ impl MysqlStore {
|
|||||||
.max_allowed_packet(config.max_allowed_packet.map(|v| v as usize))
|
.max_allowed_packet(config.max_allowed_packet.map(|v| v as usize))
|
||||||
.wait_timeout(config.timeout.map(|t| t.as_secs() as usize))
|
.wait_timeout(config.timeout.map(|t| t.as_secs() as usize))
|
||||||
.client_found_rows(true)
|
.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);
|
.tcp_port(config.port as u16);
|
||||||
|
|
||||||
if config.use_tls {
|
if config.use_tls {
|
||||||
@@ -95,7 +98,7 @@ impl MysqlStore {
|
|||||||
}
|
}
|
||||||
|
|
||||||
pub(crate) async fn create_storage_tables(&self) -> trc::Result<()> {
|
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 [
|
for table in [
|
||||||
SUBSPACE_ACL,
|
SUBSPACE_ACL,
|
||||||
@@ -169,7 +172,7 @@ impl MysqlStore {
|
|||||||
}
|
}
|
||||||
|
|
||||||
pub(crate) async fn create_search_tables(&self) -> trc::Result<()> {
|
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::<EmailSearchField>(&mut conn).await?;
|
create_search_tables::<EmailSearchField>(&mut conn).await?;
|
||||||
create_search_tables::<CalendarSearchField>(&mut conn).await?;
|
create_search_tables::<CalendarSearchField>(&mut conn).await?;
|
||||||
|
|||||||
@@ -27,6 +27,33 @@ pub struct MysqlStore {
|
|||||||
pub(crate) conn_pool: Pool,
|
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<mysql_async::Conn> {
|
||||||
|
pool_conn(&self.conn_pool, POOL_WAIT_TIMEOUT).await
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
pub(crate) async fn pool_conn(
|
||||||
|
pool: &Pool,
|
||||||
|
wait: std::time::Duration,
|
||||||
|
) -> trc::Result<mysql_async::Conn> {
|
||||||
|
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)]
|
#[inline(always)]
|
||||||
pub(crate) fn into_error(err: impl Display) -> trc::Error {
|
pub(crate) fn into_error(err: impl Display) -> trc::Error {
|
||||||
trc::StoreEvent::MysqlError.reason(err)
|
trc::StoreEvent::MysqlError.reason(err)
|
||||||
|
|||||||
@@ -2,6 +2,8 @@
|
|||||||
* SPDX-FileCopyrightText: 2020 Stalwart Labs LLC <[email protected]>
|
* SPDX-FileCopyrightText: 2020 Stalwart Labs LLC <[email protected]>
|
||||||
*
|
*
|
||||||
* SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-SEL
|
* 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};
|
use super::{MysqlStore, into_error, is_timeout_error};
|
||||||
@@ -14,7 +16,7 @@ impl MysqlStore {
|
|||||||
where
|
where
|
||||||
U: Deserialize + 'static,
|
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
|
let s = conn
|
||||||
.prep(format!(
|
.prep(format!(
|
||||||
"SELECT v FROM {} WHERE k = ?",
|
"SELECT v FROM {} WHERE k = ?",
|
||||||
@@ -36,7 +38,7 @@ impl MysqlStore {
|
|||||||
}
|
}
|
||||||
|
|
||||||
pub(crate) async fn key_exists(&self, key: impl Key) -> trc::Result<bool> {
|
pub(crate) async fn key_exists(&self, key: impl Key) -> trc::Result<bool> {
|
||||||
let mut conn = self.conn_pool.get_conn().await.map_err(into_error)?;
|
let mut conn = self.conn().await?;
|
||||||
let s = conn
|
let s = conn
|
||||||
.prep(format!(
|
.prep(format!(
|
||||||
"SELECT 1 FROM {} WHERE k = ?",
|
"SELECT 1 FROM {} WHERE k = ?",
|
||||||
@@ -56,7 +58,7 @@ impl MysqlStore {
|
|||||||
params: IterateParams<T>,
|
params: IterateParams<T>,
|
||||||
mut cb: impl for<'x> FnMut(&'x [u8], &'x [u8]) -> trc::Result<bool> + Sync + Send,
|
mut cb: impl for<'x> FnMut(&'x [u8], &'x [u8]) -> trc::Result<bool> + Sync + Send,
|
||||||
) -> trc::Result<()> {
|
) -> 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 table = char::from(params.begin.subspace());
|
||||||
let begin = params.begin.serialize(0);
|
let begin = params.begin.serialize(0);
|
||||||
let end = params.end.serialize(0);
|
let end = params.end.serialize(0);
|
||||||
@@ -155,7 +157,7 @@ impl MysqlStore {
|
|||||||
let key = key.into();
|
let key = key.into();
|
||||||
let table = char::from(key.subspace());
|
let table = char::from(key.subspace());
|
||||||
let key = key.serialize(0);
|
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
|
let s = conn
|
||||||
.prep(format!("SELECT v FROM {table} WHERE k = ?"))
|
.prep(format!("SELECT v FROM {table} WHERE k = ?"))
|
||||||
.await
|
.await
|
||||||
|
|||||||
@@ -26,7 +26,7 @@ use std::fmt::Write;
|
|||||||
|
|
||||||
impl MysqlStore {
|
impl MysqlStore {
|
||||||
pub async fn index(&self, documents: Vec<IndexDocument>) -> trc::Result<()> {
|
pub async fn index(&self, documents: Vec<IndexDocument>) -> 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();
|
let mut tx_opts = TxOpts::default();
|
||||||
tx_opts
|
tx_opts
|
||||||
.with_consistent_snapshot(false)
|
.with_consistent_snapshot(false)
|
||||||
@@ -96,7 +96,7 @@ impl MysqlStore {
|
|||||||
build_sort(&mut query, sort);
|
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)?;
|
let s = conn.prep(query).await.map_err(into_error)?;
|
||||||
|
|
||||||
conn.exec::<i64, _, _>(s, params)
|
conn.exec::<i64, _, _>(s, params)
|
||||||
@@ -110,7 +110,7 @@ impl MysqlStore {
|
|||||||
let mut query = format!("DELETE FROM {table} ");
|
let mut query = format!("DELETE FROM {table} ");
|
||||||
let params = build_filter(&mut query, &filter.filters);
|
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)?;
|
let s = conn.prep(&query).await.map_err(into_error)?;
|
||||||
|
|
||||||
match conn.exec_drop(s, params.clone()).await {
|
match conn.exec_drop(s, params.clone()).await {
|
||||||
|
|||||||
@@ -2,6 +2,8 @@
|
|||||||
* SPDX-FileCopyrightText: 2020 Stalwart Labs LLC <[email protected]>
|
* SPDX-FileCopyrightText: 2020 Stalwart Labs LLC <[email protected]>
|
||||||
*
|
*
|
||||||
* SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-SEL
|
* 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};
|
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<AssignedIds> {
|
pub(crate) async fn write(&self, mut batch: Batch<'_>) -> trc::Result<AssignedIds> {
|
||||||
let start = Instant::now();
|
let start = Instant::now();
|
||||||
let mut retry_count = 0;
|
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 {
|
loop {
|
||||||
let err = match self.write_trx(&mut conn, &mut batch).await {
|
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<()> {
|
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] {
|
for subspace in [SUBSPACE_QUOTA, SUBSPACE_COUNTER, SUBSPACE_IN_MEMORY_COUNTER] {
|
||||||
purge_table(&mut conn, char::from(subspace)).await?;
|
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<()> {
|
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 table = char::from(from.subspace());
|
||||||
let mut from = from.serialize(0);
|
let mut from = from.serialize(0);
|
||||||
let to = to.serialize(0);
|
let to = to.serialize(0);
|
||||||
|
|||||||
@@ -22,11 +22,34 @@ use crate::{
|
|||||||
use ::registry::schema::{enums::PostgreSqlRecyclingMethod, structs};
|
use ::registry::schema::{enums::PostgreSqlRecyclingMethod, structs};
|
||||||
use ahash::AHashSet;
|
use ahash::AHashSet;
|
||||||
use deadpool_postgres::{
|
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 tokio_postgres::NoTls;
|
||||||
use utils::tls::rustls_client_config;
|
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<Duration>) -> 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 {
|
impl PostgresStore {
|
||||||
pub async fn open(config: structs::PostgreSqlStore) -> Result<Store, String> {
|
pub async fn open(config: structs::PostgreSqlStore) -> Result<Store, String> {
|
||||||
// inbuxa: ST-15: where the primary is, to tell a replica from it
|
// inbuxa: ST-15: where the primary is, to tell a replica from it
|
||||||
@@ -46,9 +69,20 @@ impl PostgresStore {
|
|||||||
PostgreSqlRecyclingMethod::Clean => RecyclingMethod::Clean,
|
PostgreSqlRecyclingMethod::Clean => RecyclingMethod::Clean,
|
||||||
},
|
},
|
||||||
});
|
});
|
||||||
if let Some(max_conn) = config.pool_max_connections {
|
// inbuxa: upstream set no pool timeouts, so a request waited for a
|
||||||
cfg.pool = PoolConfig::new(max_conn as usize).into();
|
// 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 {
|
let primary_pool = if config.use_tls {
|
||||||
cfg.create_pool(
|
cfg.create_pool(
|
||||||
|
|||||||
@@ -2,6 +2,8 @@
|
|||||||
* SPDX-FileCopyrightText: 2020 Stalwart Labs LLC <[email protected]>
|
* SPDX-FileCopyrightText: 2020 Stalwart Labs LLC <[email protected]>
|
||||||
*
|
*
|
||||||
* SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-SEL
|
* SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-SEL
|
||||||
|
*
|
||||||
|
* Modified by Coffey Labs in 2026 for INBUXA.
|
||||||
*/
|
*/
|
||||||
|
|
||||||
use super::{RedisPool, RedisStore, into_error};
|
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<bool> {
|
||||||
|
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<()> {
|
pub async fn key_delete(&self, key: &[u8]) -> trc::Result<()> {
|
||||||
match &self.pool {
|
match &self.pool {
|
||||||
RedisPool::Single(pool) => {
|
RedisPool::Single(pool) => {
|
||||||
@@ -226,6 +252,22 @@ impl RedisStore {
|
|||||||
.map(|reply| reply.is_some())
|
.map(|reply| reply.is_some())
|
||||||
}
|
}
|
||||||
|
|
||||||
|
async fn renew_lock_(
|
||||||
|
conn: &mut impl AsyncCommands,
|
||||||
|
key: &[u8],
|
||||||
|
expires: u64,
|
||||||
|
) -> RedisResult<bool> {
|
||||||
|
redis::cmd("SET")
|
||||||
|
.arg(key)
|
||||||
|
.arg(now() + expires)
|
||||||
|
.arg("XX")
|
||||||
|
.arg("EX")
|
||||||
|
.arg(expires as i64)
|
||||||
|
.query_async::<Option<String>>(conn)
|
||||||
|
.await
|
||||||
|
.map(|reply| reply.is_some())
|
||||||
|
}
|
||||||
|
|
||||||
async fn key_delete_(conn: &mut impl AsyncCommands, key: &[u8]) -> RedisResult<()> {
|
async fn key_delete_(conn: &mut impl AsyncCommands, key: &[u8]) -> RedisResult<()> {
|
||||||
conn.del(key).await
|
conn.del(key).await
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -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<bool> {
|
||||||
|
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::<u64>(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<()> {
|
pub async fn remove_lock(&self, prefix: u8, key: &[u8]) -> trc::Result<()> {
|
||||||
self.key_delete(KeyValue::<()>::build_key(prefix, key))
|
self.key_delete(KeyValue::<()>::build_key(prefix, key))
|
||||||
.await
|
.await
|
||||||
|
|||||||
@@ -10,6 +10,8 @@ pub mod blob;
|
|||||||
pub mod import_export;
|
pub mod import_export;
|
||||||
pub mod lookup;
|
pub mod lookup;
|
||||||
pub mod ops;
|
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 query;
|
||||||
pub mod registry;
|
pub mod registry;
|
||||||
#[cfg(feature = "postgres")]
|
#[cfg(feature = "postgres")]
|
||||||
|
|||||||
@@ -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::<u64>(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;
|
||||||
|
}
|
||||||
@@ -79,7 +79,33 @@ pub async fn task_lock_tests() {
|
|||||||
"ran before the other node's locks expired: {elapsed:?}"
|
"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
|
// can claim those tasks at once, and this one claims nothing more
|
||||||
let ids = new_task_ids(3);
|
let ids = new_task_ids(3);
|
||||||
for id in &ids {
|
for id in &ids {
|
||||||
|
|||||||
Reference in New Issue
Block a user