1 Commits
Author SHA1 Message Date
jcoffey-dev 8a8f48c944 Coordinator: join the cluster when NATS comes up, report the connection
ci / fork-checks (pull_request) Successful in 46s
ci / build (pull_request) Successful in 4m29s
A node that started while NATS was down never got a coordinator. The
connect failed at boot, bootstrap recorded a build error and the node ran
with Coordinator::None until restarted. It had no broadcast subscriber
or publisher, so cross-node push and cache invalidation to it stayed
broken, and its healthcheck said nothing about it. Losing NATS after
startup was silent too.

- The NATS client now connects in the background
  (retry_on_initial_connect): startup never waits on NATS or fails over
  it, the node gets its coordinator, subscriber and publisher at once,
  and the client keeps trying (async-nats's backoff, at most 4 s apart)
  until NATS answers. Subscriptions made meanwhile start delivering when
  it does. A configured maxReconnects still ends the attempts.
- Three new events report the connection: cluster.coordinator-connected
  (info), cluster.coordinator-disconnected (warn: lost, closed, gave up,
  or not connected within the connection timeout at startup) and
  cluster.coordinator-error (warn: a failed attempt, reported once per
  outage rather than every retry, and server errors, slow consumers and
  lame duck mode). They are in the packaged schema, ids 644 to 646.
- GET /healthz/cluster reports the coordinator: 200
  {"coordinator":"connected"}, 503 {"coordinator":"disconnected"}, or
  200 with "none" (no coordinator) or "unknown" (a backend that doesn't
  track its connection). /healthz/live and /healthz/ready are unchanged
  on purpose: a node without its coordinator still serves mail, and
  failing those would have orchestrators restart, or pull out of
  service, every node at once whenever NATS is down.

Only NATS connects lazily; the other coordinator backends still fail at
boot as before.

cluster::coordinator::coordinator_reconnect_tests starts a node against a
NATS port with nothing behind it, checks it boots with a coordinator and
reports it disconnected, subscribes, then starts NATS on that port: the
node connects on its own and the subscription receives a message from a
second client. Stopping and restarting NATS shows disconnected, then
connected, and the same subscription keeps working.
2026-09-24 08:29:20 -07:00
41 changed files with 347 additions and 2486 deletions
+15 -69
View File
@@ -3,28 +3,11 @@
# whether a person pushed it or weekly-release.yml created it through the
# releases API.
#
# The image is multi-arch (linux/amd64, linux/arm64), built by two jobs on
# the image-build runner rather than one buildx run for both. The Dockerfile's
# builder stage runs on the build platform and cross-compiles with an aarch64
# linker, so only the small final stage (apt, setcap) goes through QEMU for
# arm64 -- but two release builds (LTO, one codegen unit) side by side on one
# machine each take twice as long. Production runs amd64, so amd64 goes first
# and on its own:
# * publish-amd64 pushes :<version>-amd64 and :<version>, a plain amd64
# image, as soon as its build is done. A deploy can start from it.
# * publish-arm64 then builds arm64, pushes :<version>-arm64, and replaces
# :<version> with the two-platform index. :latest moves only here, so it
# never names an image without arm64.
#
# Both jobs use one BuildKit builder, `gitea-builder`, whose container
# (buildx_buildkit_gitea-builder0) and state volume stay on the runner's host
# between jobs: a job container's `buildx create` finds the existing container
# and reuses it and its cache. The dependency build (`cargo chef cook`) is
# keyed on the recipe, which only a dependency change alters, so a release
# normally compiles just the workspace. Removing that container or its volume
# costs the next release a cold build, nothing more. The planner and dependency
# layers for the build platform are shared, so arm64 also reuses what amd64
# just did where it can.
# The image is multi-arch (linux/amd64, linux/arm64) as before, but built in
# one buildx run on host1 instead of one native runner per architecture: the
# Dockerfile's builder stage runs on the build platform and cross-compiles
# with an aarch64 linker, so only the small final stage (apt, setcap) goes
# through QEMU for arm64. No digest-joining job is needed.
#
# Two guards before anything is pushed:
# * the tag must be v<brand_version!>. The version is a string in
@@ -79,7 +62,7 @@ jobs:
echo "version=$V" >> "$GITHUB_OUTPUT"
echo "version $V"
publish-amd64:
publish:
needs: [version]
runs-on: docker
container:
@@ -98,15 +81,16 @@ jobs:
test -n "$REGISTRY" && test -n "$VERSION"
test -n "$PACKAGE_TOKEN" || { echo "PACKAGE_TOKEN secret is not set on this repository" >&2; exit 1; }
echo "$PACKAGE_TOKEN" | docker login -u jcoffey-dev --password-stdin "$REGISTRY"
docker run --privileged --rm tonistiigi/binfmt --install arm64
docker buildx create --use --name gitea-builder --driver docker-container || docker buildx use gitea-builder
# Attestations off, as before: they add manifests of their own, and the
# index should hold the two images and nothing else.
# Attestations off, as before: they add manifests of their own to the
# index, and the index should hold the two images and nothing else.
- run: |
docker buildx build \
--platform linux/amd64 \
--platform linux/amd64,linux/arm64 \
--provenance=false --sbom=false \
--tag "$IMAGE:$VERSION-amd64" \
--tag "$IMAGE:$VERSION" \
--tag "$IMAGE:latest" \
--push .
docker buildx imagetools inspect "$IMAGE:$VERSION"
# Gitea keeps a container package on its owner; linking it shows it on
@@ -119,47 +103,11 @@ jobs:
- if: always()
run: docker logout "$REGISTRY" || true
publish-arm64:
needs: [version, publish-amd64]
runs-on: docker
container:
image: docker:28-cli@sha256:625d9431a9f54c5a2bc90f24f0e1c3d55b1349fd857dd85035f98c2c9acbdd4d # 28-cli
volumes:
- /var/run/docker.sock:/var/run/docker.sock
env:
DOCKER_BUILDKIT: "1"
REGISTRY: ${{ vars.REGISTRY }}
IMAGE: ${{ vars.REGISTRY }}/${{ github.repository }}
VERSION: ${{ needs.version.outputs.version }}
PACKAGE_TOKEN: ${{ secrets.PACKAGE_TOKEN }}
steps:
- uses: coffey-labs/actions/checkout@fab0c4d45e0162963965f1555df27b7bed5e20ec
- run: |
echo "$PACKAGE_TOKEN" | docker login -u jcoffey-dev --password-stdin "$REGISTRY"
docker run --privileged --rm tonistiigi/binfmt --install arm64
docker buildx create --use --name gitea-builder --driver docker-container || docker buildx use gitea-builder
# The index is built from the two per-architecture tags rather than from
# :<version>, which by now is the amd64 image and would be read as such.
- run: |
docker buildx build \
--platform linux/arm64 \
--provenance=false --sbom=false \
--tag "$IMAGE:$VERSION-arm64" \
--push .
docker buildx imagetools create \
--tag "$IMAGE:$VERSION" \
--tag "$IMAGE:latest" \
"$IMAGE:$VERSION-amd64" "$IMAGE:$VERSION-arm64"
docker buildx imagetools inspect "$IMAGE:$VERSION"
- if: always()
run: docker logout "$REGISTRY" || true
# The weekly release creates its Release (and so the tag) first; a tag
# pushed by hand has none. Either way the tag ends up with exactly one
# Release, created once the amd64 image exists so its pull instructions
# work; arm64 and the binaries follow.
# Release, created after the image exists so its pull instructions work.
release:
needs: [version, publish-amd64]
needs: [version, publish]
runs-on: light
container:
image: python:3.13-slim@sha256:8d9d0b8bcf6506481eae4907c18f5e3e7902e629f5f6d684f9e7c32e85e3ddf0 # 3.13-slim
@@ -183,9 +131,7 @@ jobs:
except urllib.error.HTTPError as e:
if e.code != 404: raise
image = f"{os.environ['REGISTRY']}/{os.environ['REPO']}:{version}"
body = (f"Container image: `{image}` (linux/amd64, linux/arm64); also `:latest`. "
"amd64 is published first; arm64 is added to the same tag when its build "
"finishes, and `:latest` moves then.\n\n"
body = (f"Container image: `{image}` (linux/amd64, linux/arm64); also `:latest`.\n\n"
"Binaries for a host install are attached: `inbuxa-linux-amd64.tar.gz` and "
"`inbuxa-linux-arm64.tar.gz`, with `SHA256SUMS`. Each is the binary out of this "
"release's image for that architecture, so it is the same build. The image "
@@ -208,7 +154,7 @@ jobs:
# `docker create` does not start anything, so pulling an arm64 image on an
# amd64 runner and copying a file out of it needs no emulation.
binaries:
needs: [version, publish-arm64, release]
needs: [version, publish, release]
runs-on: docker
container:
image: docker:28-cli@sha256:625d9431a9f54c5a2bc90f24f0e1c3d55b1349fd857dd85035f98c2c9acbdd4d # 28-cli
+30 -284
View File
@@ -13,27 +13,20 @@ use crate::{
storage::Storage,
telemetry::Telemetry,
},
ipc::{BroadcastEvent, QueueEvent, RegistryChange},
ipc::{QueueEvent, RegistryChange},
network::security::{BlockedIps, IpWithTtl},
};
use ahash::AHashMap;
use directory::Directories;
use registry::{
schema::{prelude::ObjectType, structs::BlockedIp},
types::{
error::{Error, Warning},
id::ObjectId,
},
types::error::{Error, Warning},
};
use std::sync::Arc;
use store::{LookupStores, registry::bootstrap::Bootstrap, write::now};
pub struct ReloadResult {
/// Errors that kept the reload from being applied.
pub errors: Vec<Error>,
/// inbuxa: errors in objects that already failed when the running
/// settings were built; logged, but they don't refuse a reload.
pub known_errors: Vec<Error>,
pub warnings: Vec<Warning>,
pub replaced_core: bool,
}
@@ -121,60 +114,42 @@ impl Server {
directories: directory.directories,
};
// inbuxa: upstream swapped the core only when the whole build
// was free of errors, while boot runs with whatever built. So one
// object that failed (a DNS lookup that timed out, say) refused
// every later reload, cluster-wide when the reload came from
// ReloadSettings, and the running settings went stale. Now a
// reload is refused only for errors in objects that built when
// the running settings were built: those would be lost by
// applying it. Objects that already failed then are missing
// from the running settings anyway, as at boot, so their
// errors are reported but don't hold the reload back.
// Parse tracers
let tracers = Telemetry::parse(&mut bootstrap, &storage).await;
let core = Box::pin(Core::parse(&mut bootstrap, storage)).await;
let mut servers = Listeners::parse(&mut bootstrap).await;
if !self.has_new_build_errors(&bootstrap.errors) {
servers
.parse_tcp_acceptors(&mut bootstrap, self.inner.clone())
.await;
if bootstrap.errors.is_empty() {
let core = Box::pin(Core::parse(&mut bootstrap, storage)).await;
if !self.has_new_build_errors(&bootstrap.errors) {
// Update core
self.inner.shared_core.store(core.into());
if bootstrap.errors.is_empty() {
let mut servers = Listeners::parse(&mut bootstrap).await;
servers
.parse_tcp_acceptors(&mut bootstrap, self.inner.clone())
.await;
// Update tracers
tracers.update();
if bootstrap.errors.is_empty() {
// Update core
self.inner.shared_core.store(core.into());
// Reload queue settings
self.inner
.ipc
.queue_tx
.send(QueueEvent::ReloadSettings)
.await
.ok();
// Update tracers
self.record_build_errors(&bootstrap.errors);
tracers.update();
return Ok(ReloadResult {
errors: Vec::new(),
known_errors: bootstrap.errors,
warnings: bootstrap.warnings,
replaced_core: true,
});
// Reload queue settings
self.inner
.ipc
.queue_tx
.send(QueueEvent::ReloadSettings)
.await
.ok();
return Ok(ReloadResult {
errors: bootstrap.errors,
warnings: bootstrap.warnings,
replaced_core: true,
});
}
}
}
let (known_errors, errors) = std::mem::take(&mut bootstrap.errors)
.into_iter()
.partition(|error| self.is_known_build_error(error));
return Ok(ReloadResult {
errors,
known_errors,
warnings: bootstrap.warnings,
replaced_core: false,
});
}
}
@@ -188,7 +163,7 @@ impl ReloadResult {
}
pub fn log(&self) {
for error in self.errors.iter().chain(&self.known_errors) {
for error in &self.errors {
error.log();
}
for warning in &self.warnings {
@@ -201,237 +176,8 @@ impl From<Bootstrap> for ReloadResult {
fn from(bootstrap: Bootstrap) -> Self {
Self {
errors: bootstrap.errors,
known_errors: Vec::new(),
warnings: bootstrap.warnings,
replaced_core: false,
}
}
}
// inbuxa: which objects failed to build for the running settings
impl Server {
/// Records the objects that failed to build for the settings now running.
pub fn record_build_errors(&self, errors: &[Error]) {
*self.inner.data.build_errors.lock() = errors.iter().filter_map(error_object).collect();
}
fn is_known_build_error(&self, error: &Error) -> bool {
error_object(error).is_some_and(|id| self.inner.data.build_errors.lock().contains(&id))
}
fn has_new_build_errors(&self, errors: &[Error]) -> bool {
errors.iter().any(|error| !self.is_known_build_error(error))
}
}
fn error_object(error: &Error) -> Option<ObjectId> {
match error {
Error::Validation { object_id, .. }
| Error::Build { object_id, .. }
| Error::NotFound { object_id } => Some(*object_id),
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
}
-4
View File
@@ -93,11 +93,9 @@ 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"),
build_errors: Default::default(),
asn_geo_data: Default::default(),
}
}
@@ -236,11 +234,9 @@ 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(),
build_errors: Default::default(),
asn_geo_data: Default::default(),
lookup_stores: Default::default(),
}
@@ -16,6 +16,7 @@ use mail_auth::common::resolver::ToReverseName;
use nlp::classifier::model::{CcfhClassifier, FhClassifier};
use registry::schema::{
enums::{ExpressionVariable, ModelSize},
prelude::ObjectType,
structs::{
self, SpamDnsblServer, SpamDnsblSettings, SpamFileExtension, SpamPyzor, SpamRule,
SpamSettings, SpamTag,
@@ -24,10 +25,10 @@ use registry::schema::{
use sieve::SpamStatus;
use std::{
net::{IpAddr, SocketAddr},
sync::Arc,
time::{Duration, Instant},
time::Duration,
};
use store::registry::{RegistryObject, bootstrap::Bootstrap};
use tokio::net::lookup_host;
use utils::{cache::CacheItemWeight, glob::GlobMap};
#[derive(rkyv::Archive, rkyv::Deserialize, rkyv::Serialize, Debug, Default)]
@@ -156,11 +157,7 @@ pub struct FtrlParameters {
#[derive(Debug, Clone)]
pub struct PyzorConfig {
// inbuxa: the server is resolved when a message is checked, not while the
// settings are built (see PyzorConfig::address)
pub host: String,
pub port: u16,
pub resolved: Arc<parking_lot::Mutex<Option<(SocketAddr, Instant)>>>,
pub address: SocketAddr,
pub timeout: Duration,
pub min_count: u64,
pub min_wl_count: u64,
@@ -477,15 +474,31 @@ impl PyzorConfig {
return None;
}
// inbuxa: upstream resolved the host here and reported a failed lookup
// as a build error, so a DNS hiccup on one node refused every settings
// reload on it (and, from the node that ran ReloadSettings, across the
// cluster). The lookup now happens when a message is checked; a
// failure there is logged as a Pyzor error for that message.
let port = pyzor.port;
let host = pyzor.host;
let address = match lookup_host(format!("{host}:{port}"))
.await
.map(|mut a| a.next())
{
Ok(Some(address)) => address,
Ok(None) => {
bp.build_error(
ObjectType::SpamPyzor.singleton(),
"Invalid address: No addresses found.",
);
return None;
}
Err(err) => {
bp.build_error(
ObjectType::SpamPyzor.singleton(),
format!("Invalid address: {}", err),
);
return None;
}
};
PyzorConfig {
host: pyzor.host,
port: pyzor.port as u16,
resolved: Default::default(),
address,
timeout: pyzor.timeout.into_inner(),
min_count: pyzor.block_count,
min_wl_count: pyzor.allow_count,
@@ -495,35 +508,6 @@ impl PyzorConfig {
}
}
// inbuxa: how long a resolved Pyzor address is reused
const PYZOR_RESOLVE_TTL: Duration = Duration::from_secs(300);
impl PyzorConfig {
/// The server's address: the host itself when it is an IP address,
/// otherwise the first address it resolves to, reused for five minutes.
pub async fn address(&self) -> std::io::Result<SocketAddr> {
if let Ok(ip) = self.host.parse::<IpAddr>() {
return Ok(SocketAddr::new(ip, self.port));
}
if let Some((address, resolved_at)) = *self.resolved.lock()
&& resolved_at.elapsed() < PYZOR_RESOLVE_TTL
{
return Ok(address);
}
let address = tokio::net::lookup_host((self.host.as_str(), self.port))
.await?
.next()
.ok_or_else(|| {
std::io::Error::new(
std::io::ErrorKind::NotFound,
format!("{} has no addresses", self.host),
)
})?;
*self.resolved.lock() = Some((address, Instant::now()));
Ok(address)
}
}
impl ClassifierConfig {
pub async fn parse(bp: &mut Bootstrap) -> Option<Self> {
let classifier = bp.setting_infallible::<structs::SpamClassifier>().await;
+14 -13
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 self::resolver::Policy;
@@ -24,7 +22,7 @@ use registry::schema::{
};
use smtp_proto::*;
use std::{
net::{IpAddr, SocketAddr},
net::{SocketAddr, ToSocketAddrs},
str::FromStr,
time::Duration,
};
@@ -386,16 +384,19 @@ impl SessionConfig {
Some(Milter {
enable: bp.compile_expr(id, &milter.ctx_enable()),
id,
// inbuxa: upstream resolved the hostname here (a
// blocking lookup) and made a failure a build error,
// which refused the whole settings reload. An IP
// address is kept as is; a name is resolved on each
// connection (MilterClient::connect).
addrs: milter
.hostname
.parse::<IpAddr>()
.map(|ip| vec![SocketAddr::new(ip, milter.port as u16)])
.unwrap_or_default(),
addrs: format!("{}:{}", milter.hostname, milter.port)
.to_socket_addrs()
.map_err(|err| {
bp.build_error(
id,
format!(
"Unable to resolve milter hostname {}: {}",
milter.hostname, err
),
)
})
.ok()?
.collect(),
hostname: milter.hostname,
port: milter.port as u16,
timeout_connect: milter.timeout_connect.into_inner(),
-69
View File
@@ -335,72 +335,3 @@ impl EmailPush {
}
}
}
/// inbuxa: the task locks this node holds, so a graceful stop can hand them
/// back instead of leaving the tasks blocked until the locks expire.
pub struct TaskLocks {
held: parking_lot::Mutex<ahash::AHashSet<u64>>,
stopping: AtomicBool,
expiry: std::sync::atomic::AtomicU64,
}
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;
pub fn is_stopping(&self) -> bool {
self.stopping.load(Ordering::Acquire)
}
/// Stops new claims and returns the ids of every lock still held.
pub fn stop(&self) -> Vec<u64> {
self.stopping.store(true, Ordering::Release);
self.held.lock().drain().collect()
}
pub fn insert(&self, id: u64) {
self.held.lock().insert(id);
}
pub fn remove(&self, id: u64) {
self.held.lock().remove(&id);
}
pub fn held(&self) -> usize {
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)
}
/// Changes the lock lifetime; the tests shorten it.
pub fn set_expiry(&self, seconds: u64) {
self.expiry.store(seconds.max(1), Ordering::Relaxed);
}
}
impl Default for TaskLocks {
fn default() -> Self {
Self {
held: Default::default(),
stopping: AtomicBool::new(false),
expiry: std::sync::atomic::AtomicU64::new(Self::DEFAULT_EXPIRY),
}
}
}
-8
View File
@@ -161,17 +161,11 @@ 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>>,
pub smtp_connectors: TlsConnectors,
// inbuxa: the objects that failed to build when the running settings
// were built, at boot or by the last applied reload (see reload_registry)
pub build_errors: Mutex<AHashSet<registry::types::id::ObjectId>>,
}
#[derive(Clone)]
@@ -285,8 +279,6 @@ pub struct HttpAuthCache {
pub struct Ipc {
pub push_tx: mpsc::Sender<PushEvent>,
pub task_tx: Arc<Notify>,
// inbuxa: task locks held by this node, released on a graceful stop
pub task_locks: Arc<crate::ipc::TaskLocks>,
pub queue_tx: mpsc::Sender<QueueEvent>,
pub report_tx: mpsc::Sender<ReportingEvent>,
pub broadcast_tx: Option<mpsc::Sender<BroadcastEvent>>,
-4
View File
@@ -240,9 +240,6 @@ impl BootManager {
.parse_tcp_acceptors(&mut bootstrap, inner.clone())
.await;
// inbuxa: a reload isn't refused over objects that failed here
inner.build_server().record_build_errors(&bootstrap.errors);
BootManager {
inner,
bootstrap,
@@ -300,7 +297,6 @@ pub fn build_ipc(has_pubsub: bool) -> (Ipc, IpcReceivers) {
report_tx,
broadcast_tx: has_pubsub.then_some(broadcast_tx),
task_tx: Arc::new(Notify::new()),
task_locks: Arc::new(crate::ipc::TaskLocks::default()),
train_task_controller: Arc::new(TrainTaskController::default()),
},
IpcReceivers {
-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())
+1 -14
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::registry::mapping::{RegistrySetResponse, map_bootstrap_error};
@@ -101,7 +99,7 @@ pub(crate) async fn action_set(
} else {
set.response
.not_created
.append(id, reload_refused(result.errors));
.append(id, map_bootstrap_error(result.errors));
}
}
Action::InvalidateCaches => {
@@ -575,14 +573,3 @@ async fn dmarc_troubleshoot(
Some(request)
}
/// inbuxa: a refused reload names the object that stopped it and says the
/// 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)
);
map_bootstrap_error(errors).with_description(description)
}
+21 -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,30 @@ 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(_) => 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
}
-4
View File
@@ -109,10 +109,6 @@ async fn main() -> std::io::Result<()> {
// Wait for shutdown signal
wait_for_shutdown().await;
// inbuxa: hand back the task locks this node holds, so other nodes can
// run those tasks now rather than when the locks expire
services::task_manager::lock::release_task_locks(&inner.build_server()).await;
// Shutdown collector
Collector::shutdown();
+1 -9
View File
@@ -91,15 +91,7 @@ impl SearchIndexTask for Server {
build_contact_document(self, account_id, document_id).await
}
IndexDocumentType::File => {
// File indexing not implemented yet. inbuxa: still
// one result per task: update_tasks pairs them by
// position, and a missing one shifts every result
// after it onto the wrong task
results.push(IndexTaskResult {
task_type: TaskType::Insert,
index: task.document_type,
result: TaskResult::Ignored,
});
// File indexing not implemented yet
continue;
}
};
+2 -73
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::task_manager::*;
@@ -15,21 +13,13 @@ pub trait TaskLockManager: Sync + Send {
impl TaskLockManager for Server {
async fn try_lock_task(&self, id: u64) -> bool {
// inbuxa: a node that is stopping claims nothing new
let locks = &self.inner.ipc.task_locks;
if locks.is_stopping() {
return false;
}
match self
.in_memory_store()
.try_lock(KV_LOCK_TASK, &id.to_be_bytes(), locks.expiry())
.try_lock(KV_LOCK_TASK, &id.to_be_bytes(), DEFAULT_LOCK_EXPIRY)
.await
{
Ok(result) => {
if result {
locks.insert(id);
} else {
if !result {
trc::event!(
TaskManager(TaskManagerEvent::TaskLocked),
Id = id,
@@ -58,66 +48,5 @@ impl TaskLockManager for Server {
.caused_by(trc::location!())
);
}
self.inner.ipc.task_locks.remove(id);
}
}
/// inbuxa: on a graceful stop, stops claiming tasks and releases every task
/// lock this node holds, so the rest of the cluster can pick the tasks up at
/// once instead of after the lock expires. Returns how many were released.
pub async fn release_task_locks(server: &Server) -> usize {
let ids = server.inner.ipc.task_locks.stop();
for id in &ids {
if let Err(err) = server
.in_memory_store()
.remove_lock(KV_LOCK_TASK, &id.to_be_bytes())
.await
{
trc::error!(
err.details("Failed to release task lock on shutdown")
.ctx(trc::Key::Id, *id)
.caused_by(trc::location!())
);
}
}
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)
}
+161 -272
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::task_manager::acme::AcmeTask;
@@ -13,18 +11,17 @@ 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};
use crate::task_manager::restore_item::RestoreItemTask;
use crate::task_manager::spam_classifier::SpamFilterMaintenanceTask;
use crate::task_manager::{
CLAIM_RECHECK_INTERVAL, Locked, QUEUE_REFRESH_INTERVAL, TaskDetails, TaskFailureType, TaskInfo,
DEFAULT_LOCK_EXPIRY, Locked, QUEUE_REFRESH_INTERVAL, TaskDetails, TaskFailureType, TaskInfo,
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 +56,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 +70,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(),
@@ -151,47 +124,72 @@ pub fn spawn_task_manager(inner: Arc<Inner>) {
let server = inner.build_server();
let batch_size = server.core.email.index_batch_size;
let mut batch = Vec::with_capacity(batch_size);
if let Some(task) = fetch_task(&server, job).await {
batch.push(task);
match server
.store()
.get_value::<Task>(ValueKey::from(ValueClass::TaskQueue(
TaskQueueClass::Task { id: job.id },
)))
.await
{
Ok(Some(task)) => {
batch.push(TaskDetails { task, info: job });
}
Ok(None) => {
trc::event!(
TaskManager(TaskManagerEvent::TaskIgnored),
Id = job.id,
Reason = "Task not found in store, likely already processed.",
);
}
Err(err) => {
trc::error!(
err.id(job.id)
.details("Failed to retrieve task details.")
.caused_by(trc::location!())
);
}
}
while batch.len() < batch_size {
match rx.try_recv() {
Ok(job) => {
if let Some(task) = fetch_task(&server, job).await {
batch.push(task);
match server
.store()
.get_value::<Task>(ValueKey::from(ValueClass::TaskQueue(
TaskQueueClass::Task { id: job.id },
)))
.await
{
Ok(Some(task)) => {
batch.push(TaskDetails { task, info: job });
}
Ok(None) => {
trc::event!(
TaskManager(TaskManagerEvent::TaskIgnored),
Id = job.id,
Reason = "Task not found in store, likely already processed.",
);
}
Err(err) => {
trc::error!(
err.id(job.id)
.details("Failed to retrieve task details.")
.caused_by(trc::location!())
);
}
}
}
Err(_) => break,
}
}
// Dispatch. inbuxa: on a task of its own, so a panic
// releases the batch's locks and leaves this worker
// running; a dead worker would keep claiming tasks it
// can never run
// Dispatch
let mut refresh_queue = false;
let ids = batch.iter().map(|task| task.info.id).collect::<Vec<_>>();
let run = {
let server = server.clone();
tokio::spawn(async move {
let results = server.index(&batch).await;
(batch, results)
})
};
match run.await {
Ok((mut batch, results)) => {
let results = results.into_iter().map(|r| {
refresh_queue |= r.result.is_retry();
r.result
});
update_tasks(&server, &mut batch, results).await;
}
Err(err) => {
worker_failed(&server, &ids, err).await;
refresh_queue = true;
}
}
let results = server.index(&batch).await.into_iter().map(|r| {
refresh_queue |= r.result.is_retry();
r.result
});
update_tasks(&server, &mut batch, results).await;
if refresh_queue || rx.is_empty() {
server.notify_task_queue();
@@ -205,31 +203,83 @@ pub fn spawn_task_manager(inner: Arc<Inner>) {
let server = inner.build_server();
let mut refresh_queue = false;
if let Some(TaskDetails { task, info }) = fetch_task(&server, job).await {
// inbuxa: on a task of its own, as above
let run = {
let server = server.clone();
let server_instance = server_instance.clone();
tokio::spawn(async move {
let result = run_task(&server, &task, server_instance).await;
(task, result)
})
};
match run.await {
Ok((task, result)) => {
refresh_queue = result.is_retry();
match server
.store()
.get_value::<Task>(ValueKey::from(ValueClass::TaskQueue(
TaskQueueClass::Task { id: job.id },
)))
.await
{
Ok(Some(task)) => {
let result = match &task {
Task::CalendarAlarmEmail(task) => {
server.send_email_alarm(task, server_instance.clone()).await
}
Task::CalendarAlarmNotification(task) => {
server.send_display_alarm(task).await
}
Task::CalendarItipMessage(task) => {
server.send_imip(task, server_instance.clone()).await
}
Task::MergeThreads(task) => server.merge_threads(task).await,
Task::DmarcReport(task) => {
server
.submit_report(report::ReportId::Dmarc(task.report_id.id()))
.await
}
Task::TlsReport(task) => {
server
.submit_report(report::ReportId::Tls(task.report_id.id()))
.await
}
Task::RestoreArchivedItem(task) => server.restore_item(task).await,
Task::DestroyAccount(task) => server.destroy_account(task).await,
Task::AccountMaintenance(task) => {
server.account_maintenance(task).await
}
Task::TenantMaintenance(task) => {
server.tenant_maintenance(task).await
}
Task::StoreMaintenance(task) => {
server.store_maintenance(task).await
}
Task::SpamFilterMaintenance(task) => {
Box::pin(server.spam_filter_maintenance(task)).await
}
Task::AcmeRenewal(task) => server.acme_management(task).await,
Task::DkimManagement(task_dkim_rotation) => {
server.dkim_management(task_dkim_rotation).await
}
Task::DnsManagement(task_dns_management) => {
server.dns_management(task_dns_management).await
}
Task::IndexDocument(_)
| Task::UnindexDocument(_)
| Task::IndexTrace(_) => unreachable!(),
};
update_tasks(
&server,
&mut [TaskDetails { task, info }],
vec![result],
)
.await;
}
Err(err) => {
worker_failed(&server, &[info.id], err).await;
refresh_queue = true;
}
refresh_queue = result.is_retry();
update_tasks(
&server,
&mut [TaskDetails { task, info: job }],
vec![result],
)
.await;
}
Ok(None) => {
trc::event!(
TaskManager(TaskManagerEvent::TaskIgnored),
Id = job.id,
Reason = "Task not found in store, likely already processed.",
);
}
Err(err) => {
trc::error!(
err.id(job.id)
.details("Failed to retrieve task details.")
.caused_by(trc::location!())
);
}
}
@@ -268,13 +318,6 @@ pub(crate) trait TaskQueueManager: Sync + Send {
impl TaskQueueManager for Server {
async fn process_tasks(&self, ipc: &mut TaskManagerIpc) -> Duration {
// inbuxa: a node that is stopping has released its locks and claims
// nothing new
let task_locks = &self.inner.ipc.task_locks;
if task_locks.is_stopping() {
return Duration::from_secs(QUEUE_REFRESH_INTERVAL);
}
let lock_expiry = task_locks.expiry();
let now_timestamp = now();
let from_key = ValueKey::<ValueClass> {
account_id: 0,
@@ -314,13 +357,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!(
@@ -337,7 +393,9 @@ impl TaskQueueManager for Server {
let locked = entry.get_mut();
if locked.expires <= now || locked.due < task_due {
locked.expires = Instant::now()
+ std::time::Duration::from_secs(lock_expiry + 1);
+ std::time::Duration::from_secs(
DEFAULT_LOCK_EXPIRY + 1,
);
locked.due = task_due;
tasks.push((
TaskJob {
@@ -353,7 +411,9 @@ impl TaskQueueManager for Server {
Entry::Vacant(entry) => {
entry.insert(Locked {
expires: Instant::now()
+ std::time::Duration::from_secs(lock_expiry + 1),
+ std::time::Duration::from_secs(
DEFAULT_LOCK_EXPIRY + 1,
),
due: task_due,
revision: ipc.revision,
});
@@ -404,26 +464,12 @@ impl TaskQueueManager for Server {
let tx = &ipc.txs[task_type_idx as usize];
if tx.capacity() > 0 {
let id = task_job.id;
if !self.try_lock_task(id).await {
// inbuxa: another node holds the task. Look again after a
// short while rather than a full lock lifetime from now:
// the holder may have claimed it after this scan began,
// or run on a clock ahead of this one, and waiting the
// whole lifetime again would leave the task stuck for
// another hour past its lock if that holder died
if let Some(locked) = ipc.locked.get_mut(&id) {
locked.expires =
Instant::now() + Duration::from_secs(claim_recheck_interval(lock_expiry));
}
} else if tx.send(task_job).await.is_err() {
if self.try_lock_task(task_job.id).await && tx.send(task_job).await.is_err() {
trc::event!(
Server(trc::ServerEvent::ThreadError),
Details = "Error sending task.",
CausedBy = trc::location!()
);
// inbuxa: nothing will run it here, so don't hold it
self.remove_index_lock(id).await;
}
} else {
// If the channel is full, release the lock so it can be picked up in the next iteration
@@ -435,159 +481,9 @@ impl TaskQueueManager for Server {
let now = Instant::now();
ipc.locked
.retain(|_, locked| locked.expires > now && locked.revision == ipc.revision);
let sleep_for = Duration::from_secs(next_event.map_or(QUEUE_REFRESH_INTERVAL, |timestamp| {
Duration::from_secs(next_event.map_or(QUEUE_REFRESH_INTERVAL, |timestamp| {
timestamp.saturating_sub(store::write::now())
}));
// inbuxa: wake up when a claim held elsewhere is due to be tried
// again, rather than only on the next task or refresh
ipc.locked
.values()
.map(|locked| locked.expires.saturating_duration_since(now))
.min()
.map_or(sleep_for, |recheck| sleep_for.min(recheck.max(Duration::from_secs(1))))
}
}
/// 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,
server_instance: Arc<ServerInstance>,
) -> TaskResult {
match task {
Task::CalendarAlarmEmail(task) => {
server.send_email_alarm(task, server_instance.clone()).await
}
Task::CalendarAlarmNotification(task) => {
server.send_display_alarm(task).await
}
Task::CalendarItipMessage(task) => {
server.send_imip(task, server_instance.clone()).await
}
Task::MergeThreads(task) => server.merge_threads(task).await,
Task::DmarcReport(task) => {
server
.submit_report(report::ReportId::Dmarc(task.report_id.id()))
.await
}
Task::TlsReport(task) => {
server
.submit_report(report::ReportId::Tls(task.report_id.id()))
.await
}
Task::RestoreArchivedItem(task) => server.restore_item(task).await,
Task::DestroyAccount(task) => server.destroy_account(task).await,
Task::AccountMaintenance(task) => {
server.account_maintenance(task).await
}
Task::TenantMaintenance(task) => {
server.tenant_maintenance(task).await
}
Task::StoreMaintenance(task) => {
server.store_maintenance(task).await
}
Task::SpamFilterMaintenance(task) => {
Box::pin(server.spam_filter_maintenance(task)).await
}
Task::AcmeRenewal(task) => server.acme_management(task).await,
Task::DkimManagement(task_dkim_rotation) => {
server.dkim_management(task_dkim_rotation).await
}
Task::DnsManagement(task_dns_management) => {
server.dns_management(task_dns_management).await
}
Task::IndexDocument(_)
| Task::UnindexDocument(_)
| Task::IndexTrace(_) => unreachable!(),
}
}
/// 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
/// the lock expired.
async fn fetch_task(server: &Server, job: TaskJob) -> Option<TaskDetails> {
match server
.store()
.get_value::<Task>(ValueKey::from(ValueClass::TaskQueue(TaskQueueClass::Task {
id: job.id,
})))
.await
{
Ok(Some(task)) => Some(TaskDetails { task, info: job }),
Ok(None) => {
trc::event!(
TaskManager(TaskManagerEvent::TaskIgnored),
Id = job.id,
Reason = "Task not found in store, likely already processed.",
);
server.remove_index_lock(job.id).await;
None
}
Err(err) => {
trc::error!(
err.id(job.id)
.details("Failed to retrieve task details.")
.caused_by(trc::location!())
);
server.remove_index_lock(job.id).await;
None
}
}
}
/// inbuxa: a task panicked: its locks are released so it runs again, here or
/// on another node, and the worker carries on.
async fn worker_failed(server: &Server, ids: &[u64], err: tokio::task::JoinError) {
trc::event!(
Server(trc::ServerEvent::ThreadError),
Details = "Task worker failed",
Reason = err.to_string(),
CausedBy = trc::location!()
);
for id in ids {
server.remove_index_lock(*id).await;
}))
}
}
@@ -718,13 +614,6 @@ async fn update_tasks(
}
}
/// inbuxa: how long to wait before trying again to claim a task another node
/// holds: a twelfth of the lock lifetime, so five minutes for the one-hour
/// lock, never more than that and never under a second.
pub(crate) fn claim_recheck_interval(lock_expiry: u64) -> u64 {
(lock_expiry / 12).clamp(1, CLAIM_RECHECK_INTERVAL)
}
pub fn perpetual_retry_time(typ: TaskType, attempt: u64) -> Option<u64> {
matches!(
typ,
+1 -3
View File
@@ -35,9 +35,7 @@ pub mod scheduler;
pub mod spam_classifier;
const QUEUE_REFRESH_INTERVAL: u64 = 60 * 5; // 5 minutes
// inbuxa: the lock lifetime (one hour) lives in common::ipc::TaskLocks, per
// server, so a graceful stop can release the locks and the tests can shorten it
const CLAIM_RECHECK_INTERVAL: u64 = 60 * 5; // 5 minutes
const DEFAULT_LOCK_EXPIRY: u64 = 60 * 60; // 1 hour
pub(crate) struct TaskManagerIpc {
txs: [mpsc::Sender<TaskJob>; TaskType::COUNT],
+1 -15
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 common::config::smtp::session::Milter;
@@ -27,19 +25,7 @@ impl MilterClient<TcpStream> {
pub async fn connect(config: &Milter, session_id: u64) -> Result<Self> {
tokio::time::timeout(config.timeout_command, async {
let mut last_err = Error::Disconnected;
// inbuxa: a hostname is resolved here, per connection, rather
// than while the settings are built
let resolved;
let addrs = if config.addrs.is_empty() {
resolved = tokio::net::lookup_host((config.hostname.as_str(), config.port))
.await
.map_err(Error::Io)?
.collect::<Vec<_>>();
&resolved
} else {
&config.addrs
};
for addr in addrs {
for addr in &config.addrs {
match TcpStream::connect(addr).await {
Ok(stream) => {
return Ok(MilterClient {
+10 -15
View File
@@ -48,24 +48,19 @@ pub(crate) async fn pyzor_check(
// Send message to address. inbuxa: in tests, a fixed table answers
// instead of a public server (test_response).
#[cfg(not(feature = "test_mode"))]
let response = match tokio::time::timeout(config.timeout, config.address()).await {
Ok(Ok(address)) => pyzor_send_message(address, config.timeout, &request).await,
Ok(Err(err)) => Err(err),
Err(_) => Err(std::io::Error::new(
std::io::ErrorKind::TimedOut,
"Timed out resolving the Pyzor server",
)),
};
let response = pyzor_send_message(config.address, config.timeout, &request).await;
#[cfg(feature = "test_mode")]
let response = std::io::Result::Ok(test_response(&request));
response.map(Into::into).map_err(|err| {
trc::SpamEvent::PyzorError
.into_err()
.ctx(trc::Key::Url, format!("{}:{}", config.host, config.port))
.reason(err)
.details("Pyzor failed")
})
response
.map(Into::into)
.map_err(|err| {
trc::SpamEvent::PyzorError
.into_err()
.ctx(trc::Key::Url, config.address.to_string())
.reason(err)
.details("Pyzor failed")
})
}
/// inbuxa: the answers tests get, by digest, instead of a public server's,
+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);
+4 -38
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(
+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 <hello@stalw.art>
*
* 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 <hello@stalw.art>
*
* 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;
}
}
+1 -3
View File
@@ -2,8 +2,6 @@
* SPDX-FileCopyrightText: 2020 Stalwart Labs LLC <hello@stalw.art>
*
* SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-SEL
*
* Modified by Coffey Labs in 2026 for INBUXA.
*/
use crate::{
@@ -77,7 +75,7 @@ async fn milter_session() {
else_: "true".into(),
..Default::default()
},
hostname: "localhost".into(), // inbuxa: resolved when the session connects
hostname: "127.0.0.1".into(),
port: 9332,
use_tls: false,
stages: Map::new(vec![MtaStage::Data]),
+17 -40
View File
@@ -2,8 +2,6 @@
* SPDX-FileCopyrightText: 2020 Stalwart Labs LLC <hello@stalw.art>
*
* 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:?}"),
}
}
-3
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")]
@@ -23,7 +21,6 @@ pub mod replica_cluster; // inbuxa: read replicas across nodes
pub mod scaleout; // inbuxa: scale-out storage
#[cfg(any(feature = "postgres", feature = "mysql"))]
pub mod sql_timeout;
pub mod task_locks; // inbuxa: task locks across nodes
use crate::utils::server::TestServerBuilder;
use std::io::Read;
-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();
}
-210
View File
@@ -1,210 +0,0 @@
/*
* SPDX-FileCopyrightText: 2026 Coffey Labs
*
* SPDX-License-Identifier: AGPL-3.0-only
*/
//! Task locks across nodes: tasks claimed by a node that then disappears
//! run elsewhere once its locks expire, and a node that stops gracefully
//! hands its locks back at once. The other node is played by writing its
//! locks straight into the shared in-memory store, as a node that claimed
//! the tasks and died leaves them.
use crate::utils::server::TestServerBuilder;
use common::{KV_LOCK_TASK, Server};
use registry::schema::{
enums::IndexDocumentType,
structs::{Task, TaskIndexDocument, TaskStatus},
};
use services::task_manager::lock::{TaskLockManager, release_task_locks};
use std::time::{Duration, Instant};
use store::{
ValueKey,
write::{BatchBuilder, TaskQueueClass, ValueClass},
};
use utils::snowflake::SnowflakeIdGenerator;
// Short enough for a test, long enough that the recheck interval (a twelfth
// of it) is well below it
const LOCK_EXPIRY: u64 = 12;
#[tokio::test(flavor = "multi_thread")]
pub async fn task_lock_tests() {
let test = TestServerBuilder::new("task_lock_tests")
.await
.build()
.await;
let server = test.server.clone();
println!(
"Running task lock tests on {}...",
std::env::var("STORE").unwrap_or_default()
);
server.inner.ipc.task_locks.set_expiry(LOCK_EXPIRY);
// 1. Another node claimed the tasks and died. Its locks block them until
// they expire; then this node runs them, without waiting for anything
// else to wake it
let ids = new_task_ids(4);
for id in &ids {
assert!(foreign_lock(&server, *id, LOCK_EXPIRY).await);
}
schedule(&server, &ids).await;
let started = Instant::now();
server.notify_task_queue();
tokio::time::sleep(Duration::from_secs(3)).await;
assert_eq!(pending(&server, &ids).await, ids.len(), "held by the other node");
wait_until_done(&server, &ids, Duration::from_secs(LOCK_EXPIRY + 10)).await;
let elapsed = started.elapsed();
assert!(
elapsed >= Duration::from_secs(LOCK_EXPIRY - 2),
"ran before the other node's locks expired: {elapsed:?}"
);
// 2. The other node's locks outlive what this node expects: claimed just
// after this node looked, or by a node whose clock runs ahead. This node
// keeps checking at the recheck interval, so the tasks run soon after
// those locks expire, not a whole lock lifetime later
let held_for = LOCK_EXPIRY + LOCK_EXPIRY / 2;
let ids = new_task_ids(4);
for id in &ids {
assert!(foreign_lock(&server, *id, held_for).await);
}
schedule(&server, &ids).await;
let started = Instant::now();
server.notify_task_queue();
wait_until_done(&server, &ids, Duration::from_secs(held_for + 8)).await;
let elapsed = started.elapsed();
assert!(
elapsed >= Duration::from_secs(held_for - 2),
"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
// can claim those tasks at once, and this one claims nothing more
let ids = new_task_ids(3);
for id in &ids {
assert!(server.try_lock_task(*id).await, "claim {id}");
}
assert_eq!(server.inner.ipc.task_locks.held(), ids.len());
for id in &ids {
assert!(
!foreign_lock(&server, *id, LOCK_EXPIRY).await,
"held while this node runs"
);
}
assert_eq!(release_task_locks(&server).await, ids.len());
assert_eq!(server.inner.ipc.task_locks.held(), 0);
for id in &ids {
assert!(
foreign_lock(&server, *id, LOCK_EXPIRY).await,
"released on stop: {id}"
);
}
let [id] = new_task_ids(1)[..] else {
unreachable!()
};
assert!(!server.try_lock_task(id).await, "a stopping node claims nothing");
for id in ids {
let _ = server
.in_memory_store()
.remove_lock(KV_LOCK_TASK, &id.to_be_bytes())
.await;
}
if test.is_reset() {
test.temp_dir.delete();
}
}
fn new_task_ids(count: usize) -> Vec<u64> {
(0..count)
.map(|_| SnowflakeIdGenerator::global_id().unwrap())
.collect()
}
/// The other node's claim on a task, as its task manager takes it.
async fn foreign_lock(server: &Server, id: u64, seconds: u64) -> bool {
server
.in_memory_store()
.try_lock(KV_LOCK_TASK, &id.to_be_bytes(), seconds)
.await
.unwrap()
}
/// Unindex tasks for files that don't exist: files aren't search-indexed and
/// there is no undelete note, so running one only drops it from the queue.
async fn schedule(server: &Server, ids: &[u64]) {
let mut batch = BatchBuilder::new();
for (n, id) in ids.iter().enumerate() {
batch.schedule_task_with_id(
*id,
Task::UnindexDocument(TaskIndexDocument {
account_id: 0u32.into(),
document_id: (u32::MAX - n as u32).into(),
document_type: IndexDocumentType::File,
status: TaskStatus::now(),
}),
);
}
server.store().write(batch.build_all()).await.unwrap();
}
async fn pending(server: &Server, ids: &[u64]) -> usize {
let mut count = 0;
for id in ids {
if server
.store()
.get_value::<Task>(ValueKey::from(ValueClass::TaskQueue(
TaskQueueClass::Task { id: *id },
)))
.await
.unwrap()
.is_some()
{
count += 1;
}
}
count
}
async fn wait_until_done(server: &Server, ids: &[u64], within: Duration) {
let started = Instant::now();
loop {
let left = pending(server, ids).await;
if left == 0 {
return;
}
assert!(
started.elapsed() < within,
"{left} task(s) still pending after {:?}",
started.elapsed()
);
tokio::time::sleep(Duration::from_millis(250)).await;
}
}
-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)
}
-2
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;
@@ -21,7 +20,6 @@ pub mod monitoring;
pub mod oidc;
pub mod purge;
pub mod quota;
pub mod reload; // inbuxa: reloads and build errors
pub mod security;
pub mod task;
pub mod tenant;
-213
View File
@@ -1,213 +0,0 @@
/*
* SPDX-FileCopyrightText: 2026 Coffey Labs
*
* SPDX-License-Identifier: AGPL-3.0-only
*/
// inbuxa: a settings reload isn't held back by a DNS lookup, or by objects
// that already failed when the running settings were built; an error in an
// object that built then still refuses it, and says which object.
use crate::utils::server::{TestServer, TestServerBuilder};
use common::{BuildServer, config::mailstore::spamfilter::PyzorConfig, ipc::RegistryChange};
use registry::schema::{
enums::TracingLevel,
prelude::{ObjectType, Property},
structs::{Action, Expression, MtaStageAuth, SpamPyzor, Tracer, TracerStdout},
};
#[tokio::test(flavor = "multi_thread")]
pub async fn reload_tests() {
let mut test = TestServerBuilder::new("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_unresolvable_pyzor(&test).await;
test_build_errors(&test).await;
if test.is_reset() {
test.temp_dir.delete();
}
}
async fn test_unresolvable_pyzor(test: &TestServer) {
println!("Running reload with an unresolvable Pyzor host...");
let admin = test.account("[email protected]");
// Upstream resolved the host while building the settings and refused the
// reload when that failed.
admin
.registry_update_setting(
SpamPyzor {
enable: true,
host: "pyzor.invalid".into(),
port: 24441,
..Default::default()
},
&[Property::Enable, Property::Host, Property::Port],
)
.await;
admin.reload_settings().await;
let pyzor = running_pyzor(test);
assert_eq!(pyzor.host, "pyzor.invalid");
assert_eq!(pyzor.port, 24441);
assert!(pyzor.address().await.is_err());
// An IP address needs no lookup
admin
.registry_update_setting(
SpamPyzor {
host: "192.0.2.1".into(),
..Default::default()
},
&[Property::Host],
)
.await;
admin.reload_settings().await;
assert_eq!(
running_pyzor(test).address().await.unwrap().to_string(),
"192.0.2.1:24441"
);
}
async fn test_build_errors(test: &TestServer) {
println!("Running reload with build errors...");
let admin = test.account("[email protected]");
let pyzor_ratio = running_pyzor(test).ratio;
assert_ne!(pyzor_ratio, 0.25);
// Two console tracers: only one is allowed, so the build of one of them
// fails. Neither existed when the running settings were built.
let mut tracer_ids = Vec::new();
for _ in 0..2 {
tracer_ids.push(
admin
.registry_create_object(Tracer::Stdout(TracerStdout {
enable: true,
level: TracingLevel::Error,
..Default::default()
}))
.await,
);
}
admin
.registry_update_setting(
SpamPyzor {
ratio: 0.25.into(),
..Default::default()
},
&[Property::Ratio],
)
.await;
// A new error refuses the reload and names the object
let err = admin
.registry_create_object_expect_err(Action::ReloadSettings)
.await;
let description = err.description.clone().unwrap_or_default();
assert!(
description.starts_with("Settings were not reloaded. ")
&& description.contains("Tracer")
&& description.contains("Only one console tracer is allowed"),
"{err:?}"
);
assert_eq!(running_pyzor(test).ratio, pyzor_ratio);
// Had the running settings been built with that tracer failing, as a
// restart now would, the same error doesn't hold the reload back.
let result = Box::pin(
test.server
.reload_registry(RegistryChange::Reload(ObjectType::DataStore)),
)
.await
.unwrap();
assert!(!result.replaced_core);
assert_eq!(result.errors.len(), 1, "{:?}", result.errors);
test.server.record_build_errors(&result.errors);
admin.reload_settings().await;
assert_eq!(running_pyzor(test).ratio, 0.25);
let result = Box::pin(
test.server
.reload_registry(RegistryChange::Reload(ObjectType::DataStore)),
)
.await
.unwrap();
assert!(result.replaced_core);
assert!(result.errors.is_empty());
assert_eq!(result.known_errors.len(), 1);
// Once fixed, the object is no longer known to fail, so a new error
// there refuses the reload again.
admin
.registry_destroy(ObjectType::Tracer, tracer_ids.iter())
.await
.assert_destroyed(&tracer_ids);
admin.reload_settings().await;
let result = Box::pin(
test.server
.reload_registry(RegistryChange::Reload(ObjectType::DataStore)),
)
.await
.unwrap();
assert!(result.replaced_core);
assert!(result.errors.is_empty() && result.known_errors.is_empty());
for _ in 0..2 {
tracer_ids.push(
admin
.registry_create_object(Tracer::Stdout(TracerStdout {
enable: true,
level: TracingLevel::Error,
..Default::default()
}))
.await,
);
}
admin
.registry_create_object_expect_err(Action::ReloadSettings)
.await;
let tracer_ids = tracer_ids.split_off(2);
admin
.registry_destroy(ObjectType::Tracer, tracer_ids.iter())
.await
.assert_destroyed(&tracer_ids);
admin.reload_settings().await;
}
fn running_pyzor(test: &TestServer) -> PyzorConfig {
test.server
.inner
.build_server()
.core
.spam
.pyzor
.clone()
.expect("Pyzor enabled")
}