Author SHA1 Message Date
jcoffey-dev f55087dd9b Settings reload: no DNS at build time, don't refuse over old failures
ci / build (pull_request) Successful in 7m37s
ci / fork-checks (pull_request) Successful in 30s
A 3-node rehearsal found every settings reload refused, cluster-wide,
because one node couldn't resolve the Pyzor server:

- PyzorConfig::parse resolved the host while building the settings and
  made a failed lookup a build error. It now keeps the host and port and
  resolves when a message is checked (an IP address is used as is, a
  name is reused for five minutes, the lookup counts against the Pyzor
  timeout). A failure there is a Pyzor error for that message.
- A milter's hostname was resolved the same way, with a blocking
  to_socket_addrs in async code. An IP address is kept; a name is now
  resolved on each connection.

Other build-time I/O is already non-fatal: directories that can't
connect become unavailable with a warning (DIR-21), and the AI model
locality check only warns.

reload_registry swapped the core only when the whole build was free of
errors, while boot runs with whatever built. One failing object thus
refused every later reload, and the running settings went stale. Now a
reload is refused only for errors in objects that built when the
running settings were built (at boot or by the last applied reload):
applying it would lose those. Objects that already failed then are
missing from the running settings anyway, as at boot, so their errors
are logged and returned as known_errors but don't hold the reload back.
Refusing on new errors keeps a bad edit from taking a working object
out of service; the admin gets the error instead.

