5 Commits
Author SHA1 Message Date
jcoffey-dev e00978c0b4 Cluster role changes apply to delivery and tasks without a restart
ci / fork-checks (pull_request) Successful in 17s
ci / build (pull_request) Successful in 7m11s
In cluster rehearsal 3, turning outboundMta off on node1's role was
reported applied (x:settingsReload applied: true), yet node1 kept
delivering mail, a report message included, until it was restarted.
The queue and report managers were started at boot only when the
node's role included outboundMta (crates/smtp/src/lib.rs), and the task
manager only when the role had some task type (spawn_task_manager).
After that nothing looked at the role again: a queue manager that was
running kept claiming and delivering, and one that wasn't never
started.

They now start on every node (outside recovery mode) and follow the
role live:

- Queue manager: before each scan it reads the role from the running
  settings. Without outboundMta it claims nothing new; deliveries
  already running finish and report back as usual, which releases
  their locks. When the role comes back (a reload wakes the manager
  with ReloadSettings, and it looks again every 30 s regardless) it
  logs queue.started and scans the whole queue at once.
- Report scheduler: DMARC and TLS report events are handled only while
  the role has outboundMta, as at boot; events arriving without it are
  dropped, as they were on a node started without the role.
- Task manager: task_enabled already read the current role on every
  scan. It now also runs on nodes whose role has no task type (the
  scan returns at once until one is added), a job claimed before a
  role change is handed back at once rather than run or held until
  its lease lapses, and a settings reload wakes the manager so a role
  that gained task types starts claiming them straight away.

Starting the queue manager on every node also drains the queue channel
on nodes without outboundMta. Upstream left that channel unread, so
each message queued there parked a refresh in it, and by the code,
queueing would block once 1024 had piled up (not reproduced here).

A role object edit reaches the nodes that name that role in
INBUXA_ROLE. Moving a node to another role still means changing its
environment, and so a restart. Listener changes in a role still need a
restart too (listeners bind at boot); this change is about tasks and
delivery.

cluster::live_roles::live_role_tests (new; PostgreSQL, two nodes over
one store):
1. A node started with outboundMta delivers and runs a TLS report
   task; after its role loses outboundMta and the settings reload, a
   new message isn't attempted and a new report task stays pending;
   with the role back, both are taken up.
2. A node started with no task type at all gains outboundMta: a
   waiting message is attempted and a report task runs.
