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
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
15 changed files with 654 additions and 31 deletions
+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,
+57 -20
View File
@@ -567,20 +567,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 +584,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 +594,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);
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() {
if wants(SearchTracingField::EventType) && seen.insert(event.event.as_str().to_string()) {
document.index_keyword(TracingSearchField::EventType, event.event.as_str());
}
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 +685,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
+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),
+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;
+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();
}
}
+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