Author SHA1 Message Date
jcoffey-dev 95f0445d83 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 11m8s
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:38:59 -07:00
jcoffey-dev c974a0918e Merge pull request 'Task manager: release task locks on stop, recheck claims held elsewhere' (#35) from fix/task-lock-recovery into main
ci / build (push) Canceled after 6m47s
ci / fork-checks (push) Successful in 49s
2026-09-24 15:38:39 +00:00
jcoffey-dev 7c80a12d75 Task manager: release task locks on stop, recheck claims held elsewhere
ci / build (pull_request) Successful in 17m2s
ci / fork-checks (pull_request) Successful in 18s
A cluster rehearsal (PostgreSQL + NATS) left index tasks pending well
past the one-hour task lock after the node that claimed them was stopped
or killed. The exact cause there isn't confirmed; this closes every path
found in the task manager that stretches a takeover past the lock, or
keeps a task claimed without running it:

- A graceful stop never released the locks it held, so every task the
  node had claimed stayed blocked for an hour. The server now tracks the
  locks it holds (common::ipc::TaskLocks) and, once the shutdown signal
  arrives, stops claiming and releases them before exiting.
- A node that failed to claim a task (another node held it) set its own
  local hold for a full lock lifetime from that scan. If the holder
  claimed it just after the scan began, or ran on a clock ahead, that
  hold ran out a moment before the lock did and was set for another
  hour: two hours in all. Such claims are now tried again every five
  minutes (a twelfth of the lock lifetime), and the task manager wakes
  up for them: before, a node without a coordinator could sleep up to
  five minutes past the recheck, or until something else woke it.
- A worker that panicked took its task type down on that node for good,
  while the scan kept claiming that type's tasks and failing to hand them
  over, re-taking each lock as it expired and so starving every other
  node of them. Each batch now runs on a task of its own; a panic is
  logged, the batch's locks are released and the worker carries on. A
  failed hand-over releases the lock too.
- A claimed task the worker couldn't read, or found gone, kept its lock
  for the hour. It is released.
- An IndexDocument task for a file (not indexed) returned no result,
  which shifted every later result in the batch onto the wrong task in
  update_tasks. It returns Ignored. Nothing queues such a task today.

The lock lifetime stays one hour; it now lives per server so the tests
can shorten it.

store::task_locks::task_lock_tests plays a second node by writing its
locks straight into the in-memory store: tasks it claimed and abandoned
run here once its locks expire, including locks that outlive this node's
view of them, and a graceful stop hands this node's locks back at once
and claims nothing more. It passes on RocksDB, SQLite and PostgreSQL.
With the old recheck it fails.
2026-09-24 08:19:05 -07:00
jcoffey-dev 7109e67f07 Merge pull request 'Trace search: index event type and queue id as integers' (#33) from fix/pg-index-trace-types into main
ci / fork-checks (push) Successful in 30s
ci / build (push) Successful in 37m5s
2026-09-24 14:57:18 +00:00
jcoffey-dev 52b5a5f909 Mark tests/src/store/query.rs as modified by the fork
ci / fork-checks (pull_request) Successful in 1m4s
ci / build (pull_request) Successful in 4m20s
The trace document test changed an upstream file, so it carries the
AGPL section 5(a) notice (tools/fork/notice-check.py).
2026-09-24 07:39:04 -07:00
jcoffey-dev 9232662913 Trace search: index event type and queue id as integers
ci / fork-checks (pull_request) Failing after 47s
ci / build (pull_request) Successful in 4m55s
The trace index task wrote the event type (its name) and the queue id as
text, but the tracing search index types both as integers on every
backend: BIGINT on PostgreSQL and MySQL, long on Elasticsearch. On
PostgreSQL every batch holding a trace document failed with "cannot
convert between the Rust type String and the Postgres type int8", and
since a batch writes trace and email documents together, email indexing
stalled behind it.

The document is now built by trace_search_document(), which writes:

- the event type as the opening event's numeric id, the event
  x:Trace/query's event filter already matches on;
- the queue id as an integer, the first one the trace names;
- every queue id into the keywords as well, since the column holds one
  value and an SMTP session can queue several messages.

index_keyword() replaced the field on every call, so before this only the
last event type and queue id survived anyway.

x:Trace/query's queueId filter parses the id (a string, or now a number)
and matches the column or the keywords, so a session is found by any of
its queue ids on every backend. The monitoring spec says what is indexed.

Traces indexed before this on the built-in index keep their text values;
the reindexTelemetry maintenance task rebuilds them.

Tests: the search store suite builds trace documents with the index
task's code, indexes them and finds them by queue id, event type and
keyword (Sqlite, PostgreSQL, MySQL); the monitoring suite finds a real
trace by queueId through x:Trace/query.
2026-09-24 07:29:08 -07:00
jcoffey-dev 57d1c5b074 Merge pull request 'Broadcast subscriber: fix the inverted subscribe retry backoff' (#32) from fix/subscriber-backoff into main
ci / fork-checks (push) Successful in 43s
ci / build (push) Canceled after 30m24s
2026-09-24 14:26:52 +00:00
jcoffey-dev ca3abf40f0 Broadcast subscriber: fix the inverted subscribe retry backoff
ci / fork-checks (pull_request) Successful in 56s
ci / build (pull_request) Successful in 5m14s
The broadcast subscriber waited 1 << retry_count.max(6) seconds between
failed subscribe attempts. max(6) turns the cap into a floor: the first
retry waited 64 s instead of 1 s, and each later one doubled without a
bound (and would overflow the shift after enough failures).

The delay now comes from subscribe_retry_delay(), 1 s, 2 s, 4 s ... capped
at 64 s, and the retry counter saturates. A unit test pins the schedule
and the top of the range.
2026-09-24 07:01:56 -07:00
jcoffey-dev 499e4d7810 Merge pull request 'Export/import: keep archived items, spam samples and the spam model' (#31) from fix/export-all-subspaces into main
ci / fork-checks (push) Successful in 3m24s
ci / build (push) Successful in 43m50s
2026-09-23 08:50:16 +00:00
jcoffey-dev 212cd77cd3 Export/import: keep archived items, spam samples and the spam model
ci / fork-checks (pull_request) Successful in 32s
ci / build (pull_request) Successful in 7m57s
--export skipped three things, so a move from one database to another
(RocksDB to PostgreSQL, say) lost them without a word:

- archived items (subspace j), the records behind undelete;
- spam training samples (subspace w);
- the trained spam classifier and its trainer state, blobs stored under
  fixed names that no blob link points at, so the walk over links never
  reached them.

j and w now travel with the registry family, where their indexes and id
counters already were, so EXPORT_TYPES=registry keeps them consistent.
The two named blobs travel with the blob family. The file format is
unchanged and import reads any subspace it is given, so an export made
by an older binary still imports.

The full-text index (subspace z) stays out, on purpose. It belongs to one
search backend: PostgreSQL and MySQL index into their own tables and have
no z table at all, and external engines keep the index themselves. So
--import now returns the subspaces it wrote, and boot queues the
reindexAccounts and reindexTelemetry store maintenance tasks, the same
ones an administrator can queue by hand, to rebuild the index for
whichever search store the server runs with once it starts.

The round trip also turned up a loss in import itself: the SQL stores
add a negative amount with an UPDATE, which does nothing to a row that
isn't there yet, so every negative counter or quota vanished on import
into PostgreSQL, MySQL or SQLite. Import now creates the row first.

The in-memory subspaces (m, y) stay out: rate limits, locks, greylisting,
ACME challenge tokens and OAuth codes, all short-lived. Issued
certificates are registry objects and travel.

The store test now writes archived items, spam samples, directory
entries, the fork's own subspace and the named blobs, checks they come
back in place, then imports the same export into a fresh store of the
other local backend (RocksDB to SQLite, or SQLite to RocksDB), compares
it key for key and counter for counter, and checks the queued reindex.
It fails on the old export code ("Subspace j was not exported").
--help now says what an export holds.
2026-09-23 01:41:28 -07:00
jcoffey-dev 7735780807 Merge pull request 'Image build: put the vendored crate where cargo chef cooks; release 2026.9.24.3' (#30) from fix/image-vendor-before-cook into main
ci / fork-checks (push) Successful in 41s
publish / version (push) Successful in 39s
ci / build (push) Successful in 36m33s
publish / publish (push) Successful in 1h2m24s
publish / release (push) Successful in 2s
publish / binaries (push) Successful in 1m8s
Reviewed-on: #30
2026-09-23 06:10:46 +00:00
jcoffey-dev e223f7d327 Release 2026.9.24.3
ci / fork-checks (pull_request) Successful in 45s
ci / build (pull_request) Successful in 4m26s
2026.9.24.3 is 2026.9.24.2 plus the image build fix; 2026.9.24.2's tag never
published an image. Everything in 2026.9.24.2's notes applies.
2026-09-22 23:04:59 -07:00
jcoffey-dev 30df055e39 Image build: put the vendored crate where cargo chef cooks
#27 let the build context see vendor/, but the Dockerfile cooks the
dependencies before it copies the tree, from a recipe that carries only the
workspace's manifests. [patch.crates-io] points sieve-rs at vendor/, so the
cook failed the same way: failed to read /build/vendor/sieve-rs/Cargo.toml.
That's why 2026.9.24.2's publish failed. The builder stage now copies
vendor/ before cooking; a local build got past it into compiling the
dependencies.

context-check.py now also checks that each patched path is copied into the
cooking stage before the cook, and fails on the Dockerfile as it was.
2026-09-22 23:04:50 -07:00
jcoffey-dev 5393c4405a Merge pull request 'Release 2026.9.24.2' (#29) from release/2026.9.24.2 into main
ci / fork-checks (push) Successful in 50s
publish / version (push) Successful in 24s
publish / publish (push) Failing after 2m38s
publish / release (push) Skipped
publish / binaries (push) Skipped
ci / build (push) Successful in 22m35s
Reviewed-on: #29
2026-09-23 05:47:10 +00:00
jcoffey-dev cc532b914c Release 2026.9.24.2
ci / fork-checks (pull_request) Successful in 1m49s
ci / build (pull_request) Successful in 4m8s
Replaces 2026.9.24, whose tag predates the image build fix (#27) and never
published. Carries everything 2026.9.24 did -- upstream 0.16.23 and its
fixes, the scim release-profile fix -- and since then:

- identifiers renamed from the upstream name, with no aliases: the JMAP
  registry capability is urn:inbuxa:jmap:registry, WebDAV tokens
  urn:inbuxa:dav*, Sieve extensions vnd.inbuxa.*, the web interface client
  inbuxa-webui; INBUXA_* settings only. Deploy with admin and webmail
  releases that use the new names.
- the brand in lowercase where people see it.
- the spam filter rules bundled with the server; on first start they add
  the AI classifier's LLM_* scores.
- a Local AI page link in Settings › Spam Filter, for the admin release
  that draws it.
- two start-up migrations: the spam model moves to its renamed keys, and
  the web interface's old OAuth client is retired.
2026-09-22 22:40:17 -07:00
jcoffey-dev c09eff2214 Merge pull request 'Bundle the spam filter rules with the server, and link the Local AI page' (#28) from fork/bundled-spam-rules into main
ci / fork-checks (push) Successful in 20s
ci / build (push) Canceled after 7m16s
Reviewed-on: #28
2026-09-23 05:39:50 +00:00
jcoffey-dev d7c9416713 Merge pull request 'Let the image build see the dependency Cargo patches' (#27) from fix/vendor-in-build-context into main
ci / fork-checks (push) Successful in 14s
ci / build (push) Successful in 31m11s
2026-09-23 04:48:35 +00:00
jcoffey-dev 238079da66 Let the image build see the dependency Cargo patches
ci / fork-checks (pull_request) Successful in 49s
ci / build (pull_request) Successful in 4m22s
The rename pass vendored a patched sieve-rs and pointed Cargo.toml's
[patch.crates-io] at vendor/sieve-rs. .dockerignore ignores everything and
re-includes a short list that did not have vendor on it, so the image build
had no such directory and stopped at

    failed to load source for dependency `sieve-rs`
    failed to read /build/vendor/sieve-rs/Cargo.toml

CI could not have caught that: it builds from a checkout, where the
directory is simply there, and only the image build has a context to prune.
The first that was known about it was a tag that had already been pushed.

So: vendor is re-included, and tools/fork/context-check.py now asserts the
thing that was quietly assumed -- every path a [patch] section names exists
and survives .dockerignore. It runs beside the other fork checks and takes
no toolchain.

Also, the comments in .dockerignore started with // , which Docker does not
read as a comment: they were patterns that happened to match nothing. They
are # now.
2026-09-22 21:43:37 -07:00
33 changed files with 1678 additions and 211 deletions
+8 -2
View File
@@ -1,10 +1,16 @@
// Ignore everything
# Ignore everything
*
// Allow what is needed
# Allow what is needed
!crates
!tests
!resources
# The patched dependency Cargo.toml's [patch.crates-io] points at. Without
# it the build context has no vendor/, and `cargo chef cook` fails on
# "failed to load source for dependency sieve-rs" -- which CI cannot see,
# because CI builds from a checkout and only the image build has a context.
!vendor
!Cargo.lock
!Cargo.toml
+5
View File
@@ -35,6 +35,11 @@ jobs:
- run: python3 tools/fork/name-check.py
- if: always()
run: python3 tools/fork/notice-check.py
# Cargo can patch a dependency to a directory in this repository, and
# the image builds from a context .dockerignore prunes to almost
# nothing. CI never sees the difference; a release does.
- if: always()
run: python3 tools/fork/context-check.py
build:
# Either runner (host1 or host2): the build needs no docker socket.
+4
View File
@@ -19,6 +19,10 @@ RUN export DEBIAN_FRONTEND=noninteractive && \
g++-x86-64-linux-gnu binutils-x86-64-linux-gnu
RUN rustup target add "$(cat /target.txt)"
COPY --from=planner /recipe.json /recipe.json
# inbuxa: [patch.crates-io] points sieve-rs at vendor/, and the recipe only
# carries the workspace's own manifests, so cooking the dependencies needs the
# vendored crate itself (the context allows it since #27; this puts it here).
COPY vendor/ vendor/
RUN RUSTFLAGS="$(cat /flags.txt)" cargo chef cook --target "$(cat /target.txt)" --release --no-default-features --features "sqlite postgres mysql rocks s3 redis azure nats" --recipe-path /recipe.json
COPY . .
RUN RUSTFLAGS="$(cat /flags.txt)" cargo build --target "$(cat /target.txt)" --release -p inbuxa --no-default-features --features "sqlite postgres mysql rocks s3 redis azure nats"
+54
View File
@@ -335,3 +335,57 @@ 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.
pub const DEFAULT_EXPIRY: u64 = 60 * 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()
}
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),
}
}
}
+2
View File
@@ -279,6 +279,8 @@ 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>>,
+25 -6
View File
@@ -23,6 +23,13 @@ use utils::{UnwrapFailure, codec::leb128::Leb128_};
pub(super) const MAGIC_MARKER: u8 = 123;
// inbuxa: blobs kept under a fixed name instead of a content hash. Nothing
// links to them, so the export names them outright.
const NAMED_BLOBS: &[&[u8]] = &[
crate::manager::SPAM_CLASSIFIER_KEY,
crate::manager::SPAM_TRAINER_KEY,
];
#[derive(Debug, Clone, Copy, Hash, PartialEq, Eq)]
pub(super) enum Family {
Data = 0,
@@ -143,15 +150,21 @@ impl Core {
.await
.failed("Failed to iterate over data store");
for hash in blobs {
// inbuxa: the trained spam classifier and its trainer state are
// blobs stored under fixed names with no blob link, so the walk
// over links above never reaches them.
let named = NAMED_BLOBS.iter().map(|key| key.to_vec());
for key in blobs
.into_iter()
.map(|hash| hash.as_slice().to_vec())
.chain(named)
{
if let Some(blob) = blob_store
.get_blob(hash.as_slice(), 0..usize::MAX)
.get_blob(&key, 0..usize::MAX)
.await
.failed("Failed to get blob")
{
writer
.send((hash.as_slice().to_vec(), blob))
.failed("Failed to send key");
writer.send((key, blob)).failed("Failed to send key");
}
}
}),
@@ -323,7 +336,13 @@ impl Family {
SUBSPACE_REGISTRY_IDX,
SUBSPACE_REGISTRY_PK,
SUBSPACE_DIRECTORY,
store::SUBSPACE_INBUXA, // inbuxa: masked email
// inbuxa: registry objects the upstream list left out, so an
// export dropped them: archived items (undelete) and spam
// training samples. Their indexes and id counters already
// travel in this family and in `data`, so they ride along.
SUBSPACE_DELETED_ITEMS,
SUBSPACE_SPAM_SAMPLES,
store::SUBSPACE_INBUXA, // inbuxa: the fork's own data (masked email, undelete, policies)
],
Family::Changelog => &[SUBSPACE_LOGS],
Family::Queue => &[SUBSPACE_QUEUE_MESSAGE, SUBSPACE_QUEUE_EVENT],
+12 -4
View File
@@ -54,6 +54,13 @@ Options:
-o, --console Open the store console
-h, --help Print help
-V, --version Print version
An export holds everything in the data and blob stores except short-lived
in-memory state (rate limits, locks, greylisting) and the full-text search
index, which belongs to one search backend. An import into an empty store
queues the index to be rebuilt when the server next starts. EXPORT_TYPES
limits an export to some of: data, registry, blob, changelog, queue, report,
telemetry, tasks.
"#
);
@@ -256,10 +263,10 @@ impl BootManager {
telemetry.enable();
// Parse settings and restore
Box::pin(Core::parse(&mut bootstrap, storage))
.await
.restore(path)
.await;
let core = Box::pin(Core::parse(&mut bootstrap, storage)).await;
let imported = core.restore(path).await;
// inbuxa: the search index isn't exported; rebuild it
core.queue_reindex(&imported).await;
std::process::exit(0);
}
StoreOp::Console => {
@@ -290,6 +297,7 @@ 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 {
+79 -10
View File
@@ -9,15 +9,22 @@
use super::backup::MAGIC_MARKER;
use crate::{Core, DATABASE_SCHEMA_VERSION};
use lz4_flex::frame::FrameDecoder;
use registry::schema::enums::CompressionAlgo;
use registry::{
schema::{
enums::{CompressionAlgo, TaskStoreMaintenanceType},
structs::{Task, TaskStatus, TaskStoreMaintenance},
},
types::EnumImpl,
};
use std::{
fs::File,
io::{BufReader, ErrorKind, Read},
path::{Path, PathBuf},
};
use store::{
BlobStore, IterateParams, SUBSPACE_BLOBS, SUBSPACE_COUNTER, SUBSPACE_INDEXES, SUBSPACE_QUOTA,
SUBSPACE_REGISTRY_PK, Store, U32_LEN,
BlobStore, IterateParams, SUBSPACE_BLOBS, SUBSPACE_COUNTER, SUBSPACE_INDEXES,
SUBSPACE_PROPERTY, SUBSPACE_QUOTA, SUBSPACE_REGISTRY_PK, SUBSPACE_TELEMETRY_SPAN, Store,
U32_LEN,
write::{
AnyClass, AnyKey, BatchBuilder, ValueClass,
key::{DeserializeBigEndian, is_node_id_key},
@@ -27,7 +34,9 @@ use types::{collection::Collection, field::Field};
use utils::{UnwrapFailure, failed};
impl Core {
pub async fn restore(&self, src: PathBuf) {
/// Imports an export into an empty store and returns the subspaces it
/// wrote. inbuxa: the caller hands them to [`Core::queue_reindex`].
pub async fn restore(&self, src: PathBuf) -> Vec<u8> {
// Backup the core
let paths = if src.is_dir() {
let mut paths = Vec::new();
@@ -64,6 +73,13 @@ impl Core {
std::process::exit(1);
}
let mut imported = paths
.iter()
.map(|path| KeyValueReader::new(path).subspace)
.collect::<Vec<_>>();
imported.sort_unstable();
imported.dedup();
let mut tasks = Vec::new();
for path in paths {
let storage = self.storage.clone();
@@ -76,6 +92,54 @@ impl Core {
for task in tasks {
task.await.failed("Failed to wait for task");
}
imported
}
/// inbuxa: an export never carries the full-text index. It is built by
/// and for one search backend (the SQL stores index into their own
/// tables, the key-value stores into a subspace, external engines keep it
/// themselves), so it would be wrong or unreadable after a move to
/// another one. Instead, an import queues the same reindex tasks an
/// administrator can queue by hand (`reindexAccounts` and
/// `reindexTelemetry` store maintenance), and the server rebuilds the
/// index for whatever search store it is configured with once it starts.
pub async fn queue_reindex(&self, imported: &[u8]) -> Vec<TaskStoreMaintenanceType> {
let mut queued = Vec::new();
if imported.contains(&SUBSPACE_PROPERTY) {
queued.push(TaskStoreMaintenanceType::ReindexAccounts);
}
if imported.contains(&SUBSPACE_TELEMETRY_SPAN) {
queued.push(TaskStoreMaintenanceType::ReindexTelemetry);
}
if queued.is_empty() {
return queued;
}
let mut batch = BatchBuilder::new();
for maintenance_type in &queued {
batch.schedule_task(Task::StoreMaintenance(TaskStoreMaintenance {
maintenance_type: *maintenance_type,
status: TaskStatus::now(),
shard_index: None,
}));
}
self.storage
.data
.write(batch.build_all())
.await
.failed("Failed to queue the reindex tasks");
println!(
"Queued {} to rebuild the search index; it runs when the server starts.",
queued
.iter()
.map(|t| t.as_str())
.collect::<Vec<_>>()
.join(" and ")
);
queued
}
}
@@ -125,17 +189,22 @@ async fn restore_file(store: Store, blob_store: BlobStore, path: &Path) {
}
SUBSPACE_COUNTER | SUBSPACE_QUOTA => {
while let Some((key, value)) = reader.next() {
batch.add(
ValueClass::Any(AnyClass {
let class = ValueClass::Any(AnyClass {
subspace: reader.subspace,
key,
}),
u64::from_le_bytes(
});
let value = u64::from_le_bytes(
value
.try_into()
.expect("Failed to deserialize counter/quota"),
) as i64,
);
) as i64;
// inbuxa: the SQL stores add a negative amount with an UPDATE,
// which does nothing to a row that isn't there yet, so a
// negative counter vanished on import. Create the row first.
if value < 0 {
batch.add(class.clone(), 0);
}
batch.add(class, value);
if batch.is_large_batch() {
store
.write(batch.build_all())
+1 -1
View File
@@ -8,7 +8,7 @@ store = { path = "../store" }
registry = { path = "../registry" }
trc = { path = "../trc" }
futures = { version = "0.3", optional = true }
tokio = { version = "1.53", features = ["sync", "fs", "io-util"] }
tokio = { version = "1.53", features = ["sync", "fs", "io-util", "rt", "time"] }
async-nats = { version = "0.50", default-features = false, features = ["server_2_10", "server_2_11", "aws-lc-rs"], optional = true }
zenoh = { version = "1.10.0", default-features = false, features = ["auth_pubkey", "transport_multilink", "transport_compression", "transport_quic", "transport_tcp", "transport_tls", "transport_udp"], optional = true }
rdkafka = { version = "0.39", features = ["cmake-build"], optional = true }
+118 -2
View File
@@ -2,13 +2,22 @@
* 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::sync::Arc;
use std::{
sync::{
Arc,
atomic::{AtomicBool, Ordering},
},
time::Duration,
};
use crate::Coordinator;
use async_nats::Client;
use registry::schema::structs::NatsCoordinator;
use trc::ClusterEvent;
pub mod pubsub;
@@ -47,9 +56,116 @@ impl NatsPubSub {
opts = opts.token(credentials);
}
// inbuxa: connect in the background and keep trying, so a node that
// starts while NATS is down still joins the cluster once NATS is
// back, instead of running without a coordinator until restarted;
// and report the connection going and coming back
let reporter = Arc::new(Reporter::default());
opts = opts.retry_on_initial_connect().event_callback({
let reporter = reporter.clone();
move |event| {
let reporter = reporter.clone();
async move { reporter.report(event) }
}
});
let connection_timeout = config.timeout_connection.into_inner();
async_nats::connect_with_options(config.addresses.into_inner(), opts)
.await
.map(|client| Coordinator::Nats(Arc::new(NatsPubSub { client })))
.map(|client| {
reporter.watch_first_connection(client.clone(), connection_timeout);
Coordinator::Nats(Arc::new(NatsPubSub { client }))
})
.map_err(|err| format!("Failed to connect to Nats: {}", err))
}
/// inbuxa: whether the client is connected to a NATS server right now.
pub fn is_connected(&self) -> bool {
matches!(
self.client.connection_state(),
async_nats::connection::State::Connected
)
}
}
/// inbuxa: reports the client's connection events as the server's own.
#[derive(Default)]
struct Reporter {
connected_once: AtomicBool,
// A failed attempt raises an error each time the client retries, every
// few seconds while NATS is down: report the first after each change
error_reported: AtomicBool,
}
impl Reporter {
fn report(&self, event: async_nats::Event) {
match event {
async_nats::Event::Connected => {
self.connected_once.store(true, Ordering::Relaxed);
self.error_reported.store(false, Ordering::Relaxed);
trc::event!(Cluster(ClusterEvent::CoordinatorConnected), Type = "nats");
}
async_nats::Event::Disconnected => {
self.error_reported.store(false, Ordering::Relaxed);
trc::event!(
Cluster(ClusterEvent::CoordinatorDisconnected),
Type = "nats",
Details = "Connection lost; reconnecting in the background",
);
}
async_nats::Event::Closed => {
trc::event!(
Cluster(ClusterEvent::CoordinatorDisconnected),
Type = "nats",
Details = "Connection closed; no further attempts will be made",
);
}
async_nats::Event::ClientError(async_nats::ClientError::MaxReconnects) => {
trc::event!(
Cluster(ClusterEvent::CoordinatorDisconnected),
Type = "nats",
Details = "Gave up reconnecting (maxReconnects reached)",
);
}
async_nats::Event::ClientError(err) => {
if !self.error_reported.swap(true, Ordering::Relaxed) {
trc::event!(
Cluster(ClusterEvent::CoordinatorError),
Type = "nats",
Details = "Connection attempt failed; retrying",
Reason = err.to_string(),
);
}
}
event => {
trc::event!(
Cluster(ClusterEvent::CoordinatorError),
Type = "nats",
Details = event.to_string(),
);
}
}
}
/// The first connection is made in the background, so say so when it
/// hasn't been made within the connection timeout. The client keeps
/// trying, and reports the connection when it comes.
fn watch_first_connection(self: &Arc<Self>, client: Client, timeout: Duration) {
let reporter = self.clone();
tokio::spawn(async move {
tokio::time::sleep(timeout).await;
if !reporter.connected_once.load(Ordering::Relaxed)
&& !matches!(
client.connection_state(),
async_nats::connection::State::Connected
)
{
trc::event!(
Cluster(ClusterEvent::CoordinatorDisconnected),
Type = "nats",
Details = "Not connected at startup; retrying in the background",
);
}
});
}
}
+13
View File
@@ -2,6 +2,8 @@
* 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::{Coordinator, Msg, PubSubStream};
@@ -43,6 +45,17 @@ impl Coordinator {
pub fn is_none(&self) -> bool {
matches!(self, Coordinator::None)
}
/// inbuxa: whether the coordinator is connected right now, for the
/// backends that track it (NATS); `None` for the others and when no
/// coordinator is configured.
pub fn is_connected(&self) -> Option<bool> {
match self {
#[cfg(feature = "nats")]
Coordinator::Nats(store) => Some(store.is_connected()),
_ => None,
}
}
}
impl PubSubStream {
+21
View File
@@ -562,6 +562,27 @@ impl ParseHttp for Server {
})
.into_http_response());
}
// inbuxa: the cluster coordinator's connection, for
// monitoring. It stays out of live and ready on purpose:
// a node without its coordinator still serves mail, and
// failing those would have an orchestrator restart, or
// take out of service, every node at once when the
// coordinator goes down
"cluster" => {
let coordinator = &self.core.storage.coordinator;
let (status, state) = match coordinator.is_connected() {
Some(true) => (StatusCode::OK, "connected"),
Some(false) => (StatusCode::SERVICE_UNAVAILABLE, "disconnected"),
None if coordinator.is_none() => (StatusCode::OK, "none"),
None => (StatusCode::OK, "unknown"),
};
return Ok(http_proto::JsonResponse::with_status(
status,
serde_json::json!({ "coordinator": state }),
)
.no_cache()
.into_http_response());
}
_ => (),
}
}
+17 -2
View File
@@ -427,9 +427,24 @@ pub(crate) async fn trace_query(
}
None => false,
},
Property::QueueId => match value.as_str() {
// The queue id column is an integer on every search backend, and
// holds a trace's first queue id; the keywords carry all of them
Property::QueueId => match value
.as_str()
.and_then(|v| v.trim().parse::<u64>().ok())
.or_else(|| value.as_u64())
{
Some(queue_id) => {
search.push(SearchFilter::eq(TracingSearchField::QueueId, queue_id.to_string()));
search.extend([
SearchFilter::Or,
SearchFilter::eq(TracingSearchField::QueueId, queue_id),
SearchFilter::has_text(
TracingSearchField::Keywords,
queue_id.to_string(),
nlp::language::Language::None,
),
SearchFilter::End,
]);
true
}
None => false,
+4
View File
@@ -109,6 +109,10 @@ 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();
+24 -3
View File
@@ -26,7 +26,7 @@ pub fn spawn_broadcast_subscriber(inner: Arc<Inner>, mut shutdown_rx: watch::Rec
};
tokio::spawn(async move {
let mut retry_count = 0;
let mut retry_count: u32 = 0;
trc::event!(Cluster(ClusterEvent::SubscriberStart));
@@ -53,7 +53,7 @@ pub fn spawn_broadcast_subscriber(inner: Arc<Inner>, mut shutdown_rx: watch::Rec
);
match tokio::time::timeout(
Duration::from_secs(1 << retry_count.max(6)),
subscribe_retry_delay(retry_count),
shutdown_rx.changed(),
)
.await
@@ -62,7 +62,7 @@ pub fn spawn_broadcast_subscriber(inner: Arc<Inner>, mut shutdown_rx: watch::Rec
break;
}
Err(_) => {
retry_count += 1;
retry_count = retry_count.saturating_add(1);
continue;
}
}
@@ -234,6 +234,11 @@ pub fn spawn_broadcast_subscriber(inner: Arc<Inner>, mut shutdown_rx: watch::Rec
});
}
/// Delay before the next subscribe attempt: 1 s, 2 s, 4 s ... capped at 64 s.
fn subscribe_retry_delay(retry_count: u32) -> Duration {
Duration::from_secs(1u64 << retry_count.min(6))
}
fn log_event(event: &BroadcastEvent) -> trc::Value {
match event {
BroadcastEvent::PushNotification(notification) => match notification {
@@ -296,3 +301,19 @@ fn log_event(event: &BroadcastEvent) -> trc::Value {
BroadcastEvent::QueueRefresh => "QueueRefresh".into(),
}
}
#[cfg(test)]
mod tests {
use super::subscribe_retry_delay;
use std::time::Duration;
#[test]
fn subscribe_retry_backoff_grows_then_caps() {
let schedule: Vec<u64> = (0..10)
.map(|n| subscribe_retry_delay(n).as_secs())
.collect();
assert_eq!(schedule, vec![1, 2, 4, 8, 16, 32, 64, 64, 64, 64]);
// No shift overflow at the top of the range.
assert_eq!(subscribe_retry_delay(u32::MAX), Duration::from_secs(64));
}
}
+67 -22
View File
@@ -91,7 +91,15 @@ impl SearchIndexTask for Server {
build_contact_document(self, account_id, document_id).await
}
IndexDocumentType::File => {
// File indexing not implemented yet
// 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,
});
continue;
}
};
@@ -567,20 +575,14 @@ async fn build_contact_document(
}
// inbuxa: MON-16: a trace's search document, when trace search is on:
// its event types, queue ids, and addresses, their domains, hosts, IPs,
// message ids and account names as keywords
// inbuxa: MON-16: a trace's search document, when trace search is on
async fn build_tracing_span_document(
server: &Server,
span_id: u64,
) -> trc::Result<Option<IndexDocument>> {
use common::telemetry::tracers::store::MaybeTrace;
use registry::schema::{enums::SearchTracingField, structs::Search};
use store::{
search::TracingSearchField,
write::{TelemetryClass, ValueClass},
};
use trc::Key;
use registry::schema::structs::Search;
use store::write::{TelemetryClass, ValueClass};
let settings = server
.registry()
@@ -590,7 +592,6 @@ async fn build_tracing_span_document(
if !settings.index_telemetry {
return Ok(None);
}
let wants = |field: SearchTracingField| settings.index_tracing_fields.iter().any(|f| *f == field);
let Some(MaybeTrace(Some(trace))) = server
.tracing_store()
.get_value::<MaybeTrace>(ValueKey::from(ValueClass::Telemetry(TelemetryClass::Span(
@@ -601,23 +602,67 @@ async fn build_tracing_span_document(
return Ok(None);
};
Ok(Some(trace_search_document(
span_id,
&trace,
&settings
.index_tracing_fields
.iter()
.copied()
.collect::<Vec<_>>(),
)))
}
/// inbuxa: MON-16: the search document for a stored trace.
///
/// The event type and queue id columns are integers on every search backend
/// (BIGINT on PostgreSQL and MySQL, long on Elasticsearch), and each holds a
/// single value per trace: the event type is the trace's opening event, the
/// one `x:Trace/query` filters on, and the queue id is the first queue id the
/// trace mentions. Every queue id also goes into the keywords, so a session
/// that queued several messages is found by any of them.
pub fn trace_search_document(
span_id: u64,
trace: &registry::schema::structs::Trace,
fields: &[registry::schema::enums::SearchTracingField],
) -> IndexDocument {
use registry::schema::{enums::SearchTracingField, structs::TraceValue};
use store::search::TracingSearchField;
use trc::Key;
let wants = |field: SearchTracingField| fields.contains(&field);
let mut document = IndexDocument::new(SearchIndex::Tracing).with_id(span_id);
let mut seen = store::ahash::AHashSet::new();
for event in trace.events.iter() {
if wants(SearchTracingField::EventType) && seen.insert(event.event.as_str().to_string()) {
document.index_keyword(TracingSearchField::EventType, event.event.as_str());
if wants(SearchTracingField::EventType)
&& let Some(first) = trace.events.iter().next()
{
document.index_unsigned(TracingSearchField::EventType, first.event.to_id() as u64);
}
let mut seen = store::ahash::AHashSet::new();
let mut queue_id_indexed = false;
for event in trace.events.iter() {
for kv in event.key_values.iter() {
let text = match &kv.value {
registry::schema::structs::TraceValue::String(v) => v.value.clone(),
registry::schema::structs::TraceValue::UnsignedInt(v) => v.value.to_string(),
registry::schema::structs::TraceValue::IpAddr(v) => v.value.to_string(),
TraceValue::String(v) => v.value.clone(),
TraceValue::UnsignedInt(v) => v.value.to_string(),
TraceValue::IpAddr(v) => v.value.to_string(),
_ => continue,
};
match kv.key {
Key::QueueId if wants(SearchTracingField::QueueId) => {
if seen.insert(format!("q:{text}")) {
document.index_keyword(TracingSearchField::QueueId, &text);
Key::QueueId => {
let Ok(queue_id) = text.parse::<u64>() else {
continue;
};
if wants(SearchTracingField::QueueId) && !queue_id_indexed {
document.index_unsigned(TracingSearchField::QueueId, queue_id);
queue_id_indexed = true;
}
if wants(SearchTracingField::Keywords) && seen.insert(format!("k:{text}")) {
document.index_text(
TracingSearchField::Keywords,
&text,
nlp::language::Language::None,
);
}
}
Key::From
@@ -648,7 +693,7 @@ async fn build_tracing_span_document(
}
}
}
Ok(Some(document))
document
}
// inbuxa: UD-1, UD-4: archives a deleted file, event or contact noted at
+34 -2
View File
@@ -2,6 +2,8 @@
* 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::*;
@@ -13,13 +15,21 @@ 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(), DEFAULT_LOCK_EXPIRY)
.try_lock(KV_LOCK_TASK, &id.to_be_bytes(), locks.expiry())
.await
{
Ok(result) => {
if !result {
if result {
locks.insert(id);
} else {
trc::event!(
TaskManager(TaskManagerEvent::TaskLocked),
Id = id,
@@ -48,5 +58,27 @@ 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()
}
+184 -127
View File
@@ -2,6 +2,8 @@
* 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;
@@ -18,7 +20,7 @@ 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::{
DEFAULT_LOCK_EXPIRY, Locked, QUEUE_REFRESH_INTERVAL, TaskDetails, TaskFailureType, TaskInfo,
CLAIM_RECHECK_INTERVAL, Locked, QUEUE_REFRESH_INTERVAL, TaskDetails, TaskFailureType, TaskInfo,
TaskJob, TaskManagerIpc, TaskResult,
};
use common::BuildServer;
@@ -124,72 +126,47 @@ 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);
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!())
);
}
if let Some(task) = fetch_task(&server, job).await {
batch.push(task);
}
while batch.len() < batch_size {
match rx.try_recv() {
Ok(job) => {
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!())
);
}
if let Some(task) = fetch_task(&server, job).await {
batch.push(task);
}
}
Err(_) => break,
}
}
// Dispatch
// 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
let mut refresh_queue = false;
let results = server.index(&batch).await.into_iter().map(|r| {
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;
}
}
if refresh_queue || rx.is_empty() {
server.notify_task_queue();
@@ -203,83 +180,31 @@ pub fn spawn_task_manager(inner: Arc<Inner>) {
let server = inner.build_server();
let mut refresh_queue = false;
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!(),
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();
update_tasks(
&server,
&mut [TaskDetails { task, info: job }],
&mut [TaskDetails { task, info }],
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!())
);
worker_failed(&server, &[info.id], err).await;
refresh_queue = true;
}
}
}
@@ -318,6 +243,13 @@ 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,
@@ -393,9 +325,7 @@ 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(
DEFAULT_LOCK_EXPIRY + 1,
);
+ std::time::Duration::from_secs(lock_expiry + 1);
locked.due = task_due;
tasks.push((
TaskJob {
@@ -411,9 +341,7 @@ impl TaskQueueManager for Server {
Entry::Vacant(entry) => {
entry.insert(Locked {
expires: Instant::now()
+ std::time::Duration::from_secs(
DEFAULT_LOCK_EXPIRY + 1,
),
+ std::time::Duration::from_secs(lock_expiry + 1),
due: task_due,
revision: ipc.revision,
});
@@ -464,12 +392,26 @@ impl TaskQueueManager for Server {
let tx = &ipc.txs[task_type_idx as usize];
if tx.capacity() > 0 {
if self.try_lock_task(task_job.id).await && tx.send(task_job).await.is_err() {
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() {
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
@@ -481,9 +423,117 @@ impl TaskQueueManager for Server {
let now = Instant::now();
ipc.locked
.retain(|_, locked| locked.expires > now && locked.revision == ipc.revision);
Duration::from_secs(next_event.map_or(QUEUE_REFRESH_INTERVAL, |timestamp| {
let sleep_for = 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))))
}
}
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;
}
}
@@ -614,6 +664,13 @@ 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,
+3 -1
View File
@@ -35,7 +35,9 @@ pub mod scheduler;
pub mod spam_classifier;
const QUEUE_REFRESH_INTERVAL: u64 = 60 * 5; // 5 minutes
const DEFAULT_LOCK_EXPIRY: u64 = 60 * 60; // 1 hour
// 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
pub(crate) struct TaskManagerIpc {
txs: [mpsc::Sender<TaskJob>; TaskType::COUNT],
+7 -2
View File
@@ -10,8 +10,9 @@
// inbuxa: 637 to 641 are the fork's SCIM events (SCIM-54); 642 is
// auth.legacy-protocol-refused (legacy-protocols LP-6); 643 is
// security.legacy-protocols-changed (LP-8)
pub const TOTAL_EVENT_COUNT: usize = 644;
// security.legacy-protocols-changed (LP-8); 644 to 646 are the cluster
// coordinator's connection events
pub const TOTAL_EVENT_COUNT: usize = 647;
pub const TOTAL_METRIC_COUNT: usize = 369;
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
@@ -150,6 +151,10 @@ pub enum ClusterEvent {
MessageSkipped = 47,
MessageInvalid = 49,
NodeIdRenewed = 275,
// inbuxa: the coordinator's connection
CoordinatorConnected = 644,
CoordinatorDisconnected = 645,
CoordinatorError = 646,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
+32
View File
@@ -81,6 +81,10 @@ impl EventType {
b"cluster.message-skipped" => EventType::Cluster(ClusterEvent::MessageSkipped),
b"cluster.message-invalid" => EventType::Cluster(ClusterEvent::MessageInvalid),
b"cluster.node-id-renewed" => EventType::Cluster(ClusterEvent::NodeIdRenewed),
// inbuxa: coordinator connection
b"cluster.coordinator-connected" => EventType::Cluster(ClusterEvent::CoordinatorConnected),
b"cluster.coordinator-disconnected" => EventType::Cluster(ClusterEvent::CoordinatorDisconnected),
b"cluster.coordinator-error" => EventType::Cluster(ClusterEvent::CoordinatorError),
b"dane.authentication-success" => EventType::Dane(DaneEvent::AuthenticationSuccess),
b"dane.authentication-failure" => EventType::Dane(DaneEvent::AuthenticationFailure),
b"dane.no-certificates-found" => EventType::Dane(DaneEvent::NoCertificatesFound),
@@ -742,6 +746,14 @@ impl EventType {
EventType::Cluster(ClusterEvent::MessageSkipped) => "cluster.message-skipped",
EventType::Cluster(ClusterEvent::MessageInvalid) => "cluster.message-invalid",
EventType::Cluster(ClusterEvent::NodeIdRenewed) => "cluster.node-id-renewed",
// inbuxa: coordinator connection
EventType::Cluster(ClusterEvent::CoordinatorConnected) => {
"cluster.coordinator-connected"
}
EventType::Cluster(ClusterEvent::CoordinatorDisconnected) => {
"cluster.coordinator-disconnected"
}
EventType::Cluster(ClusterEvent::CoordinatorError) => "cluster.coordinator-error",
EventType::Dane(DaneEvent::AuthenticationSuccess) => "dane.authentication-success",
EventType::Dane(DaneEvent::AuthenticationFailure) => "dane.authentication-failure",
EventType::Dane(DaneEvent::NoCertificatesFound) => "dane.no-certificates-found",
@@ -1524,6 +1536,10 @@ impl EventType {
EventType::Cluster(ClusterEvent::MessageSkipped) => 47,
EventType::Cluster(ClusterEvent::MessageInvalid) => 49,
EventType::Cluster(ClusterEvent::NodeIdRenewed) => 275,
// inbuxa: coordinator connection
EventType::Cluster(ClusterEvent::CoordinatorConnected) => 644,
EventType::Cluster(ClusterEvent::CoordinatorDisconnected) => 645,
EventType::Cluster(ClusterEvent::CoordinatorError) => 646,
EventType::Dane(DaneEvent::AuthenticationSuccess) => 67,
EventType::Dane(DaneEvent::AuthenticationFailure) => 66,
EventType::Dane(DaneEvent::NoCertificatesFound) => 69,
@@ -2176,6 +2192,10 @@ impl EventType {
47 => Some(EventType::Cluster(ClusterEvent::MessageSkipped)),
49 => Some(EventType::Cluster(ClusterEvent::MessageInvalid)),
275 => Some(EventType::Cluster(ClusterEvent::NodeIdRenewed)),
// inbuxa: coordinator connection
644 => Some(EventType::Cluster(ClusterEvent::CoordinatorConnected)),
645 => Some(EventType::Cluster(ClusterEvent::CoordinatorDisconnected)),
646 => Some(EventType::Cluster(ClusterEvent::CoordinatorError)),
67 => Some(EventType::Dane(DaneEvent::AuthenticationSuccess)),
66 => Some(EventType::Dane(DaneEvent::AuthenticationFailure)),
69 => Some(EventType::Dane(DaneEvent::NoCertificatesFound)),
@@ -3114,6 +3134,10 @@ impl EventType {
EventType::Auth(AuthEvent::TooManyAttempts) => Level::Warn,
EventType::Calendar(CalendarEvent::AlarmFailed) => Level::Warn,
EventType::Cluster(ClusterEvent::SubscriberDisconnected) => Level::Warn,
// inbuxa: coordinator connection
EventType::Cluster(ClusterEvent::CoordinatorConnected) => Level::Info,
EventType::Cluster(ClusterEvent::CoordinatorDisconnected) => Level::Warn,
EventType::Cluster(ClusterEvent::CoordinatorError) => Level::Warn,
EventType::Delivery(DeliveryEvent::MissingOutboundHostname) => Level::Warn,
EventType::Delivery(DeliveryEvent::ConcurrencyLimitExceeded) => Level::Warn,
EventType::Delivery(DeliveryEvent::RateLimitExceeded) => Level::Warn,
@@ -3244,6 +3268,10 @@ impl EventType {
EventType::Cluster(ClusterEvent::MessageSkipped) => "PubSub message skipped",
EventType::Cluster(ClusterEvent::MessageInvalid) => "Invalid PubSub message",
EventType::Cluster(ClusterEvent::NodeIdRenewed) => "Node ID renewed",
// inbuxa: coordinator connection
EventType::Cluster(ClusterEvent::CoordinatorConnected) => "Coordinator connected",
EventType::Cluster(ClusterEvent::CoordinatorDisconnected) => "Coordinator unavailable",
EventType::Cluster(ClusterEvent::CoordinatorError) => "Coordinator error",
EventType::Dane(DaneEvent::AuthenticationSuccess) => "DANE authentication successful",
EventType::Dane(DaneEvent::AuthenticationFailure) => "DANE authentication failed",
EventType::Dane(DaneEvent::NoCertificatesFound) => "No certificates found for DANE",
@@ -4322,6 +4350,10 @@ impl EventType {
EventType::Cluster(ClusterEvent::MessageSkipped),
EventType::Cluster(ClusterEvent::MessageInvalid),
EventType::Cluster(ClusterEvent::NodeIdRenewed),
// inbuxa: coordinator connection
EventType::Cluster(ClusterEvent::CoordinatorConnected),
EventType::Cluster(ClusterEvent::CoordinatorDisconnected),
EventType::Cluster(ClusterEvent::CoordinatorError),
EventType::Dane(DaneEvent::AuthenticationSuccess),
EventType::Dane(DaneEvent::AuthenticationFailure),
EventType::Dane(DaneEvent::NoCertificatesFound),
+1 -1
View File
@@ -81,7 +81,7 @@ fn legacy_setting(name: &str, is_set: impl Fn(&str) -> bool) -> Option<String> {
#[macro_export]
macro_rules! brand_version {
() => {
"2026.9.24"
"2026.9.24.3"
};
}
+8 -3
View File
@@ -212,10 +212,15 @@ unchanged.
- **MON-16.** With `indexTelemetry` on, storing a trace schedules an
`IndexTrace` task. The task builds one document for `SearchIndex::Tracing`
with the fields named in `indexTracingFields`:
- `eventType`: every event type in the trace;
- `queueId`: every `queueId` value;
- `eventType`: the trace's opening event, as its numeric id;
- `queueId`: the first `queueId` value, as an integer;
- `keywords`: every address in `from` and `to`, each address's domain, every
`domain`, `hostname`, `remoteIp`, `messageId` and `accountName` value.
`domain`, `hostname`, `remoteIp`, `messageId` and `accountName` value,
and every `queueId` value.
The event type and queue id are single integer columns on every search
backend (BIGINT on PostgreSQL and MySQL), so the `queueId` filter matches
the column or any queue id in the keywords, and a session that queued
several messages is found by each of them.
So searching `example.org` finds every trace to or from that domain, as the
upstream suite expects. With `indexTelemetry` off nothing is indexed, and
the `text` and `queueId` filters are refused (see "Interfaces").
Binary file not shown.
+1 -1
View File
@@ -1 +1 @@
VbnFuwCOTBh0s2T-NuRhb2JaJr8Jl5s3LgXv4Pv2sTg
XFI3xuKC_rH1KZyaVBF0uTIiRDXRqyYboijquiGz2eg
+145
View File
@@ -0,0 +1,145 @@
/*
* SPDX-FileCopyrightText: 2026 Coffey Labs
*
* SPDX-License-Identifier: AGPL-3.0-only
*/
//! A node that starts while its NATS coordinator is down joins the cluster
//! once NATS comes up, without a restart, and reports the coordinator's
//! connection on `/healthz/cluster` as it goes and comes back.
use crate::utils::server::TestServerBuilder;
use coordinator::Coordinator;
use registry::{
schema::{
enums::NetworkListenerProtocol,
structs::{Coordinator as CoordinatorSetting, NatsCoordinator},
},
types::map::Map,
};
use serde_json::{Value, json};
use std::time::{Duration, Instant};
use testcontainers::{
GenericImage, ImageExt, core::IntoContainerPort, core::WaitFor, runners::AsyncRunner,
};
const HTTP_PORT: u16 = 11_310;
const TOPIC: &str = "inbuxa-coordinator-test";
#[tokio::test(flavor = "multi_thread")]
pub async fn coordinator_reconnect_tests() {
println!("Running coordinator reconnect tests...");
// A port with no NATS server behind it, yet
let nats_port = std::net::TcpListener::bind("127.0.0.1:0")
.unwrap()
.local_addr()
.unwrap()
.port();
let config = NatsCoordinator {
addresses: Map::new(vec![format!("127.0.0.1:{nats_port}")]),
use_tls: false,
timeout_connection: 1_000u64.into(),
..Default::default()
};
// 1. The node starts, without a build error, while NATS is down, and
// says so
let test = TestServerBuilder::new("coordinator_reconnect_tests")
.await
.with_object(CoordinatorSetting::Nats(config.clone()))
.await
.with_listener(NetworkListenerProtocol::Http, "http", HTTP_PORT, true)
.await
.build()
.await;
let coordinator = test.server.core.storage.coordinator.clone();
assert!(
coordinator.is_enabled(),
"a coordinator, though not connected"
);
assert_eq!(coordinator.is_connected(), Some(false));
assert_eq!(
cluster_health().await,
(503, json!({"coordinator": "disconnected"}))
);
// A subscription made now, as the broadcast subscriber makes it at
// startup, has to work once NATS is up
let mut stream = coordinator.subscribe(TOPIC).await.unwrap();
// 2. NATS comes up: the node connects on its own
let nats = GenericImage::new("nats", "latest")
.with_wait_for(WaitFor::message_on_stderr("Server is ready"))
.with_mapped_port(nats_port, 4222.tcp())
.start()
.await
.expect("Failed to start NATS container");
wait_for_health(200, "connected").await;
let other_node = coordinator::backend::nats::NatsPubSub::open(config.clone())
.await
.unwrap();
wait_until_connected(&other_node).await;
round_trip(&other_node, &mut stream, b"after startup").await;
// 3. NATS goes away: the node reports it; and it comes back: the node
// reconnects and the same subscription carries on
nats.stop().await.unwrap();
wait_for_health(503, "disconnected").await;
nats.start().await.unwrap();
wait_for_health(200, "connected").await;
wait_until_connected(&other_node).await;
round_trip(&other_node, &mut stream, b"after reconnect").await;
drop(nats);
if test.is_reset() {
test.temp_dir.delete();
}
}
async fn cluster_health() -> (u16, Value) {
let response = reqwest::Client::builder()
.danger_accept_invalid_certs(true)
.timeout(Duration::from_secs(5))
.build()
.unwrap()
.get(format!("https://127.0.0.1:{HTTP_PORT}/healthz/cluster"))
.send()
.await
.unwrap();
let status = response.status().as_u16();
(status, response.json().await.unwrap())
}
async fn wait_for_health(status: u16, state: &str) {
let started = Instant::now();
loop {
let health = cluster_health().await;
if health == (status, json!({"coordinator": state})) {
return;
}
assert!(
started.elapsed() < Duration::from_secs(30),
"expected {status} {state}, still {health:?}"
);
tokio::time::sleep(Duration::from_millis(250)).await;
}
}
async fn wait_until_connected(coordinator: &Coordinator) {
let started = Instant::now();
while coordinator.is_connected() != Some(true) {
assert!(started.elapsed() < Duration::from_secs(30), "not connected");
tokio::time::sleep(Duration::from_millis(100)).await;
}
}
/// Another node publishes; this one's subscription receives it.
async fn round_trip(from: &Coordinator, stream: &mut coordinator::PubSubStream, payload: &[u8]) {
from.publish(TOPIC, payload.to_vec()).await.unwrap();
let message = tokio::time::timeout(Duration::from_secs(10), stream.next())
.await
.expect("no message within 10 seconds")
.expect("subscription ended");
assert_eq!(message.payload(), payload);
}
+4
View File
@@ -2,7 +2,11 @@
* 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.
*/
pub mod broadcast;
#[cfg(feature = "nats")]
pub mod coordinator; // inbuxa: coordinator reconnects
pub mod stress;
+253 -5
View File
@@ -2,6 +2,8 @@
* SPDX-FileCopyrightText: 2020 Stalwart Labs LLC <[email protected]>
*
* SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-SEL
*
* Modified by Coffey Labs in 2026 for INBUXA.
*/
use crate::utils::{
@@ -9,14 +11,22 @@ use crate::utils::{
server::TestServer,
temp_dir::TempDir,
};
use ::registry::schema::enums::CompressionAlgo;
use ::registry::schema::{
enums::{CompressionAlgo, TaskStoreMaintenanceType},
prelude::ObjectType,
structs::Task,
};
use ahash::AHashSet;
use common::{DATABASE_SCHEMA_VERSION, manager::backup::BackupParams};
use common::{
DATABASE_SCHEMA_VERSION,
manager::{SPAM_CLASSIFIER_KEY, SPAM_TRAINER_KEY, backup::BackupParams},
};
use store::{
rand,
write::{
AnyClass, AnyKey, BatchBuilder, BlobLink, BlobOp, Operation, QueueClass, QueueEvent,
RegistryClass, ValueClass, key::KeySerializer,
RegistryClass, TaskQueueClass, ValueClass,
key::{DeserializeBigEndian, KeySerializer},
},
*,
};
@@ -167,6 +177,50 @@ pub async fn test(test: &TestServer) {
}
db.write(batch.build_all()).await.unwrap();
// inbuxa: registry objects kept outside the registry subspace (archived
// items for undelete, spam training samples, directory entries) and the
// fork's own subspace. Exports used to leave the first two behind.
println!("Creating archived items, spam samples and fork data...");
let mut batch = BatchBuilder::new();
for item_id in [1u64, 2, 3] {
for object in [
ObjectType::ArchivedItem,
ObjectType::SpamTrainingSample,
ObjectType::Account,
] {
batch.set(
ValueClass::Registry(RegistryClass::Item {
object_id: object as u16,
item_id,
}),
random_bytes(item_id as usize * 64),
);
}
batch.set(
ValueClass::Any(AnyClass {
subspace: SUBSPACE_INBUXA,
key: [b'U', b'x']
.into_iter()
.chain(item_id.to_be_bytes())
.collect(),
}),
random_bytes(32),
);
}
db.write(batch.build_all()).await.unwrap();
// inbuxa: the trained spam classifier lives in blobs with fixed names
let mut named_blobs = Vec::new();
for key in [SPAM_CLASSIFIER_KEY, SPAM_TRAINER_KEY] {
let data = random_bytes(4096);
test.server
.blob_store()
.put_blob(key, &data, CompressionAlgo::Lz4)
.await
.unwrap();
named_blobs.push((key, data));
}
// Create directory data
println!("Creating directory data...");
let mut batch = BatchBuilder::new();
@@ -185,6 +239,17 @@ pub async fn test(test: &TestServer) {
println!("Calculating store hash...");
let snapshot = Snapshot::new(&db).await;
assert!(!snapshot.keys.is_empty(), "Store hash counts are empty",);
for subspace in [
SUBSPACE_DELETED_ITEMS,
SUBSPACE_SPAM_SAMPLES,
SUBSPACE_INBUXA,
] {
assert!(
snapshot.keys.iter().any(|k| k.subspace == subspace),
"No test data in subspace {}",
char::from(subspace)
);
}
// Export store
println!("Exporting store...");
@@ -210,22 +275,188 @@ pub async fn test(test: &TestServer) {
.finalize(),
);
db.write(batch.build_all()).await.unwrap();
test.server.core.restore(temp_dir.path.clone()).await;
for (key, _) in &named_blobs {
test.server.blob_store().delete_blob(key).await.unwrap();
}
let imported = test.server.core.restore(temp_dir.path.clone()).await;
let mut batch = BatchBuilder::new();
batch.clear(ValueClass::NodeId(0));
db.write(batch.build_all()).await.unwrap();
for subspace in [
SUBSPACE_DELETED_ITEMS,
SUBSPACE_SPAM_SAMPLES,
SUBSPACE_INBUXA,
] {
assert!(
imported.contains(&subspace),
"Subspace {} was not exported",
char::from(subspace)
);
}
// Verify hash
print!("Verifying store hash...");
snapshot.assert_is_eq(&Snapshot::new(&db).await);
assert_named_blobs(test.server.blob_store(), &named_blobs).await;
println!(" GREAT SUCCESS!");
// inbuxa: import the same export into a fresh store of another backend,
// the way a move from one database to another does it
#[cfg(all(feature = "rocks", feature = "sqlite"))]
cross_backend(test, &db, &temp_dir, &named_blobs).await;
// Destroy store
for (key, _) in &named_blobs {
test.server.blob_store().delete_blob(key).await.unwrap();
}
store_destroy(&db).await;
store_assert_is_empty(&db, db.clone().into(), true).await;
temp_dir.delete();
}
#[cfg(all(feature = "rocks", feature = "sqlite"))]
async fn cross_backend(
test: &TestServer,
source: &Store,
export: &TempDir,
named_blobs: &[(&[u8], Vec<u8>)],
) {
let source_type = std::env::var("STORE").unwrap();
let target_type = if source_type.eq_ignore_ascii_case("sqlite") {
"RocksDb"
} else {
"Sqlite"
};
println!("Importing the export into a fresh {target_type} store...");
let target_dir = TempDir::new("art_vandelay_cross_backend", true);
let target = Store::build(
crate::utils::storage::build_data_store(target_type, &target_dir.path.to_string_lossy())
.await,
)
.await
.unwrap();
target.create_tables().await.unwrap();
store_destroy(&target).await;
let mut core = test.server.core.as_ref().clone();
core.storage.data = target.clone();
core.storage.blob = target.clone().into();
let imported = core.restore(export.path.clone()).await;
// Counters are stored differently by the SQL and key-value backends, so
// compare their keys here and their values through the counter API.
print!("Verifying {target_type} store hash...");
Snapshot::new_portable(source)
.await
.assert_is_eq(&Snapshot::new_portable(&target).await);
for subspace in [SUBSPACE_COUNTER, SUBSPACE_QUOTA] {
let mut keys = Vec::new();
source
.iterate(
IterateParams::new(
AnyKey {
subspace,
key: vec![0u8],
},
AnyKey {
subspace,
key: vec![u8::MAX; 10],
},
)
.no_values(),
|key, _| {
keys.push(key.to_vec());
Ok(true)
},
)
.await
.unwrap();
for key in keys {
let class = || {
ValueClass::Any(AnyClass {
subspace,
key: key.clone(),
})
};
assert_eq!(
source.get_counter(class()).await.unwrap(),
target.get_counter(class()).await.unwrap(),
"Counter mismatch in {} for {key:?}",
char::from(subspace)
);
}
}
assert_named_blobs(&core.storage.blob, named_blobs).await;
println!(" GREAT SUCCESS!");
// The search index isn't exported; the import queues its rebuild
let queued = core.queue_reindex(&imported).await;
let expected = [
TaskStoreMaintenanceType::ReindexAccounts,
TaskStoreMaintenanceType::ReindexTelemetry,
];
assert_eq!(queued, expected);
let mut task_ids = Vec::new();
target
.iterate(
IterateParams::new(
AnyKey {
subspace: SUBSPACE_TASK_QUEUE,
key: vec![0u8],
},
AnyKey {
subspace: SUBSPACE_TASK_QUEUE,
key: vec![u8::MAX; 20],
},
)
.no_values(),
|key, _| {
if key.deserialize_be_u64(0)? == 0 {
task_ids.push(key.deserialize_be_u64(U64_LEN)?);
}
Ok(true)
},
)
.await
.unwrap();
let mut found = Vec::new();
for id in task_ids {
match target
.get_value::<Task>(ValueKey::from(ValueClass::TaskQueue(
TaskQueueClass::Task { id },
)))
.await
.unwrap()
{
Some(Task::StoreMaintenance(task)) => found.push(task.maintenance_type),
other => panic!("Unexpected task {other:?}"),
}
}
found.sort_by_key(|t| *t as u16);
assert_eq!(found, expected, "Queued tasks don't match");
store_destroy(&target).await;
drop(core);
drop(target);
target_dir.delete();
}
async fn assert_named_blobs(blob_store: &BlobStore, named_blobs: &[(&[u8], Vec<u8>)]) {
for (key, data) in named_blobs {
assert_eq!(
blob_store
.get_blob(key, 0..usize::MAX)
.await
.unwrap()
.as_ref(),
Some(data),
"Blob {} was not restored",
String::from_utf8_lossy(key)
);
}
}
#[derive(Debug, PartialEq, Eq)]
struct Snapshot {
keys: AHashSet<KeyValue>,
@@ -240,7 +471,19 @@ struct KeyValue {
impl Snapshot {
async fn new(db: &Store) -> Self {
let is_sql = db.is_sql();
Self::build(db, !db.is_sql(), true).await
}
/// Comparable across backends: no counter values, which the SQL and
/// key-value stores encode differently, and no blobs, which only live in
/// the data store when it doubles as the blob store.
#[cfg(all(feature = "rocks", feature = "sqlite"))]
async fn new_portable(db: &Store) -> Self {
Self::build(db, false, false).await
}
async fn build(db: &Store, counter_values: bool, with_blobs: bool) -> Self {
let is_sql = !counter_values;
let mut keys = AHashSet::new();
@@ -265,7 +508,12 @@ impl Snapshot {
(SUBSPACE_QUOTA, !is_sql),
(SUBSPACE_REPORT_OUT, true),
(SUBSPACE_REPORT_IN, true),
(SUBSPACE_DIRECTORY, true),
(SUBSPACE_INBUXA, true),
] {
if subspace == SUBSPACE_BLOBS && !with_blobs {
continue;
}
let from_key = AnyKey {
subspace,
key: vec![0u8],
+1
View File
@@ -21,6 +21,7 @@ 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;
+163
View File
@@ -2,6 +2,8 @@
* 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::{store::deflate_test_resource, utils::server::TestServer};
@@ -122,6 +124,10 @@ pub async fn test(test: &TestServer) {
println!("Running global id filtering tests...");
test_global(store.clone()).await;
// inbuxa: trace documents as the index task builds them
println!("Running trace document tests...");
test_trace_documents(store.clone()).await;
// Large document insert test
println!("Running large document insert tests...");
let mut large_text = String::with_capacity(20 * 1024 * 1024);
@@ -809,3 +815,160 @@ async fn test_global(store: SearchStore) {
AHashSet::from_iter([3, 4, 5])
);
}
// inbuxa: MON-16: documents built by the index task from stored traces go
// into every search backend (the SQL backends type etyp and qid as BIGINT)
// and are found again by queue id and keyword.
async fn test_trace_documents(store: SearchStore) {
use registry::schema::{
enums::SearchTracingField,
structs::{
Trace, TraceEvent, TraceKeyValue, TraceValue, TraceValueString,
TraceValueUnsignedInt,
},
};
use services::task_manager::index::trace_search_document;
use trc::{DeliveryEvent, EventType, Key, SmtpEvent};
let kv_u = |key: Key, value: u64| TraceKeyValue {
key,
value: TraceValue::UnsignedInt(TraceValueUnsignedInt { value }),
};
let kv_s = |key: Key, value: &str| TraceKeyValue {
key,
value: TraceValue::String(TraceValueString {
value: value.to_string(),
}),
};
let event = |event: EventType, key_values: Vec<TraceKeyValue>| TraceEvent {
event,
key_values: key_values.into(),
..Default::default()
};
let fields = [
SearchTracingField::EventType,
SearchTracingField::QueueId,
SearchTracingField::Keywords,
];
// An SMTP session that queued two messages, and a delivery attempt
let session = Trace {
events: vec![
event(
EventType::Smtp(SmtpEvent::ConnectionStart),
vec![kv_s(Key::RemoteIp, "192.0.2.7")],
),
event(
EventType::Smtp(SmtpEvent::MailFrom),
vec![kv_s(Key::From, "[email protected]")],
),
event(
EventType::Smtp(SmtpEvent::RcptTo),
vec![kv_u(Key::QueueId, 9_000_000_001), kv_s(Key::To, "[email protected]")],
),
event(
EventType::Smtp(SmtpEvent::RcptTo),
vec![kv_u(Key::QueueId, 9_000_000_002)],
),
]
.into(),
};
let delivery = Trace {
events: vec![event(
EventType::Delivery(DeliveryEvent::AttemptStart),
vec![kv_u(Key::QueueId, 9_000_000_003), kv_s(Key::Hostname, "relay.example.net")],
)]
.into(),
};
let documents = vec![
trace_search_document(100, &session, &fields),
trace_search_document(101, &delivery, &fields),
];
assert!(
documents
.iter()
.all(|d| d.has_field(&SearchField::Tracing(TracingSearchField::QueueId))
&& d.has_field(&SearchField::Tracing(TracingSearchField::EventType))),
"trace documents carry a queue id and an event type"
);
store.index(documents).await.unwrap();
if let SearchStore::ElasticSearch(store) = &store {
store.refresh_index(SearchIndex::Tracing).await.unwrap();
}
let query = |filters: Vec<SearchFilter>| {
let store = store.clone();
async move {
store
.query_global(
SearchQuery::new(SearchIndex::Tracing)
.with_filter(SearchFilter::ge(SearchField::Id, 100u64))
.with_filters(filters),
)
.await
.unwrap()
.into_iter()
.collect::<AHashSet<_>>()
}
};
// By queue id, the way x:Trace/query asks: the queue id column, or any
// queue id in the keywords
let by_queue_id = |queue_id: u64| {
vec![
SearchFilter::Or,
SearchFilter::eq(TracingSearchField::QueueId, queue_id),
SearchFilter::has_text(
TracingSearchField::Keywords,
queue_id.to_string(),
Language::None,
),
SearchFilter::End,
]
};
assert_eq!(query(by_queue_id(9_000_000_001)).await, AHashSet::from_iter([100]));
assert_eq!(query(by_queue_id(9_000_000_002)).await, AHashSet::from_iter([100]));
assert_eq!(query(by_queue_id(9_000_000_003)).await, AHashSet::from_iter([101]));
assert_eq!(query(by_queue_id(9_000_000_004)).await, AHashSet::new());
assert_eq!(
query(vec![SearchFilter::eq(TracingSearchField::QueueId, 9_000_000_003u64)]).await,
AHashSet::from_iter([101])
);
// By opening event type
assert_eq!(
query(vec![SearchFilter::eq(
TracingSearchField::EventType,
EventType::Delivery(DeliveryEvent::AttemptStart).to_id() as u64,
)])
.await,
AHashSet::from_iter([101])
);
// By keyword: an address, lowercased, and its domain
assert_eq!(
query(vec![SearchFilter::has_text(
TracingSearchField::Keywords,
"example.org",
Language::None,
)])
.await,
AHashSet::from_iter([100])
);
assert_eq!(
query(vec![SearchFilter::has_text(
TracingSearchField::Keywords,
"relay.example.net",
Language::None,
)])
.await,
AHashSet::from_iter([101])
);
for id in [100u64, 101] {
store
.unindex(
SearchQuery::new(SearchIndex::Tracing)
.with_filter(SearchFilter::eq(SearchField::Id, id)),
)
.await
.unwrap();
}
}
+184
View File
@@ -0,0 +1,184 @@
/*
* 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 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;
}
}
+67
View File
@@ -148,6 +148,73 @@ pub async fn test(test: &mut TestServer) {
"test 9: to"
);
// MON-16: the queueId filter finds the traces that name a queue id (the
// session that queued the message and its delivery attempt) through the
// search index, given as a string or a number (the index column is an
// integer)
fn queue_ids(value: &Value, out: &mut Vec<u64>) {
match value {
Value::Object(map) => {
if map.get("key").and_then(|k| k.as_str()) == Some("queueId")
&& let Some(id) = map
.get("value")
.and_then(|v| v.get("value").unwrap_or(v).as_u64())
{
out.push(id);
}
map.values().for_each(|v| queue_ids(v, out));
}
Value::Array(list) => list.iter().for_each(|v| queue_ids(v, out)),
_ => {}
}
}
let with_ids = traces
.iter()
.map(|t| {
let mut ids = Vec::new();
queue_ids(t, &mut ids);
(t["id"].as_str().unwrap().to_string(), ids)
})
.collect::<Vec<_>>();
let queue_id = with_ids
.iter()
.find_map(|(_, ids)| ids.first().copied())
.expect("MON-16: a trace with a queue id");
let mut expected = with_ids
.iter()
.filter(|(_, ids)| ids.contains(&queue_id))
.map(|(id, _)| id.clone())
.collect::<Vec<_>>();
expected.sort();
for filter in [json!(queue_id.to_string()), json!(queue_id)] {
let response = admin
.jmap_method_call("x:Trace/query", json!({"filter": {"queueId": filter}}))
.await;
let mut found = response
.0
.pointer("/methodResponses/0/1/ids")
.and_then(|ids| ids.as_array())
.map(|ids| {
ids.iter()
.filter_map(|id| id.as_str().map(str::to_string))
.collect::<Vec<_>>()
})
.unwrap_or_default();
found.sort();
assert_eq!(found, expected, "MON-16: queueId {filter}: {response:?}");
}
let response = admin
.jmap_method_call(
"x:Trace/query",
json!({"filter": {"queueId": (queue_id ^ 0x5a5a_5a5a).to_string()}}),
)
.await;
assert_eq!(
response.0.pointer("/methodResponses/0/1/ids"),
Some(&json!([])),
"MON-16: an unknown queue id"
);
// Acceptance test 24: destroy removes a trace; create is refused
let trace_id = traces[0]["id"].as_str().unwrap().to_string();
let response = admin
+120
View File
@@ -0,0 +1,120 @@
#!/usr/bin/env python3
# SPDX-FileCopyrightText: 2026 Coffey Labs
# SPDX-License-Identifier: AGPL-3.0-or-later
"""Every path Cargo patches has to be in the image's build context.
Cargo.toml's [patch.crates-io] can point at a directory in this repository,
and the Dockerfile builds from a context that .dockerignore prunes to almost
nothing. Those two facts met on 2026-09-23: a vendored, patched sieve-rs
landed, CI stayed green -- it builds from a checkout, where the directory is
simply there -- and the release build failed on
failed to load source for dependency `sieve-rs`
failed to read /build/vendor/sieve-rs/Cargo.toml
after a tag had already been pushed. This is seconds, and it runs beside the
other fork checks rather than waiting for a release to find out.
Being in the context isn't enough on its own: the Dockerfile cooks the
dependencies (`cargo chef cook`) before it copies the tree in, from a recipe
that carries only the workspace's manifests. So each patched path must also be
copied into that stage before the cook step, or the same error comes back
there -- as it did for 2026.9.24.2, the first tag after the context fix.
"""
import re
import sys
from pathlib import Path
root = Path(__file__).resolve().parents[2]
def patched_paths(manifest: Path) -> list[str]:
"""Directories named by a [patch...] section's `path = "..."` entries."""
out, in_patch = [], False
for line in manifest.read_text().splitlines():
stripped = line.strip()
if stripped.startswith("["):
in_patch = stripped.startswith("[patch")
continue
if not in_patch:
continue
m = re.search(r'path\s*=\s*"([^"]+)"', stripped)
if m:
out.append(m.group(1))
return out
def allowed(dockerignore: Path) -> set[str]:
"""The first path segment of every re-inclusion rule."""
keep = set()
for line in dockerignore.read_text().splitlines():
stripped = line.strip()
if stripped.startswith("!"):
keep.add(stripped[1:].strip("/").split("/")[0])
return keep
def copied_before_cook(dockerfile: Path) -> list[str] | None:
"""Sources COPY'd into the stage that runs `cargo chef cook`, before it.
None when no stage cooks. A `COPY . .` covers everything.
"""
stage: list[str] = []
for line in dockerfile.read_text().splitlines():
stripped = line.strip()
if re.match(r"(?i)^FROM\s", stripped):
stage = []
continue
if "cargo chef cook" in stripped:
return stage
m = re.match(r"(?i)^COPY\s+(?!--from)(.+)$", stripped)
if m:
parts = m.group(1).split()
stage.extend(p.strip("./").split("/")[0] or "." for p in parts[:-1])
return None
def main() -> int:
paths = patched_paths(root / "Cargo.toml")
if not paths:
print("no patched paths to check")
return 0
keep = allowed(root / ".dockerignore")
bad = []
for p in paths:
top = p.strip("/").split("/")[0]
if top not in keep:
bad.append((p, top))
elif not (root / p).is_dir():
bad.append((p, None))
for path, top in bad:
if top is None:
print(f"Cargo.toml patches {path}, which does not exist", file=sys.stderr)
else:
print(
f"Cargo.toml patches {path}, but .dockerignore does not re-include {top!r}:\n"
f" the image build would not see it, and cargo would fail on it.\n"
f" Add `!{top}` to .dockerignore.",
file=sys.stderr,
)
copied = copied_before_cook(root / "Dockerfile")
if copied is not None and "." not in copied:
for p in paths:
top = p.strip("/").split("/")[0]
if top not in copied:
print(
f"Cargo.toml patches {p}, but the Dockerfile doesn't copy {top!r} into the\n"
f" stage that runs `cargo chef cook` before that step, so cooking the\n"
f" dependencies fails on it. Add `COPY {top}/ {top}/` before the cook.",
file=sys.stderr,
)
bad.append((p, top))
if bad:
return 1
print(f"build context and cook stage include every patched path: {', '.join(paths)}")
return 0
if __name__ == "__main__":
raise SystemExit(main())