On main the test fails at step 1 ("delivery attempted without
outboundMta"); with step 1 bypassed, step 2 fails (nothing picked the
message up in 20 s).
2026-09-24 17:45:30 -07:00
jcoffey-dev ad58c35f39 Merge pull request 'SQL queries time out; readiness follows the data store' (#45) from fix/query-timeouts into main
ci / build (push) Canceled after 11m21s
ci / fork-checks (push) Successful in 55s
2026-09-25 00:45:09 +00:00
jcoffey-dev 181ab1c140 Merge pull request 'PostgreSQL search GIN indexes without a pending list' (#43) from fix/pg-gin-fastupdate into main
ci / fork-checks (push) Canceled after 0s
ci / build (push) Canceled after 0s
2026-09-25 00:45:08 +00:00
jcoffey-dev 08f29926d4 SQL queries time out; readiness follows the data store
ci / fork-checks (pull_request) Successful in 17s
ci / build (pull_request) Successful in 7m13s
Cluster rehearsal 3: with PostgreSQL paused (docker pause, so its
kernel still answered TCP keepalives), requests on connections already
checked out hung until it came back, and /healthz/ready stayed 200
through the outage. #41 bounded getting a connection, not using one.

Client-side query limits (store::backend::query_timeout). Every
operation on a PostgreSQL or MySQL connection now runs under a time
limit. A server-side statement_timeout (or MySQL's MAX_EXECUTION_TIME,
which covers SELECTs only) can't do this: the server that would enforce
it is the one not answering. When an operation runs out, its connection
is closed instead of pooled, since a query may still be in flight on it
or a transaction open: deadpool's Object::take on PostgreSQL;
Conn::disconnect on MySQL, which marks the connection closed before it
sends anything, so the pool discards it even when the server never
answers.

- query, 2 minutes: reads, writes (the whole transaction with its
  retries), blobs, SQL lookups, search queries and indexing. These take
  milliseconds; two minutes leaves room for a large blob over a slow
  link and still ends a hang.
- maintenance, 30 minutes: range deletes (account removal, purges),
  unindexing, purge_store, and creating tables and indexes at startup,
  which can legitimately run long in one statement. Their existing
  chunked fallback for server-side statement timeouts is unchanged.
- iterate (exports, reindexing, maintenance scans) can run for hours,
  so the query limit bounds each wait for the database (preparing, the
  query starting, the next row) rather than the whole scan.

The limits are fixed, like the pool timeouts; the DataStore schema has
no field for them. Tests set them with Store::with_query_timeouts
(test_mode only).

Readiness. /healthz/ready answered 200 whenever a data store was
configured. It now reads one key from the data store with a 2 s limit
and reuses the answer for 2 s, so probes can't load the database;
while one probe runs, others get the last answer. The first failed
probe of an outage is logged. /healthz/live stays 200: restarting a
node doesn't bring its database back, and an orchestrator restarting on
failed liveness would restart every node at once. The container
HEALTHCHECK already uses /healthz/live.

Tests, store::pool_timeout (a proxy that stops forwarding while
keeping connections open plays the paused database):
- postgres_query_timeout, mysql_query_timeout (new): with four pooled
  connections open, a read, a scan and a write each fail with "Query
  timed out" 2.0 s after the pause (2 s test limit); once the proxy
  forwards again the store answers. With the limits set to an hour
  (upstream's behavior), the read was still waiting at the test's 20 s
  limit.
- postgres_readiness (new, STORE=PostgreSql): a node's data store
  goes through the proxy; /healthz/ready is 200, 503 about 4 s after
  the pause while /healthz/live stays 200, and 200 again about 2 s
  after it ends.
- postgres_pool_timeout, mysql_pool_timeout: pass as before.
store::store_tests (PostgreSql, MySql, including the MariaDB statement
timeout step) and store::task_locks (PostgreSql) pass;
store::search_tests (PostgreSql) fails at the same ordering assertion
(query.rs:684) as on main.
2026-09-24 16:45:54 -07:00
jcoffey-dev fde43774b4 PostgreSQL search GIN indexes without a pending list
ci / fork-checks (pull_request) Successful in 30s
ci / build (pull_request) Successful in 7m26s
A three-node rehearsal on PostgreSQL saw searches take about 185 ms
with 80 to 260 pages in the full-text indexes' pending lists, 2 to 6 ms
right after gin_clean_pending_list() or VACUUM, then creep back up as
mail came in. The search tables' GIN indexes were created with the
default fastupdate=on: new entries wait in an unindexed pending list
that every search scans in full until VACUUM (or 4 MB of backlog)
merges it, and autovacuum only visits an insert-only table after
thousands of inserts.

The search GIN indexes are now created WITH (fastupdate = off), so an
insert pays its index update at once. The schema step runs at every
startup (create_search_tables, via SearchStore::create_indexes), so
indexes made before this change are switched there: when an index's
reloptions don't already turn fastupdate off, ALTER INDEX ... SET
(fastupdate = off) and one gin_clean_pending_list() merge its backlog.
The ALTER takes a SHARE UPDATE EXCLUSIVE lock, which blocks neither
reads nor writes; after the first startup the step is one catalog read
per index. A failure is logged and startup goes on (search still
works, only slower).

Per-table autovacuum settings for the search tables are left alone.
The pending list was the only reason the insert threshold mattered for
search; dead tuples and freezing are served by the defaults, and table
settings would override whatever tuning the DBA has done.

MySQL is unaffected: InnoDB FULLTEXT keeps new entries in an in-memory
cache that queries read directly, with no setting like fastupdate.

store::search_gin::postgres_gin_fastupdate (new, PostgreSQL) builds
the search schema in a schema of its own and checks pg_class.reloptions:
fastupdate=off on every GIN index of a fresh schema; then, with the
option reset to the default and 500 rows pending, one startup turns it
off everywhere and leaves no pending tuples (pgstatginindex); a second
startup changes nothing. On main it fails at the first check.
2026-09-24 15:51:06 -07:00
31 changed files with 2154 additions and 810 deletions
+6
View File
@@ -155,6 +155,12 @@ impl Server {
.await .await
.ok(); .ok();
// inbuxa: the task manager reads the node's role on
// every scan; scan now, so a role that gained task
// types starts claiming them without waiting out the
// refresh interval
self.inner.ipc.task_tx.notify_one();
self.record_build_errors(&bootstrap.errors); self.record_build_errors(&bootstrap.errors);
return Ok(ReloadResult { return Ok(ReloadResult {
+2
View File
@@ -94,6 +94,7 @@ impl Data {
span_id_gen: id_generator, span_id_gen: id_generator,
queue_status: true.into(), queue_status: true.into(),
settings_reload: Default::default(), settings_reload: Default::default(),
store_health: Default::default(),
applications, applications,
logos: Default::default(), logos: Default::default(),
smtp_connectors: TlsConnectors::try_new().failed("Failed to build TLS connectors"), smtp_connectors: TlsConnectors::try_new().failed("Failed to build TLS connectors"),
@@ -237,6 +238,7 @@ impl Default for Data {
registry_id_gen: Default::default(), registry_id_gen: Default::default(),
queue_status: true.into(), queue_status: true.into(),
settings_reload: Default::default(), settings_reload: Default::default(),
store_health: Default::default(),
applications: WebApplications::new(), applications: WebApplications::new(),
logos: Default::default(), logos: Default::default(),
smtp_connectors: TlsConnectors::try_new().unwrap(), smtp_connectors: TlsConnectors::try_new().unwrap(),
+2
View File
@@ -163,6 +163,8 @@ pub struct Data {
pub queue_status: AtomicBool, pub queue_status: AtomicBool,
// inbuxa: coalesces the settings reloads registry writes trigger // inbuxa: coalesces the settings reloads registry writes trigger
pub settings_reload: cache::reload::SettingsReloadGate, pub settings_reload: cache::reload::SettingsReloadGate,
// inbuxa: the readiness probe's cached answer
pub store_health: storage::ready::StoreHealth,
pub applications: WebApplications, pub applications: WebApplications,
pub logos: Mutex<AHashMap<Box<str>, LogoCache>>, pub logos: Mutex<AHashMap<Box<str>, LogoCache>>,
+1
View File
@@ -26,6 +26,7 @@ pub mod document;
pub mod encryption; pub mod encryption;
pub mod index; pub mod index;
pub mod quota; pub mod quota;
pub mod ready; // inbuxa: readiness follows the data store
pub mod state; pub mod state;
pub mod transaction; pub mod transaction;
+83
View File
@@ -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<Option<(Instant, bool)>>,
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::<u64>(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
}
}
+3 -1
View File
@@ -553,8 +553,10 @@ impl ParseHttp for Server {
return Ok(JsonProblemResponse(StatusCode::OK).into_http_response()); return Ok(JsonProblemResponse(StatusCode::OK).into_http_response());
} }
"ready" => { "ready" => {
// inbuxa: ready only while the data store answers
// (a cached, time-limited read); liveness stays 200
return Ok(JsonProblemResponse({ return Ok(JsonProblemResponse({
if !self.core.storage.data.is_none() { if self.is_data_store_ready().await {
StatusCode::OK StatusCode::OK
} else { } else {
StatusCode::SERVICE_UNAVAILABLE StatusCode::SERVICE_UNAVAILABLE
+43 -21
View File
@@ -55,23 +55,12 @@ const PERPETUAL_RETRY_MIN_DELAY: u64 = 3600;
const PERPETUAL_RETRY_MAX_DELAY: u64 = 21600; const PERPETUAL_RETRY_MAX_DELAY: u64 = 21600;
pub fn spawn_task_manager(inner: Arc<Inner>) { pub fn spawn_task_manager(inner: Arc<Inner>) {
let is_clustered = { // inbuxa: upstream didn't start the task manager on a node whose role
let server = inner.build_server(); // had no task types at boot, so adding one later did nothing until a
let roles = &server.core.network.roles; // restart. It now always runs and reads the role on every scan and
// before every job (task_enabled), so a role change applies at the next
// inbuxa: outbound_mta too, which now governs report tasks // settings reload.
if !roles.account_maintenance let is_clustered = inner.build_server().core.storage.coordinator.is_enabled();
&& !roles.store_maintenance
&& !roles.search_indexing
&& !roles.spam_training
&& !roles.outbound_mta
&& !roles.task_manager
{
return;
}
server.core.storage.coordinator.is_enabled()
};
trc::event!(TaskManager(TaskManagerEvent::ManagerStarted)); trc::event!(TaskManager(TaskManagerEvent::ManagerStarted));
@@ -151,20 +140,23 @@ pub fn spawn_task_manager(inner: Arc<Inner>) {
let server = inner.build_server(); let server = inner.build_server();
let batch_size = server.core.email.index_batch_size; let batch_size = server.core.email.index_batch_size;
let mut batch = Vec::with_capacity(batch_size); let mut batch = Vec::with_capacity(batch_size);
if let Some(task) = fetch_task(&server, job).await { if let Some(task) = fetch_enabled_task(&server, job).await {
batch.push(task); batch.push(task);
} }
while batch.len() < batch_size { while batch.len() < batch_size {
match rx.try_recv() { match rx.try_recv() {
Ok(job) => { Ok(job) => {
if let Some(task) = fetch_task(&server, job).await { if let Some(task) = fetch_enabled_task(&server, job).await {
batch.push(task); batch.push(task);
} }
} }
Err(_) => break, Err(_) => break,
} }
} }
if batch.is_empty() {
continue;
}
// Dispatch. inbuxa: on a task of its own, so a panic // Dispatch. inbuxa: on a task of its own, so a panic
// releases the batch's locks and leaves this worker // releases the batch's locks and leaves this worker
@@ -205,7 +197,8 @@ pub fn spawn_task_manager(inner: Arc<Inner>) {
let server = inner.build_server(); let server = inner.build_server();
let mut refresh_queue = false; let mut refresh_queue = false;
if let Some(TaskDetails { task, info }) = fetch_task(&server, job).await { if let Some(TaskDetails { task, info }) = fetch_enabled_task(&server, job).await
{
// inbuxa: on a task of its own, as above // inbuxa: on a task of its own, as above
let run = { let run = {
let server = server.clone(); let server = server.clone();
@@ -274,6 +267,17 @@ impl TaskQueueManager for Server {
if task_locks.is_stopping() { if task_locks.is_stopping() {
return Duration::from_secs(QUEUE_REFRESH_INTERVAL); return Duration::from_secs(QUEUE_REFRESH_INTERVAL);
} }
// inbuxa: with no task type enabled by this node's role there is
// nothing to claim; a settings reload wakes the manager when that
// changes
let roles = &self.core.network.roles;
if !(0..TaskType::COUNT as u16)
.filter_map(TaskType::from_id)
.any(|task_type| task_enabled(roles, task_type))
{
ipc.locked.clear();
return Duration::from_secs(QUEUE_REFRESH_INTERVAL);
}
let lock_expiry = task_locks.expiry(); let lock_expiry = task_locks.expiry();
let now_timestamp = now(); let now_timestamp = now();
let from_key = ValueKey::<ValueClass> { let from_key = ValueKey::<ValueClass> {
@@ -296,7 +300,6 @@ impl TaskQueueManager for Server {
let mut tasks = Vec::new(); let mut tasks = Vec::new();
let now = Instant::now(); let now = Instant::now();
let mut next_event = None; let mut next_event = None;
let roles = &self.core.network.roles;
ipc.revision += 1; ipc.revision += 1;
let _ = self let _ = self
.store() .store()
@@ -544,6 +547,25 @@ async fn run_task(
} }
} }
/// inbuxa: reads a claimed task when this node's role still allows its type.
/// The role may have changed since the task was claimed (a settings reload in
/// between); the claim is then handed back at once for a node that may run
/// it, rather than held until the lease runs out.
async fn fetch_enabled_task(server: &Server, job: TaskJob) -> Option<TaskDetails> {
if task_enabled(&server.core.network.roles, job.typ) {
fetch_task(server, job).await
} else {
trc::event!(
TaskManager(TaskManagerEvent::TaskIgnored),
Id = job.id,
Details = job.typ.as_str(),
Reason = "Task type was disabled by cluster roles after it was claimed.",
);
server.remove_index_lock(job.id).await;
None
}
}
/// Reads a claimed task. When it is gone or can't be read, the claim is /// Reads a claimed task. When it is gone or can't be read, the claim is
/// released: inbuxa: holding it would block the task, everywhere, until /// released: inbuxa: holding it would block the task, everywhere, until
/// the lock expired. /// the lock expired.
+8 -1
View File
@@ -44,7 +44,14 @@ impl StartQueueManager for BootManager {
impl SpawnQueueManager for IpcReceivers { impl SpawnQueueManager for IpcReceivers {
fn spawn_queue_manager(&mut self, inner: Arc<Inner>) { fn spawn_queue_manager(&mut self, inner: Arc<Inner>) {
let core = inner.shared_core.load(); let core = inner.shared_core.load();
if !core.storage.registry.is_recovery_mode() && core.network.roles.outbound_mta { // inbuxa: upstream started these only when the node's role included
// outboundMta at boot, so turning the role on later did nothing and
// turning it off left them delivering until a restart. They now run
// on every node and follow the role live (see Queue::start and the
// report scheduler). This also drains the queue channel on nodes
// without the role, where every queued message's refresh used to sit
// in a channel nobody read until it filled and queueing blocked.
if !core.storage.registry.is_recovery_mode() {
// Spawn queue manager // Spawn queue manager
self.queue_rx.take().unwrap().spawn(inner.clone()); self.queue_rx.take().unwrap().spawn(inner.clone());
+30
View File
@@ -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::{Message, QueueId, Status, spool::SmtpSpool}; use super::{Message, QueueId, Status, spool::SmtpSpool};
@@ -39,6 +41,9 @@ pub struct Queue {
pub urgent_refresh: bool, pub urgent_refresh: bool,
pub last_scan: Instant, pub last_scan: Instant,
pub last_full_scan: Instant, pub last_full_scan: Instant,
/// inbuxa: whether this node's role included outboundMta when last
/// checked (None before the first check)
pub role_enabled: Option<bool>,
} }
#[derive(Debug)] #[derive(Debug)]
@@ -67,6 +72,9 @@ impl SpawnQueue for mpsc::Receiver<QueueEvent> {
const BACK_PRESSURE_WARN_INTERVAL: Duration = Duration::from_secs(60); const BACK_PRESSURE_WARN_INTERVAL: Duration = Duration::from_secs(60);
const MIN_SCAN_INTERVAL: Duration = Duration::from_millis(100); const MIN_SCAN_INTERVAL: Duration = Duration::from_millis(100);
const FULL_SCAN_INTERVAL: Duration = Duration::from_secs(QUEUE_REFRESH / 2); const FULL_SCAN_INTERVAL: Duration = Duration::from_secs(QUEUE_REFRESH / 2);
/// inbuxa: how often a node without the outbound MTA role looks at its role
/// again when nothing else wakes it (a settings reload does)
const ROLE_RECHECK_INTERVAL: Duration = Duration::from_secs(30);
impl Queue { impl Queue {
pub fn new(core: Arc<Inner>, rx: mpsc::Receiver<QueueEvent>) -> Self { pub fn new(core: Arc<Inner>, rx: mpsc::Receiver<QueueEvent>) -> Self {
@@ -87,6 +95,7 @@ impl Queue {
urgent_refresh: false, urgent_refresh: false,
last_scan: now.checked_sub(MIN_SCAN_INTERVAL).unwrap_or(now), last_scan: now.checked_sub(MIN_SCAN_INTERVAL).unwrap_or(now),
last_full_scan: now, last_full_scan: now,
role_enabled: None,
} }
} }
@@ -123,6 +132,27 @@ impl Queue {
continue; continue;
} }
// inbuxa: follow the node's role live. Without outboundMta the
// queue claims nothing new; deliveries already running finish
// and report back as usual, releasing their locks. When the role
// comes back, the whole queue is scanned at once.
let role_enabled = self.core.shared_core.load().network.roles.outbound_mta;
if self.role_enabled.replace(role_enabled) == Some(false) && role_enabled {
trc::event!(
Queue(trc::QueueEvent::Started),
Details = "This node's cluster role now includes outboundMta",
);
self.scan_from = 0;
self.pending_refresh = true;
self.urgent_refresh = true;
}
if !role_enabled {
self.pending_refresh = false;
self.urgent_refresh = false;
self.next_refresh = Instant::now() + ROLE_RECHECK_INTERVAL;
continue;
}
self.pending_refresh |= refresh_queue; self.pending_refresh |= refresh_queue;
if !self.pending_refresh && self.next_refresh > Instant::now() { if !self.pending_refresh && self.next_refresh > Instant::now() {
continue; continue;
+10
View File
@@ -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::{dmarc::DmarcReporting, tls::TlsReporting}; use super::{dmarc::DmarcReporting, tls::TlsReporting};
@@ -18,6 +20,14 @@ impl SpawnReport for mpsc::Receiver<ReportingEvent> {
tokio::spawn(async move { tokio::spawn(async move {
while let Some(event) = self.recv().await { while let Some(event) = self.recv().await {
let server = inner.build_server(); let server = inner.build_server();
// inbuxa: reports are the outbound MTA's business, as at
// boot, but the role is read per event so a change applies
// without a restart. Events that arrive while the role is
// off are dropped, as they were on a node started without it
if !matches!(event, ReportingEvent::Stop) && !server.core.network.roles.outbound_mta
{
continue;
}
match event { match event {
ReportingEvent::Dmarc(event) => server.schedule_dmarc(event).await, ReportingEvent::Dmarc(event) => server.schedule_dmarc(event).await,
ReportingEvent::Tls(event) => server.schedule_tls(event).await, ReportingEvent::Tls(event) => server.schedule_tls(event).await,
+3
View File
@@ -30,6 +30,9 @@ pub mod s3;
pub mod sqlite; pub mod sqlite;
// inbuxa: scale-out storage (sharded stores) // inbuxa: scale-out storage (sharded stores)
pub mod scaleout; 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; pub const MAX_TOKEN_LENGTH: usize = (u8::MAX >> 1) as usize;
+16 -1
View File
@@ -10,7 +10,7 @@ use std::ops::Range;
use mysql_async::prelude::Queryable; use mysql_async::prelude::Queryable;
use super::{MysqlStore, into_error}; use super::{MysqlStore, bounded, into_error};
impl MysqlStore { impl MysqlStore {
pub(crate) async fn get_blob( pub(crate) async fn get_blob(
@@ -19,6 +19,8 @@ impl MysqlStore {
range: Range<usize>, range: Range<usize>,
) -> trc::Result<Option<Vec<u8>>> { ) -> trc::Result<Option<Vec<u8>>> {
let mut conn = self.conn().await?; let mut conn = self.conn().await?;
let limit = self.timeouts.query;
let result = tokio::time::timeout(limit, async {
let s = conn let s = conn
.prep("SELECT v FROM t WHERE k = ?") .prep("SELECT v FROM t WHERE k = ?")
.await .await
@@ -38,10 +40,15 @@ impl MysqlStore {
} }
}) })
.map_err(into_error) .map_err(into_error)
})
.await;
bounded(conn, result, limit)
} }
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().await?; let mut conn = self.conn().await?;
let limit = self.timeouts.query;
let result = tokio::time::timeout(limit, async {
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
@@ -50,10 +57,15 @@ impl MysqlStore {
.await .await
.map_err(into_error) .map_err(into_error)
.map(|_| ()) .map(|_| ())
})
.await;
bounded(conn, result, limit)
} }
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().await?; let mut conn = self.conn().await?;
let limit = self.timeouts.query;
let result = tokio::time::timeout(limit, async {
let s = conn let s = conn
.prep("DELETE FROM t WHERE k = ?") .prep("DELETE FROM t WHERE k = ?")
.await .await
@@ -62,5 +74,8 @@ impl MysqlStore {
.await .await
.map_err(into_error) .map_err(into_error)
.map(|hits| hits.affected_rows() > 0) .map(|hits| hits.affected_rows() > 0)
})
.await;
bounded(conn, result, limit)
} }
} }
+6 -1
View File
@@ -10,7 +10,7 @@ use mysql_async::{Params, Row, prelude::Queryable};
use crate::{IntoRows, QueryResult, QueryType, Value}; use crate::{IntoRows, QueryResult, QueryType, Value};
use super::{MysqlStore, into_error}; use super::{MysqlStore, bounded, into_error};
impl MysqlStore { impl MysqlStore {
pub(crate) async fn sql_query<T: QueryResult>( pub(crate) async fn sql_query<T: QueryResult>(
@@ -19,6 +19,8 @@ impl MysqlStore {
params: &[Value<'_>], params: &[Value<'_>],
) -> trc::Result<T> { ) -> trc::Result<T> {
let mut conn = self.conn().await?; let mut conn = self.conn().await?;
let limit = self.timeouts.query;
let result = tokio::time::timeout(limit, async {
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());
@@ -40,6 +42,9 @@ impl MysqlStore {
.await .await
.map_or_else(|e| Err(into_error(e)), |r| Ok(T::from_query_all(r))), .map_or_else(|e| Err(into_error(e)), |r| Ok(T::from_query_all(r))),
} }
})
.await;
bounded(conn, result, limit)
} }
} }
+13 -3
View File
@@ -6,7 +6,7 @@
* Modified by Coffey Labs in 2026 for INBUXA. * Modified by Coffey Labs in 2026 for INBUXA.
*/ */
use super::{MysqlStore, into_error}; use super::{MysqlStore, bounded, into_error};
use crate::{ use crate::{
backend::mysql::MysqlSearchField, backend::mysql::MysqlSearchField,
search::{ search::{
@@ -72,6 +72,7 @@ impl MysqlStore {
.db_name(Some(replica.database.clone())) .db_name(Some(replica.database.clone()))
.tcp_port(replica.port as u16), .tcp_port(replica.port as u16),
), ),
timeouts: Default::default(),
})), })),
replica.host, replica.host,
replica.port as u16, replica.port as u16,
@@ -81,6 +82,7 @@ impl MysqlStore {
let primary = Store::MySQL(Arc::new(MysqlStore { let primary = Store::MySQL(Arc::new(MysqlStore {
conn_pool: Pool::new(opts), conn_pool: Pool::new(opts),
timeouts: Default::default(),
})); }));
// ST-1: no replicas, no change // ST-1: no replicas, no change
@@ -99,7 +101,8 @@ 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().await?; let mut conn = self.conn().await?;
let limit = self.timeouts.maintenance;
let result = tokio::time::timeout(limit, async {
for table in [ for table in [
SUBSPACE_ACL, SUBSPACE_ACL,
SUBSPACE_TASK_QUEUE, SUBSPACE_TASK_QUEUE,
@@ -169,11 +172,15 @@ impl MysqlStore {
} }
Ok(()) Ok(())
})
.await;
bounded(conn, result, limit)
} }
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().await?; let mut conn = self.conn().await?;
let limit = self.timeouts.maintenance;
let result = tokio::time::timeout(limit, async {
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?;
create_search_tables::<ContactSearchField>(&mut conn).await?; create_search_tables::<ContactSearchField>(&mut conn).await?;
@@ -181,6 +188,9 @@ impl MysqlStore {
create_search_tables::<TracingSearchField>(&mut conn).await?; create_search_tables::<TracingSearchField>(&mut conn).await?;
Ok(()) Ok(())
})
.await;
bounded(conn, result, limit)
} }
} }
+41 -1
View File
@@ -6,6 +6,7 @@
* Modified by Coffey Labs in 2026 for INBUXA. * Modified by Coffey Labs in 2026 for INBUXA.
*/ */
use crate::backend::query_timeout::QueryTimeouts;
use crate::{ use crate::{
search::{ search::{
CalendarSearchField, ContactSearchField, EmailSearchField, FileSearchField, SearchField, CalendarSearchField, ContactSearchField, EmailSearchField, FileSearchField, SearchField,
@@ -14,7 +15,7 @@ use crate::{
write::SearchIndex, write::SearchIndex,
}; };
use mysql_async::Pool; use mysql_async::Pool;
use std::fmt::Display; use std::{fmt::Display, time::Duration};
pub mod blob; pub mod blob;
pub mod lookup; pub mod lookup;
@@ -25,6 +26,8 @@ pub mod write;
pub struct MysqlStore { pub struct MysqlStore {
pub(crate) conn_pool: Pool, 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 /// 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<T>(
conn: mysql_async::Conn,
result: Result<trc::Result<T>, tokio::time::error::Elapsed>,
limit: Duration,
) -> trc::Result<T> {
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)] #[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)
+56 -13
View File
@@ -6,7 +6,7 @@
* Modified by Coffey Labs in 2026 for INBUXA. * 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 crate::{Deserialize, IterateParams, Key, ValueKey, write::ValueClass};
use futures::TryStreamExt; use futures::TryStreamExt;
use mysql_async::{Row, prelude::Queryable}; use mysql_async::{Row, prelude::Queryable};
@@ -17,6 +17,8 @@ impl MysqlStore {
U: Deserialize + 'static, U: Deserialize + 'static,
{ {
let mut conn = self.conn().await?; let mut conn = self.conn().await?;
let limit = self.timeouts.query;
let result = tokio::time::timeout(limit, async {
let s = conn let s = conn
.prep(format!( .prep(format!(
"SELECT v FROM {} WHERE k = ?", "SELECT v FROM {} WHERE k = ?",
@@ -35,10 +37,15 @@ impl MysqlStore {
Ok(None) Ok(None)
} }
}) })
})
.await;
bounded(conn, result, limit)
} }
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().await?; let mut conn = self.conn().await?;
let limit = self.timeouts.query;
let result = tokio::time::timeout(limit, async {
let s = conn let s = conn
.prep(format!( .prep(format!(
"SELECT 1 FROM {} WHERE k = ?", "SELECT 1 FROM {} WHERE k = ?",
@@ -51,6 +58,9 @@ impl MysqlStore {
.await .await
.map_err(into_error) .map_err(into_error)
.map(|r| r.is_some()) .map(|r| r.is_some())
})
.await;
bounded(conn, result, limit)
} }
pub(crate) async fn iterate<T: Key>( pub(crate) async fn iterate<T: Key>(
@@ -64,12 +74,14 @@ impl MysqlStore {
let end = params.end.serialize(0); let end = params.end.serialize(0);
let keys = if params.values { "k, v" } else { "k" }; let keys = if params.values { "k, v" } else { "k" };
let s = conn // inbuxa: a scan may run for hours, so the query limit bounds each
.prep(&match (params.first, params.ascending) { // 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) => { (true, true) => {
format!( format!("SELECT {keys} FROM {table} WHERE k >= ? AND k <= ? ORDER BY k ASC LIMIT 1")
"SELECT {keys} FROM {table} WHERE k >= ? AND k <= ? ORDER BY k ASC LIMIT 1"
)
} }
(true, false) => { (true, false) => {
format!( format!(
@@ -82,10 +94,16 @@ impl MysqlStore {
(false, false) => { (false, false) => {
format!("SELECT {keys} FROM {table} WHERE k >= ? AND k <= ? ORDER BY k DESC") format!("SELECT {keys} FROM {table} WHERE k >= ? AND k <= ? ORDER BY k DESC")
} }
}) };
.await let s = match tokio::time::timeout(limit, conn.prep(&query)).await {
.map_err(into_error)?; Ok(s) => s.map_err(into_error)?,
Err(_) => {
discard(conn);
return Err(query_timeout_error(limit));
}
};
let mut from = begin; let mut from = begin;
let mut stalled = false;
let mut to = end; let mut to = end;
let mut resume_key = None; let mut resume_key = None;
@@ -94,13 +112,26 @@ impl MysqlStore {
let mut timed_out = false; let mut timed_out = false;
{ {
let mut rows = conn let mut rows = match tokio::time::timeout(
.exec_stream::<Row, _, _>(&s, (from.clone(), to.clone())) limit,
conn.exec_stream::<Row, _, _>(&s, (from.clone(), to.clone())),
)
.await .await
.map_err(into_error)?; {
Ok(rows) => rows.map_err(into_error)?,
// Leaves the scan loop for the timeout below
Err(_) => break,
};
loop { 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)) => { Ok(Some(mut row)) => {
let value = if params.values { let value = if params.values {
row.take_opt::<Vec<u8>, _>(1) row.take_opt::<Vec<u8>, _>(1)
@@ -136,6 +167,10 @@ impl MysqlStore {
} }
} }
if stalled {
break;
}
match last_key { match last_key {
Some(last_key) if timed_out => { Some(last_key) if timed_out => {
if params.ascending { if params.ascending {
@@ -148,6 +183,9 @@ impl MysqlStore {
_ => return Ok(()), _ => return Ok(()),
} }
} }
discard(conn);
Err(query_timeout_error(limit))
} }
pub(crate) async fn get_counter( pub(crate) async fn get_counter(
@@ -158,6 +196,8 @@ impl MysqlStore {
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().await?; let mut conn = self.conn().await?;
let limit = self.timeouts.query;
let result = tokio::time::timeout(limit, async {
let s = conn let s = conn
.prep(format!("SELECT v FROM {table} WHERE k = ?")) .prep(format!("SELECT v FROM {table} WHERE k = ?"))
.await .await
@@ -167,5 +207,8 @@ impl MysqlStore {
Ok(None) => Ok(0), Ok(None) => Ok(0),
Err(e) => Err(into_error(e)), Err(e) => Err(into_error(e)),
} }
})
.await;
bounded(conn, result, limit)
} }
} }
+20 -3
View File
@@ -10,8 +10,8 @@ use crate::{
backend::{ backend::{
MAX_TOKEN_LENGTH, MAX_TOKEN_LENGTH,
mysql::{ mysql::{
DELETE_CHUNK_SIZE, MIN_DELETE_CHUNK_SIZE, MysqlSearchField, MysqlStore, into_error, DELETE_CHUNK_SIZE, MIN_DELETE_CHUNK_SIZE, MysqlSearchField, MysqlStore, bounded,
is_timeout_error, into_error, is_timeout_error,
}, },
}, },
search::{ search::{
@@ -27,6 +27,8 @@ 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().await?; let mut conn = self.conn().await?;
let limit = self.timeouts.query;
let result = tokio::time::timeout(limit, async {
let mut tx_opts = TxOpts::default(); let mut tx_opts = TxOpts::default();
tx_opts tx_opts
.with_consistent_snapshot(false) .with_consistent_snapshot(false)
@@ -78,6 +80,9 @@ impl MysqlStore {
} }
trx.commit().await.map_err(into_error) trx.commit().await.map_err(into_error)
})
.await;
bounded(conn, result, limit)
} }
pub async fn query<R: SearchDocumentId>( pub async fn query<R: SearchDocumentId>(
@@ -97,12 +102,17 @@ impl MysqlStore {
} }
let mut conn = self.conn().await?; let mut conn = self.conn().await?;
let limit = self.timeouts.query;
let result = tokio::time::timeout(limit, async {
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)
.await .await
.map(|r| r.into_iter().map(|r| R::from_u64(r as u64)).collect()) .map(|r| r.into_iter().map(|r| R::from_u64(r as u64)).collect())
.map_err(into_error) .map_err(into_error)
})
.await;
bounded(conn, result, limit)
} }
pub async fn unindex(&self, filter: SearchQuery) -> trc::Result<u64> { pub async fn unindex(&self, filter: SearchQuery) -> trc::Result<u64> {
@@ -111,6 +121,8 @@ impl MysqlStore {
let params = build_filter(&mut query, &filter.filters); let params = build_filter(&mut query, &filter.filters);
let mut conn = self.conn().await?; let mut conn = self.conn().await?;
let limit = self.timeouts.maintenance;
let result = tokio::time::timeout(limit, async {
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 {
@@ -137,7 +149,9 @@ impl MysqlStore {
} }
deleted += affected; deleted += affected;
} }
Err(err) if is_timeout_error(&err) && chunk_size > MIN_DELETE_CHUNK_SIZE => { Err(err)
if is_timeout_error(&err) && chunk_size > MIN_DELETE_CHUNK_SIZE =>
{
chunk_size = (chunk_size / 2).max(MIN_DELETE_CHUNK_SIZE); chunk_size = (chunk_size / 2).max(MIN_DELETE_CHUNK_SIZE);
break; break;
} }
@@ -145,6 +159,9 @@ impl MysqlStore {
} }
} }
} }
})
.await;
bounded(conn, result, limit)
} }
} }
+18 -2
View File
@@ -6,7 +6,9 @@
* Modified by Coffey Labs in 2026 for INBUXA. * 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::{ use crate::{
IndexKey, Key, LogKey, SUBSPACE_COUNTER, SUBSPACE_IN_MEMORY_COUNTER, SUBSPACE_QUOTA, IndexKey, Key, LogKey, SUBSPACE_COUNTER, SUBSPACE_IN_MEMORY_COUNTER, SUBSPACE_QUOTA,
SUBSPACE_REGISTRY_IDX, SUBSPACE_REGISTRY_IDX,
@@ -32,7 +34,8 @@ impl MysqlStore {
let start = Instant::now(); let start = Instant::now();
let mut retry_count = 0; let mut retry_count = 0;
let mut conn = self.conn().await?; let mut conn = self.conn().await?;
let limit = self.timeouts.query;
let result = tokio::time::timeout(limit, async {
loop { loop {
let err = match self.write_trx(&mut conn, &mut batch).await { let err = match self.write_trx(&mut conn, &mut batch).await {
Ok(result) => { Ok(result) => {
@@ -67,6 +70,9 @@ impl MysqlStore {
tokio::time::sleep(Duration::from_millis(backoff)).await; tokio::time::sleep(Duration::from_millis(backoff)).await;
retry_count += 1; retry_count += 1;
} }
})
.await;
bounded(conn, result, limit)
} }
async fn write_trx( async fn write_trx(
@@ -385,15 +391,22 @@ 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().await?; let mut conn = self.conn().await?;
let limit = self.timeouts.maintenance;
let result = tokio::time::timeout(limit, async {
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?;
} }
Ok(()) Ok(())
})
.await;
bounded(conn, result, limit)
} }
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().await?; let mut conn = self.conn().await?;
let limit = self.timeouts.maintenance;
let result = tokio::time::timeout(limit, async {
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);
@@ -450,6 +463,9 @@ impl MysqlStore {
} }
} }
} }
})
.await;
bounded(conn, result, limit)
} }
} }
+18 -1
View File
@@ -2,13 +2,15 @@
* 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;
use crate::backend::postgres::into_pool_error; use crate::backend::postgres::into_pool_error;
use super::{PostgresStore, into_error}; use super::{PostgresStore, bounded, into_error};
impl PostgresStore { impl PostgresStore {
pub(crate) async fn get_blob( pub(crate) async fn get_blob(
@@ -17,6 +19,8 @@ impl PostgresStore {
range: Range<usize>, range: Range<usize>,
) -> trc::Result<Option<Vec<u8>>> { ) -> trc::Result<Option<Vec<u8>>> {
let conn = self.conn_pool.get().await.map_err(into_pool_error)?; let conn = self.conn_pool.get().await.map_err(into_pool_error)?;
let limit = self.timeouts.query;
let result = tokio::time::timeout(limit, async {
let s = conn let s = conn
.prepare_cached("SELECT v FROM t WHERE k = $1") .prepare_cached("SELECT v FROM t WHERE k = $1")
.await .await
@@ -39,10 +43,15 @@ impl PostgresStore {
} }
}) })
.map_err(into_error) .map_err(into_error)
})
.await;
bounded(conn, result, limit)
} }
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 conn = self.conn_pool.get().await.map_err(into_pool_error)?; let conn = self.conn_pool.get().await.map_err(into_pool_error)?;
let limit = self.timeouts.query;
let result = tokio::time::timeout(limit, async {
let s = conn let s = conn
.prepare_cached( .prepare_cached(
"INSERT INTO t (k, v) VALUES ($1, $2) ON CONFLICT (k) DO UPDATE SET v = EXCLUDED.v", "INSERT INTO t (k, v) VALUES ($1, $2) ON CONFLICT (k) DO UPDATE SET v = EXCLUDED.v",
@@ -53,10 +62,15 @@ impl PostgresStore {
.await .await
.map_err(into_error) .map_err(into_error)
.map(|_| ()) .map(|_| ())
})
.await;
bounded(conn, result, limit)
} }
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 conn = self.conn_pool.get().await.map_err(into_pool_error)?; let conn = self.conn_pool.get().await.map_err(into_pool_error)?;
let limit = self.timeouts.query;
let result = tokio::time::timeout(limit, async {
let s = conn let s = conn
.prepare_cached("DELETE FROM t WHERE k = $1") .prepare_cached("DELETE FROM t WHERE k = $1")
.await .await
@@ -65,5 +79,8 @@ impl PostgresStore {
.await .await
.map_err(into_error) .map_err(into_error)
.map(|hits| hits > 0) .map(|hits| hits > 0)
})
.await;
bounded(conn, result, limit)
} }
} }
+8 -1
View File
@@ -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 crate::{QueryResult, QueryType, backend::postgres::into_pool_error}; use crate::{QueryResult, QueryType, backend::postgres::into_pool_error};
@@ -12,7 +14,7 @@ use tokio_postgres::types::{FromSql, ToSql, Type};
use crate::IntoRows; use crate::IntoRows;
use super::{PostgresStore, into_error}; use super::{PostgresStore, bounded, into_error};
impl PostgresStore { impl PostgresStore {
pub(crate) async fn sql_query<T: QueryResult>( pub(crate) async fn sql_query<T: QueryResult>(
@@ -21,6 +23,8 @@ impl PostgresStore {
params_: &[crate::Value<'_>], params_: &[crate::Value<'_>],
) -> trc::Result<T> { ) -> trc::Result<T> {
let conn = self.conn_pool.get().await.map_err(into_pool_error)?; let conn = self.conn_pool.get().await.map_err(into_pool_error)?;
let limit = self.timeouts.query;
let result = tokio::time::timeout(limit, async {
let s = conn.prepare_cached(query).await.map_err(into_error)?; let s = conn.prepare_cached(query).await.map_err(into_error)?;
let params = params_ let params = params_
.iter() .iter()
@@ -48,6 +52,9 @@ impl PostgresStore {
.await .await
.map_or_else(|e| Err(into_error(e)), |r| Ok(T::from_query_all(r))), .map_or_else(|e| Err(into_error(e)), |r| Ok(T::from_query_all(r))),
} }
})
.await;
bounded(conn, result, limit)
} }
} }
+86 -4
View File
@@ -6,7 +6,7 @@
* Modified by Coffey Labs in 2026 for INBUXA. * Modified by Coffey Labs in 2026 for INBUXA.
*/ */
use super::{PostgresStore, into_error}; use super::{PostgresStore, bounded, into_error};
use crate::{ use crate::{
backend::postgres::{ backend::postgres::{
PsqlSearchField, into_pool_error, PsqlSearchField, into_pool_error,
@@ -119,6 +119,7 @@ impl PostgresStore {
Store::PostgreSQL(Arc::new(PostgresStore { Store::PostgreSQL(Arc::new(PostgresStore {
conn_pool: pool, conn_pool: pool,
ts_configs: ts_configs.clone(), ts_configs: ts_configs.clone(),
timeouts: Default::default(),
})), })),
replica.host, replica.host,
replica.port as u16, replica.port as u16,
@@ -129,6 +130,7 @@ impl PostgresStore {
let primary = Store::PostgreSQL(Arc::new(PostgresStore { let primary = Store::PostgreSQL(Arc::new(PostgresStore {
conn_pool: primary_pool, conn_pool: primary_pool,
ts_configs, ts_configs,
timeouts: Default::default(),
})); }));
// ST-1: no replicas, no change // ST-1: no replicas, no change
@@ -147,7 +149,8 @@ impl PostgresStore {
pub(crate) async fn create_storage_tables(&self) -> trc::Result<()> { pub(crate) async fn create_storage_tables(&self) -> trc::Result<()> {
let conn = self.conn_pool.get().await.map_err(into_pool_error)?; 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 [ for table in [
SUBSPACE_ACL, SUBSPACE_ACL,
SUBSPACE_TASK_QUEUE, SUBSPACE_TASK_QUEUE,
@@ -213,11 +216,15 @@ impl PostgresStore {
} }
Ok(()) Ok(())
})
.await;
bounded(conn, result, limit)
} }
pub(crate) async fn create_search_tables(&self) -> trc::Result<()> { pub(crate) async fn create_search_tables(&self) -> trc::Result<()> {
let conn = self.conn_pool.get().await.map_err(into_pool_error)?; 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::<EmailSearchField>(&conn).await?; create_search_tables::<EmailSearchField>(&conn).await?;
create_search_tables::<CalendarSearchField>(&conn).await?; create_search_tables::<CalendarSearchField>(&conn).await?;
create_search_tables::<ContactSearchField>(&conn).await?; create_search_tables::<ContactSearchField>(&conn).await?;
@@ -225,6 +232,9 @@ impl PostgresStore {
create_search_tables::<TracingSearchField>(&conn).await?; create_search_tables::<TracingSearchField>(&conn).await?;
Ok(()) Ok(())
})
.await;
bounded(conn, result, limit)
} }
} }
@@ -265,12 +275,21 @@ async fn create_search_tables<T: SearchableField + PsqlSearchField + 'static>(
for field in T::all_fields() { for field in T::all_fields() {
if field.is_text() || field.is_json() { if field.is_text() || field.is_json() {
let column_name = field.column(); let column_name = field.column();
// inbuxa: with GIN's default fastupdate=on, new entries wait in
// an unindexed pending list that every search scans in full
// until a VACUUM (or 4 MB of backlog) merges it. On a mailbox
// taking steady mail that list never drains and searches slow
// from milliseconds to hundreds of them. Pay the index update
// at insert time instead.
let index_name = format!("gin_{table_name}_{column_name}");
let create_index_query = format!( let create_index_query = format!(
"CREATE INDEX IF NOT EXISTS gin_{table_name}_{column_name} ON {table_name} USING GIN({column_name})", "CREATE INDEX IF NOT EXISTS {index_name} ON {table_name} USING GIN({column_name}) WITH (fastupdate = off)",
); );
conn.execute(&create_index_query, &[]) conn.execute(&create_index_query, &[])
.await .await
.map_err(into_error)?; .map_err(into_error)?;
// Indexes made before this change keep fastupdate=on
disable_gin_fastupdate(conn, &index_name).await;
} }
if field.is_indexed() { if field.is_indexed() {
@@ -287,6 +306,69 @@ async fn create_search_tables<T: SearchableField + PsqlSearchField + 'static>(
Ok(()) Ok(())
} }
/// inbuxa: turns fastupdate off on a GIN index made with the default and
/// merges the pending list it has built up. Idempotent: an index that already
/// has the option is left alone, so this costs one catalog read per index at
/// startup. A failure is logged and startup goes on, since search still works,
/// only slower.
async fn disable_gin_fastupdate(conn: &Object, index_name: &str) {
if let Err(err) = try_disable_gin_fastupdate(conn, index_name).await {
trc::event!(
Store(trc::StoreEvent::PostgresqlError),
Details = format!("Failed to turn off fastupdate on search index {index_name}"),
Reason = err.to_string(),
);
}
}
async fn try_disable_gin_fastupdate(conn: &Object, index_name: &str) -> trc::Result<()> {
let options = conn
.query_opt(
"SELECT COALESCE(reloptions, '{}')::text[] FROM pg_class WHERE oid = to_regclass($1)",
&[&index_name],
)
.await
.map_err(into_error)?
.map(|row| row.try_get::<_, Vec<String>>(0))
.transpose()
.map_err(into_error)?;
let Some(options) = options else {
return Ok(());
};
if gin_fastupdate_is_off(&options) {
return Ok(());
}
// SET (fastupdate) takes a SHARE UPDATE EXCLUSIVE lock, which doesn't
// block reads or writes. Turning it off stops new entries going to the
// pending list but doesn't flush the entries already there.
conn.execute(
&format!("ALTER INDEX {index_name} SET (fastupdate = off)"),
&[],
)
.await
.map_err(into_error)?;
conn.query_one(
"SELECT gin_clean_pending_list($1::text::regclass)",
&[&index_name],
)
.await
.map_err(into_error)?;
Ok(())
}
/// Whether a relation's reloptions turn GIN's fastupdate off.
fn gin_fastupdate_is_off(options: &[String]) -> bool {
options.iter().any(|option| {
option.split_once('=').is_some_and(|(name, value)| {
name.trim().eq_ignore_ascii_case("fastupdate")
&& matches!(
value.trim().to_ascii_lowercase().as_str(),
"off" | "false" | "no" | "0" | "f" | "n"
)
})
})
}
async fn discover_ts_configs(pool: &Pool) -> AHashSet<&'static str> { async fn discover_ts_configs(pool: &Pool) -> AHashSet<&'static str> {
let mut ts_configs = AHashSet::from_iter([PG_FALLBACK_LANG, PG_UNSTEMMED_LANG]); let mut ts_configs = AHashSet::from_iter([PG_FALLBACK_LANG, PG_UNSTEMMED_LANG]);
+33 -1
View File
@@ -6,6 +6,7 @@
* Modified by Coffey Labs in 2026 for INBUXA. * Modified by Coffey Labs in 2026 for INBUXA.
*/ */
use crate::backend::query_timeout::QueryTimeouts;
use crate::{ use crate::{
search::{ search::{
CalendarSearchField, ContactSearchField, EmailSearchField, FileSearchField, SearchField, CalendarSearchField, ContactSearchField, EmailSearchField, FileSearchField, SearchField,
@@ -14,7 +15,8 @@ use crate::{
write::SearchIndex, write::SearchIndex,
}; };
use ahash::AHashSet; use ahash::AHashSet;
use deadpool_postgres::Pool; use deadpool_postgres::{Object, Pool};
use std::time::Duration;
use tokio_postgres::error::SqlState; use tokio_postgres::error::SqlState;
pub mod blob; pub mod blob;
@@ -28,6 +30,8 @@ pub mod write;
pub struct PostgresStore { pub struct PostgresStore {
pub(crate) conn_pool: Pool, pub(crate) conn_pool: Pool,
pub(crate) ts_configs: AHashSet<&'static str>, pub(crate) ts_configs: AHashSet<&'static str>,
/// inbuxa: client-side query limits (see backend::query_timeout)
pub(crate) timeouts: QueryTimeouts,
} }
#[inline(always)] #[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<T>(
conn: Object,
result: Result<trc::Result<T>, tokio::time::error::Elapsed>,
limit: Duration,
) -> trc::Result<T> {
match result {
Ok(result) => result,
Err(_) => {
drop(Object::take(conn));
Err(query_timeout_error(limit))
}
}
}
#[inline(always)] #[inline(always)]
pub(crate) fn into_pool_error(err: deadpool_postgres::PoolError) -> trc::Error { pub(crate) fn into_pool_error(err: deadpool_postgres::PoolError) -> trc::Error {
match err { match err {
+55 -10
View File
@@ -2,9 +2,11 @@
* 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::{PostgresStore, into_error, is_timeout_error}; use super::{PostgresStore, bounded, into_error, is_timeout_error, query_timeout_error};
use crate::{ use crate::{
Deserialize, IterateParams, Key, ValueKey, backend::postgres::into_pool_error, Deserialize, IterateParams, Key, ValueKey, backend::postgres::into_pool_error,
write::ValueClass, write::ValueClass,
@@ -17,6 +19,8 @@ impl PostgresStore {
U: Deserialize + 'static, U: Deserialize + 'static,
{ {
let conn = self.conn_pool.get().await.map_err(into_pool_error)?; let conn = self.conn_pool.get().await.map_err(into_pool_error)?;
let limit = self.timeouts.query;
let result = tokio::time::timeout(limit, async {
let s = conn let s = conn
.prepare_cached(&format!( .prepare_cached(&format!(
"SELECT v FROM {} WHERE k = $1", "SELECT v FROM {} WHERE k = $1",
@@ -35,10 +39,15 @@ impl PostgresStore {
Ok(None) Ok(None)
} }
}) })
})
.await;
bounded(conn, result, limit)
} }
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 conn = self.conn_pool.get().await.map_err(into_pool_error)?; let conn = self.conn_pool.get().await.map_err(into_pool_error)?;
let limit = self.timeouts.query;
let result = tokio::time::timeout(limit, async {
let s = conn let s = conn
.prepare_cached(&format!( .prepare_cached(&format!(
"SELECT 1 FROM {} WHERE k = $1", "SELECT 1 FROM {} WHERE k = $1",
@@ -51,6 +60,9 @@ impl PostgresStore {
.await .await
.map_err(into_error) .map_err(into_error)
.map(|r| r.is_some()) .map(|r| r.is_some())
})
.await;
bounded(conn, result, limit)
} }
pub(crate) async fn iterate<T: Key>( pub(crate) async fn iterate<T: Key>(
@@ -64,8 +76,12 @@ impl PostgresStore {
let end = params.end.serialize(0); let end = params.end.serialize(0);
let keys = if params.values { "k, v" } else { "k" }; let keys = if params.values { "k, v" } else { "k" };
let s = conn // inbuxa: a scan may run for hours, so the query limit bounds each
.prepare_cached(&match (params.first, params.ascending) { // 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) => { (true, true) => {
format!( format!(
"SELECT {keys} FROM {table} WHERE k >= $1 AND k <= $2 ORDER BY k ASC LIMIT 1" "SELECT {keys} FROM {table} WHERE k >= $1 AND k <= $2 ORDER BY k ASC LIMIT 1"
@@ -82,26 +98,43 @@ impl PostgresStore {
(false, false) => { (false, false) => {
format!("SELECT {keys} FROM {table} WHERE k >= $1 AND k <= $2 ORDER BY k DESC") format!("SELECT {keys} FROM {table} WHERE k >= $1 AND k <= $2 ORDER BY k DESC")
} }
}) };
.await.map_err(into_error)?; 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 from = begin;
let mut to = end; let mut to = end;
let mut resume_key: Option<Vec<u8>> = None; let mut resume_key: Option<Vec<u8>> = None;
let mut stalled = false;
loop { loop {
let mut last_key = None; let mut last_key = None;
let mut timed_out = false; let mut timed_out = false;
{ {
let rows = conn let rows =
.query_raw(&s, &[&from, &to]) match tokio::time::timeout(limit, conn.query_raw(&s, &[&from, &to])).await {
.await Ok(rows) => rows.map_err(into_error)?,
.map_err(into_error)?; // Leaves the scan loop for the timeout below
Err(_) => break,
};
pin_mut!(rows); pin_mut!(rows);
loop { 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)) => { Ok(Some(row)) => {
let key = row.try_get::<_, &[u8]>(0).map_err(into_error)?; let key = row.try_get::<_, &[u8]>(0).map_err(into_error)?;
let value = if params.values { let value = if params.values {
@@ -132,6 +165,10 @@ impl PostgresStore {
} }
} }
if stalled {
break;
}
match last_key { match last_key {
Some(last_key) if timed_out => { Some(last_key) if timed_out => {
if params.ascending { if params.ascending {
@@ -144,6 +181,9 @@ impl PostgresStore {
_ => return Ok(()), _ => return Ok(()),
} }
} }
drop(deadpool_postgres::Object::take(conn));
Err(query_timeout_error(limit))
} }
pub(crate) async fn get_counter( pub(crate) async fn get_counter(
@@ -155,6 +195,8 @@ impl PostgresStore {
let key = key.serialize(0); let key = key.serialize(0);
let conn = self.conn_pool.get().await.map_err(into_pool_error)?; let conn = self.conn_pool.get().await.map_err(into_pool_error)?;
let limit = self.timeouts.query;
let result = tokio::time::timeout(limit, async {
let s = conn let s = conn
.prepare_cached(&format!("SELECT v FROM {table} WHERE k = $1")) .prepare_cached(&format!("SELECT v FROM {table} WHERE k = $1"))
.await .await
@@ -164,5 +206,8 @@ impl PostgresStore {
Ok(None) => Ok(0), Ok(None) => Ok(0),
Err(e) => Err(into_error(e)), Err(e) => Err(into_error(e)),
} }
})
.await;
bounded(conn, result, limit)
} }
} }
+19 -4
View File
@@ -10,8 +10,8 @@ use crate::{
backend::{ backend::{
MAX_TOKEN_LENGTH, MAX_TOKEN_LENGTH,
postgres::{ postgres::{
DELETE_CHUNK_SIZE, MIN_DELETE_CHUNK_SIZE, PostgresStore, PsqlSearchField, into_error, DELETE_CHUNK_SIZE, MIN_DELETE_CHUNK_SIZE, PostgresStore, PsqlSearchField, bounded,
into_pool_error, is_timeout_error, into_error, into_pool_error, is_timeout_error,
}, },
}, },
search::{ search::{
@@ -36,6 +36,8 @@ impl PostgresStore {
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().await.map_err(into_pool_error)?; let mut conn = self.conn_pool.get().await.map_err(into_pool_error)?;
let limit = self.timeouts.query;
let result = tokio::time::timeout(limit, async {
let trx = conn let trx = conn
.build_transaction() .build_transaction()
.isolation_level(IsolationLevel::ReadCommitted) .isolation_level(IsolationLevel::ReadCommitted)
@@ -85,8 +87,8 @@ impl PostgresStore {
if let Some(value) = fields.get(field) { if let Some(value) = fields.get(field) {
let value_ref = format!("${}", values.len() + 1); let value_ref = format!("${}", values.len() + 1);
let (text_len, language) = if let SearchValue::Text { value, language } = value let (text_len, language) =
{ if let SearchValue::Text { value, language } = value {
(value.len(), self.ts_config(language)) (value.len(), self.ts_config(language))
} else { } else {
(0, PG_UNSTEMMED_LANG) (0, PG_UNSTEMMED_LANG)
@@ -155,6 +157,9 @@ impl PostgresStore {
} }
trx.commit().await.map_err(into_error) trx.commit().await.map_err(into_error)
})
.await;
bounded(conn, result, limit)
} }
pub async fn query<R: SearchDocumentId>( pub async fn query<R: SearchDocumentId>(
@@ -170,6 +175,8 @@ impl PostgresStore {
build_sort(&mut query, sort); build_sort(&mut query, sort);
} }
let conn = self.conn_pool.get().await.map_err(into_pool_error)?; let conn = self.conn_pool.get().await.map_err(into_pool_error)?;
let limit = self.timeouts.query;
let result = tokio::time::timeout(limit, async {
let s = conn.prepare_cached(&query).await.map_err(into_error)?; let s = conn.prepare_cached(&query).await.map_err(into_error)?;
conn.query(&s, params.as_slice()) conn.query(&s, params.as_slice())
@@ -180,6 +187,9 @@ impl PostgresStore {
.collect::<Result<Vec<R>, _>>() .collect::<Result<Vec<R>, _>>()
}) })
.map_err(into_error) .map_err(into_error)
})
.await;
bounded(conn, result, limit)
} }
pub async fn unindex(&self, filter: SearchQuery) -> trc::Result<u64> { pub async fn unindex(&self, filter: SearchQuery) -> trc::Result<u64> {
@@ -189,6 +199,8 @@ impl PostgresStore {
let params = self.build_filter(&mut where_clause, &filter.filters); let params = self.build_filter(&mut where_clause, &filter.filters);
let params = params.iter().map(SqlParam::as_sql).collect::<Vec<_>>(); let params = params.iter().map(SqlParam::as_sql).collect::<Vec<_>>();
let conn = self.conn_pool.get().await.map_err(into_pool_error)?; let conn = self.conn_pool.get().await.map_err(into_pool_error)?;
let limit = self.timeouts.maintenance;
let result = tokio::time::timeout(limit, async {
let s = conn let s = conn
.prepare_cached(&format!("DELETE FROM {table}{where_clause}")) .prepare_cached(&format!("DELETE FROM {table}{where_clause}"))
.await .await
@@ -223,6 +235,9 @@ impl PostgresStore {
} }
} }
} }
})
.await;
bounded(conn, result, limit)
} }
fn build_filter<'x>( fn build_filter<'x>(
+18 -2
View File
@@ -2,9 +2,11 @@
* 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::{PostgresStore, into_error, is_timeout_error}; use super::{PostgresStore, bounded, into_error, is_timeout_error};
use crate::{ use crate::{
IndexKey, Key, LogKey, SUBSPACE_COUNTER, SUBSPACE_IN_MEMORY_COUNTER, SUBSPACE_QUOTA, IndexKey, Key, LogKey, SUBSPACE_COUNTER, SUBSPACE_IN_MEMORY_COUNTER, SUBSPACE_QUOTA,
SUBSPACE_REGISTRY_IDX, SUBSPACE_REGISTRY_IDX,
@@ -30,6 +32,8 @@ enum CommitError {
impl PostgresStore { impl PostgresStore {
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 mut conn = self.conn_pool.get().await.map_err(into_pool_error)?; let mut conn = self.conn_pool.get().await.map_err(into_pool_error)?;
let limit = self.timeouts.query;
let result = tokio::time::timeout(limit, async {
let start = Instant::now(); let start = Instant::now();
let mut retry_count = 0; let mut retry_count = 0;
@@ -72,6 +76,9 @@ impl PostgresStore {
} }
} }
} }
})
.await;
bounded(conn, result, limit)
} }
async fn write_trx( async fn write_trx(
@@ -393,16 +400,22 @@ impl PostgresStore {
pub(crate) async fn purge_store(&self) -> trc::Result<()> { pub(crate) async fn purge_store(&self) -> trc::Result<()> {
let conn = self.conn_pool.get().await.map_err(into_pool_error)?; 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] { for subspace in [SUBSPACE_QUOTA, SUBSPACE_COUNTER, SUBSPACE_IN_MEMORY_COUNTER] {
purge_table(&conn, char::from(subspace)).await?; 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<()> { 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 conn = self.conn_pool.get().await.map_err(into_pool_error)?;
let limit = self.timeouts.maintenance;
let result = tokio::time::timeout(limit, async {
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);
@@ -459,6 +472,9 @@ impl PostgresStore {
} }
} }
} }
})
.await;
bounded(conn, result, limit)
} }
} }
+77
View File
@@ -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,
}
}
}
+297
View File
@@ -0,0 +1,297 @@
/*
* SPDX-FileCopyrightText: 2026 Coffey Labs
*
* SPDX-License-Identifier: AGPL-3.0-only
*/
//! A node follows edits to its cluster role without a restart: outbound
//! delivery and report tasks start when the role gains outboundMta and stop
//! when it loses it. Upstream decided at boot whether the queue, report and
//! task managers ran at all. Needs a store the seed and the node can share
//! (STORE=PostgreSql or MySql).
use crate::utils::server::{TestServer, TestServerBuilder};
use common::Server;
use registry::{
schema::{
enums::ClusterTaskType,
prelude::{Object, ObjectType},
structs::{
ClusterListenerGroup, ClusterRole, ClusterTaskGroup, ClusterTaskGroupProperties, Task,
TaskStatus, TaskTlsReport,
},
},
types::{id::ObjectId, map::Map},
};
use smtp::{
queue::{Message, Status},
reporting::send::MtaReportSend,
};
use std::time::{Duration, Instant};
use store::{
Deserialize, IterateParams, ValueKey,
registry::write::{RegistryWrite, RegistryWriteResult},
write::{AlignedBytes, Archive, BatchBuilder, QueueClass, TaskQueueClass, ValueClass},
};
use types::id::Id;
use utils::snowflake::SnowflakeIdGenerator;
const BUSY_ROLE: &str = "live_role_busy";
const IDLE_ROLE: &str = "live_role_idle";
const WITH: &[ClusterTaskType] = &[
ClusterTaskType::PushNotifications,
ClusterTaskType::OutboundMta,
];
const WITHOUT: &[ClusterTaskType] = &[ClusterTaskType::PushNotifications];
const RCPT_DOMAIN: &str = "live-role.invalid";
#[tokio::test(flavor = "multi_thread")]
pub async fn live_role_tests() {
if matches!(
std::env::var("STORE").as_deref(),
Ok("RocksDb" | "Sqlite") | Err(_)
) {
println!("Skipping live role tests: they need a store the nodes can share.");
return;
}
println!(
"Running live role tests on {}...",
std::env::var("STORE").unwrap_or_default()
);
// Two roles: one with outboundMta, one with no task type at all
let seed = TestServerBuilder::new("live_roles_seed").await;
let busy_id = seed.insert_object(role(BUSY_ROLE, WITH)).await;
let idle_id = seed.insert_object(role(IDLE_ROLE, WITHOUT)).await;
let seed = seed.disable_services().build().await;
let registry = seed.server.clone();
// 1. The rehearsal case: a node started with outboundMta has it taken
// away. Upstream kept delivering, report messages included, until a
// restart.
let node = start_node("live_roles_busy", BUSY_ROLE).await;
let server = node.server.clone();
assert!(server.core.network.roles.outbound_mta);
let (msg, task) = queue_work(&server, "busy-before").await;
assert_runs(&server, &msg, task).await;
set_role(&registry, &server, busy_id, role(BUSY_ROLE, WITHOUT)).await;
let (msg, task) = queue_work(&server, "busy-off").await;
assert_idle(&server, &msg, task).await;
// Given back, it takes up the work left waiting
set_role(&registry, &server, busy_id, role(BUSY_ROLE, WITH)).await;
assert_runs(&server, &msg, task).await;
// Off again, so it leaves the next node's work alone
set_role(&registry, &server, busy_id, role(BUSY_ROLE, WITHOUT)).await;
// 2. A node started with no task type at all gains outboundMta.
// Upstream never started its queue, report or task manager, so the
// role did nothing until a restart.
let node2 = start_node("live_roles_idle", IDLE_ROLE).await;
let server2 = node2.server.clone();
assert!(!server2.core.network.roles.outbound_mta);
let (msg, task) = queue_work(&server2, "idle-off").await;
assert_idle(&server2, &msg, task).await;
set_role(&registry, &server2, idle_id, role(IDLE_ROLE, WITH)).await;
assert_runs(&server2, &msg, task).await;
set_role(&registry, &server2, idle_id, role(IDLE_ROLE, WITHOUT)).await;
if seed.is_reset() {
seed.temp_dir.delete();
node.temp_dir.delete();
node2.temp_dir.delete();
}
}
async fn start_node(name: &str, role: &str) -> TestServer {
TestServerBuilder::new_with_role(
name,
format!("{name}.example.com").replace('_', "-"),
Some(role.into()),
false,
)
.await
.build_with_opts(false)
.await
}
/// Neither the message nor the report task is touched.
async fn assert_idle(server: &Server, msg: &str, task: u64) {
tokio::time::sleep(Duration::from_secs(4)).await;
server.notify_task_queue();
tokio::time::sleep(Duration::from_secs(1)).await;
assert!(
!attempted(server, msg).await,
"delivery attempted without outboundMta"
);
assert!(
is_pending(server, task).await,
"report task claimed without outboundMta"
);
}
/// Delivery of the message is attempted and the report task runs.
async fn assert_runs(server: &Server, msg: &str, task: u64) {
wait_for(Duration::from_secs(20), "message delivery attempt", || {
attempted(server, msg)
})
.await;
wait_for(Duration::from_secs(20), "report task to run", || async {
!is_pending(server, task).await
})
.await;
}
fn role(name: &str, tasks: &[ClusterTaskType]) -> ClusterRole {
ClusterRole {
name: name.into(),
description: None,
listeners: ClusterListenerGroup::EnableAll,
tasks: ClusterTaskGroup::EnableSome(ClusterTaskGroupProperties {
task_types: Map::new(tasks.to_vec()),
}),
}
}
/// Stores a new version of a role and reloads the node's settings, as a
/// JMAP write to the role does.
async fn set_role(registry: &Server, node: &Server, id: Id, new: ClusterRole) {
let enabled = matches!(&new.tasks, ClusterTaskGroup::EnableSome(group)
if group.task_types.iter().any(|t| *t == ClusterTaskType::OutboundMta));
let old = registry
.registry()
.get(ObjectId::new(ObjectType::ClusterRole, id))
.await
.unwrap()
.expect("role not found");
let new = Object::from(new);
let result = registry
.registry()
.write(RegistryWrite::update(id, &new, &old))
.await
.unwrap();
assert!(
matches!(result, RegistryWriteResult::Success(_)),
"role update refused"
);
assert_eq!(
node.reload_after_write(ObjectType::ClusterRole).await,
Some(Ok(()))
);
assert_eq!(
node.inner.shared_core.load().network.roles.outbound_mta,
enabled
);
}
/// Queues a message to an unreachable domain and schedules a TLS report
/// task, both due now. Returns the recipient's local part and the task id.
async fn queue_work(server: &Server, name: &str) -> (String, u64) {
let local = format!("{name}-{}", SnowflakeIdGenerator::global_id().unwrap());
let rcpt = format!("{local}@{RCPT_DOMAIN}");
server
.send_autogenerated(
"[email protected]",
[rcpt.as_str()].into_iter(),
format!(
"From: [email protected]\r\nTo: {rcpt}\r\n\
Subject: live role test\r\n\r\nTest\r\n"
)
.into_bytes(),
None,
0,
)
.await;
assert!(
queued_recipient(server, &rcpt).await.is_some(),
"message to {rcpt} was not queued"
);
let task = SnowflakeIdGenerator::global_id().unwrap();
let mut batch = BatchBuilder::new();
batch.schedule_task_with_id(
task,
Task::TlsReport(TaskTlsReport {
report_id: u64::MAX.into(),
status: TaskStatus::now(),
}),
);
server.store().write(batch.build_all()).await.unwrap();
server.notify_task_queue();
(rcpt, task)
}
/// Whether delivery to `rcpt` was tried: the message is gone, or its
/// recipient is no longer scheduled or has a retry count.
async fn attempted(server: &Server, rcpt: &str) -> bool {
match queued_recipient(server, rcpt).await {
None => true,
Some((status_scheduled, retries)) => !status_scheduled || retries > 0,
}
}
/// The queued recipient `rcpt`: whether it is still scheduled, and how many
/// times delivery was retried.
async fn queued_recipient(server: &Server, rcpt: &str) -> Option<(bool, u32)> {
let mut found = None;
server
.store()
.iterate(
IterateParams::new(
ValueKey::from(ValueClass::Queue(QueueClass::Message(0))),
ValueKey::from(ValueClass::Queue(QueueClass::Message(u64::MAX))),
),
|_, value| {
let message = <Archive<AlignedBytes> as Deserialize>::deserialize(value)?
.deserialize::<Message>()?;
if let Some(recipient) = message
.recipients
.iter()
.find(|recipient| recipient.address.as_ref() == rcpt)
{
found = Some((
matches!(recipient.status, Status::Scheduled),
recipient.retry.inner,
));
return Ok(false);
}
Ok(true)
},
)
.await
.unwrap();
found
}
async fn is_pending(server: &Server, id: u64) -> bool {
matches!(
server
.store()
.get_value::<Task>(ValueKey::from(ValueClass::TaskQueue(
TaskQueueClass::Task { id },
)))
.await
.unwrap()
.map(|task| task.status().clone()),
Some(TaskStatus::Pending(_))
)
}
async fn wait_for<F, Fut>(within: Duration, what: &str, mut check: F)
where
F: FnMut() -> Fut,
Fut: Future<Output = bool>,
{
let started = Instant::now();
while !check().await {
assert!(
started.elapsed() < within,
"still waiting for the {what} after {:?}",
started.elapsed()
);
tokio::time::sleep(Duration::from_millis(250)).await;
}
}
+1
View File
@@ -7,6 +7,7 @@
*/ */
pub mod broadcast; pub mod broadcast;
pub mod live_roles; // inbuxa: role edits apply without a restart
#[cfg(feature = "nats")] #[cfg(feature = "nats")]
pub mod coordinator; // inbuxa: coordinator reconnects pub mod coordinator; // inbuxa: coordinator reconnects
pub mod stress; pub mod stress;
+2
View File
@@ -21,6 +21,8 @@ pub mod replica_mysql; // inbuxa: read replicas on MySQL
#[cfg(all(feature = "postgres", feature = "redis"))] #[cfg(all(feature = "postgres", feature = "redis"))]
pub mod replica_cluster; // inbuxa: read replicas across nodes pub mod replica_cluster; // inbuxa: read replicas across nodes
pub mod scaleout; // inbuxa: scale-out storage pub mod scaleout; // inbuxa: scale-out storage
#[cfg(feature = "postgres")]
pub mod search_gin; // inbuxa: GIN indexes without a pending list
#[cfg(any(feature = "postgres", feature = "mysql"))] #[cfg(any(feature = "postgres", feature = "mysql"))]
pub mod sql_timeout; pub mod sql_timeout;
pub mod task_locks; // inbuxa: task locks across nodes pub mod task_locks; // inbuxa: task locks across nodes
+283 -3
View File
@@ -9,11 +9,31 @@
//! the pool's timeouts. Upstream's pools had none, so the worker waited for //! 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 //! good. No database is needed: a local listener that never answers plays
//! the server. //! 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 registry::schema::structs::DataStore;
use std::time::{Duration, Instant}; use std::{
use store::{Store, ValueKey, write::ValueClass}; sync::{
use tokio::net::TcpListener; 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. /// Accepts connections on a local port and never sends a byte.
async fn silent_server() -> u16 { async fn silent_server() -> u16 {
@@ -94,3 +114,263 @@ pub async fn mysql_pool_timeout() {
) )
.await; .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<AtomicBool>,
}
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<AtomicBool>,
) {
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::<u64>(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::<u64>(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::<u64>(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();
}
}
+159
View File
@@ -0,0 +1,159 @@
/*
* SPDX-FileCopyrightText: 2026 Coffey Labs
*
* SPDX-License-Identifier: AGPL-3.0-only
*/
//! PostgreSQL full-text GIN indexes are built with fastupdate off, and an
//! index made earlier with the default is switched over at startup. With
//! fastupdate on, new entries wait in a pending list that every search scans
//! in full until VACUUM merges it.
use crate::utils::storage::build_data_store;
use registry::schema::structs::DataStore;
use store::{Rows, SearchStore, Store};
const SCHEMA: &str = "gin_fastupdate_test";
#[tokio::test(flavor = "multi_thread")]
pub async fn postgres_gin_fastupdate() {
println!("Running PostgreSQL GIN fastupdate test...");
// Work in a schema of our own so the shared search tables are untouched
let admin = Store::build(build_data_store("PostgreSql", "").await)
.await
.expect("Failed to connect to PostgreSQL");
for query in [
format!("DROP SCHEMA IF EXISTS {SCHEMA} CASCADE"),
format!("CREATE SCHEMA {SCHEMA}"),
] {
admin.sql_query::<usize>(&query, vec![]).await.unwrap();
}
let DataStore::PostgreSql(mut config) = build_data_store("PostgreSql", "").await else {
unreachable!()
};
config.options = Some(format!("-c search_path={SCHEMA}"));
let store = Store::build(DataStore::PostgreSql(config))
.await
.expect("Failed to connect to PostgreSQL");
let search = SearchStore::Store(store.clone());
// A fresh schema
search.create_indexes().await.unwrap();
let indexes = gin_indexes(&admin).await;
assert!(
indexes.len() >= 4,
"expected the search GIN indexes, found {indexes:?}"
);
for (name, options) in &indexes {
assert!(
options.contains("fastupdate=off"),
"fresh index {name} has options {options:?}"
);
}
// A schema from before the change: the same indexes, made with the
// default fastupdate=on, and a pending list with something in it
for (name, _) in &indexes {
admin
.sql_query::<usize>(
&format!("ALTER INDEX {SCHEMA}.{name} RESET (fastupdate)"),
vec![],
)
.await
.unwrap();
}
for (name, options) in gin_indexes(&admin).await {
assert!(
!options.contains("fastupdate"),
"index {name} still has options {options:?}"
);
}
admin
.sql_query::<usize>(
&format!(
"INSERT INTO {SCHEMA}.s_email (accid, docid, subj, body) \
SELECT 1, n, to_tsvector('simple', 'pending subject ' || n), \
to_tsvector('simple', 'pending body text ' || n) \
FROM generate_series(1, 500) n"
),
vec![],
)
.await
.unwrap();
assert!(
pending_tuples(&admin, "gin_s_email_body").await > 0,
"no pending list to merge"
);
// Startup on the existing schema switches every index over and merges
// what was pending
search.create_indexes().await.unwrap();
for (name, options) in gin_indexes(&admin).await {
assert!(
options.contains("fastupdate=off"),
"existing index {name} has options {options:?} after startup"
);
}
assert_eq!(pending_tuples(&admin, "gin_s_email_body").await, 0);
// And a second startup changes nothing
search.create_indexes().await.unwrap();
for (name, options) in gin_indexes(&admin).await {
assert!(options.contains("fastupdate=off"), "{name}: {options:?}");
}
admin
.sql_query::<usize>(&format!("DROP SCHEMA {SCHEMA} CASCADE"), vec![])
.await
.unwrap();
}
/// The GIN indexes in the test schema with their reloptions.
async fn gin_indexes(admin: &Store) -> Vec<(String, String)> {
admin
.sql_query::<Rows>(
&format!(
"SELECT c.relname::text, COALESCE(array_to_string(c.reloptions, ','), '') \
FROM pg_class c JOIN pg_namespace n ON n.oid = c.relnamespace \
JOIN pg_am a ON a.oid = c.relam \
WHERE n.nspname = '{SCHEMA}' AND c.relkind = 'i' AND a.amname = 'gin' \
ORDER BY 1"
),
vec![],
)
.await
.unwrap()
.rows
.into_iter()
.map(|row| {
let mut values = row.values.into_iter();
(
values.next().unwrap().to_str().into_owned(),
values.next().unwrap().to_str().into_owned(),
)
})
.collect()
}
/// Tuples waiting in a GIN index's pending list (pgstattuple is a contrib
/// extension the test database has).
async fn pending_tuples(admin: &Store, index: &str) -> i64 {
admin
.sql_query::<usize>("CREATE EXTENSION IF NOT EXISTS pgstattuple", vec![])
.await
.unwrap();
admin
.sql_query::<Rows>(
&format!("SELECT pending_tuples FROM pgstatginindex('{SCHEMA}.{index}'::regclass)"),
vec![],
)
.await
.unwrap()
.rows
.into_iter()
.next()
.and_then(|row| row.values.into_iter().next())
.map(|value| value.to_str().parse::<i64>().unwrap())
.unwrap()
}