ReloadSettings now says "Settings were not reloaded." and names the
object and its error ("Tracer with id ...: Only one console tracer is
allowed"), with a count of any further errors. A refused reload after a
directory change logs its errors too.

system::reload::reload_tests (new): with Pyzor enabled on an
unresolvable host, ReloadSettings succeeds (on main it fails with
"Invalid address: failed to lookup address information"); an IP host
needs no lookup; a new build error refuses the reload, names the object
and leaves the running settings unchanged; the same error, once known
from the running settings' build, no longer blocks; once fixed, a new
error there blocks again. smtp::inbound::milter's session test now
names its milter "localhost", so the connect-time lookup is exercised.
2026-09-24 11:58:52 -07:00
31 changed files with 136 additions and 1762 deletions
+1 -204
View File
@@ -13,7 +13,7 @@ use crate::{
storage::Storage,
telemetry::Telemetry,
},
ipc::{BroadcastEvent, QueueEvent, RegistryChange},
ipc::{QueueEvent, RegistryChange},
network::security::{BlockedIps, IpWithTtl},
};
use ahash::AHashMap;
@@ -232,206 +232,3 @@ fn error_object(error: &Error) -> Option<ObjectId> {
Error::Internal { object_id, .. } => *object_id,
}
}
// inbuxa: upstream applied a registry write to the running settings only on
// an explicit x:Action ReloadSettings (Directory and Authentication aside), so
// a new MtaDeliverySchedule, say, stayed unknown ("Queue strategy not found")
// until someone reloaded. Writes to objects the settings are built from now
// reload them, here and across the cluster, as ReloadSettings does.
/// Coalesces the full reloads that registry writes trigger: a write waits for
/// a reload that started after it was stored, and joins one if it can, so a
/// burst of writes costs a reload or two rather than one each.
#[derive(Default)]
pub struct SettingsReloadGate {
requested: std::sync::atomic::AtomicU64,
state: tokio::sync::Mutex<SettingsReloadState>,
}
#[derive(Default)]
struct SettingsReloadState {
completed: u64,
refused: Option<String>,
}
/// The reload a write to `object` calls for: the object to reload, or None
/// when the running settings don't hold that object (accounts, domains and
/// other data read as needed, stores, which take a restart, and objects with
/// reload actions of their own, such as applications). Blocked IPs have a
/// reload of their own; allowed IPs take the full one.
pub fn write_reload_target(object: ObjectType) -> Option<ObjectType> {
match object {
ObjectType::Certificate => Some(ObjectType::Certificate),
ObjectType::MemoryLookupKey
| ObjectType::MemoryLookupKeyValue
| ObjectType::HttpLookup
| ObjectType::StoreLookup => Some(ObjectType::StoreLookup),
ObjectType::BlockedIp => Some(ObjectType::BlockedIp),
// Allowed IPs are part of the core's security settings
// (Security::parse), which only a full reload rebuilds; the blocked-IP
// reload doesn't touch them
ObjectType::AllowedIp
| ObjectType::AcmeProvider
| ObjectType::AddressBook
| ObjectType::AiModel
| ObjectType::Asn
| ObjectType::Authentication
| ObjectType::Cache
| ObjectType::Calendar
| ObjectType::CalendarAlarm
| ObjectType::CalendarScheduling
| ObjectType::ClusterRole
| ObjectType::DataRetention
| ObjectType::Directory
| ObjectType::DkimReportSettings
| ObjectType::DmarcReportSettings
| ObjectType::DnsResolver
| ObjectType::DsnReportSettings
| ObjectType::Email
| ObjectType::EventTracingLevel
| ObjectType::FileStorage
| ObjectType::Http
| ObjectType::HttpForm
| ObjectType::Imap
| ObjectType::Jmap
| ObjectType::Metrics
| ObjectType::MtaConnectionStrategy
| ObjectType::MtaDeliverySchedule
| ObjectType::MtaExtensions
| ObjectType::MtaHook
| ObjectType::MtaInboundSession
| ObjectType::MtaInboundThrottle
| ObjectType::MtaMilter
| ObjectType::MtaOutboundStrategy
| ObjectType::MtaOutboundThrottle
| ObjectType::MtaQueueQuota
| ObjectType::MtaRoute
| ObjectType::MtaStageAuth
| ObjectType::MtaStageConnect
| ObjectType::MtaStageData
| ObjectType::MtaStageEhlo
| ObjectType::MtaStageMail
| ObjectType::MtaStageRcpt
| ObjectType::MtaSts
| ObjectType::MtaTlsStrategy
| ObjectType::MtaVirtualQueue
| ObjectType::NetworkListener
| ObjectType::OidcProvider
| ObjectType::ReportSettings
| ObjectType::Search
| ObjectType::Security
| ObjectType::SenderAuth
| ObjectType::Sharing
| ObjectType::SieveSystemInterpreter
| ObjectType::SieveSystemScript
| ObjectType::SieveUserInterpreter
| ObjectType::SieveUserScript
| ObjectType::SpamClassifier
| ObjectType::SpamDnsblServer
| ObjectType::SpamDnsblSettings
| ObjectType::SpamFileExtension
| ObjectType::SpamPyzor
| ObjectType::SpamRule
| ObjectType::SpamSettings
| ObjectType::SpamTag
| ObjectType::SpfReportSettings
| ObjectType::SystemSettings
| ObjectType::TaskManager
| ObjectType::TlsReportSettings
| ObjectType::Tracer
| ObjectType::WebDav
| ObjectType::WebHook => Some(object),
_ => None,
}
}
impl Server {
/// Applies a stored registry write to `object` to the running settings,
/// and on success tells the other nodes to do the same. Returns None when
/// the write needs no reload, Some(Ok(())) when it was applied, and
/// Some(Err(reason)) when the reload was refused (the write stays stored;
/// ReloadSettings reports the same errors).
pub async fn reload_after_write(&self, object: ObjectType) -> Option<Result<(), String>> {
let target = write_reload_target(object)?;
let change = RegistryChange::Reload(target);
if matches!(
target,
ObjectType::Certificate | ObjectType::StoreLookup | ObjectType::BlockedIp
) {
// Cheap, and limited to their own objects
let result = self.reload_and_broadcast(change).await;
return Some(result);
}
let gate = &self.inner.data.settings_reload;
let ticket = gate
.requested
.fetch_add(1, std::sync::atomic::Ordering::SeqCst)
+ 1;
let mut state = gate.state.lock().await;
if state.completed >= ticket {
// A reload that started after this write was stored has run
return Some(state.refused.clone().map_or(Ok(()), Err));
}
let covers = gate.requested.load(std::sync::atomic::Ordering::SeqCst);
let result = self.reload_and_broadcast(change).await;
state.completed = covers;
state.refused = result.clone().err();
Some(result)
}
async fn reload_and_broadcast(&self, change: RegistryChange) -> Result<(), String> {
match Box::pin(self.reload_registry(change)).await {
Ok(reload) if !reload.has_errors() => {
reload.log();
self.cluster_broadcast(BroadcastEvent::RegistryChange(change))
.await;
Ok(())
}
Ok(reload) => {
reload.log();
let reason = describe_reload_errors(&reload.errors);
trc::event!(
Registry(trc::RegistryEvent::BuildWarning),
Details = "Settings didn't reload after a registry write",
Reason = reason.clone(),
);
Err(reason)
}
Err(err) => {
let reason = err.to_string();
trc::error!(err.details("Failed to reload settings after a registry write"));
Err(reason)
}
}
}
}
/// inbuxa: a refused reload's errors in a sentence: the first one, naming its
/// object, and how many more there are.
pub fn describe_reload_errors(errors: &[Error]) -> String {
let mut description = match errors.first() {
Some(Error::Build { object_id, message }) => format!("{object_id}: {message}"),
Some(Error::Validation { object_id, errors }) => format!(
"{object_id}: {}",
errors
.iter()
.map(|err| err.to_string())
.collect::<Vec<_>>()
.join("; ")
),
Some(Error::Internal {
object_id: Some(object_id),
error,
}) => format!("{object_id}: {error}"),
Some(Error::Internal { error, .. }) => error.to_string(),
Some(Error::NotFound { object_id }) => format!("{object_id} was not found"),
None => String::new(),
};
let more = errors.len().saturating_sub(1);
if more > 0 {
description.push_str(&format!(" ({more} more in the server log.)"));
}
description
}
-2
View File
@@ -93,7 +93,6 @@ impl Data {
registry_id_gen: id_generator.clone(),
span_id_gen: id_generator,
queue_status: true.into(),
settings_reload: Default::default(),
applications,
logos: Default::default(),
smtp_connectors: TlsConnectors::try_new().failed("Failed to build TLS connectors"),
@@ -236,7 +235,6 @@ impl Default for Data {
span_id_gen: Default::default(),
registry_id_gen: Default::default(),
queue_status: true.into(),
settings_reload: Default::default(),
applications: WebApplications::new(),
logos: Default::default(),
smtp_connectors: TlsConnectors::try_new().unwrap(),
+2 -17
View File
@@ -345,13 +345,8 @@ pub struct TaskLocks {
}
impl TaskLocks {
/// How long a task lock lasts, in seconds, unless it is released first
/// or renewed. inbuxa: upstream held a lock for an hour, so a killed
/// node's tasks waited that long; the lock is now a five-minute lease
/// that the task manager renews every third of it while the task runs
/// (renew_task_locks), so a dead node's tasks run elsewhere within
/// minutes.
pub const DEFAULT_EXPIRY: u64 = 5 * 60;
/// How long a task lock lasts, in seconds, unless it is released first.
pub const DEFAULT_EXPIRY: u64 = 60 * 60;
pub fn is_stopping(&self) -> bool {
self.stopping.load(Ordering::Acquire)
@@ -375,16 +370,6 @@ impl TaskLocks {
self.held.lock().len()
}
/// inbuxa: the tasks this node holds, to renew their locks.
pub fn held_ids(&self) -> Vec<u64> {
self.held.lock().iter().copied().collect()
}
/// inbuxa: whether this node holds (and is running) the task.
pub fn is_held(&self, id: u64) -> bool {
self.held.lock().contains(&id)
}
pub fn expiry(&self) -> u64 {
self.expiry.load(Ordering::Relaxed)
}
-2
View File
@@ -161,8 +161,6 @@ pub struct Data {
pub span_id_gen: SnowflakeIdGenerator,
pub registry_id_gen: SnowflakeIdGenerator,
pub queue_status: AtomicBool,
// inbuxa: coalesces the settings reloads registry writes trigger
pub settings_reload: cache::reload::SettingsReloadGate,
pub applications: WebApplications,
pub logos: Mutex<AHashMap<Box<str>, LogoCache>>,
-20
View File
@@ -2,8 +2,6 @@
* SPDX-FileCopyrightText: 2020 Stalwart Labs LLC <[email protected]>
*
* SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-SEL
*
* Modified by Coffey Labs in 2026 for INBUXA.
*/
use super::ahash_is_empty;
@@ -73,23 +71,6 @@ pub struct SetResponse<T: JmapObject> {
#[serde(rename = "notDestroyed")]
#[serde(skip_serializing_if = "VecMap::is_empty")]
pub not_destroyed: VecMap<MaybeInvalid<Id>, SetError<T::Property>>,
// inbuxa: on a registry write that changes the running settings, whether
// the server applied it
#[serde(rename = "x:settingsReload")]
#[serde(skip_serializing_if = "Option::is_none")]
pub settings_reload: Option<SettingsReload>,
}
/// inbuxa: the settings reload that followed a registry write.
#[derive(Debug, Clone, serde::Serialize)]
pub struct SettingsReload {
/// The running settings (here and, through the cluster, on every node)
/// include the write.
pub applied: bool,
/// Why they don't, when they don't.
#[serde(skip_serializing_if = "Option::is_none")]
pub description: Option<String>,
}
impl<'de, T: JmapObject> DeserializeArguments<'de> for SetRequest<'de, T> {
@@ -218,7 +199,6 @@ impl<T: JmapObject> SetResponse<T> {
not_created: VecMap::new(),
not_updated: VecMap::new(),
not_destroyed: VecMap::new(),
settings_reload: None,
})
} else {
Err(trc::JmapEvent::RequestTooLarge.into_err())
+24 -4
View File
@@ -580,9 +580,29 @@ async fn dmarc_troubleshoot(
/// settings weren't applied; upstream passed on the first error's bare message
/// ("Invalid address: ..."), which read like a problem with the request.
fn reload_refused(errors: Vec<registry::types::error::Error>) -> SetError<Property> {
let description = format!(
"Settings were not reloaded. {}",
common::cache::reload::describe_reload_errors(&errors)
);
use registry::types::error::Error;
let more = errors.len().saturating_sub(1);
let mut description = match errors.first() {
Some(Error::Build { object_id, message }) => format!("{object_id}: {message}"),
Some(Error::Validation { object_id, errors }) => format!(
"{object_id}: {}",
errors
.iter()
.map(|err| err.to_string())
.collect::<Vec<_>>()
.join("; ")
),
Some(Error::Internal {
object_id: Some(object_id),
error,
}) => format!("{object_id}: {error}"),
Some(Error::Internal { error, .. }) => error.to_string(),
Some(Error::NotFound { object_id }) => format!("{object_id} was not found"),
None => String::new(),
};
description.insert_str(0, "Settings were not reloaded. ");
if more > 0 {
description.push_str(&format!(" ({more} more in the server log.)"));
}
map_bootstrap_error(errors).with_description(description)
}
+25 -19
View File
@@ -38,7 +38,7 @@ use directory::core::secret::{hash_secret, is_password_hash};
use http_proto::HttpSessionData;
use jmap_proto::{
error::set::{SetError, SetErrorType},
method::set::{SetRequest, SetResponse, SettingsReload},
method::set::{SetRequest, SetResponse},
object::registry::Registry,
references::resolve::ResolveCreatedReference,
request::{IntoValid, MaybeInvalid},
@@ -931,28 +931,34 @@ impl RegistrySet for Server {
}
};
// inbuxa: a write to an object the running settings are built from
// applies at once, here and on every node (DIR-17 did this for
// directories and the server default; now it covers every such object)
let mut result = result;
if let Ok(response) = &mut result
// inbuxa: DIR-17: a directory or the server default applies on the
// next request, here and on every node
if matches!(
object_type,
ObjectType::Directory | ObjectType::Authentication
) && let Ok(response) = &result
&& (!response.created.is_empty()
|| !response.updated.is_empty()
|| !response.destroyed.is_empty())
&& let Some(reload) = self.reload_after_write(object_type).await
{
response.settings_reload = Some(match reload {
Ok(()) => SettingsReload {
applied: true,
description: None,
},
Err(reason) => SettingsReload {
applied: false,
description: Some(format!(
"Saved, but the running settings were not reloaded. {reason}"
)),
},
});
let change = common::ipc::RegistryChange::Reload(ObjectType::Directory);
match Box::pin(self.reload_registry(change)).await {
Ok(reload) if !reload.has_errors() => {
self.cluster_broadcast(common::ipc::BroadcastEvent::RegistryChange(change))
.await;
}
Ok(reload) => {
// inbuxa: name what stopped it
reload.log();
trc::event!(
Registry(trc::RegistryEvent::BuildWarning),
Details = "Settings didn't reload after a directory change",
)
}
Err(err) => {
trc::error!(err.details("Failed to reload directories"));
}
}
}
result
}
-39
View File
@@ -82,42 +82,3 @@ pub async fn release_task_locks(server: &Server) -> usize {
}
ids.len()
}
/// inbuxa: renews the lease on every task this node is running, so it stays
/// claimed for as long as it runs while a node that dies loses its claims
/// within one lock lifetime. Returns how many leases were renewed and how
/// many were found lost (expired, perhaps taken by another node).
pub async fn renew_task_locks(server: &Server) -> (usize, usize) {
let locks = &server.inner.ipc.task_locks;
let expiry = locks.expiry();
let (mut renewed, mut lost) = (0, 0);
for id in locks.held_ids() {
match server
.in_memory_store()
.renew_lock(KV_LOCK_TASK, &id.to_be_bytes(), expiry)
.await
{
Ok(true) => renewed += 1,
Ok(false) => {
// Still held here as far as this node knows; the task
// finishes and its lock is removed as usual
if locks.is_held(id) {
lost += 1;
trc::event!(
TaskManager(TaskManagerEvent::TaskLocked),
Id = id,
Details = "Task lock expired while the task was running",
);
}
}
Err(err) => {
trc::error!(
err.details("Failed to renew task lock")
.ctx(trc::Key::Id, id)
.caused_by(trc::location!())
);
}
}
}
(renewed, lost)
}
+21 -75
View File
@@ -13,7 +13,7 @@ use crate::task_manager::dkim::DkimManagementTask;
use crate::task_manager::dns::DnsManagementTask;
use crate::task_manager::imip::SendImipTask;
use crate::task_manager::index::SearchIndexTask;
use crate::task_manager::lock::{TaskLockManager, renew_task_locks};
use crate::task_manager::lock::TaskLockManager;
use crate::task_manager::maintenance::MaintenanceTask;
use crate::task_manager::merge_threads::MergeThreadsTask;
use crate::task_manager::report::{self, SubmitReportTask};
@@ -24,7 +24,6 @@ use crate::task_manager::{
TaskJob, TaskManagerIpc, TaskResult,
};
use common::BuildServer;
use common::config::network::ClusterRoles;
use common::config::server::{DEFAULT_TLS_TIMEOUT, ServerProtocol};
use common::network::limiter::ConcurrencyLimiter;
use common::network::{ServerInstance, TcpAcceptor};
@@ -59,12 +58,10 @@ pub fn spawn_task_manager(inner: Arc<Inner>) {
let server = inner.build_server();
let roles = &server.core.network.roles;
// inbuxa: outbound_mta too, which now governs report tasks
if !roles.account_maintenance
&& !roles.store_maintenance
&& !roles.search_indexing
&& !roles.spam_training
&& !roles.outbound_mta
&& !roles.task_manager
{
return;
@@ -75,28 +72,6 @@ pub fn spawn_task_manager(inner: Arc<Inner>) {
trc::event!(TaskManager(TaskManagerEvent::ManagerStarted));
// inbuxa: keep the leases of running tasks alive, every third of a lock
// lifetime, until the node stops
{
let inner = inner.clone();
tokio::spawn(async move {
let mut renewed_at = Instant::now();
loop {
tokio::time::sleep(Duration::from_secs(1)).await;
let locks = &inner.ipc.task_locks;
if locks.is_stopping() {
break;
}
if renewed_at.elapsed() >= Duration::from_secs((locks.expiry() / 3).max(1)) {
renewed_at = Instant::now();
if locks.held() > 0 {
renew_task_locks(&inner.build_server()).await;
}
}
}
});
}
// Create dummy server instance for alarms
let server_instance = Arc::new(ServerInstance {
id: "_local".to_string(),
@@ -314,13 +289,26 @@ impl TaskQueueManager for Server {
.caused_by(trc::location!())
.ctx(trc::Key::Value, value)
})?;
// inbuxa: running here under a lease this node
// renews; don't hand it to a worker again
if task_locks.is_held(task_id) {
return Ok(true);
}
let enabled = task_enabled(roles, task_type);
let enabled = match task_type {
TaskType::IndexDocument
| TaskType::UnindexDocument
| TaskType::IndexTrace => roles.search_indexing,
TaskType::AccountMaintenance
| TaskType::TenantMaintenance
| TaskType::DestroyAccount => roles.account_maintenance,
TaskType::StoreMaintenance => roles.store_maintenance,
TaskType::SpamFilterMaintenance => roles.spam_training,
TaskType::CalendarAlarmEmail
| TaskType::CalendarAlarmNotification
| TaskType::CalendarItipMessage
| TaskType::MergeThreads
| TaskType::DmarcReport
| TaskType::TlsReport
| TaskType::RestoreArchivedItem
| TaskType::AcmeRenewal
| TaskType::DkimManagement
| TaskType::DnsManagement => true,
};
if !enabled {
trc::event!(
@@ -449,48 +437,6 @@ impl TaskQueueManager for Server {
}
}
/// inbuxa: whether this node's cluster role lets it run a task type. Upstream
/// checked the dedicated roles (search indexing, account and store
/// maintenance, spam training) and let every node with a task manager run
/// the rest, whatever its taskQueueProcessing setting. Every task type now
/// answers to one ClusterTaskType:
///
/// - IndexDocument, UnindexDocument, IndexTrace: searchIndexing
/// - AccountMaintenance, TenantMaintenance, DestroyAccount: accountMaintenance
/// - StoreMaintenance: storeMaintenance
/// - SpamFilterMaintenance: spamClassifierTraining
/// - DmarcReport, TlsReport: outboundMta. They build and send reports to
/// other domains (TLS reports can go straight to an HTTPS endpoint), which
/// is the outbound MTA's business.
/// - CalendarAlarmEmail, CalendarAlarmNotification, CalendarItipMessage,
/// MergeThreads, RestoreArchivedItem, AcmeRenewal, DkimManagement,
/// DnsManagement: taskQueueProcessing, the role for queue tasks with no
/// role of their own.
///
/// A node that may not run a task leaves it unclaimed, so a node that may
/// picks it up.
pub fn task_enabled(roles: &ClusterRoles, task_type: TaskType) -> bool {
match task_type {
TaskType::IndexDocument | TaskType::UnindexDocument | TaskType::IndexTrace => {
roles.search_indexing
}
TaskType::AccountMaintenance | TaskType::TenantMaintenance | TaskType::DestroyAccount => {
roles.account_maintenance
}
TaskType::StoreMaintenance => roles.store_maintenance,
TaskType::SpamFilterMaintenance => roles.spam_training,
TaskType::DmarcReport | TaskType::TlsReport => roles.outbound_mta,
TaskType::CalendarAlarmEmail
| TaskType::CalendarAlarmNotification
| TaskType::CalendarItipMessage
| TaskType::MergeThreads
| TaskType::RestoreArchivedItem
| TaskType::AcmeRenewal
| TaskType::DkimManagement
| TaskType::DnsManagement => roles.task_manager,
}
}
async fn run_task(
server: &Server,
task: &Task,
+3 -5
View File
@@ -2,8 +2,6 @@
* SPDX-FileCopyrightText: 2020 Stalwart Labs LLC <[email protected]>
*
* SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-SEL
*
* Modified by Coffey Labs in 2026 for INBUXA.
*/
use std::ops::Range;
@@ -18,7 +16,7 @@ impl MysqlStore {
key: &[u8],
range: Range<usize>,
) -> trc::Result<Option<Vec<u8>>> {
let mut conn = self.conn().await?;
let mut conn = self.conn_pool.get_conn().await.map_err(into_error)?;
let s = conn
.prep("SELECT v FROM t WHERE k = ?")
.await
@@ -41,7 +39,7 @@ impl MysqlStore {
}
pub(crate) async fn put_blob(&self, key: &[u8], data: &[u8]) -> trc::Result<()> {
let mut conn = self.conn().await?;
let mut conn = self.conn_pool.get_conn().await.map_err(into_error)?;
let s = conn
.prep("INSERT INTO t (k, v) VALUES (?, ?) ON DUPLICATE KEY UPDATE v = VALUES(v)")
.await
@@ -53,7 +51,7 @@ impl MysqlStore {
}
pub(crate) async fn delete_blob(&self, key: &[u8]) -> trc::Result<bool> {
let mut conn = self.conn().await?;
let mut conn = self.conn_pool.get_conn().await.map_err(into_error)?;
let s = conn
.prep("DELETE FROM t WHERE k = ?")
.await
+1 -3
View File
@@ -2,8 +2,6 @@
* SPDX-FileCopyrightText: 2020 Stalwart Labs LLC <[email protected]>
*
* SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-SEL
*
* Modified by Coffey Labs in 2026 for INBUXA.
*/
use mysql_async::{Params, Row, prelude::Queryable};
@@ -18,7 +16,7 @@ impl MysqlStore {
query: &str,
params: &[Value<'_>],
) -> trc::Result<T> {
let mut conn = self.conn().await?;
let mut conn = self.conn_pool.get_conn().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());
+2 -5
View File
@@ -32,9 +32,6 @@ impl MysqlStore {
.max_allowed_packet(config.max_allowed_packet.map(|v| v as usize))
.wait_timeout(config.timeout.map(|t| t.as_secs() as usize))
.client_found_rows(true)
// inbuxa: notice a server that went away without closing the
// connection in minutes, not the system default of two hours
.tcp_keepalive(Some(super::POOL_KEEPALIVE_IDLE))
.tcp_port(config.port as u16);
if config.use_tls {
@@ -98,7 +95,7 @@ impl MysqlStore {
}
pub(crate) async fn create_storage_tables(&self) -> trc::Result<()> {
let mut conn = self.conn().await?;
let mut conn = self.conn_pool.get_conn().await.map_err(into_error)?;
for table in [
SUBSPACE_ACL,
@@ -172,7 +169,7 @@ impl MysqlStore {
}
pub(crate) async fn create_search_tables(&self) -> trc::Result<()> {
let mut conn = self.conn().await?;
let mut conn = self.conn_pool.get_conn().await.map_err(into_error)?;
create_search_tables::<EmailSearchField>(&mut conn).await?;
create_search_tables::<CalendarSearchField>(&mut conn).await?;
-27
View File
@@ -27,33 +27,6 @@ pub struct MysqlStore {
pub(crate) conn_pool: Pool,
}
/// inbuxa: how long a request waits for a pooled connection (including
/// opening one). mysql_async's pool has no wait timeout, so upstream waited
/// forever when the server stopped answering.
pub(crate) const POOL_WAIT_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(30);
/// inbuxa: idle time before TCP keepalive probes start.
pub(crate) const POOL_KEEPALIVE_IDLE: std::time::Duration = std::time::Duration::from_secs(60);
impl MysqlStore {
/// inbuxa: a pooled connection, or an error once POOL_WAIT_TIMEOUT has
/// passed without one.
pub(crate) async fn conn(&self) -> trc::Result<mysql_async::Conn> {
pool_conn(&self.conn_pool, POOL_WAIT_TIMEOUT).await
}
}
pub(crate) async fn pool_conn(
pool: &Pool,
wait: std::time::Duration,
) -> trc::Result<mysql_async::Conn> {
match tokio::time::timeout(wait, pool.get_conn()).await {
Ok(result) => result.map_err(into_error),
Err(_) => Err(trc::StoreEvent::MysqlError
.reason("Timed out waiting for a database connection")
.details(format!("No connection within {} s", wait.as_secs()))),
}
}
#[inline(always)]
pub(crate) fn into_error(err: impl Display) -> trc::Error {
trc::StoreEvent::MysqlError.reason(err)
+4 -6
View File
@@ -2,8 +2,6 @@
* SPDX-FileCopyrightText: 2020 Stalwart Labs LLC <[email protected]>
*
* SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-SEL
*
* Modified by Coffey Labs in 2026 for INBUXA.
*/
use super::{MysqlStore, into_error, is_timeout_error};
@@ -16,7 +14,7 @@ impl MysqlStore {
where
U: Deserialize + 'static,
{
let mut conn = self.conn().await?;
let mut conn = self.conn_pool.get_conn().await.map_err(into_error)?;
let s = conn
.prep(format!(
"SELECT v FROM {} WHERE k = ?",
@@ -38,7 +36,7 @@ impl MysqlStore {
}
pub(crate) async fn key_exists(&self, key: impl Key) -> trc::Result<bool> {
let mut conn = self.conn().await?;
let mut conn = self.conn_pool.get_conn().await.map_err(into_error)?;
let s = conn
.prep(format!(
"SELECT 1 FROM {} WHERE k = ?",
@@ -58,7 +56,7 @@ impl MysqlStore {
params: IterateParams<T>,
mut cb: impl for<'x> FnMut(&'x [u8], &'x [u8]) -> trc::Result<bool> + Sync + Send,
) -> trc::Result<()> {
let mut conn = self.conn().await?;
let mut conn = self.conn_pool.get_conn().await.map_err(into_error)?;
let table = char::from(params.begin.subspace());
let begin = params.begin.serialize(0);
let end = params.end.serialize(0);
@@ -157,7 +155,7 @@ impl MysqlStore {
let key = key.into();
let table = char::from(key.subspace());
let key = key.serialize(0);
let mut conn = self.conn().await?;
let mut conn = self.conn_pool.get_conn().await.map_err(into_error)?;
let s = conn
.prep(format!("SELECT v FROM {table} WHERE k = ?"))
.await
+16 -79
View File
@@ -2,8 +2,6 @@
* SPDX-FileCopyrightText: 2020 Stalwart Labs LLC <[email protected]>
*
* SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-SEL
*
* Modified by Coffey Labs in 2026 for INBUXA.
*/
use crate::{
@@ -21,12 +19,12 @@ use crate::{
write::SearchIndex,
};
use mysql_async::{IsolationLevel, TxOpts, Value, prelude::Queryable};
use nlp::{language::Language, tokenizers::word::WordTokenizer};
use nlp::tokenizers::word::WordTokenizer;
use std::fmt::Write;
impl MysqlStore {
pub async fn index(&self, documents: Vec<IndexDocument>) -> trc::Result<()> {
let mut conn = self.conn().await?;
let mut conn = self.conn_pool.get_conn().await.map_err(into_error)?;
let mut tx_opts = TxOpts::default();
tx_opts
.with_consistent_snapshot(false)
@@ -96,7 +94,7 @@ impl MysqlStore {
build_sort(&mut query, sort);
}
let mut conn = self.conn().await?;
let mut conn = self.conn_pool.get_conn().await.map_err(into_error)?;
let s = conn.prep(query).await.map_err(into_error)?;
conn.exec::<i64, _, _>(s, params)
@@ -110,7 +108,7 @@ impl MysqlStore {
let mut query = format!("DELETE FROM {table} ");
let params = build_filter(&mut query, &filter.filters);
let mut conn = self.conn().await?;
let mut conn = self.conn_pool.get_conn().await.map_err(into_error)?;
let s = conn.prep(&query).await.map_err(into_error)?;
match conn.exec_drop(s, params.clone()).await {
@@ -148,20 +146,6 @@ impl MysqlStore {
}
}
// inbuxa: InnoDB's default full-text stopword list
// (INFORMATION_SCHEMA.INNODB_FT_DEFAULT_STOPWORD) and innodb_ft_min_token_size
// default; words outside these are not in a FULLTEXT index.
const FT_STOPWORDS: &[&str] = &[
"a", "about", "an", "are", "as", "at", "be", "by", "com", "de", "en", "for", "from", "how",
"i", "in", "is", "it", "la", "of", "on", "or", "that", "the", "this", "to", "was", "what",
"when", "where", "who", "will", "with", "und", "www",
];
const FT_MIN_TOKEN_SIZE: usize = 3;
fn is_ft_indexed(word: &str) -> bool {
word.chars().count() >= FT_MIN_TOKEN_SIZE && !FT_STOPWORDS.contains(&word)
}
fn build_filter(query: &mut String, filters: &[SearchFilter]) -> Vec<Value> {
if filters.is_empty() {
return Vec::new();
@@ -187,77 +171,30 @@ fn build_filter(query: &mut String, filters: &[SearchFilter]) -> Vec<Value> {
if field.is_text() && matches!(op, SearchOperator::Equal | SearchOperator::Contains)
{
let (value, mode, unindexed) = match (value, op) {
(SearchValue::Text { value, .. }, SearchOperator::Equal) => (
Value::Bytes(format!("{value:?}").into_bytes()),
"BOOLEAN",
Vec::new(),
),
(SearchValue::Text { value, language }, ..) => {
let (value, mode) = match (value, op) {
(SearchValue::Text { value, .. }, SearchOperator::Equal) => {
(Value::Bytes(format!("{value:?}").into_bytes()), "BOOLEAN")
}
(SearchValue::Text { value, .. }, ..) => {
let mut text_query = String::with_capacity(value.len() + 1);
let mut unindexed = Vec::new();
for item in WordTokenizer::new(value, MAX_TOKEN_LENGTH) {
// inbuxa: InnoDB never indexes stopwords ("com",
// "de", "www", ...) or words under
// innodb_ft_min_token_size, and a required
// (+word) term it has not indexed matches no row,
// so "example.com" or "[email protected]" found
// nothing. Such words are matched with a
// word-boundary REGEXP instead.
if is_ft_indexed(&item.word) {
if !text_query.is_empty() {
text_query.push(' ');
}
text_query.push('+');
text_query.push_str(&item.word);
} else {
unindexed.push(item.word);
if !text_query.is_empty() {
text_query.push(' ');
}
text_query.push('+');
text_query.push_str(&item.word);
}
// For language text (bodies, subjects) the unindexed
// words are noise words and only checked when nothing
// else is left to match; keyword text (addresses,
// contact fields) checks every word, as the other
// backends do.
if !text_query.is_empty() && !matches!(language, Language::None) {
unindexed.clear();
}
(Value::Bytes(text_query.into_bytes()), "BOOLEAN", unindexed)
(Value::Bytes(text_query.into_bytes()), "BOOLEAN")
}
_ => {
debug_assert!(false, "Invalid search value for text field");
continue;
}
};
if unindexed.is_empty() {
let _ =
write!(query, "MATCH({}) AGAINST(? IN {mode} MODE)", field.column());
values.push(value);
} else {
query.push('(');
let is_empty = matches!(&value, Value::Bytes(v) if v.is_empty());
if !is_empty {
let _ = write!(
query,
"MATCH({}) AGAINST(? IN {mode} MODE) AND ",
field.column()
);
values.push(value);
}
for (i, word) in unindexed.iter().enumerate() {
if i > 0 {
query.push_str(" AND ");
}
let _ = write!(query, "{} REGEXP ?", field.column());
values.push(Value::Bytes(
format!("(^|[^[:alnum:]]){word}([^[:alnum:]]|$)").into_bytes(),
));
}
query.push(')');
}
let _ = write!(query, "MATCH({}) AGAINST(? IN {mode} MODE)", field.column());
values.push(value);
} else if let SearchValue::KeyValues(kv) = value {
let (key, value) = kv.iter().next().unwrap();
+3 -5
View File
@@ -2,8 +2,6 @@
* SPDX-FileCopyrightText: 2020 Stalwart Labs LLC <[email protected]>
*
* SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-SEL
*
* Modified by Coffey Labs in 2026 for INBUXA.
*/
use super::{DELETE_CHUNK_SIZE, MIN_DELETE_CHUNK_SIZE, MysqlStore, into_error, is_timeout_error};
@@ -31,7 +29,7 @@ impl MysqlStore {
pub(crate) async fn write(&self, mut batch: Batch<'_>) -> trc::Result<AssignedIds> {
let start = Instant::now();
let mut retry_count = 0;
let mut conn = self.conn().await?;
let mut conn = self.conn_pool.get_conn().await.map_err(into_error)?;
loop {
let err = match self.write_trx(&mut conn, &mut batch).await {
@@ -384,7 +382,7 @@ impl MysqlStore {
}
pub(crate) async fn purge_store(&self) -> trc::Result<()> {
let mut conn = self.conn().await?;
let mut conn = self.conn_pool.get_conn().await.map_err(into_error)?;
for subspace in [SUBSPACE_QUOTA, SUBSPACE_COUNTER, SUBSPACE_IN_MEMORY_COUNTER] {
purge_table(&mut conn, char::from(subspace)).await?;
}
@@ -393,7 +391,7 @@ impl MysqlStore {
}
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_pool.get_conn().await.map_err(into_error)?;
let table = char::from(from.subspace());
let mut from = from.serialize(0);
let to = to.serialize(0);
+5 -111
View File
@@ -22,34 +22,11 @@ use crate::{
use ::registry::schema::{enums::PostgreSqlRecyclingMethod, structs};
use ahash::AHashSet;
use deadpool_postgres::{
Config, ManagerConfig, Object, Pool, PoolConfig, RecyclingMethod, Runtime, Timeouts,
Config, ManagerConfig, Object, Pool, PoolConfig, RecyclingMethod, Runtime,
};
use std::time::Duration;
use tokio_postgres::NoTls;
use utils::tls::rustls_client_config;
/// inbuxa: how long a request waits for a pooled connection.
pub(crate) const POOL_WAIT_TIMEOUT: Duration = Duration::from_secs(30);
/// inbuxa: how long opening a connection may take when the store sets no
/// timeout of its own.
pub(crate) const POOL_CREATE_TIMEOUT: Duration = Duration::from_secs(15);
/// inbuxa: how long checking a pooled connection before reuse may take.
pub(crate) const POOL_RECYCLE_TIMEOUT: Duration = Duration::from_secs(10);
/// inbuxa: idle time before TCP keepalive probes start.
pub(crate) const POOL_KEEPALIVE_IDLE: Duration = Duration::from_secs(60);
/// inbuxa: the pool's timeouts. Opening a connection is bounded by the
/// store's own timeout when it has one; waiting for one covers at least that
/// long, so a slow connect isn't cut short by the wait.
pub(crate) fn pool_timeouts(connect_timeout: Option<Duration>) -> Timeouts {
let create = connect_timeout.unwrap_or(POOL_CREATE_TIMEOUT);
Timeouts {
wait: POOL_WAIT_TIMEOUT.max(create).into(),
create: create.into(),
recycle: POOL_RECYCLE_TIMEOUT.into(),
}
}
impl PostgresStore {
pub async fn open(config: structs::PostgreSqlStore) -> Result<Store, String> {
// inbuxa: ST-15: where the primary is, to tell a replica from it
@@ -69,20 +46,9 @@ impl PostgresStore {
PostgreSqlRecyclingMethod::Clean => RecyclingMethod::Clean,
},
});
// inbuxa: upstream set no pool timeouts, so a request waited for a
// free connection, or for one to be made or recycled, for as long as
// it took: forever when the server stopped answering. A worker now
// gets an error instead and the task or request is retried.
let mut pool = config
.pool_max_connections
.map(|max_conn| PoolConfig::new(max_conn as usize))
.unwrap_or_default();
pool.timeouts = pool_timeouts(cfg.connect_timeout);
cfg.pool = pool.into();
// Notice a server that went away without closing the connection in
// minutes rather than the system default of two hours
cfg.keepalives = true.into();
cfg.keepalives_idle = POOL_KEEPALIVE_IDLE.into();
if let Some(max_conn) = config.pool_max_connections {
cfg.pool = PoolConfig::new(max_conn as usize).into();
}
let primary_pool = if config.use_tls {
cfg.create_pool(
@@ -265,21 +231,12 @@ async fn create_search_tables<T: SearchableField + PsqlSearchField + 'static>(
for field in T::all_fields() {
if field.is_text() || field.is_json() {
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!(
"CREATE INDEX IF NOT EXISTS {index_name} ON {table_name} USING GIN({column_name}) WITH (fastupdate = off)",
"CREATE INDEX IF NOT EXISTS gin_{table_name}_{column_name} ON {table_name} USING GIN({column_name})",
);
conn.execute(&create_index_query, &[])
.await
.map_err(into_error)?;
// Indexes made before this change keep fastupdate=on
disable_gin_fastupdate(conn, &index_name).await;
}
if field.is_indexed() {
@@ -296,69 +253,6 @@ async fn create_search_tables<T: SearchableField + PsqlSearchField + 'static>(
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> {
let mut ts_configs = AHashSet::from_iter([PG_FALLBACK_LANG, PG_UNSTEMMED_LANG]);
+10 -79
View File
@@ -2,17 +2,12 @@
* SPDX-FileCopyrightText: 2020 Stalwart Labs LLC <[email protected]>
*
* SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-SEL
*
* Modified by Coffey Labs in 2026 for INBUXA.
*/
use crate::{
backend::{
MAX_TOKEN_LENGTH,
postgres::{
DELETE_CHUNK_SIZE, MIN_DELETE_CHUNK_SIZE, PostgresStore, PsqlSearchField, into_error,
into_pool_error, is_timeout_error,
},
backend::postgres::{
DELETE_CHUNK_SIZE, MIN_DELETE_CHUNK_SIZE, PostgresStore, PsqlSearchField, into_error,
into_pool_error, is_timeout_error,
},
search::{
IndexDocument, SearchComparator, SearchDocumentId, SearchFilter, SearchOperator,
@@ -20,7 +15,7 @@ use crate::{
},
write::SearchIndex,
};
use nlp::{language::Language, tokenizers::space::SpaceTokenizer};
use nlp::language::Language;
use std::fmt::Write;
use tokio_postgres::{
IsolationLevel,
@@ -48,19 +43,6 @@ impl PostgresStore {
let primary_keys = index.primary_keys();
let all_fields = index.all_fields();
let fields = document.fields;
// inbuxa: keyword text (addresses, contact fields, ...) is split into
// words before it reaches the text parser, see keyword_terms().
let keywords = primary_keys
.iter()
.chain(all_fields)
.map(|field| match fields.get(field) {
Some(SearchValue::Text {
value,
language: Language::None,
}) if field.is_text() => Some(keyword_terms(value)),
_ => None,
})
.collect::<Vec<_>>();
let mut values = Vec::with_capacity(fields.len() + 2);
let mut query = format!("INSERT INTO {} (", index.psql_table());
@@ -92,20 +74,7 @@ impl PostgresStore {
(0, PG_UNSTEMMED_LANG)
};
if let Some(keywords) = &keywords[i] {
let _ = write!(&mut query, "to_tsvector('{language}',{value_ref})");
values.push(keywords as &(dyn ToSql + Sync));
if field.sort_column().is_some() {
let value_ref = format!("${}", values.len() + 1);
if text_len > 255 {
let _ = write!(&mut query, ",left({value_ref},255)");
} else {
let _ = write!(&mut query, ",{value_ref}");
}
values.push(value as &(dyn ToSql + Sync));
}
continue;
} else if field.is_text() {
if field.is_text() {
let _ = write!(&mut query, "to_tsvector('{language}',{value_ref})");
} else if text_len > 512 {
query.push_str("left(");
@@ -165,7 +134,6 @@ impl PostgresStore {
) -> trc::Result<Vec<R>> {
let mut query = format!("SELECT {} FROM {}", R::field().column(), index.psql_table());
let params = self.build_filter(&mut query, filters);
let params = params.iter().map(SqlParam::as_sql).collect::<Vec<_>>();
if !sort.is_empty() {
build_sort(&mut query, sort);
}
@@ -187,7 +155,6 @@ impl PostgresStore {
let table = filter.index.psql_table();
let mut where_clause = String::new();
let params = self.build_filter(&mut where_clause, &filter.filters);
let params = params.iter().map(SqlParam::as_sql).collect::<Vec<_>>();
let conn = self.conn_pool.get().await.map_err(into_pool_error)?;
let s = conn
.prepare_cached(&format!("DELETE FROM {table}{where_clause}"))
@@ -229,7 +196,7 @@ impl PostgresStore {
&self,
query: &mut String,
filters: &'x [SearchFilter],
) -> Vec<SqlParam<'x>> {
) -> Vec<&'x (dyn ToSql + Sync)> {
if filters.is_empty() {
return Vec::new();
}
@@ -270,10 +237,6 @@ impl PostgresStore {
if matches!(language, Language::None) {
let _ = write!(query, "@@ {method}('{config}', ${value_pos})");
if let SearchValue::Text { value, .. } = value {
values.push(SqlParam::Owned(keyword_terms(value)));
continue;
}
} else {
let _ = write!(query, "@@ ({method}('{config}', ${value_pos})");
for fallback in [PG_FALLBACK_LANG, PG_UNSTEMMED_LANG] {
@@ -284,18 +247,18 @@ impl PostgresStore {
}
query.push(')');
}
values.push(SqlParam::Ref(value));
values.push(value as &(dyn ToSql + Sync));
} else if let SearchValue::KeyValues(kv) = value {
query.push_str(field.column());
query.push(' ');
let (key, value) = kv.iter().next().unwrap();
values.push(SqlParam::Ref(key));
values.push(key as &(dyn ToSql + Sync));
if !value.is_empty() {
let _ = write!(query, "->> ${value_pos} ");
op.write_pqsql(query, values.len() + 1);
values.push(SqlParam::Ref(value));
values.push(value as &(dyn ToSql + Sync));
} else {
let _ = write!(query, " ? ${value_pos}");
}
@@ -304,7 +267,7 @@ impl PostgresStore {
query.push(' ');
op.write_pqsql(query, value_pos);
values.push(SqlParam::Ref(value));
values.push(value as &(dyn ToSql + Sync));
}
}
SearchFilter::And | SearchFilter::Or => {
@@ -358,38 +321,6 @@ impl PostgresStore {
}
}
// inbuxa: PostgreSQL's text parser keeps "[email protected]" (and host names,
// URLs, file paths, ...) as a single token, so a search for "user" or
// "example.com" never matched an address. Keyword text is split into words the
// same way the built-in index splits it (SpaceTokenizer: lowercase runs of
// alphanumerics) on both the indexing and the query side, so a full address,
// its local part, its domain and the display-name words all match, as they do
// on the other backends.
pub(crate) fn keyword_terms(value: &str) -> String {
let mut terms = String::with_capacity(value.len());
for token in SpaceTokenizer::new(value, MAX_TOKEN_LENGTH) {
if !terms.is_empty() {
terms.push(' ');
}
terms.push_str(&token);
}
terms
}
pub(super) enum SqlParam<'x> {
Ref(&'x (dyn ToSql + Sync)),
Owned(String),
}
impl SqlParam<'_> {
fn as_sql(&self) -> &(dyn ToSql + Sync) {
match self {
SqlParam::Ref(value) => *value,
SqlParam::Owned(value) => value,
}
}
}
fn build_sort(query: &mut String, sort: &[SearchComparator]) {
query.push_str(" ORDER BY ");
for (i, comparator) in sort.iter().enumerate() {
-42
View File
@@ -2,8 +2,6 @@
* SPDX-FileCopyrightText: 2020 Stalwart Labs LLC <[email protected]>
*
* SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-SEL
*
* Modified by Coffey Labs in 2026 for INBUXA.
*/
use super::{RedisPool, RedisStore, into_error};
@@ -81,30 +79,6 @@ impl RedisStore {
}
}
// inbuxa: see InMemoryStore::renew_lock
pub async fn renew_lock(&self, key: &[u8], expires: u64) -> trc::Result<bool> {
match &self.pool {
RedisPool::Single(pool) => {
with_conn(pool, async |conn| {
Self::renew_lock_(conn, key, expires).await
})
.await
}
RedisPool::Cluster(pool) => {
with_conn(pool, async |conn| {
Self::renew_lock_(conn, key, expires).await
})
.await
}
RedisPool::Sentinel(pool) => {
with_conn(pool, async |conn| {
Self::renew_lock_(conn, key, expires).await
})
.await
}
}
}
pub async fn key_delete(&self, key: &[u8]) -> trc::Result<()> {
match &self.pool {
RedisPool::Single(pool) => {
@@ -252,22 +226,6 @@ impl RedisStore {
.map(|reply| reply.is_some())
}
async fn renew_lock_(
conn: &mut impl AsyncCommands,
key: &[u8],
expires: u64,
) -> RedisResult<bool> {
redis::cmd("SET")
.arg(key)
.arg(now() + expires)
.arg("XX")
.arg("EX")
.arg(expires as i64)
.query_async::<Option<String>>(conn)
.await
.map(|reply| reply.is_some())
}
async fn key_delete_(conn: &mut impl AsyncCommands, key: &[u8]) -> RedisResult<()> {
conn.del(key).await
}
-51
View File
@@ -401,57 +401,6 @@ impl InMemoryStore {
}
}
/// inbuxa: extends a lock this node holds to `duration` seconds from now.
/// Returns false when the lock is gone or has expired: it may have been
/// taken by someone else since, so it is left alone.
pub async fn renew_lock(&self, prefix: u8, key: &[u8], duration: u64) -> trc::Result<bool> {
match self {
InMemoryStore::Store(store) => {
let key = KeyValue::<()>::build_key(prefix, key);
let key = ValueClass::InMemory(InMemoryClass::Key(key));
let Some(lock_expiry) = store
.get_value::<u64>(ValueKey::from(key.clone()))
.await
.caused_by(trc::location!())?
else {
return Ok(false);
};
let now = now();
if lock_expiry <= now {
return Ok(false);
}
let mut batch = BatchBuilder::new();
batch.assert_value(key.clone(), AssertValue::U64(lock_expiry));
batch.set(key, (now + duration).serialize());
match store.write(batch.build_all()).await {
Ok(_) => Ok(true),
Err(err) if err.is_assertion_failure() => Ok(false),
Err(err) => Err(err
.details("Failed to renew lock.")
.caused_by(trc::location!())),
}
}
InMemoryStore::Sharded(store) => {
Box::pin(
store
.member(&KeyValue::<()>::build_key(prefix, key))
.renew_lock(prefix, key, duration),
)
.await
}
#[cfg(feature = "redis")]
InMemoryStore::Redis(store) => {
store
.renew_lock(&KeyValue::<()>::build_key(prefix, key), duration)
.await
}
InMemoryStore::Static(_) | InMemoryStore::Http(_) => {
Err(trc::StoreEvent::NotSupported.into_err())
}
}
}
pub async fn remove_lock(&self, prefix: u8, key: &[u8]) -> trc::Result<()> {
self.key_delete(KeyValue::<()>::build_key(prefix, key))
.await
+1 -44
View File
@@ -2,8 +2,6 @@
* SPDX-FileCopyrightText: 2020 Stalwart Labs LLC <[email protected]>
*
* SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-SEL
*
* Modified by Coffey Labs in 2026 for INBUXA.
*/
use crate::{
@@ -13,7 +11,6 @@ use crate::{
server::TestServerBuilder,
},
};
use common::BuildServer;
use imap_proto::ResponseType;
use registry::{
schema::{
@@ -21,8 +18,7 @@ use registry::{
prelude::{ObjectType, Property, SocketAddr},
structs::{
ClusterListenerGroup, ClusterListenerGroupProperties, ClusterRole, ClusterTaskGroup,
Coordinator, Imap, MtaDeliverySchedule, MtaVirtualQueue, NatsCoordinator,
NetworkListener, RedisStore,
Coordinator, Imap, NatsCoordinator, NetworkListener, RedisStore,
},
},
types::map::Map,
@@ -213,45 +209,6 @@ pub async fn cluster_tests() {
Some("John Doe")
);
// inbuxa: a settings write applies on every node, no ReloadSettings
let queue_id = admin
.registry_create_object(MtaVirtualQueue {
name: "clusterq".into(),
threads_per_node: 1,
description: None,
})
.await;
admin
.registry_create_object(MtaDeliverySchedule {
name: "cluster-autoreload".into(),
queue_id,
..Default::default()
})
.await;
for (node_id, test) in servers.iter().enumerate() {
let started = std::time::Instant::now();
while !test
.server
.inner
.build_server()
.core
.smtp
.queue
.queue_strategy
.contains_key("cluster-autoreload")
{
assert!(
started.elapsed() < std::time::Duration::from_secs(5),
"node {node_id} didn't pick up the new delivery schedule"
);
tokio::time::sleep(std::time::Duration::from_millis(50)).await;
}
println!(
"Node {node_id} has the new delivery schedule after {} ms",
started.elapsed().as_millis()
);
}
// Run IMAP idle tests across nodes
let mut node1_client = imap_client("[email protected]", "this is john's secret", 1).await;
let mut node2_client = imap_client("[email protected]", "this is john's secret", 2).await;
-1
View File
@@ -10,4 +10,3 @@ pub mod broadcast;
#[cfg(feature = "nats")]
pub mod coordinator; // inbuxa: coordinator reconnects
pub mod stress;
pub mod task_roles; // inbuxa: task types follow cluster roles
-218
View File
@@ -1,218 +0,0 @@
/*
* SPDX-FileCopyrightText: 2026 Coffey Labs
*
* SPDX-License-Identifier: AGPL-3.0-only
*/
//! Two task managers with different cluster roles over one shared store:
//! each runs only the task types its role allows, and a task one node may
//! not run is left for the node that may. Needs a store both nodes can open
//! (STORE=PostgreSql or MySql).
use crate::utils::server::TestServerBuilder;
use common::Server;
use registry::{
schema::{
enums::{ClusterTaskType, IndexDocumentType},
structs::{
ClusterListenerGroup, ClusterRole, ClusterTaskGroup, ClusterTaskGroupProperties, Task,
TaskDnsManagement, TaskIndexDocument, TaskStatus, TaskTlsReport,
},
},
types::map::Map,
};
use std::time::{Duration, Instant};
use store::{
ValueKey,
write::{BatchBuilder, TaskQueueClass, ValueClass},
};
use utils::snowflake::SnowflakeIdGenerator;
const QUEUE_ROLE: &str = "tasks_queue";
const INDEX_MTA_ROLE: &str = "tasks_index_mta";
#[tokio::test(flavor = "multi_thread")]
pub async fn task_role_tests() {
if matches!(
std::env::var("STORE").as_deref(),
Ok("RocksDb" | "Sqlite") | Err(_)
) {
println!("Skipping task role tests: they need a store both nodes can open.");
return;
}
println!(
"Running task role tests on {}...",
std::env::var("STORE").unwrap_or_default()
);
// The roles, stored by a node that runs no services of its own (a node
// looks its role up when it starts)
let seed = TestServerBuilder::new("task_roles_seed")
.await
.with_object(role(QUEUE_ROLE, &[ClusterTaskType::TaskQueueProcessing]))
.await
.with_object(role(
INDEX_MTA_ROLE,
&[
ClusterTaskType::SearchIndexing,
ClusterTaskType::OutboundMta,
],
))
.await
.disable_services()
.build()
.await;
// Node A runs queue tasks (taskQueueProcessing) only
let node_a = TestServerBuilder::new_with_role(
"task_roles_a",
"node-a.example.com".into(),
Some(QUEUE_ROLE.into()),
false,
)
.await
.build_with_opts(false)
.await;
let server_a = node_a.server.clone();
let roles = &server_a.core.network.roles;
assert!(roles.task_manager && !roles.search_indexing && !roles.outbound_mta);
// A DNS task (taskQueueProcessing), an unindex task (searchIndexing) and
// a TLS report (outboundMta), all due now
let [dns, unindex, report] = new_task_ids();
let mut batch = BatchBuilder::new();
batch
.schedule_task_with_id(
dns,
Task::DnsManagement(TaskDnsManagement {
status: TaskStatus::now(),
..Default::default()
}),
)
.schedule_task_with_id(
unindex,
Task::UnindexDocument(TaskIndexDocument {
account_id: 0u32.into(),
document_id: u32::MAX.into(),
document_type: IndexDocumentType::File,
status: TaskStatus::now(),
}),
)
.schedule_task_with_id(
report,
Task::TlsReport(TaskTlsReport {
report_id: u64::MAX.into(),
status: TaskStatus::now(),
}),
);
server_a.store().write(batch.build_all()).await.unwrap();
server_a.notify_task_queue();
// Node A runs the DNS task and leaves the other two alone. Upstream ran
// the TLS report here too: report tasks ran on any node with a task
// manager.
wait_until_run(&server_a, &[dns], Duration::from_secs(20)).await;
tokio::time::sleep(Duration::from_secs(3)).await;
server_a.notify_task_queue();
tokio::time::sleep(Duration::from_secs(2)).await;
assert!(
is_pending(&server_a, unindex).await,
"unindex ran on node A"
);
assert!(
is_pending(&server_a, report).await,
"TLS report ran on node A"
);
// Node B (search indexing and outbound MTA) comes up and picks up what
// node A left
let node_b = TestServerBuilder::new_with_role(
"task_roles_b",
"node-b.example.com".into(),
Some(INDEX_MTA_ROLE.into()),
false,
)
.await
.build_with_opts(false)
.await;
let server_b = node_b.server.clone();
let roles = &server_b.core.network.roles;
assert!(!roles.task_manager && roles.search_indexing && roles.outbound_mta);
server_b.notify_task_queue();
wait_until_run(&server_b, &[unindex, report], Duration::from_secs(20)).await;
// A queue task scheduled now still runs, on node A: node B may not
// claim it
let [dns] = new_task_ids();
let mut batch = BatchBuilder::new();
batch.schedule_task_with_id(
dns,
Task::DnsManagement(TaskDnsManagement {
status: TaskStatus::now(),
..Default::default()
}),
);
server_b.store().write(batch.build_all()).await.unwrap();
server_b.notify_task_queue();
tokio::time::sleep(Duration::from_secs(3)).await;
assert!(is_pending(&server_b, dns).await, "DNS task ran on node B");
server_a.notify_task_queue();
wait_until_run(&server_a, &[dns], Duration::from_secs(20)).await;
if seed.is_reset() {
seed.temp_dir.delete();
node_a.temp_dir.delete();
node_b.temp_dir.delete();
}
}
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()),
}),
}
}
fn new_task_ids<const N: usize>() -> [u64; N] {
std::array::from_fn(|_| SnowflakeIdGenerator::global_id().unwrap())
}
/// Still due and never run: present, and pending.
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_until_run(server: &Server, ids: &[u64], within: Duration) {
let started = Instant::now();
loop {
let mut left = 0;
for id in ids {
if is_pending(server, *id).await {
left += 1;
}
}
if left == 0 {
return;
}
assert!(
started.elapsed() < within,
"{left} task(s) still pending after {:?}",
started.elapsed()
);
tokio::time::sleep(Duration::from_millis(250)).await;
}
}
+17 -40
View File
@@ -2,8 +2,6 @@
* SPDX-FileCopyrightText: 2020 Stalwart Labs LLC <[email protected]>
*
* SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-SEL
*
* Modified by Coffey Labs in 2026 for INBUXA.
*/
use crate::utils::server::TestServer;
@@ -42,23 +40,15 @@ pub mod vrfy;
const EVENT_TIMEOUT: Duration = Duration::from_secs(5);
impl TestServer {
// inbuxa: registry writes reload the settings, and each reload sends the
// queue a ReloadSettings; read_event, try_read_event and assert_no_events
// pass over those (expect_reload_settings still waits for one)
pub async fn read_event(&mut self) -> QueueEvent {
while let Some(event) = self.queue_events.pop_front() {
if !event.is_reload_settings() {
return event;
}
if let Some(event) = self.queue_events.pop_front() {
return event;
}
loop {
match tokio::time::timeout(EVENT_TIMEOUT, self.queue_rx.recv()).await {
Ok(Some(event)) if event.is_reload_settings() => (),
Ok(Some(event)) => return event,
Ok(None) => panic!("Channel closed."),
Err(_) => panic!("No queue event received."),
}
match tokio::time::timeout(EVENT_TIMEOUT, self.queue_rx.recv()).await {
Ok(Some(event)) => event,
Ok(None) => panic!("Channel closed."),
Err(_) => panic!("No queue event received."),
}
}
@@ -88,39 +78,26 @@ impl TestServer {
}
pub async fn try_read_event(&mut self) -> Option<QueueEvent> {
while let Some(event) = self.queue_events.pop_front() {
if !event.is_reload_settings() {
return Some(event);
}
if let Some(event) = self.queue_events.pop_front() {
return Some(event);
}
loop {
match tokio::time::timeout(EVENT_TIMEOUT, self.queue_rx.recv()).await {
Ok(Some(event)) if event.is_reload_settings() => (),
Ok(Some(event)) => return Some(event),
Ok(None) => panic!("Channel closed."),
Err(_) => return None,
}
match tokio::time::timeout(EVENT_TIMEOUT, self.queue_rx.recv()).await {
Ok(Some(event)) => Some(event),
Ok(None) => panic!("Channel closed."),
Err(_) => None,
}
}
pub fn assert_no_events(&mut self) {
if let Some(event) = self
.queue_events
.iter()
.find(|event| !event.is_reload_settings())
{
if let Some(event) = self.queue_events.pop_front() {
panic!("Expected empty queue but got {event:?}");
}
self.queue_events.clear();
loop {
match self.queue_rx.try_recv() {
Ok(event) if event.is_reload_settings() => (),
Err(TryRecvError::Empty) => break,
Ok(event) => panic!("Expected empty queue but got {event:?}"),
Err(err) => panic!("Queue error: {err:?}"),
}
match self.queue_rx.try_recv() {
Err(TryRecvError::Empty) => (),
Ok(event) => panic!("Expected empty queue but got {event:?}"),
Err(err) => panic!("Queue error: {err:?}"),
}
}
-4
View File
@@ -10,8 +10,6 @@ pub mod blob;
pub mod import_export;
pub mod lookup;
pub mod ops;
#[cfg(any(feature = "postgres", feature = "mysql"))]
pub mod pool_timeout; // inbuxa: SQL pools give up instead of hanging
pub mod query;
pub mod registry;
#[cfg(feature = "postgres")]
@@ -21,8 +19,6 @@ pub mod replica_mysql; // inbuxa: read replicas on MySQL
#[cfg(all(feature = "postgres", feature = "redis"))]
pub mod replica_cluster; // inbuxa: read replicas across nodes
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"))]
pub mod sql_timeout;
pub mod task_locks; // inbuxa: task locks across nodes
-96
View File
@@ -1,96 +0,0 @@
/*
* SPDX-FileCopyrightText: 2026 Coffey Labs
*
* SPDX-License-Identifier: AGPL-3.0-only
*/
//! A database that accepts connections and then says nothing (a hung or
//! half-dead server, a black-holed failover) gives a worker an error within
//! the pool's timeouts. Upstream's pools had none, so the worker waited for
//! good. No database is needed: a local listener that never answers plays
//! the server.
use registry::schema::structs::DataStore;
use std::time::{Duration, Instant};
use store::{Store, ValueKey, write::ValueClass};
use tokio::net::TcpListener;
/// Accepts connections on a local port and never sends a byte.
async fn silent_server() -> u16 {
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
let port = listener.local_addr().unwrap().port();
tokio::spawn(async move {
let mut held = Vec::new();
while let Ok((socket, _)) = listener.accept().await {
held.push(socket);
}
});
port
}
/// Builds the store and reads a key; both must end, with an error for the
/// read, well within `limit`.
async fn assert_times_out(data_store: DataStore, limit: Duration) {
let started = Instant::now();
let result = tokio::time::timeout(limit, async {
match Store::build(data_store).await {
Ok(store) => store
.get_value::<u64>(ValueKey::from(ValueClass::Property(0)))
.await
.map(|_| ())
.map_err(|err| err.to_string()),
Err(err) => Err(err.to_string()),
}
})
.await;
let elapsed = started.elapsed();
match result {
Ok(Err(err)) => println!("Got {err} after {elapsed:?}"),
Ok(Ok(())) => panic!("a silent server answered?"),
Err(_) => panic!("still waiting for a connection after {elapsed:?}"),
}
}
#[cfg(feature = "postgres")]
#[tokio::test(flavor = "multi_thread")]
pub async fn postgres_pool_timeout() {
use registry::schema::structs::PostgreSqlStore;
let port = silent_server().await;
println!("Running PostgreSQL pool timeout test...");
// The store's own timeout bounds opening a connection, handshake
// included (tokio-postgres's connect_timeout covers only the TCP connect)
assert_times_out(
DataStore::PostgreSql(PostgreSqlStore {
host: "127.0.0.1".into(),
port: port as u64,
database: "none".into(),
timeout: Some(Duration::from_secs(2).into()),
use_tls: false,
..Default::default()
}),
Duration::from_secs(20),
)
.await;
}
#[cfg(feature = "mysql")]
#[tokio::test(flavor = "multi_thread")]
pub async fn mysql_pool_timeout() {
use registry::schema::structs::MySqlStore;
let port = silent_server().await;
println!("Running MySQL pool timeout test...");
// mysql_async has no pool timeout; the store waits 30 s for a connection
assert_times_out(
DataStore::MySql(MySqlStore {
host: "127.0.0.1".into(),
port: port as u64,
database: "none".into(),
use_tls: false,
..Default::default()
}),
Duration::from_secs(60),
)
.await;
}
-157
View File
@@ -128,11 +128,6 @@ pub async fn test(test: &TestServer) {
println!("Running trace document tests...");
test_trace_documents(store.clone()).await;
// inbuxa: address fields match by full address, local part, domain and
// display name on every backend
println!("Running address search tests...");
test_address_search(store.clone()).await;
// Large document insert test
println!("Running large document insert tests...");
let mut large_text = String::with_capacity(20 * 1024 * 1024);
@@ -977,155 +972,3 @@ async fn test_trace_documents(store: SearchStore) {
.unwrap();
}
}
// inbuxa: the message indexer passes each display name and each address of
// From/To/Cc/Bcc as keyword text (Language::None). The built-in index splits
// that text into words, so an address is found by its full form, its local
// part, its domain or a display-name word; PostgreSQL kept the whole address
// as one token and MySQL dropped stopwords such as "com" and words under three
// characters. The expected results below are the built-in (RocksDB/SQLite)
// results and must be the same on every backend.
async fn test_address_search(store: SearchStore) {
const ACCOUNT_ID: u32 = 7;
let messages: [[&[(&str, &str)]; 4]; 5] = [
// From, To, Cc, Bcc
[
&[("Amazon.com", "[email protected]")],
&[("Jane Doe", "[email protected]")],
&[],
&[],
],
[
&[("", "[email protected]")],
&[("", "[email protected]")],
&[("Jane Doe", "[email protected]")],
&[],
],
[
&[("GitHub", "[email protected]")],
&[("Jo Li", "[email protected]")],
&[],
&[("Audit", "[email protected]")],
],
[
&[("Jane Doe", "[email protected]")],
&[("Amazon Web Services", "[email protected]")],
&[("Bob", "[email protected]")],
&[("", "[email protected]")],
],
[
&[("Newsletter", "[email protected]")],
&[("", "[email protected]")],
&[],
&[],
],
];
let fields = [
EmailSearchField::From,
EmailSearchField::To,
EmailSearchField::Cc,
EmailSearchField::Bcc,
];
let mut documents = Vec::new();
let mut mask = RoaringBitmap::new();
for (document_id, message) in messages.iter().enumerate() {
let mut document = IndexDocument::new(SearchIndex::Email)
.with_account_id(ACCOUNT_ID)
.with_document_id(document_id as u32);
for (field, addresses) in fields.iter().zip(message.iter()) {
for (name, address) in addresses.iter() {
if !name.is_empty() {
document.index_text(field.clone(), name, Language::None);
}
document.index_text(field.clone(), address, Language::None);
}
}
document.index_unsigned(EmailSearchField::ReceivedAt, document_id as u64);
documents.push(document);
mask.insert(document_id as u32);
}
store.index(documents).await.unwrap();
if let SearchStore::ElasticSearch(store) = &store {
store.refresh_index(SearchIndex::Email).await.unwrap();
}
for (field, text, expected) in [
// full address
(EmailSearchField::From, "[email protected]", vec![0u32]),
(EmailSearchField::To, "[email protected]", vec![0]),
(EmailSearchField::Cc, "[email protected]", vec![1]),
(EmailSearchField::Bcc, "[email protected]", vec![3]),
(EmailSearchField::To, "[email protected]", vec![1, 2]),
// local part
(EmailSearchField::From, "noreply", vec![0, 2]),
(EmailSearchField::To, "jo", vec![1, 2]),
(EmailSearchField::Cc, "bob", vec![3]),
(EmailSearchField::Bcc, "audit", vec![2]),
// domain
(EmailSearchField::From, "amazon.com", vec![0, 1]),
(EmailSearchField::From, "amazon", vec![0, 1]),
(EmailSearchField::To, "example.org", vec![0, 4]),
(EmailSearchField::To, "io.de", vec![1, 2]),
(EmailSearchField::Cc, "example.net", vec![3]),
(EmailSearchField::Bcc, "example.org", vec![2]),
(EmailSearchField::From, "www.example.com", vec![4]),
(EmailSearchField::From, "com", vec![0, 1, 2, 4]),
// display name
(EmailSearchField::From, "Jane", vec![3]),
(EmailSearchField::From, "jane doe", vec![3]),
(EmailSearchField::To, "Web Services", vec![3]),
(EmailSearchField::To, "Li", vec![2]),
(EmailSearchField::Cc, "Doe", vec![1]),
(EmailSearchField::Bcc, "Audit", vec![2]),
// hyphenated local part
(EmailSearchField::From, "shipment-tracking", vec![1]),
(EmailSearchField::From, "tracking", vec![1]),
// no match
(EmailSearchField::From, "amazon.org", vec![]),
(EmailSearchField::To, "noreply", vec![]),
(EmailSearchField::Bcc, "jane", vec![]),
] {
let ids = store
.query_account(
SearchQuery::new(SearchIndex::Email)
.with_filters(vec![
SearchFilter::eq(SearchField::AccountId, ACCOUNT_ID),
SearchFilter::has_keyword(field.clone(), text),
])
.with_comparator(SearchComparator::ascending(EmailSearchField::ReceivedAt))
.with_mask(mask.clone()),
)
.await
.unwrap();
assert_eq!(ids, expected, "{field:?} {text:?}");
}
// TEXT-style search across all address fields
let ids = store
.query_account(
SearchQuery::new(SearchIndex::Email)
.with_filters(vec![
SearchFilter::eq(SearchField::AccountId, ACCOUNT_ID),
SearchFilter::Or,
SearchFilter::has_keyword(EmailSearchField::From, "example.org"),
SearchFilter::has_keyword(EmailSearchField::To, "example.org"),
SearchFilter::has_keyword(EmailSearchField::Cc, "example.org"),
SearchFilter::has_keyword(EmailSearchField::Bcc, "example.org"),
SearchFilter::End,
])
.with_comparator(SearchComparator::ascending(EmailSearchField::ReceivedAt))
.with_mask(mask.clone()),
)
.await
.unwrap();
assert_eq!(ids, vec![0, 1, 2, 3, 4]);
store
.unindex(
SearchQuery::new(SearchIndex::Email)
.with_filter(SearchFilter::eq(SearchField::AccountId, ACCOUNT_ID)),
)
.await
.unwrap();
}
-159
View File
@@ -1,159 +0,0 @@
/*
* 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()
}
+1 -27
View File
@@ -79,33 +79,7 @@ pub async fn task_lock_tests() {
"ran before the other node's locks expired: {elapsed:?}"
);
// 3. A task that runs longer than a lock lifetime keeps its claim: the
// task manager renews the lease while this node holds it, and the claim
// ends when the task does. (Before, a lock simply lasted an hour.)
let [id] = new_task_ids(1)[..] else {
unreachable!()
};
assert!(server.try_lock_task(id).await, "claim {id}");
tokio::time::sleep(Duration::from_secs(LOCK_EXPIRY + LOCK_EXPIRY / 2)).await;
assert!(
!foreign_lock(&server, id, LOCK_EXPIRY).await,
"lease lapsed while the task ran"
);
server.remove_index_lock(id).await;
assert!(
foreign_lock(&server, id, LOCK_EXPIRY).await,
"released when the task ended"
);
let _ = server
.in_memory_store()
.remove_lock(KV_LOCK_TASK, &id.to_be_bytes())
.await;
assert!(
common::ipc::TaskLocks::DEFAULT_EXPIRY <= 5 * 60,
"a dead node's tasks wait no more than a few minutes"
);
// 4. A graceful stop releases the locks this node holds: another node
// 3. A graceful stop releases the locks this node holds: another node
// can claim those tasks at once, and this one claims nothing more
let ids = new_task_ids(3);
for id in &ids {
-220
View File
@@ -1,220 +0,0 @@
/*
* SPDX-FileCopyrightText: 2026 Coffey Labs
*
* SPDX-License-Identifier: AGPL-3.0-only
*/
// inbuxa: a registry write to an object the running settings are built from
// applies without an x:Action ReloadSettings, and the set response says so.
use crate::utils::{
jmap::JmapResponse,
server::{TestServer, TestServerBuilder},
};
use common::BuildServer;
use registry::{
schema::{
enums::TracingLevel,
prelude::ObjectType,
structs::{
AllowedIp, CertificateManagement, DkimManagement, DnsManagement, Domain, Expression,
MtaDeliverySchedule, MtaStageAuth, MtaVirtualQueue, Tracer, TracerStdout,
},
},
types::ipmask::IpAddrOrMask,
};
use serde_json::Value;
#[tokio::test(flavor = "multi_thread")]
pub async fn settings_reload_tests() {
let mut test = TestServerBuilder::new("settings_reload_tests")
.await
.with_default_listeners()
.await
.with_object(MtaStageAuth {
require: Expression {
else_: "false".to_string(),
..Default::default()
},
..Default::default()
})
.await
.build()
.await;
let admin = test
.create_user_account(
"admin",
"[email protected]",
"these_pretzels_are_making_me_thirsty",
&[],
"Admin",
)
.await;
test.account("admin")
.assign_roles_to_account(admin.id(), &["user", "system"])
.await;
test.insert_account(admin);
test_write_applies(&test).await;
if test.is_reset() {
test.temp_dir.delete();
}
}
async fn test_write_applies(test: &TestServer) {
println!("Running settings reload after registry writes...");
let admin = test.account("[email protected]");
// A delivery schedule is in use as soon as it is saved
let response = admin
.registry_create([MtaVirtualQueue {
name: "autorld".into(),
threads_per_node: 2,
description: None,
}])
.await;
assert_applied(&response);
let queue_id = response.created_id(0);
assert!(!has_schedule(test, "autoreload-schedule"));
let response = admin
.registry_create([MtaDeliverySchedule {
name: "autoreload-schedule".into(),
queue_id,
..Default::default()
}])
.await;
assert_applied(&response);
assert!(has_schedule(test, "autoreload-schedule"));
// Destroyed, it's gone at once too
let schedule_id = response.created_id(0);
let response = admin
.registry_destroy(ObjectType::MtaDeliverySchedule, [schedule_id])
.await;
assert_applied(&response);
assert!(!has_schedule(test, "autoreload-schedule"));
// Concurrent writes all end up in the running settings
let names = (0..8)
.map(|i| format!("autoreload-{i}"))
.collect::<Vec<_>>();
let mut writes = Vec::new();
for name in &names {
writes.push(admin.registry_create([MtaDeliverySchedule {
name: name.clone(),
queue_id,
..Default::default()
}]));
}
let mut schedule_ids = Vec::new();
for response in futures::future::join_all(writes).await {
assert_applied(&response);
schedule_ids.push(response.created_id(0));
}
for name in &names {
assert!(has_schedule(test, name), "{name} missing");
}
// Several objects in one request: one reload
let response = admin
.registry_destroy(ObjectType::MtaDeliverySchedule, schedule_ids.iter())
.await;
assert_applied(&response);
for name in &names {
assert!(!has_schedule(test, name), "{name} still present");
}
// A write whose reload fails is stored, and the response says the
// settings weren't reloaded: only one console tracer is allowed.
let response = admin
.registry_create([
Tracer::Stdout(TracerStdout {
enable: true,
level: TracingLevel::Error,
..Default::default()
}),
Tracer::Stdout(TracerStdout {
enable: true,
level: TracingLevel::Error,
..Default::default()
}),
])
.await;
let reload = settings_reload(&response).expect("x:settingsReload missing");
assert_eq!(reload["applied"], Value::Bool(false), "{response:?}");
let description = reload["description"].as_str().unwrap_or_default();
assert!(
description.starts_with("Saved, but the running settings were not reloaded. ")
&& description.contains("Only one console tracer is allowed"),
"{description}"
);
let tracer_ids = [response.created_id(0), response.created_id(1)];
let response = admin
.registry_destroy(ObjectType::Tracer, tracer_ids.iter())
.await;
assert_applied(&response);
// An allowed IP is live as soon as it is saved, and gone once
// destroyed. It lives in the core's security settings, which the
// blocked-IP reload it used to get doesn't rebuild.
let ip: std::net::IpAddr = "198.51.100.7".parse().unwrap();
assert!(!is_allowed(test, ip));
let response = admin
.registry_create([AllowedIp {
address: IpAddrOrMask::from_ip(ip),
reason: Some("autoreload".into()),
..Default::default()
}])
.await;
assert_applied(&response);
assert!(
is_allowed(test, ip),
"allowed IP not in the running settings"
);
let allowed_id = response.created_id(0);
let response = admin
.registry_destroy(ObjectType::AllowedIp, [allowed_id])
.await;
assert_applied(&response);
assert!(!is_allowed(test, ip), "destroyed allowed IP still live");
// Data that isn't part of the running settings doesn't reload them
let response = admin
.registry_create([Domain {
name: "autoreload.example.org".into(),
certificate_management: CertificateManagement::Manual,
dns_management: DnsManagement::Manual,
dkim_management: DkimManagement::Manual,
..Default::default()
}])
.await;
assert!(settings_reload(&response).is_none(), "{response:?}");
}
fn settings_reload(response: &JmapResponse) -> Option<&Value> {
response.pointer("/methodResponses/0/1/x:settingsReload")
}
fn assert_applied(response: &JmapResponse) {
assert_eq!(
settings_reload(response),
Some(&serde_json::json!({"applied": true})),
"{response:?}"
);
}
fn has_schedule(test: &TestServer, name: &str) -> bool {
test.server
.inner
.build_server()
.core
.smtp
.queue
.queue_strategy
.contains_key(name)
}
fn is_allowed(test: &TestServer, ip: std::net::IpAddr) -> bool {
test.server.inner.build_server().is_ip_allowed(ip)
}
-1
View File
@@ -11,7 +11,6 @@ pub mod authentication;
pub mod ai;
pub mod ai_calibration;
pub mod authorization;
pub mod auto_reload; // inbuxa: registry writes apply at once
pub mod branding;
pub mod crypto;
pub mod delivery;