diff --git a/crates/common/src/config/inner.rs b/crates/common/src/config/inner.rs index 325aacc..b1f32d8 100644 --- a/crates/common/src/config/inner.rs +++ b/crates/common/src/config/inner.rs @@ -94,6 +94,7 @@ impl Data { span_id_gen: id_generator, queue_status: true.into(), settings_reload: Default::default(), + store_health: Default::default(), applications, logos: Default::default(), smtp_connectors: TlsConnectors::try_new().failed("Failed to build TLS connectors"), @@ -237,6 +238,7 @@ impl Default for Data { registry_id_gen: Default::default(), queue_status: true.into(), settings_reload: Default::default(), + store_health: Default::default(), applications: WebApplications::new(), logos: Default::default(), smtp_connectors: TlsConnectors::try_new().unwrap(), diff --git a/crates/common/src/lib.rs b/crates/common/src/lib.rs index d9538da..3f84f1d 100644 --- a/crates/common/src/lib.rs +++ b/crates/common/src/lib.rs @@ -163,6 +163,8 @@ pub struct Data { pub queue_status: AtomicBool, // inbuxa: coalesces the settings reloads registry writes trigger pub settings_reload: cache::reload::SettingsReloadGate, + // inbuxa: the readiness probe's cached answer + pub store_health: storage::ready::StoreHealth, pub applications: WebApplications, pub logos: Mutex, LogoCache>>, diff --git a/crates/common/src/storage/mod.rs b/crates/common/src/storage/mod.rs index 2acc061..d192067 100644 --- a/crates/common/src/storage/mod.rs +++ b/crates/common/src/storage/mod.rs @@ -26,6 +26,7 @@ pub mod document; pub mod encryption; pub mod index; pub mod quota; +pub mod ready; // inbuxa: readiness follows the data store pub mod state; pub mod transaction; diff --git a/crates/common/src/storage/ready.rs b/crates/common/src/storage/ready.rs new file mode 100644 index 0000000..9d7838b --- /dev/null +++ b/crates/common/src/storage/ready.rs @@ -0,0 +1,83 @@ +/* + * SPDX-FileCopyrightText: 2026 Coffey Labs + * + * SPDX-License-Identifier: AGPL-3.0-only + */ + +//! Readiness that reflects the data store. +//! +//! /healthz/ready used to answer 200 whenever a data store was configured, +//! so a load balancer kept sending traffic to a node through a database +//! outage. It now reads one key from the data store, with a short time +//! limit, and caches the answer for a couple of seconds so probes can't load +//! the database. Liveness stays 200: restarting a node doesn't bring its +//! database back, and an orchestrator that restarts on failed liveness would +//! otherwise restart every node at once. + +use crate::Server; +use parking_lot::Mutex; +use std::{ + sync::atomic::{AtomicBool, Ordering}, + time::{Duration, Instant}, +}; +use store::{ValueKey, write::ValueClass}; + +/// How long a probe's answer is reused. +pub const READY_CACHE: Duration = Duration::from_secs(2); +/// How long a probe waits for the data store. +pub const READY_PROBE_TIMEOUT: Duration = Duration::from_secs(2); + +#[derive(Default)] +pub struct StoreHealth { + last: Mutex>, + probing: AtomicBool, +} + +/// Clears the probing flag even when the request is dropped mid-probe. +struct ProbeGuard<'x>(&'x AtomicBool); + +impl Drop for ProbeGuard<'_> { + fn drop(&mut self) { + self.0.store(false, Ordering::Release); + } +} + +impl Server { + /// Whether the data store answers: a cached result younger than + /// READY_CACHE, or a fresh read bounded by READY_PROBE_TIMEOUT. While + /// one probe is running, other callers get the last answer. + pub async fn is_data_store_ready(&self) -> bool { + let store = &self.core.storage.data; + if store.is_none() { + return false; + } + let health = &self.inner.data.store_health; + let last = *health.last.lock(); + if let Some((at, ready)) = last + && at.elapsed() < READY_CACHE + { + return ready; + } + if health.probing.swap(true, Ordering::AcqRel) { + return last.is_none_or(|(_, ready)| ready); + } + let _guard = ProbeGuard(&health.probing); + + let ready = tokio::time::timeout( + READY_PROBE_TIMEOUT, + store.get_value::(ValueKey::from(ValueClass::Property(0))), + ) + .await + .is_ok_and(|result| result.is_ok()); + // Say so once per outage, not on every probe + if !ready && last.is_none_or(|(_, ready)| ready) { + trc::event!( + Store(trc::StoreEvent::UnexpectedError), + Details = "Readiness probe: the data store didn't answer", + Limit = READY_PROBE_TIMEOUT, + ); + } + *health.last.lock() = Some((Instant::now(), ready)); + ready + } +} diff --git a/crates/http/src/request.rs b/crates/http/src/request.rs index 576b003..a09f42c 100644 --- a/crates/http/src/request.rs +++ b/crates/http/src/request.rs @@ -553,8 +553,10 @@ impl ParseHttp for Server { return Ok(JsonProblemResponse(StatusCode::OK).into_http_response()); } "ready" => { + // inbuxa: ready only while the data store answers + // (a cached, time-limited read); liveness stays 200 return Ok(JsonProblemResponse({ - if !self.core.storage.data.is_none() { + if self.is_data_store_ready().await { StatusCode::OK } else { StatusCode::SERVICE_UNAVAILABLE diff --git a/crates/store/src/backend/mod.rs b/crates/store/src/backend/mod.rs index b817a75..5af6fb2 100644 --- a/crates/store/src/backend/mod.rs +++ b/crates/store/src/backend/mod.rs @@ -30,6 +30,9 @@ pub mod s3; pub mod sqlite; // inbuxa: scale-out storage (sharded stores) pub mod scaleout; +// inbuxa: client-side SQL query limits +#[cfg(any(feature = "postgres", feature = "mysql"))] +pub mod query_timeout; pub const MAX_TOKEN_LENGTH: usize = (u8::MAX >> 1) as usize; diff --git a/crates/store/src/backend/mysql/blob.rs b/crates/store/src/backend/mysql/blob.rs index 4866229..65a4e06 100644 --- a/crates/store/src/backend/mysql/blob.rs +++ b/crates/store/src/backend/mysql/blob.rs @@ -10,7 +10,7 @@ use std::ops::Range; use mysql_async::prelude::Queryable; -use super::{MysqlStore, into_error}; +use super::{MysqlStore, bounded, into_error}; impl MysqlStore { pub(crate) async fn get_blob( @@ -19,48 +19,63 @@ impl MysqlStore { range: Range, ) -> trc::Result>> { let mut conn = self.conn().await?; - let s = conn - .prep("SELECT v FROM t WHERE k = ?") - .await - .map_err(into_error)?; - conn.exec_first::, _, _>(&s, (key,)) - .await - .map(|bytes| { - if range.start == 0 && range.end == usize::MAX { - bytes - } else { - bytes.map(|bytes| { + let limit = self.timeouts.query; + let result = tokio::time::timeout(limit, async { + let s = conn + .prep("SELECT v FROM t WHERE k = ?") + .await + .map_err(into_error)?; + conn.exec_first::, _, _>(&s, (key,)) + .await + .map(|bytes| { + if range.start == 0 && range.end == usize::MAX { bytes - .get(range.start..std::cmp::min(bytes.len(), range.end)) - .unwrap_or_default() - .to_vec() - }) - } - }) - .map_err(into_error) + } else { + bytes.map(|bytes| { + bytes + .get(range.start..std::cmp::min(bytes.len(), range.end)) + .unwrap_or_default() + .to_vec() + }) + } + }) + .map_err(into_error) + }) + .await; + bounded(conn, result, limit) } pub(crate) async fn put_blob(&self, key: &[u8], data: &[u8]) -> trc::Result<()> { let mut conn = self.conn().await?; - let s = conn - .prep("INSERT INTO t (k, v) VALUES (?, ?) ON DUPLICATE KEY UPDATE v = VALUES(v)") - .await - .map_err(into_error)?; - conn.exec_drop(&s, (key, data)) - .await - .map_err(into_error) - .map(|_| ()) + let limit = self.timeouts.query; + let result = tokio::time::timeout(limit, async { + let s = conn + .prep("INSERT INTO t (k, v) VALUES (?, ?) ON DUPLICATE KEY UPDATE v = VALUES(v)") + .await + .map_err(into_error)?; + conn.exec_drop(&s, (key, data)) + .await + .map_err(into_error) + .map(|_| ()) + }) + .await; + bounded(conn, result, limit) } pub(crate) async fn delete_blob(&self, key: &[u8]) -> trc::Result { let mut conn = self.conn().await?; - let s = conn - .prep("DELETE FROM t WHERE k = ?") - .await - .map_err(into_error)?; - conn.exec_iter(&s, (key,)) - .await - .map_err(into_error) - .map(|hits| hits.affected_rows() > 0) + let limit = self.timeouts.query; + let result = tokio::time::timeout(limit, async { + let s = conn + .prep("DELETE FROM t WHERE k = ?") + .await + .map_err(into_error)?; + conn.exec_iter(&s, (key,)) + .await + .map_err(into_error) + .map(|hits| hits.affected_rows() > 0) + }) + .await; + bounded(conn, result, limit) } } diff --git a/crates/store/src/backend/mysql/lookup.rs b/crates/store/src/backend/mysql/lookup.rs index a9bedff..5789e08 100644 --- a/crates/store/src/backend/mysql/lookup.rs +++ b/crates/store/src/backend/mysql/lookup.rs @@ -10,7 +10,7 @@ use mysql_async::{Params, Row, prelude::Queryable}; use crate::{IntoRows, QueryResult, QueryType, Value}; -use super::{MysqlStore, into_error}; +use super::{MysqlStore, bounded, into_error}; impl MysqlStore { pub(crate) async fn sql_query( @@ -19,27 +19,32 @@ impl MysqlStore { params: &[Value<'_>], ) -> trc::Result { 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()); + let limit = self.timeouts.query; + let result = tokio::time::timeout(limit, async { + let s = conn.prep(query).await.map_err(into_error)?; + let params = Params::Positional(params.iter().map(Into::into).collect()); - match T::query_type() { - QueryType::Execute => conn.exec_drop(s, params).await.map_or_else( - |e| Err(into_error(e)), - |_| Ok(T::from_exec(conn.affected_rows() as usize)), - ), - QueryType::Exists => conn - .exec_first::(s, params) - .await - .map_or_else(|e| Err(into_error(e)), |r| Ok(T::from_exists(r.is_some()))), - QueryType::QueryOne => conn - .exec_first::(s, params) - .await - .map_or_else(|e| Err(into_error(e)), |r| Ok(T::from_query_one(r))), - QueryType::QueryAll => conn - .exec::(s, params) - .await - .map_or_else(|e| Err(into_error(e)), |r| Ok(T::from_query_all(r))), - } + match T::query_type() { + QueryType::Execute => conn.exec_drop(s, params).await.map_or_else( + |e| Err(into_error(e)), + |_| Ok(T::from_exec(conn.affected_rows() as usize)), + ), + QueryType::Exists => conn + .exec_first::(s, params) + .await + .map_or_else(|e| Err(into_error(e)), |r| Ok(T::from_exists(r.is_some()))), + QueryType::QueryOne => conn + .exec_first::(s, params) + .await + .map_or_else(|e| Err(into_error(e)), |r| Ok(T::from_query_one(r))), + QueryType::QueryAll => conn + .exec::(s, params) + .await + .map_or_else(|e| Err(into_error(e)), |r| Ok(T::from_query_all(r))), + } + }) + .await; + bounded(conn, result, limit) } } diff --git a/crates/store/src/backend/mysql/main.rs b/crates/store/src/backend/mysql/main.rs index 0e8a6dd..deb8aa2 100644 --- a/crates/store/src/backend/mysql/main.rs +++ b/crates/store/src/backend/mysql/main.rs @@ -6,7 +6,7 @@ * Modified by Coffey Labs in 2026 for INBUXA. */ -use super::{MysqlStore, into_error}; +use super::{MysqlStore, bounded, into_error}; use crate::{ backend::mysql::MysqlSearchField, search::{ @@ -72,6 +72,7 @@ impl MysqlStore { .db_name(Some(replica.database.clone())) .tcp_port(replica.port as u16), ), + timeouts: Default::default(), })), replica.host, replica.port as u16, @@ -81,6 +82,7 @@ impl MysqlStore { let primary = Store::MySQL(Arc::new(MysqlStore { conn_pool: Pool::new(opts), + timeouts: Default::default(), })); // ST-1: no replicas, no change @@ -99,88 +101,96 @@ impl MysqlStore { pub(crate) async fn create_storage_tables(&self) -> trc::Result<()> { let mut conn = self.conn().await?; + let limit = self.timeouts.maintenance; + let result = tokio::time::timeout(limit, async { + for table in [ + SUBSPACE_ACL, + SUBSPACE_TASK_QUEUE, + SUBSPACE_DELETED_ITEMS, + SUBSPACE_SPAM_SAMPLES, + crate::SUBSPACE_INBUXA, // inbuxa: masked email + SUBSPACE_BLOB_LINK, + SUBSPACE_IN_MEMORY_VALUE, + SUBSPACE_PROPERTY, + SUBSPACE_REGISTRY, + SUBSPACE_REGISTRY_PK, + SUBSPACE_DIRECTORY, + SUBSPACE_QUEUE_MESSAGE, + SUBSPACE_QUEUE_EVENT, + SUBSPACE_REPORT_OUT, + SUBSPACE_REPORT_IN, + SUBSPACE_LOGS, + SUBSPACE_TELEMETRY_SPAN, + SUBSPACE_TELEMETRY_METRIC, + ] { + let table = char::from(table); + conn.query_drop(format!( + "CREATE TABLE IF NOT EXISTS {table} ( + k VARBINARY(255) NOT NULL, + v MEDIUMBLOB NOT NULL, + PRIMARY KEY (k) + ) ENGINE=InnoDB" + )) + .await + .map_err(into_error)?; + } - for table in [ - SUBSPACE_ACL, - SUBSPACE_TASK_QUEUE, - SUBSPACE_DELETED_ITEMS, - SUBSPACE_SPAM_SAMPLES, - crate::SUBSPACE_INBUXA, // inbuxa: masked email - SUBSPACE_BLOB_LINK, - SUBSPACE_IN_MEMORY_VALUE, - SUBSPACE_PROPERTY, - SUBSPACE_REGISTRY, - SUBSPACE_REGISTRY_PK, - SUBSPACE_DIRECTORY, - SUBSPACE_QUEUE_MESSAGE, - SUBSPACE_QUEUE_EVENT, - SUBSPACE_REPORT_OUT, - SUBSPACE_REPORT_IN, - SUBSPACE_LOGS, - SUBSPACE_TELEMETRY_SPAN, - SUBSPACE_TELEMETRY_METRIC, - ] { - let table = char::from(table); - conn.query_drop(format!( - "CREATE TABLE IF NOT EXISTS {table} ( - k VARBINARY(255) NOT NULL, - v MEDIUMBLOB NOT NULL, - PRIMARY KEY (k) - ) ENGINE=InnoDB" - )) - .await - .map_err(into_error)?; - } - - conn.query_drop(format!( - "CREATE TABLE IF NOT EXISTS {} ( - k VARBINARY(255) NOT NULL, - v LONGBLOB NOT NULL, - PRIMARY KEY (k) - ) ENGINE=InnoDB", - char::from(SUBSPACE_BLOBS), - )) - .await - .map_err(into_error)?; - - for table in [SUBSPACE_INDEXES, SUBSPACE_REGISTRY_IDX] { - let table = char::from(table); - conn.query_drop(format!( - "CREATE TABLE IF NOT EXISTS {table} ( - k BLOB, - PRIMARY KEY (k(400)) - ) ENGINE=InnoDB" - )) - .await - .map_err(into_error)?; - } - - for table in [SUBSPACE_COUNTER, SUBSPACE_QUOTA, SUBSPACE_IN_MEMORY_COUNTER] { conn.query_drop(format!( "CREATE TABLE IF NOT EXISTS {} ( - k VARBINARY(255) NOT NULL, - v BIGINT NOT NULL DEFAULT 0, - PRIMARY KEY (k) - ) ENGINE=InnoDB", - char::from(table) + k VARBINARY(255) NOT NULL, + v LONGBLOB NOT NULL, + PRIMARY KEY (k) + ) ENGINE=InnoDB", + char::from(SUBSPACE_BLOBS), )) .await .map_err(into_error)?; - } - Ok(()) + for table in [SUBSPACE_INDEXES, SUBSPACE_REGISTRY_IDX] { + let table = char::from(table); + conn.query_drop(format!( + "CREATE TABLE IF NOT EXISTS {table} ( + k BLOB, + PRIMARY KEY (k(400)) + ) ENGINE=InnoDB" + )) + .await + .map_err(into_error)?; + } + + for table in [SUBSPACE_COUNTER, SUBSPACE_QUOTA, SUBSPACE_IN_MEMORY_COUNTER] { + conn.query_drop(format!( + "CREATE TABLE IF NOT EXISTS {} ( + k VARBINARY(255) NOT NULL, + v BIGINT NOT NULL DEFAULT 0, + PRIMARY KEY (k) + ) ENGINE=InnoDB", + char::from(table) + )) + .await + .map_err(into_error)?; + } + + Ok(()) + }) + .await; + bounded(conn, result, limit) } pub(crate) async fn create_search_tables(&self) -> trc::Result<()> { let mut conn = self.conn().await?; + let limit = self.timeouts.maintenance; + let result = tokio::time::timeout(limit, async { + create_search_tables::(&mut conn).await?; + create_search_tables::(&mut conn).await?; + create_search_tables::(&mut conn).await?; + //create_search_tables::(&mut conn).await?; + create_search_tables::(&mut conn).await?; - create_search_tables::(&mut conn).await?; - create_search_tables::(&mut conn).await?; - create_search_tables::(&mut conn).await?; - //create_search_tables::(&mut conn).await?; - create_search_tables::(&mut conn).await?; - - Ok(()) + Ok(()) + }) + .await; + bounded(conn, result, limit) } } diff --git a/crates/store/src/backend/mysql/mod.rs b/crates/store/src/backend/mysql/mod.rs index 023d63c..d1c0b8c 100644 --- a/crates/store/src/backend/mysql/mod.rs +++ b/crates/store/src/backend/mysql/mod.rs @@ -6,6 +6,7 @@ * Modified by Coffey Labs in 2026 for INBUXA. */ +use crate::backend::query_timeout::QueryTimeouts; use crate::{ search::{ CalendarSearchField, ContactSearchField, EmailSearchField, FileSearchField, SearchField, @@ -14,7 +15,7 @@ use crate::{ write::SearchIndex, }; use mysql_async::Pool; -use std::fmt::Display; +use std::{fmt::Display, time::Duration}; pub mod blob; pub mod lookup; @@ -25,6 +26,8 @@ pub mod write; pub struct MysqlStore { pub(crate) conn_pool: Pool, + /// inbuxa: client-side query limits (see backend::query_timeout) + pub(crate) timeouts: QueryTimeouts, } /// inbuxa: how long a request waits for a pooled connection (including @@ -54,6 +57,43 @@ pub(crate) async fn pool_conn( } } +/// inbuxa: the error for an operation that ran past its time limit. +pub(crate) fn query_timeout_error(limit: Duration) -> trc::Error { + trc::StoreEvent::MysqlError + .reason("Query timed out") + .details(format!( + "No answer from the database within {} s", + limit.as_secs() + )) +} + +/// inbuxa: ends an operation run on `conn` under `limit`. When it ran out, +/// the connection is closed rather than returned to the pool: a query may +/// still be in flight on it, or a transaction open. Conn::disconnect marks +/// the connection closed before it sends anything, so even when the server +/// doesn't answer and the attempt is dropped, the pool discards it instead +/// of waiting to clean it up. +pub(crate) fn bounded( + conn: mysql_async::Conn, + result: Result, tokio::time::error::Elapsed>, + limit: Duration, +) -> trc::Result { + match result { + Ok(result) => result, + Err(_) => { + discard(conn); + Err(query_timeout_error(limit)) + } + } +} + +/// inbuxa: closes a connection whose state is unknown (see bounded). +pub(crate) fn discard(conn: mysql_async::Conn) { + tokio::spawn(async move { + let _ = tokio::time::timeout(Duration::from_secs(1), conn.disconnect()).await; + }); +} + #[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 74f03d4..ce42bac 100644 --- a/crates/store/src/backend/mysql/read.rs +++ b/crates/store/src/backend/mysql/read.rs @@ -6,7 +6,7 @@ * Modified by Coffey Labs in 2026 for INBUXA. */ -use super::{MysqlStore, into_error, is_timeout_error}; +use super::{MysqlStore, bounded, discard, into_error, is_timeout_error, query_timeout_error}; use crate::{Deserialize, IterateParams, Key, ValueKey, write::ValueClass}; use futures::TryStreamExt; use mysql_async::{Row, prelude::Queryable}; @@ -17,40 +17,50 @@ impl MysqlStore { U: Deserialize + 'static, { let mut conn = self.conn().await?; - let s = conn - .prep(format!( - "SELECT v FROM {} WHERE k = ?", - char::from(key.subspace()) - )) - .await - .map_err(into_error)?; - let key = key.serialize(0); - conn.exec_first::, _, _>(&s, (&key,)) - .await - .map_err(into_error) - .and_then(|r| { - if let Some(r) = r { - Ok(Some(U::deserialize_owned_with_key(&key, r)?)) - } else { - Ok(None) - } - }) + let limit = self.timeouts.query; + let result = tokio::time::timeout(limit, async { + let s = conn + .prep(format!( + "SELECT v FROM {} WHERE k = ?", + char::from(key.subspace()) + )) + .await + .map_err(into_error)?; + let key = key.serialize(0); + conn.exec_first::, _, _>(&s, (&key,)) + .await + .map_err(into_error) + .and_then(|r| { + if let Some(r) = r { + Ok(Some(U::deserialize_owned_with_key(&key, r)?)) + } else { + Ok(None) + } + }) + }) + .await; + bounded(conn, result, limit) } pub(crate) async fn key_exists(&self, key: impl Key) -> trc::Result { let mut conn = self.conn().await?; - let s = conn - .prep(format!( - "SELECT 1 FROM {} WHERE k = ?", - char::from(key.subspace()) - )) - .await - .map_err(into_error)?; - let key = key.serialize(0); - conn.exec_first::(&s, (&key,)) - .await - .map_err(into_error) - .map(|r| r.is_some()) + let limit = self.timeouts.query; + let result = tokio::time::timeout(limit, async { + let s = conn + .prep(format!( + "SELECT 1 FROM {} WHERE k = ?", + char::from(key.subspace()) + )) + .await + .map_err(into_error)?; + let key = key.serialize(0); + conn.exec_first::(&s, (&key,)) + .await + .map_err(into_error) + .map(|r| r.is_some()) + }) + .await; + bounded(conn, result, limit) } pub(crate) async fn iterate( @@ -64,28 +74,36 @@ impl MysqlStore { let end = params.end.serialize(0); let keys = if params.values { "k, v" } else { "k" }; - let s = conn - .prep(&match (params.first, params.ascending) { - (true, true) => { - format!( - "SELECT {keys} FROM {table} WHERE k >= ? AND k <= ? ORDER BY k ASC LIMIT 1" - ) - } - (true, false) => { - format!( - "SELECT {keys} FROM {table} WHERE k >= ? AND k <= ? ORDER BY k DESC LIMIT 1" - ) - } - (false, true) => { - format!("SELECT {keys} FROM {table} WHERE k >= ? AND k <= ? ORDER BY k ASC") - } - (false, false) => { - format!("SELECT {keys} FROM {table} WHERE k >= ? AND k <= ? ORDER BY k DESC") - } - }) - .await - .map_err(into_error)?; + // inbuxa: a scan may run for hours, so the query limit bounds each + // wait for the database (preparing, the query starting, the next + // row) rather than the scan. A wait that runs out closes the + // connection. + let limit = self.timeouts.query; + let query = match (params.first, params.ascending) { + (true, true) => { + format!("SELECT {keys} FROM {table} WHERE k >= ? AND k <= ? ORDER BY k ASC LIMIT 1") + } + (true, false) => { + format!( + "SELECT {keys} FROM {table} WHERE k >= ? AND k <= ? ORDER BY k DESC LIMIT 1" + ) + } + (false, true) => { + format!("SELECT {keys} FROM {table} WHERE k >= ? AND k <= ? ORDER BY k ASC") + } + (false, false) => { + format!("SELECT {keys} FROM {table} WHERE k >= ? AND k <= ? ORDER BY k DESC") + } + }; + let s = match tokio::time::timeout(limit, conn.prep(&query)).await { + Ok(s) => s.map_err(into_error)?, + Err(_) => { + discard(conn); + return Err(query_timeout_error(limit)); + } + }; let mut from = begin; + let mut stalled = false; let mut to = end; let mut resume_key = None; @@ -94,13 +112,26 @@ impl MysqlStore { let mut timed_out = false; { - let mut rows = conn - .exec_stream::(&s, (from.clone(), to.clone())) - .await - .map_err(into_error)?; + let mut rows = match tokio::time::timeout( + limit, + conn.exec_stream::(&s, (from.clone(), to.clone())), + ) + .await + { + Ok(rows) => rows.map_err(into_error)?, + // Leaves the scan loop for the timeout below + Err(_) => break, + }; loop { - match rows.try_next().await { + let next = match tokio::time::timeout(limit, rows.try_next()).await { + Ok(next) => next, + Err(_) => { + stalled = true; + break; + } + }; + match next { Ok(Some(mut row)) => { let value = if params.values { row.take_opt::, _>(1) @@ -136,6 +167,10 @@ impl MysqlStore { } } + if stalled { + break; + } + match last_key { Some(last_key) if timed_out => { if params.ascending { @@ -148,6 +183,9 @@ impl MysqlStore { _ => return Ok(()), } } + + discard(conn); + Err(query_timeout_error(limit)) } pub(crate) async fn get_counter( @@ -158,14 +196,19 @@ impl MysqlStore { let table = char::from(key.subspace()); let key = key.serialize(0); let mut conn = self.conn().await?; - let s = conn - .prep(format!("SELECT v FROM {table} WHERE k = ?")) - .await - .map_err(into_error)?; - match conn.exec_first::(&s, (key,)).await { - Ok(Some(num)) => Ok(num), - Ok(None) => Ok(0), - Err(e) => Err(into_error(e)), - } + let limit = self.timeouts.query; + let result = tokio::time::timeout(limit, async { + let s = conn + .prep(format!("SELECT v FROM {table} WHERE k = ?")) + .await + .map_err(into_error)?; + match conn.exec_first::(&s, (key,)).await { + Ok(Some(num)) => Ok(num), + Ok(None) => Ok(0), + Err(e) => Err(into_error(e)), + } + }) + .await; + bounded(conn, result, limit) } } diff --git a/crates/store/src/backend/mysql/search.rs b/crates/store/src/backend/mysql/search.rs index f8c1a67..906ce66 100644 --- a/crates/store/src/backend/mysql/search.rs +++ b/crates/store/src/backend/mysql/search.rs @@ -10,8 +10,8 @@ use crate::{ backend::{ MAX_TOKEN_LENGTH, mysql::{ - DELETE_CHUNK_SIZE, MIN_DELETE_CHUNK_SIZE, MysqlSearchField, MysqlStore, into_error, - is_timeout_error, + DELETE_CHUNK_SIZE, MIN_DELETE_CHUNK_SIZE, MysqlSearchField, MysqlStore, bounded, + into_error, is_timeout_error, }, }, search::{ @@ -27,57 +27,62 @@ use std::fmt::Write; impl MysqlStore { pub async fn index(&self, documents: Vec) -> trc::Result<()> { let mut conn = self.conn().await?; - let mut tx_opts = TxOpts::default(); - tx_opts - .with_consistent_snapshot(false) - .with_isolation_level(IsolationLevel::ReadCommitted); - let mut trx = conn.start_transaction(tx_opts).await.map_err(into_error)?; + let limit = self.timeouts.query; + let result = tokio::time::timeout(limit, async { + let mut tx_opts = TxOpts::default(); + tx_opts + .with_consistent_snapshot(false) + .with_isolation_level(IsolationLevel::ReadCommitted); + let mut trx = conn.start_transaction(tx_opts).await.map_err(into_error)?; - for document in documents { - let index = document.index; - let primary_keys = index.primary_keys(); - let all_fields = index.all_fields(); - let mut fields = document.fields; - let mut values = Vec::with_capacity(fields.len() + 2); - let mut query = format!("INSERT INTO {} (", index.mysql_table()); + for document in documents { + let index = document.index; + let primary_keys = index.primary_keys(); + let all_fields = index.all_fields(); + let mut fields = document.fields; + let mut values = Vec::with_capacity(fields.len() + 2); + let mut query = format!("INSERT INTO {} (", index.mysql_table()); - for (i, field) in primary_keys.iter().chain(all_fields).enumerate() { - if i > 0 { - query.push(','); + for (i, field) in primary_keys.iter().chain(all_fields).enumerate() { + if i > 0 { + query.push(','); + } + query.push_str(field.column()); } - query.push_str(field.column()); + + query.push_str(") VALUES ("); + + for (i, field) in primary_keys.iter().chain(all_fields).enumerate() { + if i > 0 { + query.push(','); + } + + if let Some(value) = fields.remove(field) { + query.push('?'); + values.push(value); + } else { + query.push_str("NULL"); + } + } + + query.push_str(") ON DUPLICATE KEY UPDATE "); + for (i, field) in all_fields.iter().enumerate() { + if i > 0 { + query.push(','); + } + let column = field.column(); + let _ = write!(&mut query, "{column} = VALUES({column})"); + } + + let s = trx.prep(&query).await.map_err(into_error)?; + + trx.exec_drop(&s, values).await.map_err(into_error)?; } - query.push_str(") VALUES ("); - - for (i, field) in primary_keys.iter().chain(all_fields).enumerate() { - if i > 0 { - query.push(','); - } - - if let Some(value) = fields.remove(field) { - query.push('?'); - values.push(value); - } else { - query.push_str("NULL"); - } - } - - query.push_str(") ON DUPLICATE KEY UPDATE "); - for (i, field) in all_fields.iter().enumerate() { - if i > 0 { - query.push(','); - } - let column = field.column(); - let _ = write!(&mut query, "{column} = VALUES({column})"); - } - - let s = trx.prep(&query).await.map_err(into_error)?; - - trx.exec_drop(&s, values).await.map_err(into_error)?; - } - - trx.commit().await.map_err(into_error) + trx.commit().await.map_err(into_error) + }) + .await; + bounded(conn, result, limit) } pub async fn query( @@ -97,12 +102,17 @@ impl MysqlStore { } let mut conn = self.conn().await?; - let s = conn.prep(query).await.map_err(into_error)?; + let limit = self.timeouts.query; + let result = tokio::time::timeout(limit, async { + let s = conn.prep(query).await.map_err(into_error)?; - conn.exec::(s, params) - .await - .map(|r| r.into_iter().map(|r| R::from_u64(r as u64)).collect()) - .map_err(into_error) + conn.exec::(s, params) + .await + .map(|r| r.into_iter().map(|r| R::from_u64(r as u64)).collect()) + .map_err(into_error) + }) + .await; + bounded(conn, result, limit) } pub async fn unindex(&self, filter: SearchQuery) -> trc::Result { @@ -111,40 +121,47 @@ impl MysqlStore { let params = build_filter(&mut query, &filter.filters); let mut conn = self.conn().await?; - let s = conn.prep(&query).await.map_err(into_error)?; + let limit = self.timeouts.maintenance; + let result = tokio::time::timeout(limit, async { + let s = conn.prep(&query).await.map_err(into_error)?; - match conn.exec_drop(s, params.clone()).await { - Ok(_) => return Ok(conn.affected_rows()), - Err(err) if is_timeout_error(&err) => (), - Err(err) => return Err(into_error(err)), - } + match conn.exec_drop(s, params.clone()).await { + Ok(_) => return Ok(conn.affected_rows()), + Err(err) if is_timeout_error(&err) => (), + Err(err) => return Err(into_error(err)), + } - let mut chunk_size = DELETE_CHUNK_SIZE; - let mut deleted = 0; - - loop { - let s = conn - .prep(format!("{query} LIMIT {chunk_size}")) - .await - .map_err(into_error)?; + let mut chunk_size = DELETE_CHUNK_SIZE; + let mut deleted = 0; loop { - match conn.exec_drop(&s, params.clone()).await { - Ok(_) => { - let affected = conn.affected_rows(); - if affected == 0 { - return Ok(deleted); + let s = conn + .prep(format!("{query} LIMIT {chunk_size}")) + .await + .map_err(into_error)?; + + loop { + match conn.exec_drop(&s, params.clone()).await { + Ok(_) => { + let affected = conn.affected_rows(); + if affected == 0 { + return Ok(deleted); + } + deleted += affected; } - deleted += affected; + Err(err) + if is_timeout_error(&err) && chunk_size > MIN_DELETE_CHUNK_SIZE => + { + chunk_size = (chunk_size / 2).max(MIN_DELETE_CHUNK_SIZE); + break; + } + Err(err) => return Err(into_error(err)), } - Err(err) if is_timeout_error(&err) && chunk_size > MIN_DELETE_CHUNK_SIZE => { - chunk_size = (chunk_size / 2).max(MIN_DELETE_CHUNK_SIZE); - break; - } - Err(err) => return Err(into_error(err)), } } - } + }) + .await; + bounded(conn, result, limit) } } diff --git a/crates/store/src/backend/mysql/write.rs b/crates/store/src/backend/mysql/write.rs index 150ad3b..d869bc0 100644 --- a/crates/store/src/backend/mysql/write.rs +++ b/crates/store/src/backend/mysql/write.rs @@ -6,7 +6,9 @@ * 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, bounded, into_error, is_timeout_error, +}; use crate::{ IndexKey, Key, LogKey, SUBSPACE_COUNTER, SUBSPACE_IN_MEMORY_COUNTER, SUBSPACE_QUOTA, SUBSPACE_REGISTRY_IDX, @@ -32,41 +34,45 @@ impl MysqlStore { let start = Instant::now(); let mut retry_count = 0; let mut conn = self.conn().await?; - - loop { - let err = match self.write_trx(&mut conn, &mut batch).await { - Ok(result) => { - return Ok(result); - } - Err(err) => err, - }; - - let _ = conn.query_drop("ROLLBACK;").await; - - match err { - CommitError::Mysql(Error::Server(err)) - if [1062, 1213].contains(&err.code) - && retry_count < MAX_COMMIT_ATTEMPTS - && start.elapsed() < MAX_COMMIT_TIME => {} - /*CommitError::Retry => { - if retry_count > MAX_COMMIT_ATTEMPTS || start.elapsed() > MAX_COMMIT_TIME { - return Err(trc::StoreEvent::AssertValueFailed - .into_err() - .caused_by(trc::location!())); + let limit = self.timeouts.query; + let result = tokio::time::timeout(limit, async { + loop { + let err = match self.write_trx(&mut conn, &mut batch).await { + Ok(result) => { + return Ok(result); } - }*/ - CommitError::Mysql(err) => { - return Err(into_error(err)); - } - CommitError::Internal(err) => { - return Err(err); - } - } + Err(err) => err, + }; - let backoff = rand::rng().random_range(50..=300); - tokio::time::sleep(Duration::from_millis(backoff)).await; - retry_count += 1; - } + let _ = conn.query_drop("ROLLBACK;").await; + + match err { + CommitError::Mysql(Error::Server(err)) + if [1062, 1213].contains(&err.code) + && retry_count < MAX_COMMIT_ATTEMPTS + && start.elapsed() < MAX_COMMIT_TIME => {} + /*CommitError::Retry => { + if retry_count > MAX_COMMIT_ATTEMPTS || start.elapsed() > MAX_COMMIT_TIME { + return Err(trc::StoreEvent::AssertValueFailed + .into_err() + .caused_by(trc::location!())); + } + }*/ + CommitError::Mysql(err) => { + return Err(into_error(err)); + } + CommitError::Internal(err) => { + return Err(err); + } + } + + let backoff = rand::rng().random_range(50..=300); + tokio::time::sleep(Duration::from_millis(backoff)).await; + retry_count += 1; + } + }) + .await; + bounded(conn, result, limit) } async fn write_trx( @@ -385,71 +391,81 @@ impl MysqlStore { pub(crate) async fn purge_store(&self) -> trc::Result<()> { 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?; - } + let limit = self.timeouts.maintenance; + let result = tokio::time::timeout(limit, async { + for subspace in [SUBSPACE_QUOTA, SUBSPACE_COUNTER, SUBSPACE_IN_MEMORY_COUNTER] { + purge_table(&mut conn, char::from(subspace)).await?; + } - Ok(()) + Ok(()) + }) + .await; + bounded(conn, result, limit) } pub(crate) async fn delete_range(&self, from: impl Key, to: impl Key) -> trc::Result<()> { let mut conn = self.conn().await?; - let table = char::from(from.subspace()); - let mut from = from.serialize(0); - let to = to.serialize(0); + let limit = self.timeouts.maintenance; + let result = tokio::time::timeout(limit, async { + let table = char::from(from.subspace()); + let mut from = from.serialize(0); + let to = to.serialize(0); - let delete = conn - .prep(format!("DELETE FROM {table} WHERE k >= ? AND k < ?")) - .await - .map_err(into_error)?; - - match conn.exec_drop(&delete, (&from, &to)).await { - Ok(_) => return Ok(()), - Err(err) if is_timeout_error(&err) => (), - Err(err) => return Err(into_error(err)), - } - - let mut chunk_size = DELETE_CHUNK_SIZE; - - loop { - let boundary = conn - .prep(format!( - "SELECT k FROM {table} WHERE k >= ? AND k < ? ORDER BY k ASC LIMIT 1 OFFSET {chunk_size}" - )) + let delete = conn + .prep(format!("DELETE FROM {table} WHERE k >= ? AND k < ?")) .await .map_err(into_error)?; + match conn.exec_drop(&delete, (&from, &to)).await { + Ok(_) => return Ok(()), + Err(err) if is_timeout_error(&err) => (), + Err(err) => return Err(into_error(err)), + } + + let mut chunk_size = DELETE_CHUNK_SIZE; + loop { - let next = match conn - .exec_first::, _, _>(&boundary, (&from, &to)) + let boundary = conn + .prep(format!( + "SELECT k FROM {table} WHERE k >= ? AND k < ? ORDER BY k ASC LIMIT 1 OFFSET {chunk_size}" + )) .await - { - Ok(next) => next, - Err(err) if is_timeout_error(&err) && chunk_size > MIN_DELETE_CHUNK_SIZE => { - chunk_size = (chunk_size / 2).max(MIN_DELETE_CHUNK_SIZE); - break; - } - Err(err) => return Err(into_error(err)), - }; + .map_err(into_error)?; - match conn - .exec_drop(&delete, (&from, next.as_ref().unwrap_or(&to))) - .await - { - Ok(_) => (), - Err(err) if is_timeout_error(&err) && chunk_size > MIN_DELETE_CHUNK_SIZE => { - chunk_size = (chunk_size / 2).max(MIN_DELETE_CHUNK_SIZE); - break; - } - Err(err) => return Err(into_error(err)), - } + loop { + let next = match conn + .exec_first::, _, _>(&boundary, (&from, &to)) + .await + { + Ok(next) => next, + Err(err) if is_timeout_error(&err) && chunk_size > MIN_DELETE_CHUNK_SIZE => { + chunk_size = (chunk_size / 2).max(MIN_DELETE_CHUNK_SIZE); + break; + } + Err(err) => return Err(into_error(err)), + }; - match next { - Some(next) => from = next, - None => return Ok(()), + match conn + .exec_drop(&delete, (&from, next.as_ref().unwrap_or(&to))) + .await + { + Ok(_) => (), + Err(err) if is_timeout_error(&err) && chunk_size > MIN_DELETE_CHUNK_SIZE => { + chunk_size = (chunk_size / 2).max(MIN_DELETE_CHUNK_SIZE); + break; + } + Err(err) => return Err(into_error(err)), + } + + match next { + Some(next) => from = next, + None => return Ok(()), + } } } - } + }) + .await; + bounded(conn, result, limit) } } diff --git a/crates/store/src/backend/postgres/blob.rs b/crates/store/src/backend/postgres/blob.rs index 66d437c..716f946 100644 --- a/crates/store/src/backend/postgres/blob.rs +++ b/crates/store/src/backend/postgres/blob.rs @@ -2,13 +2,15 @@ * 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; use crate::backend::postgres::into_pool_error; -use super::{PostgresStore, into_error}; +use super::{PostgresStore, bounded, into_error}; impl PostgresStore { pub(crate) async fn get_blob( @@ -17,53 +19,68 @@ impl PostgresStore { range: Range, ) -> trc::Result>> { let conn = self.conn_pool.get().await.map_err(into_pool_error)?; - let s = conn - .prepare_cached("SELECT v FROM t WHERE k = $1") - .await - .map_err(into_error)?; - conn.query_opt(&s, &[&key]) - .await - .and_then(|row| { - if let Some(row) = row { - Ok(Some(if range.start == 0 && range.end == usize::MAX { - row.try_get::<_, Vec>(0)? + let limit = self.timeouts.query; + let result = tokio::time::timeout(limit, async { + let s = conn + .prepare_cached("SELECT v FROM t WHERE k = $1") + .await + .map_err(into_error)?; + conn.query_opt(&s, &[&key]) + .await + .and_then(|row| { + if let Some(row) = row { + Ok(Some(if range.start == 0 && range.end == usize::MAX { + row.try_get::<_, Vec>(0)? + } else { + let bytes = row.try_get::<_, &[u8]>(0)?; + bytes + .get(range.start..std::cmp::min(bytes.len(), range.end)) + .unwrap_or_default() + .to_vec() + })) } else { - let bytes = row.try_get::<_, &[u8]>(0)?; - bytes - .get(range.start..std::cmp::min(bytes.len(), range.end)) - .unwrap_or_default() - .to_vec() - })) - } else { - Ok(None) - } - }) - .map_err(into_error) + Ok(None) + } + }) + .map_err(into_error) + }) + .await; + bounded(conn, result, limit) } pub(crate) async fn put_blob(&self, key: &[u8], data: &[u8]) -> trc::Result<()> { let conn = self.conn_pool.get().await.map_err(into_pool_error)?; - let s = conn - .prepare_cached( - "INSERT INTO t (k, v) VALUES ($1, $2) ON CONFLICT (k) DO UPDATE SET v = EXCLUDED.v", - ) - .await - .map_err(into_error)?; - conn.execute(&s, &[&key, &data]) - .await - .map_err(into_error) - .map(|_| ()) + let limit = self.timeouts.query; + let result = tokio::time::timeout(limit, async { + let s = conn + .prepare_cached( + "INSERT INTO t (k, v) VALUES ($1, $2) ON CONFLICT (k) DO UPDATE SET v = EXCLUDED.v", + ) + .await + .map_err(into_error)?; + conn.execute(&s, &[&key, &data]) + .await + .map_err(into_error) + .map(|_| ()) + }) + .await; + bounded(conn, result, limit) } pub(crate) async fn delete_blob(&self, key: &[u8]) -> trc::Result { let conn = self.conn_pool.get().await.map_err(into_pool_error)?; - let s = conn - .prepare_cached("DELETE FROM t WHERE k = $1") - .await - .map_err(into_error)?; - conn.execute(&s, &[&key]) - .await - .map_err(into_error) - .map(|hits| hits > 0) + let limit = self.timeouts.query; + let result = tokio::time::timeout(limit, async { + let s = conn + .prepare_cached("DELETE FROM t WHERE k = $1") + .await + .map_err(into_error)?; + conn.execute(&s, &[&key]) + .await + .map_err(into_error) + .map(|hits| hits > 0) + }) + .await; + bounded(conn, result, limit) } } diff --git a/crates/store/src/backend/postgres/lookup.rs b/crates/store/src/backend/postgres/lookup.rs index 30b052f..8e647b2 100644 --- a/crates/store/src/backend/postgres/lookup.rs +++ b/crates/store/src/backend/postgres/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 crate::{QueryResult, QueryType, backend::postgres::into_pool_error}; @@ -12,7 +14,7 @@ use tokio_postgres::types::{FromSql, ToSql, Type}; use crate::IntoRows; -use super::{PostgresStore, into_error}; +use super::{PostgresStore, bounded, into_error}; impl PostgresStore { pub(crate) async fn sql_query( @@ -21,33 +23,38 @@ impl PostgresStore { params_: &[crate::Value<'_>], ) -> trc::Result { let conn = self.conn_pool.get().await.map_err(into_pool_error)?; - let s = conn.prepare_cached(query).await.map_err(into_error)?; - let params = params_ - .iter() - .map(|v| v as &(dyn tokio_postgres::types::ToSql + Sync)) - .collect::>(); + let limit = self.timeouts.query; + let result = tokio::time::timeout(limit, async { + let s = conn.prepare_cached(query).await.map_err(into_error)?; + let params = params_ + .iter() + .map(|v| v as &(dyn tokio_postgres::types::ToSql + Sync)) + .collect::>(); - match T::query_type() { - QueryType::Execute => conn - .execute(&s, params.as_slice()) - .await - .map_or_else(|e| Err(into_error(e)), |r| Ok(T::from_exec(r as usize))), - QueryType::Exists => { - let rows = conn.query_raw(&s, params).await.map_err(into_error)?; - pin_mut!(rows); - rows.try_next() + match T::query_type() { + QueryType::Execute => conn + .execute(&s, params.as_slice()) .await - .map_or_else(|e| Err(into_error(e)), |r| Ok(T::from_exists(r.is_some()))) + .map_or_else(|e| Err(into_error(e)), |r| Ok(T::from_exec(r as usize))), + QueryType::Exists => { + let rows = conn.query_raw(&s, params).await.map_err(into_error)?; + pin_mut!(rows); + rows.try_next() + .await + .map_or_else(|e| Err(into_error(e)), |r| Ok(T::from_exists(r.is_some()))) + } + QueryType::QueryOne => conn + .query_opt(&s, params.as_slice()) + .await + .map_or_else(|e| Err(into_error(e)), |r| Ok(T::from_query_one(r))), + QueryType::QueryAll => conn + .query(&s, params.as_slice()) + .await + .map_or_else(|e| Err(into_error(e)), |r| Ok(T::from_query_all(r))), } - QueryType::QueryOne => conn - .query_opt(&s, params.as_slice()) - .await - .map_or_else(|e| Err(into_error(e)), |r| Ok(T::from_query_one(r))), - QueryType::QueryAll => conn - .query(&s, params.as_slice()) - .await - .map_or_else(|e| Err(into_error(e)), |r| Ok(T::from_query_all(r))), - } + }) + .await; + bounded(conn, result, limit) } } diff --git a/crates/store/src/backend/postgres/main.rs b/crates/store/src/backend/postgres/main.rs index 6dbd853..00c8ed0 100644 --- a/crates/store/src/backend/postgres/main.rs +++ b/crates/store/src/backend/postgres/main.rs @@ -6,7 +6,7 @@ * Modified by Coffey Labs in 2026 for INBUXA. */ -use super::{PostgresStore, into_error}; +use super::{PostgresStore, bounded, into_error}; use crate::{ backend::postgres::{ PsqlSearchField, into_pool_error, @@ -119,6 +119,7 @@ impl PostgresStore { Store::PostgreSQL(Arc::new(PostgresStore { conn_pool: pool, ts_configs: ts_configs.clone(), + timeouts: Default::default(), })), replica.host, replica.port as u16, @@ -129,6 +130,7 @@ impl PostgresStore { let primary = Store::PostgreSQL(Arc::new(PostgresStore { conn_pool: primary_pool, ts_configs, + timeouts: Default::default(), })); // ST-1: no replicas, no change @@ -147,84 +149,92 @@ impl PostgresStore { pub(crate) async fn create_storage_tables(&self) -> trc::Result<()> { let conn = self.conn_pool.get().await.map_err(into_pool_error)?; + let limit = self.timeouts.maintenance; + let result = tokio::time::timeout(limit, async { + for table in [ + SUBSPACE_ACL, + SUBSPACE_TASK_QUEUE, + SUBSPACE_DELETED_ITEMS, + SUBSPACE_SPAM_SAMPLES, + crate::SUBSPACE_INBUXA, // inbuxa: masked email + SUBSPACE_BLOB_LINK, + SUBSPACE_IN_MEMORY_VALUE, + SUBSPACE_PROPERTY, + SUBSPACE_REGISTRY, + SUBSPACE_REGISTRY_PK, + SUBSPACE_QUEUE_MESSAGE, + SUBSPACE_QUEUE_EVENT, + SUBSPACE_REPORT_OUT, + SUBSPACE_REPORT_IN, + SUBSPACE_LOGS, + SUBSPACE_BLOBS, + SUBSPACE_DIRECTORY, + SUBSPACE_TELEMETRY_SPAN, + SUBSPACE_TELEMETRY_METRIC, + ] { + let table = char::from(table); + conn.execute( + &format!( + "CREATE TABLE IF NOT EXISTS {table} ( + k BYTEA PRIMARY KEY, + v BYTEA NOT NULL + )" + ), + &[], + ) + .await + .map_err(into_error)?; + } - for table in [ - SUBSPACE_ACL, - SUBSPACE_TASK_QUEUE, - SUBSPACE_DELETED_ITEMS, - SUBSPACE_SPAM_SAMPLES, - crate::SUBSPACE_INBUXA, // inbuxa: masked email - SUBSPACE_BLOB_LINK, - SUBSPACE_IN_MEMORY_VALUE, - SUBSPACE_PROPERTY, - SUBSPACE_REGISTRY, - SUBSPACE_REGISTRY_PK, - SUBSPACE_QUEUE_MESSAGE, - SUBSPACE_QUEUE_EVENT, - SUBSPACE_REPORT_OUT, - SUBSPACE_REPORT_IN, - SUBSPACE_LOGS, - SUBSPACE_BLOBS, - SUBSPACE_DIRECTORY, - SUBSPACE_TELEMETRY_SPAN, - SUBSPACE_TELEMETRY_METRIC, - ] { - let table = char::from(table); - conn.execute( - &format!( - "CREATE TABLE IF NOT EXISTS {table} ( + for table in [SUBSPACE_INDEXES, SUBSPACE_REGISTRY_IDX] { + let table = char::from(table); + conn.execute( + &format!( + "CREATE TABLE IF NOT EXISTS {table} ( + k BYTEA PRIMARY KEY + )" + ), + &[], + ) + .await + .map_err(into_error)?; + } + + for table in [SUBSPACE_COUNTER, SUBSPACE_QUOTA, SUBSPACE_IN_MEMORY_COUNTER] { + conn.execute( + &format!( + "CREATE TABLE IF NOT EXISTS {} ( k BYTEA PRIMARY KEY, - v BYTEA NOT NULL - )" - ), - &[], - ) - .await - .map_err(into_error)?; - } + v BIGINT NOT NULL DEFAULT 0 + )", + char::from(table) + ), + &[], + ) + .await + .map_err(into_error)?; + } - for table in [SUBSPACE_INDEXES, SUBSPACE_REGISTRY_IDX] { - let table = char::from(table); - conn.execute( - &format!( - "CREATE TABLE IF NOT EXISTS {table} ( - k BYTEA PRIMARY KEY - )" - ), - &[], - ) - .await - .map_err(into_error)?; - } - - for table in [SUBSPACE_COUNTER, SUBSPACE_QUOTA, SUBSPACE_IN_MEMORY_COUNTER] { - conn.execute( - &format!( - "CREATE TABLE IF NOT EXISTS {} ( - k BYTEA PRIMARY KEY, - v BIGINT NOT NULL DEFAULT 0 - )", - char::from(table) - ), - &[], - ) - .await - .map_err(into_error)?; - } - - Ok(()) + Ok(()) + }) + .await; + bounded(conn, result, limit) } pub(crate) async fn create_search_tables(&self) -> trc::Result<()> { let conn = self.conn_pool.get().await.map_err(into_pool_error)?; + let limit = self.timeouts.maintenance; + let result = tokio::time::timeout(limit, async { + create_search_tables::(&conn).await?; + create_search_tables::(&conn).await?; + create_search_tables::(&conn).await?; + //create_search_tables::(&conn).await?; + create_search_tables::(&conn).await?; - create_search_tables::(&conn).await?; - create_search_tables::(&conn).await?; - create_search_tables::(&conn).await?; - //create_search_tables::(&conn).await?; - create_search_tables::(&conn).await?; - - Ok(()) + Ok(()) + }) + .await; + bounded(conn, result, limit) } } diff --git a/crates/store/src/backend/postgres/mod.rs b/crates/store/src/backend/postgres/mod.rs index 19b1dc9..ebe425e 100644 --- a/crates/store/src/backend/postgres/mod.rs +++ b/crates/store/src/backend/postgres/mod.rs @@ -6,6 +6,7 @@ * Modified by Coffey Labs in 2026 for INBUXA. */ +use crate::backend::query_timeout::QueryTimeouts; use crate::{ search::{ CalendarSearchField, ContactSearchField, EmailSearchField, FileSearchField, SearchField, @@ -14,7 +15,8 @@ use crate::{ write::SearchIndex, }; use ahash::AHashSet; -use deadpool_postgres::Pool; +use deadpool_postgres::{Object, Pool}; +use std::time::Duration; use tokio_postgres::error::SqlState; pub mod blob; @@ -28,6 +30,8 @@ pub mod write; pub struct PostgresStore { pub(crate) conn_pool: Pool, pub(crate) ts_configs: AHashSet<&'static str>, + /// inbuxa: client-side query limits (see backend::query_timeout) + pub(crate) timeouts: QueryTimeouts, } #[inline(always)] @@ -72,6 +76,34 @@ pub(crate) fn is_timeout_error(err: &tokio_postgres::Error) -> bool { }) } +/// inbuxa: the error for an operation that ran past its time limit. +pub(crate) fn query_timeout_error(limit: Duration) -> trc::Error { + trc::StoreEvent::PostgresqlError + .reason("Query timed out") + .details(format!( + "No answer from the database within {} s", + limit.as_secs() + )) +} + +/// inbuxa: ends an operation run on `conn` under `limit`. When it ran out, +/// the connection is taken out of the pool and closed: a query may still be +/// in flight on it, or a transaction open, so it can't be handed to the +/// next caller. +pub(crate) fn bounded( + conn: Object, + result: Result, tokio::time::error::Elapsed>, + limit: Duration, +) -> trc::Result { + match result { + Ok(result) => result, + Err(_) => { + drop(Object::take(conn)); + Err(query_timeout_error(limit)) + } + } +} + #[inline(always)] pub(crate) fn into_pool_error(err: deadpool_postgres::PoolError) -> trc::Error { match err { diff --git a/crates/store/src/backend/postgres/read.rs b/crates/store/src/backend/postgres/read.rs index 813d853..681c3d4 100644 --- a/crates/store/src/backend/postgres/read.rs +++ b/crates/store/src/backend/postgres/read.rs @@ -2,9 +2,11 @@ * 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::{PostgresStore, into_error, is_timeout_error}; +use super::{PostgresStore, bounded, into_error, is_timeout_error, query_timeout_error}; use crate::{ Deserialize, IterateParams, Key, ValueKey, backend::postgres::into_pool_error, write::ValueClass, @@ -17,40 +19,50 @@ impl PostgresStore { U: Deserialize + 'static, { let conn = self.conn_pool.get().await.map_err(into_pool_error)?; - let s = conn - .prepare_cached(&format!( - "SELECT v FROM {} WHERE k = $1", - char::from(key.subspace()) - )) - .await - .map_err(into_error)?; - let key = key.serialize(0); - conn.query_opt(&s, &[&key]) - .await - .map_err(into_error) - .and_then(|r| { - if let Some(r) = r { - Ok(Some(U::deserialize_with_key(&key, r.get(0))?)) - } else { - Ok(None) - } - }) + let limit = self.timeouts.query; + let result = tokio::time::timeout(limit, async { + let s = conn + .prepare_cached(&format!( + "SELECT v FROM {} WHERE k = $1", + char::from(key.subspace()) + )) + .await + .map_err(into_error)?; + let key = key.serialize(0); + conn.query_opt(&s, &[&key]) + .await + .map_err(into_error) + .and_then(|r| { + if let Some(r) = r { + Ok(Some(U::deserialize_with_key(&key, r.get(0))?)) + } else { + Ok(None) + } + }) + }) + .await; + bounded(conn, result, limit) } pub(crate) async fn key_exists(&self, key: impl Key) -> trc::Result { let conn = self.conn_pool.get().await.map_err(into_pool_error)?; - let s = conn - .prepare_cached(&format!( - "SELECT 1 FROM {} WHERE k = $1", - char::from(key.subspace()) - )) - .await - .map_err(into_error)?; - let key = key.serialize(0); - conn.query_opt(&s, &[&key]) - .await - .map_err(into_error) - .map(|r| r.is_some()) + let limit = self.timeouts.query; + let result = tokio::time::timeout(limit, async { + let s = conn + .prepare_cached(&format!( + "SELECT 1 FROM {} WHERE k = $1", + char::from(key.subspace()) + )) + .await + .map_err(into_error)?; + let key = key.serialize(0); + conn.query_opt(&s, &[&key]) + .await + .map_err(into_error) + .map(|r| r.is_some()) + }) + .await; + bounded(conn, result, limit) } pub(crate) async fn iterate( @@ -64,44 +76,65 @@ impl PostgresStore { let end = params.end.serialize(0); let keys = if params.values { "k, v" } else { "k" }; - let s = conn - .prepare_cached(&match (params.first, params.ascending) { - (true, true) => { - format!( - "SELECT {keys} FROM {table} WHERE k >= $1 AND k <= $2 ORDER BY k ASC LIMIT 1" - ) - } - (true, false) => { - format!( + // inbuxa: a scan may run for hours, so the query limit bounds each + // wait for the database (preparing, the query starting, the next + // row) rather than the scan. A wait that runs out closes the + // connection. + let limit = self.timeouts.query; + let query = match (params.first, params.ascending) { + (true, true) => { + format!( + "SELECT {keys} FROM {table} WHERE k >= $1 AND k <= $2 ORDER BY k ASC LIMIT 1" + ) + } + (true, false) => { + format!( "SELECT {keys} FROM {table} WHERE k >= $1 AND k <= $2 ORDER BY k DESC LIMIT 1" ) - } - (false, true) => { - format!("SELECT {keys} FROM {table} WHERE k >= $1 AND k <= $2 ORDER BY k ASC") - } - (false, false) => { - format!("SELECT {keys} FROM {table} WHERE k >= $1 AND k <= $2 ORDER BY k DESC") - } - }) - .await.map_err(into_error)?; + } + (false, true) => { + format!("SELECT {keys} FROM {table} WHERE k >= $1 AND k <= $2 ORDER BY k ASC") + } + (false, false) => { + format!("SELECT {keys} FROM {table} WHERE k >= $1 AND k <= $2 ORDER BY k DESC") + } + }; + let s = match tokio::time::timeout(limit, conn.prepare_cached(&query)).await { + Ok(s) => s.map_err(into_error)?, + Err(_) => { + drop(deadpool_postgres::Object::take(conn)); + return Err(query_timeout_error(limit)); + } + }; let mut from = begin; let mut to = end; let mut resume_key: Option> = None; + let mut stalled = false; + loop { let mut last_key = None; let mut timed_out = false; { - let rows = conn - .query_raw(&s, &[&from, &to]) - .await - .map_err(into_error)?; + let rows = + match tokio::time::timeout(limit, conn.query_raw(&s, &[&from, &to])).await { + Ok(rows) => rows.map_err(into_error)?, + // Leaves the scan loop for the timeout below + Err(_) => break, + }; pin_mut!(rows); loop { - match rows.try_next().await { + let next = match tokio::time::timeout(limit, rows.try_next()).await { + Ok(next) => next, + Err(_) => { + stalled = true; + break; + } + }; + match next { Ok(Some(row)) => { let key = row.try_get::<_, &[u8]>(0).map_err(into_error)?; let value = if params.values { @@ -132,6 +165,10 @@ impl PostgresStore { } } + if stalled { + break; + } + match last_key { Some(last_key) if timed_out => { if params.ascending { @@ -144,6 +181,9 @@ impl PostgresStore { _ => return Ok(()), } } + + drop(deadpool_postgres::Object::take(conn)); + Err(query_timeout_error(limit)) } pub(crate) async fn get_counter( @@ -155,14 +195,19 @@ impl PostgresStore { let key = key.serialize(0); let conn = self.conn_pool.get().await.map_err(into_pool_error)?; - let s = conn - .prepare_cached(&format!("SELECT v FROM {table} WHERE k = $1")) - .await - .map_err(into_error)?; - match conn.query_opt(&s, &[&key]).await { - Ok(Some(row)) => row.try_get(0).map_err(into_error), - Ok(None) => Ok(0), - Err(e) => Err(into_error(e)), - } + let limit = self.timeouts.query; + let result = tokio::time::timeout(limit, async { + let s = conn + .prepare_cached(&format!("SELECT v FROM {table} WHERE k = $1")) + .await + .map_err(into_error)?; + match conn.query_opt(&s, &[&key]).await { + Ok(Some(row)) => row.try_get(0).map_err(into_error), + Ok(None) => Ok(0), + Err(e) => Err(into_error(e)), + } + }) + .await; + bounded(conn, result, limit) } } diff --git a/crates/store/src/backend/postgres/search.rs b/crates/store/src/backend/postgres/search.rs index a626894..92ba457 100644 --- a/crates/store/src/backend/postgres/search.rs +++ b/crates/store/src/backend/postgres/search.rs @@ -10,8 +10,8 @@ use crate::{ backend::{ MAX_TOKEN_LENGTH, postgres::{ - DELETE_CHUNK_SIZE, MIN_DELETE_CHUNK_SIZE, PostgresStore, PsqlSearchField, into_error, - into_pool_error, is_timeout_error, + DELETE_CHUNK_SIZE, MIN_DELETE_CHUNK_SIZE, PostgresStore, PsqlSearchField, bounded, + into_error, into_pool_error, is_timeout_error, }, }, search::{ @@ -36,125 +36,130 @@ impl PostgresStore { pub async fn index(&self, documents: Vec) -> trc::Result<()> { let mut conn = self.conn_pool.get().await.map_err(into_pool_error)?; - let trx = conn - .build_transaction() - .isolation_level(IsolationLevel::ReadCommitted) - .start() - .await - .map_err(into_error)?; + let limit = self.timeouts.query; + let result = tokio::time::timeout(limit, async { + let trx = conn + .build_transaction() + .isolation_level(IsolationLevel::ReadCommitted) + .start() + .await + .map_err(into_error)?; - for document in documents { - let index = document.index; - let primary_keys = index.primary_keys(); - let all_fields = index.all_fields(); - let fields = document.fields; - // inbuxa: keyword text (addresses, contact fields, ...) is split into - // words before it reaches the text parser, see keyword_terms(). - let keywords = primary_keys - .iter() - .chain(all_fields) - .map(|field| match fields.get(field) { - Some(SearchValue::Text { - value, - language: Language::None, - }) if field.is_text() => Some(keyword_terms(value)), - _ => None, - }) - .collect::>(); - let mut values = Vec::with_capacity(fields.len() + 2); - let mut query = format!("INSERT INTO {} (", index.psql_table()); + for document in documents { + let index = document.index; + let primary_keys = index.primary_keys(); + let all_fields = index.all_fields(); + let fields = document.fields; + // inbuxa: keyword text (addresses, contact fields, ...) is split into + // words before it reaches the text parser, see keyword_terms(). + let keywords = primary_keys + .iter() + .chain(all_fields) + .map(|field| match fields.get(field) { + Some(SearchValue::Text { + value, + language: Language::None, + }) if field.is_text() => Some(keyword_terms(value)), + _ => None, + }) + .collect::>(); + let mut values = Vec::with_capacity(fields.len() + 2); + let mut query = format!("INSERT INTO {} (", index.psql_table()); - for (i, field) in primary_keys.iter().chain(all_fields).enumerate() { - if i > 0 { - query.push(','); - } - query.push_str(field.column()); + for (i, field) in primary_keys.iter().chain(all_fields).enumerate() { + if i > 0 { + query.push(','); + } + query.push_str(field.column()); - if let Some(sort_column) = field.sort_column() { - query.push(','); - query.push_str(sort_column); - } - } - - query.push_str(") VALUES ("); - - for (i, field) in primary_keys.iter().chain(all_fields).enumerate() { - if i > 0 { - query.push(','); + if let Some(sort_column) = field.sort_column() { + query.push(','); + query.push_str(sort_column); + } } - if let Some(value) = fields.get(field) { - let value_ref = format!("${}", values.len() + 1); - let (text_len, language) = if let SearchValue::Text { value, language } = value - { - (value.len(), self.ts_config(language)) - } else { - (0, PG_UNSTEMMED_LANG) - }; + query.push_str(") VALUES ("); - if let Some(keywords) = &keywords[i] { - let _ = write!(&mut query, "to_tsvector('{language}',{value_ref})"); - values.push(keywords as &(dyn ToSql + Sync)); - if field.sort_column().is_some() { - let value_ref = format!("${}", values.len() + 1); - if text_len > 255 { - let _ = write!(&mut query, ",left({value_ref},255)"); + for (i, field) in primary_keys.iter().chain(all_fields).enumerate() { + if i > 0 { + query.push(','); + } + + if let Some(value) = fields.get(field) { + let value_ref = format!("${}", values.len() + 1); + let (text_len, language) = + if let SearchValue::Text { value, language } = value { + (value.len(), self.ts_config(language)) } else { - let _ = write!(&mut query, ",{value_ref}"); + (0, PG_UNSTEMMED_LANG) + }; + + if let Some(keywords) = &keywords[i] { + let _ = write!(&mut query, "to_tsvector('{language}',{value_ref})"); + values.push(keywords as &(dyn ToSql + Sync)); + if field.sort_column().is_some() { + let value_ref = format!("${}", values.len() + 1); + if text_len > 255 { + let _ = write!(&mut query, ",left({value_ref},255)"); + } else { + let _ = write!(&mut query, ",{value_ref}"); + } + values.push(value as &(dyn ToSql + Sync)); } - values.push(value as &(dyn ToSql + Sync)); - } - continue; - } else if field.is_text() { - let _ = write!(&mut query, "to_tsvector('{language}',{value_ref})"); - } else if text_len > 512 { - query.push_str("left("); - query.push_str(&value_ref); - query.push_str(",512)"); - } else { - query.push_str(&value_ref); - } - - if field.sort_column().is_some() { - if text_len > 255 { - query.push_str(",left("); + continue; + } else if field.is_text() { + let _ = write!(&mut query, "to_tsvector('{language}',{value_ref})"); + } else if text_len > 512 { + query.push_str("left("); query.push_str(&value_ref); - query.push_str(",255)"); + query.push_str(",512)"); } else { - query.push(','); query.push_str(&value_ref); } - } - values.push(value as &(dyn ToSql + Sync)); - } else { - query.push_str("NULL"); - if field.sort_column().is_some() { - query.push_str(",NULL"); + if field.sort_column().is_some() { + if text_len > 255 { + query.push_str(",left("); + query.push_str(&value_ref); + query.push_str(",255)"); + } else { + query.push(','); + query.push_str(&value_ref); + } + } + + values.push(value as &(dyn ToSql + Sync)); + } else { + query.push_str("NULL"); + if field.sort_column().is_some() { + query.push_str(",NULL"); + } } } - } - query.push_str(") ON CONFLICT ("); - for (i, pkey) in primary_keys.iter().enumerate() { - if i > 0 { - query.push(','); + query.push_str(") ON CONFLICT ("); + for (i, pkey) in primary_keys.iter().enumerate() { + if i > 0 { + query.push(','); + } + query.push_str(pkey.column()); } - query.push_str(pkey.column()); - } - query.push_str(") DO UPDATE SET "); - for (i, field) in all_fields.iter().enumerate() { - if i > 0 { - query.push(','); + query.push_str(") DO UPDATE SET "); + for (i, field) in all_fields.iter().enumerate() { + if i > 0 { + query.push(','); + } + let column = field.column(); + let _ = write!(&mut query, "{column} = EXCLUDED.{column}"); } - let column = field.column(); - let _ = write!(&mut query, "{column} = EXCLUDED.{column}"); + + trx.execute(&query, &values).await.map_err(into_error)?; } - trx.execute(&query, &values).await.map_err(into_error)?; - } - - trx.commit().await.map_err(into_error) + trx.commit().await.map_err(into_error) + }) + .await; + bounded(conn, result, limit) } pub async fn query( @@ -170,16 +175,21 @@ impl PostgresStore { build_sort(&mut query, sort); } let conn = self.conn_pool.get().await.map_err(into_pool_error)?; - let s = conn.prepare_cached(&query).await.map_err(into_error)?; + let limit = self.timeouts.query; + let result = tokio::time::timeout(limit, async { + let s = conn.prepare_cached(&query).await.map_err(into_error)?; - conn.query(&s, params.as_slice()) - .await - .and_then(|rows| { - rows.into_iter() - .map(|row| row.try_get::<_, DocId>(0).map(|v| R::from_u64(v.0))) - .collect::, _>>() - }) - .map_err(into_error) + conn.query(&s, params.as_slice()) + .await + .and_then(|rows| { + rows.into_iter() + .map(|row| row.try_get::<_, DocId>(0).map(|v| R::from_u64(v.0))) + .collect::, _>>() + }) + .map_err(into_error) + }) + .await; + bounded(conn, result, limit) } pub async fn unindex(&self, filter: SearchQuery) -> trc::Result { @@ -189,40 +199,45 @@ impl PostgresStore { let params = self.build_filter(&mut where_clause, &filter.filters); let params = params.iter().map(SqlParam::as_sql).collect::>(); let conn = self.conn_pool.get().await.map_err(into_pool_error)?; - let s = conn - .prepare_cached(&format!("DELETE FROM {table}{where_clause}")) - .await - .map_err(into_error)?; - - match conn.execute(&s, params.as_slice()).await { - Ok(deleted) => return Ok(deleted), - Err(err) if is_timeout_error(&err) => (), - Err(err) => return Err(into_error(err)), - } - - let mut chunk_size = DELETE_CHUNK_SIZE; - let mut deleted = 0; - - loop { + let limit = self.timeouts.maintenance; + let result = tokio::time::timeout(limit, async { let s = conn - .prepare_cached(&format!( - "DELETE FROM {table} WHERE ctid IN (SELECT ctid FROM {table}{where_clause} LIMIT {chunk_size})" - )) + .prepare_cached(&format!("DELETE FROM {table}{where_clause}")) .await .map_err(into_error)?; + match conn.execute(&s, params.as_slice()).await { + Ok(deleted) => return Ok(deleted), + Err(err) if is_timeout_error(&err) => (), + Err(err) => return Err(into_error(err)), + } + + let mut chunk_size = DELETE_CHUNK_SIZE; + let mut deleted = 0; + loop { - match conn.execute(&s, params.as_slice()).await { - Ok(0) => return Ok(deleted), - Ok(affected) => deleted += affected, - Err(err) if is_timeout_error(&err) && chunk_size > MIN_DELETE_CHUNK_SIZE => { - chunk_size = (chunk_size / 2).max(MIN_DELETE_CHUNK_SIZE); - break; + let s = conn + .prepare_cached(&format!( + "DELETE FROM {table} WHERE ctid IN (SELECT ctid FROM {table}{where_clause} LIMIT {chunk_size})" + )) + .await + .map_err(into_error)?; + + loop { + match conn.execute(&s, params.as_slice()).await { + Ok(0) => return Ok(deleted), + Ok(affected) => deleted += affected, + Err(err) if is_timeout_error(&err) && chunk_size > MIN_DELETE_CHUNK_SIZE => { + chunk_size = (chunk_size / 2).max(MIN_DELETE_CHUNK_SIZE); + break; + } + Err(err) => return Err(into_error(err)), } - Err(err) => return Err(into_error(err)), } } - } + }) + .await; + bounded(conn, result, limit) } fn build_filter<'x>( diff --git a/crates/store/src/backend/postgres/write.rs b/crates/store/src/backend/postgres/write.rs index d91e7f6..77f2af0 100644 --- a/crates/store/src/backend/postgres/write.rs +++ b/crates/store/src/backend/postgres/write.rs @@ -2,9 +2,11 @@ * 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::{PostgresStore, into_error, is_timeout_error}; +use super::{PostgresStore, bounded, into_error, is_timeout_error}; use crate::{ IndexKey, Key, LogKey, SUBSPACE_COUNTER, SUBSPACE_IN_MEMORY_COUNTER, SUBSPACE_QUOTA, SUBSPACE_REGISTRY_IDX, @@ -30,48 +32,53 @@ enum CommitError { impl PostgresStore { pub(crate) async fn write(&self, mut batch: Batch<'_>) -> trc::Result { let mut conn = self.conn_pool.get().await.map_err(into_pool_error)?; - let start = Instant::now(); - let mut retry_count = 0; + let limit = self.timeouts.query; + let result = tokio::time::timeout(limit, async { + let start = Instant::now(); + let mut retry_count = 0; - loop { - match self.write_trx(&mut conn, &mut batch).await { - Ok(result) => { - return Ok(result); - } - Err(err) => { - match err { - CommitError::Postgres(err) => match err.code() { - Some( - &SqlState::T_R_SERIALIZATION_FAILURE - | &SqlState::T_R_DEADLOCK_DETECTED, - ) if retry_count < MAX_COMMIT_ATTEMPTS - && start.elapsed() < MAX_COMMIT_TIME => {} - Some(&SqlState::UNIQUE_VIOLATION) => { - return Err(trc::StoreEvent::AssertValueFailed - .into_err() - .reason("Unique violation") - .caused_by(trc::location!())); - } - _ => return Err(into_error(err)), - }, - CommitError::Internal(err) => return Err(err), - /*CommitError::Retry => { - if retry_count > MAX_COMMIT_ATTEMPTS - || start.elapsed() > MAX_COMMIT_TIME - { - return Err(trc::StoreEvent::AssertValueFailed - .into_err() - .caused_by(trc::location!())); - } - }*/ + loop { + match self.write_trx(&mut conn, &mut batch).await { + Ok(result) => { + return Ok(result); } + Err(err) => { + match err { + CommitError::Postgres(err) => match err.code() { + Some( + &SqlState::T_R_SERIALIZATION_FAILURE + | &SqlState::T_R_DEADLOCK_DETECTED, + ) if retry_count < MAX_COMMIT_ATTEMPTS + && start.elapsed() < MAX_COMMIT_TIME => {} + Some(&SqlState::UNIQUE_VIOLATION) => { + return Err(trc::StoreEvent::AssertValueFailed + .into_err() + .reason("Unique violation") + .caused_by(trc::location!())); + } + _ => return Err(into_error(err)), + }, + CommitError::Internal(err) => return Err(err), + /*CommitError::Retry => { + if retry_count > MAX_COMMIT_ATTEMPTS + || start.elapsed() > MAX_COMMIT_TIME + { + return Err(trc::StoreEvent::AssertValueFailed + .into_err() + .caused_by(trc::location!())); + } + }*/ + } - let backoff = rand::rng().random_range(50..=300); - tokio::time::sleep(Duration::from_millis(backoff)).await; - retry_count += 1; + let backoff = rand::rng().random_range(50..=300); + tokio::time::sleep(Duration::from_millis(backoff)).await; + retry_count += 1; + } } } - } + }) + .await; + bounded(conn, result, limit) } async fn write_trx( @@ -393,72 +400,81 @@ impl PostgresStore { pub(crate) async fn purge_store(&self) -> trc::Result<()> { let conn = self.conn_pool.get().await.map_err(into_pool_error)?; + let limit = self.timeouts.maintenance; + let result = tokio::time::timeout(limit, async { + for subspace in [SUBSPACE_QUOTA, SUBSPACE_COUNTER, SUBSPACE_IN_MEMORY_COUNTER] { + purge_table(&conn, char::from(subspace)).await?; + } - for subspace in [SUBSPACE_QUOTA, SUBSPACE_COUNTER, SUBSPACE_IN_MEMORY_COUNTER] { - purge_table(&conn, char::from(subspace)).await?; - } - - Ok(()) + Ok(()) + }) + .await; + bounded(conn, result, limit) } pub(crate) async fn delete_range(&self, from: impl Key, to: impl Key) -> trc::Result<()> { let conn = self.conn_pool.get().await.map_err(into_pool_error)?; - let table = char::from(from.subspace()); - let mut from = from.serialize(0); - let to = to.serialize(0); + let limit = self.timeouts.maintenance; + let result = tokio::time::timeout(limit, async { + let table = char::from(from.subspace()); + let mut from = from.serialize(0); + let to = to.serialize(0); - let delete = conn - .prepare_cached(&format!("DELETE FROM {table} WHERE k >= $1 AND k < $2")) - .await - .map_err(into_error)?; - - match conn.execute(&delete, &[&from, &to]).await { - Ok(_) => return Ok(()), - Err(err) if is_timeout_error(&err) => (), - Err(err) => return Err(into_error(err)), - } - - let mut chunk_size = DELETE_CHUNK_SIZE; - - loop { - let boundary = conn - .prepare_cached(&format!( - "SELECT k FROM {table} WHERE k >= $1 AND k < $2 ORDER BY k ASC LIMIT 1 OFFSET {chunk_size}" - )) + let delete = conn + .prepare_cached(&format!("DELETE FROM {table} WHERE k >= $1 AND k < $2")) .await .map_err(into_error)?; + match conn.execute(&delete, &[&from, &to]).await { + Ok(_) => return Ok(()), + Err(err) if is_timeout_error(&err) => (), + Err(err) => return Err(into_error(err)), + } + + let mut chunk_size = DELETE_CHUNK_SIZE; + loop { - let next = match conn.query_opt(&boundary, &[&from, &to]).await { - Ok(next) => match next { - Some(row) => Some(row.try_get::<_, Vec>(0).map_err(into_error)?), - None => None, - }, - Err(err) if is_timeout_error(&err) && chunk_size > MIN_DELETE_CHUNK_SIZE => { - chunk_size = (chunk_size / 2).max(MIN_DELETE_CHUNK_SIZE); - break; - } - Err(err) => return Err(into_error(err)), - }; - - match conn - .execute(&delete, &[&from, next.as_ref().unwrap_or(&to)]) + let boundary = conn + .prepare_cached(&format!( + "SELECT k FROM {table} WHERE k >= $1 AND k < $2 ORDER BY k ASC LIMIT 1 OFFSET {chunk_size}" + )) .await - { - Ok(_) => (), - Err(err) if is_timeout_error(&err) && chunk_size > MIN_DELETE_CHUNK_SIZE => { - chunk_size = (chunk_size / 2).max(MIN_DELETE_CHUNK_SIZE); - break; - } - Err(err) => return Err(into_error(err)), - } + .map_err(into_error)?; - match next { - Some(next) => from = next, - None => return Ok(()), + loop { + let next = match conn.query_opt(&boundary, &[&from, &to]).await { + Ok(next) => match next { + Some(row) => Some(row.try_get::<_, Vec>(0).map_err(into_error)?), + None => None, + }, + Err(err) if is_timeout_error(&err) && chunk_size > MIN_DELETE_CHUNK_SIZE => { + chunk_size = (chunk_size / 2).max(MIN_DELETE_CHUNK_SIZE); + break; + } + Err(err) => return Err(into_error(err)), + }; + + match conn + .execute(&delete, &[&from, next.as_ref().unwrap_or(&to)]) + .await + { + Ok(_) => (), + Err(err) if is_timeout_error(&err) && chunk_size > MIN_DELETE_CHUNK_SIZE => { + chunk_size = (chunk_size / 2).max(MIN_DELETE_CHUNK_SIZE); + break; + } + Err(err) => return Err(into_error(err)), + } + + match next { + Some(next) => from = next, + None => return Ok(()), + } } } - } + }) + .await; + bounded(conn, result, limit) } } diff --git a/crates/store/src/backend/query_timeout.rs b/crates/store/src/backend/query_timeout.rs new file mode 100644 index 0000000..1b1c576 --- /dev/null +++ b/crates/store/src/backend/query_timeout.rs @@ -0,0 +1,77 @@ +/* + * SPDX-FileCopyrightText: 2026 Coffey Labs + * + * SPDX-License-Identifier: AGPL-3.0-only + */ + +//! Client-side limits on SQL queries. +//! +//! The pool timeouts bound getting a connection, not using one. A database +//! that stops answering while the TCP connection stays up (a paused +//! container, a hung server whose kernel still acknowledges keepalives) +//! left a query on a checked-out connection waiting for as long as it took. +//! A server-side statement_timeout can't help there: the server that would +//! enforce it is the one not answering. So each operation on a PostgreSQL +//! or MySQL connection runs under a time limit here, and a connection whose +//! operation ran out is closed rather than put back in the pool, since its +//! protocol state is unknown. +//! +//! Two limits: +//! - `query`, two minutes, for request-path work: reads, writes, blob +//! transfers, search queries and document indexing. Those take +//! milliseconds; two minutes leaves room for a large blob over a slow +//! link and still ends a hang. +//! - `maintenance`, thirty minutes, for work that legitimately runs long in +//! one statement: range deletes (account removal, purges), unindexing, +//! and creating tables and indexes at startup. +//! +//! Iterating over a range (exports, reindexing, maintenance scans) can run +//! for hours, so there the `query` limit applies to each wait for the next +//! row instead of the whole scan. + +use std::time::Duration; + +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub struct QueryTimeouts { + pub query: Duration, + pub maintenance: Duration, +} + +impl QueryTimeouts { + pub const QUERY: Duration = Duration::from_secs(120); + pub const MAINTENANCE: Duration = Duration::from_secs(30 * 60); +} + +impl Default for QueryTimeouts { + fn default() -> Self { + Self { + query: Self::QUERY, + maintenance: Self::MAINTENANCE, + } + } +} + +#[cfg(feature = "test_mode")] +impl crate::Store { + /// Sets the query limits of a SQL store that was just built (tests only: + /// the limits aren't configurable). + pub fn with_query_timeouts(self, timeouts: QueryTimeouts) -> Self { + match self { + #[cfg(feature = "postgres")] + crate::Store::PostgreSQL(mut store) => { + std::sync::Arc::get_mut(&mut store) + .expect("store already shared") + .timeouts = timeouts; + crate::Store::PostgreSQL(store) + } + #[cfg(feature = "mysql")] + crate::Store::MySQL(mut store) => { + std::sync::Arc::get_mut(&mut store) + .expect("store already shared") + .timeouts = timeouts; + crate::Store::MySQL(store) + } + store => store, + } + } +} diff --git a/tests/src/store/pool_timeout.rs b/tests/src/store/pool_timeout.rs index b9b829a..b6511c6 100644 --- a/tests/src/store/pool_timeout.rs +++ b/tests/src/store/pool_timeout.rs @@ -9,11 +9,31 @@ //! 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. +//! +//! inbuxa: the same for a database that stops answering while connections +//! are already open (a paused container): a query on a checked-out +//! connection ends within the query limit, the store works again once the +//! database is back, and /healthz/ready says 503 in between while +//! /healthz/live stays 200. These need the local test databases; a proxy +//! that can stop forwarding plays the pause. use registry::schema::structs::DataStore; -use std::time::{Duration, Instant}; -use store::{Store, ValueKey, write::ValueClass}; -use tokio::net::TcpListener; +use std::{ + sync::{ + Arc, + atomic::{AtomicBool, Ordering}, + }, + time::{Duration, Instant}, +}; +use store::{ + IterateParams, Store, ValueKey, + backend::query_timeout::QueryTimeouts, + write::{BatchBuilder, ValueClass}, +}; +use tokio::{ + io::{AsyncReadExt, AsyncWriteExt}, + net::{TcpListener, TcpStream}, +}; /// Accepts connections on a local port and never sends a byte. async fn silent_server() -> u16 { @@ -94,3 +114,263 @@ pub async fn mysql_pool_timeout() { ) .await; } + +/// A TCP proxy to a local port that can stop forwarding, in both +/// directions, while keeping every connection open: a paused server whose +/// kernel still keeps the connections up. +struct PausableProxy { + port: u16, + paused: Arc, +} + +impl PausableProxy { + async fn start(upstream: u16) -> Self { + let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); + let port = listener.local_addr().unwrap().port(); + let paused = Arc::new(AtomicBool::new(false)); + let paused_ = paused.clone(); + tokio::spawn(async move { + while let Ok((client, _)) = listener.accept().await { + let Ok(server) = TcpStream::connect(("127.0.0.1", upstream)).await else { + continue; + }; + let (client_rx, client_tx) = client.into_split(); + let (server_rx, server_tx) = server.into_split(); + tokio::spawn(forward(client_rx, server_tx, paused_.clone())); + tokio::spawn(forward(server_rx, client_tx, paused_.clone())); + } + }); + PausableProxy { port, paused } + } + + fn pause(&self, paused: bool) { + self.paused.store(paused, Ordering::SeqCst); + } +} + +async fn forward( + mut from: tokio::net::tcp::OwnedReadHalf, + mut to: tokio::net::tcp::OwnedWriteHalf, + paused: Arc, +) { + let mut buf = vec![0u8; 16384]; + loop { + while paused.load(Ordering::SeqCst) { + tokio::time::sleep(Duration::from_millis(20)).await; + } + let n = match from.read(&mut buf).await { + Ok(0) | Err(_) => return, + Ok(n) => n, + }; + // Hold what arrived while paused until the pause ends + while paused.load(Ordering::SeqCst) { + tokio::time::sleep(Duration::from_millis(20)).await; + } + if to.write_all(&buf[..n]).await.is_err() { + return; + } + } +} + +const TEST_LIMITS: QueryTimeouts = QueryTimeouts { + query: Duration::from_secs(2), + maintenance: Duration::from_secs(3), +}; + +/// Opens `connections` pooled connections at once, so the operations that +/// follow find one idle and check it out. +async fn warm(store: &Store, connections: usize) { + let reads = (0..connections).map(|_| async { + store + .get_value::(ValueKey::from(ValueClass::Property(0))) + .await + .unwrap(); + }); + futures::future::join_all(reads).await; +} + +/// With the database paused, reads, scans and writes on connections the +/// pool already holds end in an error within the query limit; once it is +/// back, the store works again. +async fn assert_queries_time_out(store: Store, proxy: &PausableProxy) { + store.create_tables().await.unwrap(); + warm(&store, 4).await; + // mysql_async resets a connection on its way back to the pool; let + // those finish, or the connections are stuck in the reset when the + // pause starts and the pool's own wait timeout answers instead + tokio::time::sleep(Duration::from_secs(1)).await; + proxy.pause(true); + + let key = || ValueKey::from(ValueClass::Property(0)); + let limit = TEST_LIMITS.query; + for (what, op) in [("read", 0), ("scan", 1), ("write", 2)] { + let started = Instant::now(); + let result = tokio::time::timeout(Duration::from_secs(20), async { + match op { + 0 => store.get_value::(key()).await.map(|_| ()), + 1 => { + store + .iterate( + IterateParams::new( + ValueKey::from(ValueClass::Property(0)), + ValueKey::from(ValueClass::Property(u8::MAX)), + ), + |_, _| Ok(true), + ) + .await + } + _ => { + let mut batch = BatchBuilder::new(); + batch + .with_account_id(u32::MAX - 7) + .with_collection(types::collection::Collection::Email) + .with_document(0) + .set(ValueClass::Property(0), 1u64.to_be_bytes().to_vec()); + store.write(batch.build_all()).await.map(|_| ()) + } + } + }) + .await; + let elapsed = started.elapsed(); + match result { + Ok(Err(err)) => { + let err = format!("{err:?}"); + println!("Paused database, {what}: {err} after {elapsed:?}"); + assert!(err.contains("Query timed out"), "{what}: {err}"); + assert!( + elapsed >= limit && elapsed < limit * 3, + "{what} ended after {elapsed:?}" + ); + } + Ok(Ok(())) => panic!("{what} succeeded against a paused database"), + Err(_) => panic!("{what} still waiting after {elapsed:?}"), + } + } + + proxy.pause(false); + tokio::time::timeout(Duration::from_secs(20), store.get_value::(key())) + .await + .expect("still waiting after the database came back") + .expect("the store didn't recover"); +} + +#[cfg(feature = "postgres")] +#[tokio::test(flavor = "multi_thread")] +pub async fn postgres_query_timeout() { + println!("Running PostgreSQL query timeout test..."); + let DataStore::PostgreSql(mut config) = + crate::utils::storage::build_data_store("PostgreSql", "").await + else { + unreachable!() + }; + let proxy = PausableProxy::start(config.port as u16).await; + config.host = "127.0.0.1".into(); + config.port = proxy.port as u64; + // New connections through the paused proxy give up as quickly + config.timeout = Some(TEST_LIMITS.query.into()); + let store = Store::build(DataStore::PostgreSql(config)) + .await + .unwrap() + .with_query_timeouts(TEST_LIMITS); + assert_queries_time_out(store, &proxy).await; +} + +#[cfg(feature = "mysql")] +#[tokio::test(flavor = "multi_thread")] +pub async fn mysql_query_timeout() { + println!("Running MySQL query timeout test..."); + let DataStore::MySql(mut config) = crate::utils::storage::build_data_store("MySql", "").await + else { + unreachable!() + }; + let proxy = PausableProxy::start(config.port as u16).await; + config.host = "127.0.0.1".into(); + config.port = proxy.port as u64; + let store = Store::build(DataStore::MySql(config)) + .await + .unwrap() + .with_query_timeouts(TEST_LIMITS); + assert_queries_time_out(store, &proxy).await; +} + +/// /healthz/ready follows the data store; /healthz/live doesn't. +#[cfg(feature = "postgres")] +#[tokio::test(flavor = "multi_thread")] +pub async fn postgres_readiness() { + use crate::utils::server::TestServerBuilder; + use registry::schema::enums::NetworkListenerProtocol; + + const HTTP_PORT: u16 = 11_320; + if std::env::var("STORE").as_deref() != Ok("PostgreSql") { + println!("Skipping the readiness test: it runs with STORE=PostgreSql."); + return; + } + println!("Running readiness test..."); + + let test = TestServerBuilder::new("postgres_readiness") + .await + .with_listener(NetworkListenerProtocol::Http, "http", HTTP_PORT, true) + .await + .build() + .await; + + // Point the running node's data store at the database through the proxy + let DataStore::PostgreSql(mut config) = + crate::utils::storage::build_data_store("PostgreSql", "").await + else { + unreachable!() + }; + let proxy = PausableProxy::start(config.port as u16).await; + config.host = "127.0.0.1".into(); + config.port = proxy.port as u64; + config.timeout = Some(TEST_LIMITS.query.into()); + let store = Store::build(DataStore::PostgreSql(config)) + .await + .unwrap() + .with_query_timeouts(TEST_LIMITS); + let inner = &test.server.inner; + let mut core = inner.shared_core.load_full().as_ref().clone(); + core.storage.data = store; + inner.shared_core.store(Arc::new(core)); + + let health = |path: &'static str| async move { + reqwest::Client::builder() + .danger_accept_invalid_certs(true) + .timeout(Duration::from_secs(10)) + .build() + .unwrap() + .get(format!("https://127.0.0.1:{HTTP_PORT}/healthz/{path}")) + .send() + .await + .unwrap() + .status() + .as_u16() + }; + let wait_for = |path: &'static str, status: u16| async move { + let started = Instant::now(); + loop { + let got = health(path).await; + if got == status { + println!("/healthz/{path}: {got} after {:?}", started.elapsed()); + return; + } + assert!( + started.elapsed() < Duration::from_secs(20), + "/healthz/{path} still {got}, expected {status}" + ); + tokio::time::sleep(Duration::from_millis(250)).await; + } + }; + + wait_for("ready", 200).await; + proxy.pause(true); + wait_for("ready", 503).await; + assert_eq!(health("live").await, 200); + proxy.pause(false); + wait_for("ready", 200).await; + assert_eq!(health("live").await, 200); + + if test.is_reset() { + test.temp_dir.delete(); + } +}