15 Commits
Author SHA1 Message Date
jcoffey-dev 9fa5433665 Merge pull request 'Release 2026.9.25' (#53) from bump/2026.9.25 into main
ci / fork-checks (push) Successful in 1m9s
publish / version (push) Successful in 1m1s
publish / publish-amd64 (push) Successful in 28m2s
publish / release (push) Successful in 1s
ci / build (push) Successful in 39m42s
publish / publish-arm64 (push) Successful in 40m42s
publish / binaries (push) Successful in 1m6s
2026-09-25 05:35:39 +00:00
jcoffey-dev f2605877f7 Release 2026.9.25
ci / fork-checks (pull_request) Successful in 21s
ci / build (pull_request) Successful in 7m36s
2026-09-24 22:27:11 -07:00
jcoffey-dev c521f060ba Merge pull request 'PostgreSQL search: find words inside URLs and file names in body text' (#52) from fix/pg-url-body-tokens into main
ci / fork-checks (push) Successful in 47s
ci / build (push) Canceled after 53m30s
2026-09-25 04:42:05 +00:00
jcoffey-dev 5927dda7e2 PostgreSQL search: find words inside URLs and file names in body text
ci / fork-checks (pull_request) Successful in 48s
ci / build (pull_request) Successful in 3m26s
After #37, address fields on PostgreSQL are split into words as the
built-in index splits them, but language text (subject, body,
attachments) still goes straight to PostgreSQL's parser, which keeps a
URL, host, path or file name as tokens of its own:
"https://x.example/shipping-support/" becomes a url, a host and a
url_path, "invoice-2024.pdf" a file. So TEXT/BODY "shipping" missed
messages where the word appears only inside a link, while RocksDB and
the other built-in backends found them: 8 messages across a handful
of searches in the rehearsal.

On insert, language text is now indexed as it was, followed by the
word parts of each token that holds a URL separator (/ . @ : ? = & # _
% + ~ \), split with SpaceTokenizer as keyword_terms() splits addresses.
The parts go through the same text search configuration as the rest of
the text, so they are stemmed like the words around them. Plain words,
words that only carry punctuation ("end.", "(see") and hyphenated words
(the parser already splits those) add nothing, so text without links
is indexed exactly as before. Each part is added once per document.
On sample mail, the text vector of a short order notice with three
links grows from 546 to 716 bytes, a newsletter with 25 tracking links
from 5586 to 6430, and a plain letter not at all.

On search, a query word written as a URL, host, file or hyphenated word
also matches as its word parts, ORed with the query as written, so
"shipping-support" or "invoice-2024.pdf" match the new parts and
documents indexed before this change still match as they did.

Existing messages keep their old vectors until they are reindexed (the
reindexAccounts task); new and reindexed messages match at once.

store::search_tests gains test_url_word_search: five bodies, 19 body
searches for words found only in a URL path, query string, host or
file name, the tokens as written, plain words and non-matches, with the
same expected ids on every backend. It passes on RocksDB, SQLite,
MySQL and PostgreSQL; on main PostgreSQL fails at the first ("shipping"
finds [3], not [0, 3]). On PostgreSQL the suite then stops at the
account sort assertion (query.rs:689) exactly as it does on main.
2026-09-24 21:35:12 -07:00
jcoffey-dev 9e49597ae4 Merge pull request 'Settings writes: wait for a burst to settle before reloading' (#51) from fix/settings-write-debounce into main
ci / fork-checks (push) Successful in 25s
ci / build (push) Canceled after 6m56s
2026-09-25 04:35:08 +00:00
jcoffey-dev 71ce11c57d Settings writes: wait for a burst to settle before reloading
ci / fork-checks (pull_request) Successful in 56s
ci / build (pull_request) Successful in 12m43s
A cluster rehearsal sent ten x:<Object>/set requests at once and got
ten full reloads on every node. #39's coalescing only joined writes
that queued behind a running reload, but the requests reached the
server about 33 ms apart and a reload takes tens of milliseconds, so
none overlapped one.

A full reload after a registry write now waits for writes to settle:
75 ms after the last one, and at most 250 ms after the first it
covers, so a steady stream still reloads at least four times a
second. 75 ms is a little over twice the gap the rehearsal saw between
requests. A single write pays it once: in the tests a settings write
takes about 140 ms instead of 60. The reload runs in a task of its
own, so a request that goes away doesn't cancel it for the others.
Each write takes the result of the first reload that started after it
was stored (the gate keeps the last 64 results), so applied true or
false still describes the reload that covered that write.

The 33 ms gap was a queue on the server, not password hashing: Basic
credentials are cached per Authorization header, so they are checked
once. Every authenticated HTTP request counted itself against the
account's rate limit by incrementing one counter per account in the
in-memory store, so parallel requests from one account queued on that
key: a row lock on PostgreSQL (a few round trips to the database
each) and conflict retries with a 50-300 ms backoff on RocksDB. An
account with the unlimitedRequests permission (administrators, by
default) passes the rate and concurrency limits anyway, so its
requests are no longer counted. Ten parallel Core/echo calls as the
admin now finish in 1-4 ms; before, they finished one after another
over 20 ms on a local PostgreSQL and 300-450 ms on RocksDB. Other
accounts still count every request.

system::auto_reload::settings_reload_tests: ten concurrent writes now
take one reload (the gate counts them; at most two allowed), all are
applied: true and in the running settings, and a single write takes
exactly one reload. RocksDB and PostgreSQL, 1 reload in 141-196 ms.
With the old behavior (no wait, requests counted) the same writes
took 5 reloads; without the wait but with the rate fix, 2.
cluster::broadcast (3 nodes, PostgreSQL + NATS) and system::reload
still pass.
2026-09-24 21:15:44 -07:00
jcoffey-dev 51b159a1a2 Merge pull request 'Tracers whose settings change start over on reload' (#50) from fix/tracer-live-reload into main
ci / fork-checks (push) Successful in 24s
ci / build (push) Successful in 34m49s
2026-09-25 03:57:52 +00:00
jcoffey-dev a891667149 Tracers whose settings change start over on reload
ci / fork-checks (pull_request) Successful in 47s
ci / build (pull_request) Successful in 3m49s
A cluster rehearsal moved a Log tracer to another directory: the write
was reported x:settingsReload applied:true, but the tracer kept writing
to the old file until a restart. Telemetry::update only refreshed each
running tracer's events, level and lossiness; a tracer's own settings
(path, prefix, rotation, format, endpoint, headers, ...) stayed as built.

Each tracer now carries a hash of the registry object it was built
from, less the fields that change in place. The reload compares it with
the running tracer's: unchanged ones are updated in place as before,
changed ones are started over, new ones started and removed ones
stopped. Only tracers this server started are removed; upstream removed
every subscriber not in the settings, which also cut off live-tracing
streams on each reload.

Starting over is a swap in the collector, so no event is lost or
written twice: a subscriber registered under a running one's id
replaces it between two collection passes. The old one's batch is sent
first (what its full channel can't take moves to the new one), and
dropping it closes its channel, so its task writes what is queued and
ends. Per tracer kind:

- Log: a tracer started over on the same files (rotation or format
  changed) waits for the old one to finish, so lines don't interleave.
- Webhook: the task held a sender of its own channel for retries, so
  it never ended; retries now use a weak sender, and pending events are
  posted when the channel closes.
- OpenTelemetry: pending logs and spans are exported when the channel
  closes instead of dropped, and a span that was open across the swap
  is exported by the new tracer with the events it saw.
- Console and journal: nothing kept between batches.
- Trace history: built from the tracing store, which takes a restart,
  so it is never started over.

No kind needs a restart, so x:settingsReload doesn't gain one.

system::tracer_reload::tracer_reload_tests (new): a Log tracer created
over JMAP writes to its directory; its path is changed over JMAP while
2000 numbered events are emitted; after the reload, events land in the
new file and not the old one, each numbered event is in exactly one of
the two files, and a destroyed tracer writes nothing. On main the new
file never appears.
2026-09-24 20:52:46 -07:00
jcoffey-dev 59e631eded Merge pull request 'Every node records DMARC and TLS results for the aggregate reports' (#47) from fix/front-node-dmarc into main
ci / fork-checks (push) Successful in 1m27s
ci / build (push) Successful in 40m18s
2026-09-25 01:46:19 +00:00
jcoffey-dev 716800d681 Merge pull request 'Publish: accept tags on release/* branches for hotfix releases' (#48) from ci/publish-release-branches into main
ci / fork-checks (push) Successful in 19s
ci / build (push) Canceled after 7m17s
2026-09-25 01:38:56 +00:00
jcoffey-dev 4b85113262 Publish: accept tags on release/* branches for hotfix releases
ci / fork-checks (pull_request) Successful in 20s
ci / build (pull_request) Successful in 7m16s
The publish workflow only built a tag whose commit is on main. That keeps
every image tied to reviewed code, but it means production can only get a
fix together with everything that has landed on main since its release.

A tag on a release/* branch is now accepted too. A hotfix branch starts at
an earlier release tag, takes fixes through pull requests into it (so the
code is still reviewed and CI-tested before it is tagged), bumps
brand_version! and is tagged there. The tag must still equal
v<brand_version!>, and the step prints which branch it was found on.

A tag runs the workflow file from its own commit, so a hotfix branch that
starts before this change needs this commit cherry-picked onto it before
its tag is pushed.
2026-09-24 18:31:24 -07:00
jcoffey-dev 9cc9951428 Merge pull request 'Report reschedules keep the task queue readable' (#46) from fix/report-reschedule into main
ci / fork-checks (push) Successful in 32s
ci / build (push) Canceled after 8m20s
2026-09-25 01:30:35 +00:00
jcoffey-dev 5dde9793eb Every node records DMARC and TLS results for the aggregate reports
ci / fork-checks (pull_request) Successful in 43s
ci / build (pull_request) Successful in 17m25s
The report scheduler dropped DMARC and TLS events on a node whose role
lacks outboundMta (upstream never started it there, so they sat in a
channel nobody read). Mail received on a front node therefore never
reached an aggregate report, which is meant to cover all of a domain's
inbound mail, whichever node received it. In rehearsal, five messages
received on port 25 on a front node were missing from every report.

- The report scheduler records on every node. Recording is a store write
  the nodes already share, so it needs nothing from the outbound MTA.
  Building and sending a report (the DmarcReport and TlsReport tasks) stay
  with outboundMta nodes, as the task manager already enforces.
- More nodes now append to one report at once. Appends already guard the
  report's versioned primary key; a write that loses now retries up to ten
  times after a short random pause, not three times at once.
- The node sending a report deletes it only if it is unchanged since it
  was read, and reads it again otherwise, so a record another node appends
  meanwhile goes out with the report instead of being deleted unsent.

Test: cluster::front_reports (PostgreSQL and MySQL). A front node's
results appear in the report the MTA node sends, alongside eight appended
at once from both nodes, and the front node never runs the report task.
It fails on main: the front node's results are never recorded.
2026-09-24 18:28:00 -07:00
jcoffey-dev 1a7859a8cc Report reschedules keep the task queue readable
ci / fork-checks (pull_request) Successful in 52s
ci / build (pull_request) Successful in 4m0s
Setting deliverAt on an internal DMARC or TLS report wrote the new task
queue row with the report's object type (0x21, 0x6e) instead of the task
type (7, 8), and left the task row at its old due. The task manager's scan
failed on that row with store.data-corruption ("Failed to iterate over task
queue"), and because the error ended the whole scan, every task due after
the row stopped running on every node.

- reschedule_ops writes the new queue row through schedule_task_with_id, so
  it carries the task type and the task row gets the new due. It removes
  the row the task is actually queued under (the task's due, which differs
  from deliverAt once the task has been retried) and any row an earlier
  reschedule left at deliverAt.
- x:DmarcInternalReport/set and x:TlsInternalReport/set lock the report's
  task while they move it, as x:Task/set does, refuse while the report is
  being sent, release the locks however the request ends, and wake the task
  manager.
- The task manager logs a queue row it can't read (id, due, key, value) and
  skips it instead of ending the scan. It then repairs the row from its task:
  the row is rewritten with the task's type, and a row with no task behind
  it is removed. A row holding a report's object type for a report task is
  what the old reschedule wrote: the task is moved to that row's time, as
  the reschedule intended, and its old queue row is removed. Stores that
  already hold such a row recover on their own once it comes due.
- x:Task/query with a type filter skips an unreadable row instead of
  failing.

Test: smtp::reporting::reschedule (RocksDB and PostgreSQL). It fails on
main: x:Task/get shows the old due, and with that check removed, neither
report nor a later task ever runs.
2026-09-24 18:09:44 -07:00
jcoffey-dev b90a7f173e Merge pull request 'Cluster role changes apply to delivery and tasks without a restart' (#44) from fix/live-role-changes into main
ci / fork-checks (push) Successful in 17s
ci / build (push) Canceled after 34m3s
2026-09-25 00:56:30 +00:00
30 changed files with 1832 additions and 127 deletions
+12 -4
View File
@@ -31,8 +31,11 @@
# crates/types/src/branding.rs, not Cargo.toml, and the image is tagged
# with it, so a tag beside an unbumped macro would publish an image that
# reports a different version from its tag.
# * the tag must be on main, so an image never describes code that was never
# reviewed onto the default branch.
# * the tag must be on main or on a release/* branch, so an image never
# describes code that was never reviewed onto one of them. A release/*
# branch carries a hotfix: it starts at an earlier release tag, takes
# fixes through pull requests into it, and is tagged there, so production
# can get a fix without everything that has landed on main since.
#
# :latest moves with every published tag: tags are cut by the weekly release
# (or by hand for a real release); there are no prerelease tags here.
@@ -74,8 +77,13 @@ jobs:
echo "Refusing to publish an image that would report the wrong version." >&2
exit 1
fi
git merge-base --is-ancestor "$(git rev-parse "${TAG}^{commit}")" origin/main \
|| { echo "$TAG is not on main" >&2; exit 1; }
commit="$(git rev-parse "${TAG}^{commit}")"
on=""
for ref in origin/main $(git for-each-ref --format='%(refname:short)' 'refs/remotes/origin/release/*'); do
if git merge-base --is-ancestor "$commit" "$ref"; then on="$ref"; break; fi
done
[ -n "$on" ] || { echo "$TAG is not on main or a release/* branch" >&2; exit 1; }
echo "$TAG is on $on"
echo "version=$V" >> "$GITHUB_OUTPUT"
echo "version $V"
+12
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::auth::AccessToken;
@@ -18,6 +20,16 @@ impl Server {
access_token: &AccessToken,
addr: IpAddr,
) -> trc::Result<Option<InFlight>> {
// inbuxa: an account with unlimited requests passes both limits
// below anyway, so don't count its requests. The count is a write to
// one counter per account in the in-memory store, and concurrent
// requests from one account queue on that key (a row lock on SQL,
// conflict retries on RocksDB): in a cluster rehearsal ten parallel
// admin writes were accepted one after another, about 33 ms apart.
if access_token.has_permission(Permission::UnlimitedRequests) {
return Ok(None);
}
let rate_reset = if let Some(rate) = &self.core.network.http.rate_authenticated {
if self.is_ip_allowed(addr) {
None
+133 -17
View File
@@ -7,7 +7,7 @@
*/
use crate::{
Core, Server,
BuildServer, Core, Server,
config::{
server::{Listeners, tls::parse_certificates},
storage::Storage,
@@ -245,21 +245,73 @@ fn error_object(error: &Error) -> Option<ObjectId> {
// until someone reloaded. Writes to objects the settings are built from now
// reload them, here and across the cluster, as ReloadSettings does.
/// Coalesces the full reloads that registry writes trigger: a write waits for
/// a reload that started after it was stored, and joins one if it can, so a
/// burst of writes costs a reload or two rather than one each.
#[derive(Default)]
/// Coalesces the full reloads that registry writes trigger. A write waits
/// for more writes before a reload starts (see [`WRITE_QUIET`]), then
/// takes the result of the first reload that started after it was stored,
/// so a burst of writes, or a request with many objects, costs one reload
/// or two rather than one each.
pub struct SettingsReloadGate {
requested: std::sync::atomic::AtomicU64,
state: tokio::sync::Mutex<SettingsReloadState>,
reloads: std::sync::atomic::AtomicU64,
state: parking_lot::Mutex<SettingsReloadState>,
completed: tokio::sync::watch::Sender<u64>,
}
#[derive(Default)]
struct SettingsReloadState {
completed: u64,
refused: Option<String>,
/// A reload is waiting for writes to settle, or running.
scheduled: bool,
/// When the oldest write not yet covered by a reload was stored, and
/// the newest.
first_write: Option<std::time::Instant>,
last_write: Option<std::time::Instant>,
/// Recent reloads, oldest first: the last write each covered, and why
/// it was refused, if it was.
results: std::collections::VecDeque<(u64, Option<String>)>,
}
impl Default for SettingsReloadGate {
fn default() -> Self {
Self {
requested: Default::default(),
reloads: Default::default(),
state: Default::default(),
completed: tokio::sync::watch::Sender::new(0),
}
}
}
impl SettingsReloadGate {
/// How many full reloads registry writes have run.
pub fn reloads(&self) -> u64 {
self.reloads.load(std::sync::atomic::Ordering::Relaxed)
}
}
impl SettingsReloadState {
/// The result of the reload that covered write `ticket`, once it ran.
fn result_for(&self, ticket: u64) -> Option<Result<(), String>> {
self.results
.iter()
.find(|(covers, _)| *covers >= ticket)
.map(|(_, refused)| refused.clone().map_or(Ok(()), Err))
}
}
/// How long a full reload waits after the last registry write for another.
/// Parallel requests reach the server tens of milliseconds apart (in a
/// cluster rehearsal, ten x:<Object>/set requests sent at once arrived about
/// 33 ms apart and each got a reload of its own), so the window is a little
/// over twice that. A single write pays it once, on top of the reload.
pub const WRITE_QUIET: std::time::Duration = std::time::Duration::from_millis(75);
/// The longest a full reload waits after the first write it covers, so a
/// steady stream of writes still reloads at least this often.
pub const WRITE_MAX_WAIT: std::time::Duration = std::time::Duration::from_millis(250);
/// How many past reload results a waiting write can look up.
const RELOAD_RESULTS: usize = 64;
/// The reload a write to `object` calls for: the object to reload, or None
/// when the running settings don't hold that object (accounts, domains and
/// other data read as needed, stores, which take a restart, and objects with
@@ -370,21 +422,85 @@ impl Server {
return Some(result);
}
// inbuxa: #39 joined only writes that queued behind a running
// reload; requests that arrive tens of milliseconds apart never
// overlapped one, so each got a reload of its own. The reload now
// waits until writes settle (WRITE_QUIET after the last one, at
// most WRITE_MAX_WAIT after the first) and covers them all. It runs
// in a task of its own, so a request that goes away doesn't take
// it with it; each write then takes the result of the reload that
// started after it was stored.
let gate = &self.inner.data.settings_reload;
let ticket = gate
.requested
.fetch_add(1, std::sync::atomic::Ordering::SeqCst)
+ 1;
let mut state = gate.state.lock().await;
if state.completed >= ticket {
// A reload that started after this write was stored has run
return Some(state.refused.clone().map_or(Ok(()), Err));
let now = std::time::Instant::now();
{
let mut state = gate.state.lock();
state.first_write.get_or_insert(now);
state.last_write = Some(now);
}
let covers = gate.requested.load(std::sync::atomic::Ordering::SeqCst);
let result = self.reload_and_broadcast(change).await;
state.completed = covers;
state.refused = result.clone().err();
Some(result)
loop {
let mut completed = {
let mut state = gate.state.lock();
if let Some(result) = state.result_for(ticket) {
return Some(result);
}
if !state.scheduled {
state.scheduled = true;
let server = self.clone();
tokio::spawn(async move {
server.run_write_reload(change).await;
});
}
gate.completed.subscribe()
};
if completed.changed().await.is_err() {
return Some(Err("The settings reload was interrupted".to_string()));
}
}
}
/// Waits for registry writes to settle, then reloads the settings once
/// for all the writes stored so far.
async fn run_write_reload(&self, change: RegistryChange) {
let gate = &self.inner.data.settings_reload;
loop {
let deadline = {
let state = gate.state.lock();
let now = std::time::Instant::now();
let first = state.first_write.unwrap_or(now);
let last = state.last_write.unwrap_or(now);
(last + WRITE_QUIET).min(first + WRITE_MAX_WAIT)
};
if deadline <= std::time::Instant::now() {
break;
}
tokio::time::sleep_until(deadline.into()).await;
}
// Writes stored from here on wait for the next reload
let covers = {
let mut state = gate.state.lock();
state.first_write = None;
state.last_write = None;
gate.requested.load(std::sync::atomic::Ordering::SeqCst)
};
gate.reloads
.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
let result = self.inner.build_server().reload_and_broadcast(change).await;
{
let mut state = gate.state.lock();
if state.results.len() == RELOAD_RESULTS {
state.results.pop_front();
}
state.results.push_back((covers, result.err()));
state.scheduled = false;
}
gate.completed.send_replace(covers);
}
async fn reload_and_broadcast(&self, change: RegistryChange) -> Result<(), String> {
+48
View File
@@ -31,6 +31,10 @@ pub struct TelemetrySubscriber {
pub interests: Interests,
pub typ: TelemetrySubscriberType,
pub lossy: bool,
/// inbuxa: a hash of the settings the running tracer is built from
/// (everything but its events, level and lossiness, which change in
/// place), so a reload can tell which tracers to start over.
pub settings: u64,
}
#[allow(clippy::large_enum_variant)]
@@ -167,6 +171,7 @@ impl Tracers {
for tracer in bp.list_infallible::<Tracer>().await {
let id = tracer.id;
let tracer = tracer.object;
let settings = tracer_settings(&tracer);
let level;
let lossy;
let events;
@@ -379,6 +384,7 @@ impl Tracers {
interests: Default::default(),
lossy,
typ,
settings,
};
// Parse disabled events
@@ -426,6 +432,7 @@ impl Tracers {
for hook in bp.list_infallible::<WebHook>().await {
let id = hook.id;
let hook = hook.object;
let settings = webhook_settings(&hook);
if !hook.enable {
continue;
@@ -448,6 +455,7 @@ impl Tracers {
id: format!("w_{}", id.id()),
interests: Default::default(),
lossy: hook.lossy,
settings,
typ: TelemetrySubscriberType::Webhook(WebhookTracer {
url: hook.url,
timeout: hook.timeout.into_inner(),
@@ -516,6 +524,8 @@ impl Tracers {
data: storage.data.clone(),
}),
lossy: true,
// Stores take a restart
settings: 0,
});
}
@@ -541,6 +551,7 @@ impl Tracers {
buffered: true,
}),
lossy: false,
settings: 0,
});
}
} else {
@@ -568,6 +579,7 @@ impl Tracers {
buffered: true,
}),
lossy: false,
settings: 0,
});
}
@@ -701,6 +713,42 @@ impl Metrics {
}
}
// inbuxa: what a tracer is built from, less what changes in place
macro_rules! in_place_reset {
($tracer:expr) => {{
$tracer.enable = true;
$tracer.level = Default::default();
$tracer.lossy = false;
$tracer.events = Default::default();
$tracer.events_policy = Default::default();
}};
}
fn settings_hash(settings: &impl std::fmt::Debug) -> u64 {
use std::hash::{Hash, Hasher};
let mut hasher = std::collections::hash_map::DefaultHasher::new();
format!("{settings:?}").hash(&mut hasher);
hasher.finish()
}
fn tracer_settings(tracer: &Tracer) -> u64 {
let mut tracer = tracer.clone();
match &mut tracer {
Tracer::Log(tracer) => in_place_reset!(tracer),
Tracer::Stdout(tracer) => in_place_reset!(tracer),
Tracer::Journal(tracer) => in_place_reset!(tracer),
Tracer::OtelHttp(tracer) => in_place_reset!(tracer),
Tracer::OtelGrpc(tracer) => in_place_reset!(tracer),
}
settings_hash(&tracer)
}
fn webhook_settings(hook: &WebHook) -> u64 {
let mut hook = hook.clone();
in_place_reset!(hook);
settings_hash(&hook)
}
fn apply_events(
event_types: impl IntoIterator<Item = EventType>,
policy: EventPolicy,
+34 -9
View File
@@ -14,15 +14,26 @@ pub mod webhooks;
use tracers::log::spawn_log_tracer;
use tracers::otel::spawn_otel_tracer;
use tracers::stdout::spawn_console_tracer;
use ahash::AHashMap;
use parking_lot::Mutex;
use trc::{Collector, ipc::subscriber::SubscriberBuilder};
use webhooks::spawn_webhook_tracer;
use crate::config::telemetry::{Telemetry, TelemetrySubscriberType};
/// inbuxa: the tracers this server started, by subscriber id, with the
/// settings each was built from. Live-tracing streams and other subscribers
/// registered elsewhere aren't listed, so a reload leaves them running.
static RUNNING_TRACERS: Mutex<Option<AHashMap<String, u64>>> = Mutex::new(None);
impl Telemetry {
pub fn enable(self) {
let mut running = RUNNING_TRACERS.lock();
let running = running.get_or_insert_with(AHashMap::new);
// Spawn tracers
for tracer in self.tracers.subscribers {
running.insert(tracer.id.clone(), tracer.settings);
tracer.typ.spawn(
SubscriberBuilder::new(tracer.id)
.with_interests(tracer.interests)
@@ -37,25 +48,39 @@ impl Telemetry {
Collector::reload();
}
// inbuxa: upstream only refreshed the events, level and lossiness of a
// tracer that was already running, so a Log tracer moved to another
// path (or any tracer whose own settings changed) kept going as it was
// built until a restart, while the reload reported the change applied.
// A tracer whose settings changed is now started over: the new one is
// registered under the same id and the collector swaps it in at an
// event boundary, so no event is lost or written twice (see
// Update::RegisterSubscriber); the old one writes what it has queued
// and stops.
pub fn update(self) {
let mut running = RUNNING_TRACERS.lock();
let running = running.get_or_insert_with(AHashMap::new);
// Remove tracers that are no longer active
let active_subscribers = Collector::get_subscribers();
for subscribed_id in &active_subscribers {
if !self
running.retain(|id, _| {
let keep = self
.tracers
.subscribers
.iter()
.any(|tracer| tracer.id == *subscribed_id)
{
Collector::remove_subscriber(subscribed_id.clone());
.any(|tracer| tracer.id == *id);
if !keep {
Collector::remove_subscriber(id.clone());
}
}
keep
});
// Activate new tracers or update existing ones
// Start new tracers, start over those whose settings changed and
// update the rest in place
for tracer in self.tracers.subscribers {
if active_subscribers.contains(&tracer.id) {
if running.get(&tracer.id) == Some(&tracer.settings) {
Collector::update_subscriber(tracer.id, tracer.interests, tracer.lossy);
} else {
running.insert(tracer.id.clone(), tracer.settings);
tracer.typ.spawn(
SubscriberBuilder::new(tracer.id)
.with_interests(tracer.interests)
@@ -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 std::{path::PathBuf, time::SystemTime};
@@ -15,9 +17,27 @@ use tokio::{
};
use trc::{TelemetryEvent, ipc::subscriber::SubscriberBuilder, serializers::text::FmtWriter};
// inbuxa: when a Log tracer is started over on the same files (its rotation
// or format changed), the new one waits for the old one to write what it
// has queued, so their lines don't interleave. Keyed by path and prefix;
// each entry is the last tracer's "done" signal, sent when it ends.
type LogFileOwners = ahash::AHashMap<(String, String), tokio::sync::oneshot::Receiver<()>>;
static LOG_FILE_OWNERS: parking_lot::Mutex<Option<LogFileOwners>> = parking_lot::Mutex::new(None);
pub(crate) fn spawn_log_tracer(builder: SubscriberBuilder, settings: LogTracer) {
let (done_tx, done_rx) = tokio::sync::oneshot::channel::<()>();
let previous = LOG_FILE_OWNERS
.lock()
.get_or_insert_with(Default::default)
.insert((settings.path.clone(), settings.prefix.clone()), done_rx);
let (_, mut rx) = builder.register();
tokio::spawn(async move {
// Dropped when this tracer ends, however it ends
let _done = done_tx;
if let Some(previous) = previous {
let _ = previous.await;
}
if let Some(writer) = settings.build_writer().await {
let mut buf = FmtWriter::new(writer)
.with_ansi(settings.ansi)
+22 -1
View File
@@ -47,6 +47,10 @@ pub(crate) fn spawn_otel_tracer(builder: SubscriberBuilder, mut otel: OtelTracer
let mut pending_spans = Vec::new();
let mut active_spans = AHashMap::new();
let mut closing = false;
let started = std::time::SystemTime::now()
.duration_since(std::time::SystemTime::UNIX_EPOCH)
.map_or(0, |d| d.as_secs());
loop {
// Wait for the next event or timeout
@@ -75,12 +79,26 @@ pub(crate) fn spawn_otel_tracer(builder: SubscriberBuilder, mut otel: OtelTracer
events.iter().chain(std::iter::once(&event)),
&instrumentation,
));
} else if span.inner.timestamp < started {
// inbuxa: a span that was open when this
// tracer replaced another one (its settings
// changed) is exported with its end event
// rather than dropped
pending_spans.push(build_span_data(
span,
&event,
std::iter::once(&event),
&instrumentation,
));
}
}
}
}
Ok(None) => {
break;
// inbuxa: the tracer was removed or replaced; export
// what is pending now rather than drop it
closing = true;
next_delivery = Instant::now();
}
Err(_) => (),
}
@@ -131,6 +149,9 @@ pub(crate) fn spawn_otel_tracer(builder: SubscriberBuilder, mut otel: OtelTracer
}
}
}
if closing {
break;
}
wakeup_time = next_retry.unwrap_or(LONG_1Y_SLUMBER);
}
});
+22 -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::{LONG_1Y_SLUMBER, config::telemetry::WebhookTracer};
@@ -25,6 +27,11 @@ use trc::{
pub(crate) fn spawn_webhook_tracer(builder: SubscriberBuilder, settings: WebhookTracer) {
let (tx, mut rx) = builder.register();
// inbuxa: failed deliveries come back through a weak sender, so the
// channel closes when the collector drops this webhook (removed, or
// replaced after a settings change) and the task ends; upstream held a
// sender here and the task outlived its subscription
let tx = tx.downgrade();
tokio::spawn(async move {
let settings = Arc::new(settings);
let mut wakeup_time = LONG_1Y_SLUMBER;
@@ -58,6 +65,15 @@ pub(crate) fn spawn_webhook_tracer(builder: SubscriberBuilder, settings: Webhook
}
}
Ok(None) => {
// inbuxa: deliver what is pending rather than drop it
if !pending_events.is_empty() {
spawn_webhook_handler(
settings.clone(),
in_flight.clone(),
std::mem::take(&mut pending_events),
tx.clone(),
);
}
break;
}
Err(_) => (),
@@ -102,7 +118,7 @@ fn spawn_webhook_handler(
settings: Arc<WebhookTracer>,
in_flight: Arc<AtomicBool>,
events: EventBatch,
webhook_tx: mpsc::Sender<EventBatch>,
webhook_tx: mpsc::WeakSender<EventBatch>,
) {
tokio::spawn(async move {
in_flight.store(true, Ordering::Relaxed);
@@ -113,7 +129,11 @@ fn spawn_webhook_handler(
if let Err(err) = post_webhook_events(&settings, &wrapper).await {
trc::event!(Telemetry(TelemetryEvent::WebhookError), Details = err);
if webhook_tx.send(wrapper.events.into_inner()).await.is_err() {
let sent = match webhook_tx.upgrade() {
Some(webhook_tx) => webhook_tx.send(wrapper.events.into_inner()).await.is_ok(),
None => false,
};
if !sent {
trc::event!(
Server(ServerEvent::ThreadError),
Details = "Failed to send failed webhook events back to main thread",
+63 -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::{
@@ -15,22 +17,42 @@ use jmap_proto::{error::set::SetError, types::state::State};
use jmap_tools::{Key, Value};
use registry::{
jmap::IntoValue,
schema::prelude::{Object, ObjectInner, ObjectType, Property},
schema::{
prelude::{Object, ObjectInner, ObjectType, Property},
structs::Task,
},
types::{EnumImpl, datetime::UTCDateTime},
};
use services::task_manager::lock::TaskLockManager;
use smtp::reporting::index::{ExternalReportIndex, InternalReportIndex};
use std::str::FromStr;
use store::{
U64_LEN, ValueKey,
registry::{RegistryFilter, RegistryFilterValue, RegistryQuery},
write::{BatchBuilder, RegistryClass, ValueClass, key::KeySerializer},
write::{BatchBuilder, RegistryClass, TaskQueueClass, ValueClass, key::KeySerializer},
};
use trc::AddContext;
use types::id::Id;
pub(crate) async fn report_set(
mut set: RegistrySetResponse<'_>,
set: RegistrySetResponse<'_>,
) -> trc::Result<RegistrySetResponse<'_>> {
// inbuxa: task locks taken to reschedule reports are released however
// the request ends; a held lock is renewed, so a leaked one would keep
// the report's task from ever running
let server = set.server;
let mut locked_tasks = Vec::new();
let result = report_set_locked(set, &mut locked_tasks).await;
for task_id in locked_tasks {
server.remove_index_lock(task_id).await;
}
result
}
async fn report_set_locked<'x>(
mut set: RegistrySetResponse<'x>,
locked_tasks: &mut Vec<u64>,
) -> trc::Result<RegistrySetResponse<'x>> {
let object_id = set.object_type.to_id();
// Reports cannot be created
@@ -89,12 +111,45 @@ pub(crate) async fn report_set(
.get_value::<Object>(ValueKey::from(key.clone()))
.await?
{
// inbuxa: the report's task shares its id. Hold the task
// while its queue rows move, as x:Task/set does, and move the
// row the task is actually queued under
if !set.server.try_lock_task(item_id).await {
set.response.not_updated.append(
id,
SetError::forbidden().with_description(
"The report is being sent and cannot be rescheduled".to_string(),
),
);
continue;
}
locked_tasks.push(item_id);
let queued = set
.server
.store()
.get_value::<Task>(ValueKey::from(ValueClass::TaskQueue(
TaskQueueClass::Task { id: item_id },
)))
.await?;
match &mut report_obj.inner {
ObjectInner::DmarcInternalReport(report) => {
report.reschedule_ops(&mut batch, item_id, report_obj.revision, deliver_at);
report.reschedule_ops(
&mut batch,
item_id,
report_obj.revision,
deliver_at,
queued.as_ref(),
);
}
ObjectInner::TlsInternalReport(report) => {
report.reschedule_ops(&mut batch, item_id, report_obj.revision, deliver_at);
report.reschedule_ops(
&mut batch,
item_id,
report_obj.revision,
deliver_at,
queued.as_ref(),
);
}
_ => {}
}
@@ -156,6 +211,9 @@ pub(crate) async fn report_set(
.write(batch.build_all())
.await
.caused_by(trc::location!())?;
// inbuxa: a rescheduled report may now be due sooner than the task
// manager's next scan
set.server.notify_task_queue();
}
Ok(set)
+4 -9
View File
@@ -463,15 +463,10 @@ pub(crate) async fn task_query(
.set_values(typ.is_some()),
|key, value| {
if let Some(typ) = typ {
let task_type =
TaskType::from_id(value.deserialize_be_u16(0)?).ok_or_else(|| {
trc::StoreEvent::DataCorruption
.into_err()
.ctx(trc::Key::Key, key.to_vec())
.ctx(trc::Key::Value, value.to_vec())
.caused_by(trc::location!())
})?;
if task_type != typ {
// inbuxa: a row whose type can't be read matches no type
// filter; the task manager logs and repairs it
let task_type = value.deserialize_be_u16(0).ok().and_then(TaskType::from_id);
if task_type != Some(typ) {
return Ok(true);
}
}
+133 -6
View File
@@ -30,6 +30,7 @@ use common::network::limiter::ConcurrencyLimiter;
use common::network::{ServerInstance, TcpAcceptor};
use common::{Inner, Server};
use registry::schema::enums::TaskType;
use registry::schema::prelude::ObjectType;
use registry::schema::structs::{
Task, TaskManager, TaskRetryStrategy, TaskStatus, TaskStatusFailed, TaskStatusRetry,
};
@@ -298,6 +299,7 @@ impl TaskQueueManager for Server {
// Retrieve tasks pending to be processed
let mut tasks = Vec::new();
let mut unreadable = Vec::new();
let now = Instant::now();
let mut next_event = None;
ipc.revision += 1;
@@ -311,12 +313,21 @@ impl TaskQueueManager for Server {
let task_id = key.deserialize_be_u64(U64_LEN)?;
if task_due <= now_timestamp {
let task_type_idx = value.deserialize_be_u16(0)?;
let task_type = TaskType::from_id(task_type_idx).ok_or_else(|| {
trc::StoreEvent::DataCorruption
.caused_by(trc::location!())
.ctx(trc::Key::Value, value)
})?;
// inbuxa: a row whose task type can't be read is
// set aside, not allowed to end the scan: every
// task due after it would wait behind it
let Some((task_type_idx, task_type)) = value
.deserialize_be_u16(0)
.ok()
.and_then(|idx| TaskType::from_id(idx).map(|typ| (idx, typ)))
else {
unreadable.push(UnreadableDueRow {
due: task_due,
id: task_id,
value: value.to_vec(),
});
return Ok(true);
};
// inbuxa: running here under a lease this node
// renews; don't hand it to a worker again
if task_locks.is_held(task_id) {
@@ -389,6 +400,11 @@ impl TaskQueueManager for Server {
);
});
if !unreadable.is_empty() && repair_due_rows(self, unreadable).await {
// Look again at once for the rows that were rewritten
self.notify_task_queue();
}
if !tasks.is_empty() {
trc::event!(
TaskManager(TaskManagerEvent::TaskAcquired),
@@ -819,3 +835,114 @@ impl TaskResult {
)
}
}
/// inbuxa: a task queue row whose task type could not be read.
struct UnreadableDueRow {
due: u64,
id: u64,
value: Vec<u8>,
}
/// inbuxa: logs each unreadable queue row and repairs it from the task it
/// schedules. The task row says what the task is, so the queue row is
/// rewritten with that task's type; a row with no task behind it is removed.
///
/// Rescheduling an internal DMARC or TLS report wrote the report's object
/// type into the queue row instead of the task type. Such a row is the time
/// an administrator chose, so the task is moved to it as the reschedule
/// meant to do: the task row takes that due, and a queue row left at the
/// task's previous due is removed. Returns whether any row was repaired.
async fn repair_due_rows(server: &Server, rows: Vec<UnreadableDueRow>) -> bool {
let mut repaired = false;
for row in rows {
let UnreadableDueRow { due, id, value } = row;
trc::error!(
trc::StoreEvent::DataCorruption
.into_err()
.id(id)
.ctx(trc::Key::Due, trc::Value::Timestamp(due))
.ctx(
trc::Key::Key,
[due.to_be_bytes(), id.to_be_bytes()].concat()
)
.ctx(trc::Key::Value, value.clone())
.details("Unreadable task queue row skipped")
.caused_by(trc::location!())
);
let task_key = ValueClass::TaskQueue(TaskQueueClass::Task { id });
let due_key = ValueClass::TaskQueue(TaskQueueClass::Due { id, due });
let task = match server
.store()
.get_value::<Task>(ValueKey::from(task_key.clone()))
.await
{
Ok(task) => task,
Err(err) => {
trc::error!(
err.id(id)
.details("Failed to read the task of an unreadable queue row.")
.caused_by(trc::location!())
);
continue;
}
};
let mut batch = BatchBuilder::new();
let action = if let Some(mut task) = task {
let task_type = task.object_type();
batch.assert_value(task_key.clone(), AssertValue::Some);
if rescheduled_report_type(&value) == Some(task_type) {
let old_due = task.due_timestamp();
if old_due != due {
batch.clear(ValueClass::TaskQueue(TaskQueueClass::Due {
id,
due: old_due,
}));
}
task.set_status(TaskStatus::at(due as i64));
}
batch
.set(due_key, task_type.to_id().serialize())
.set(task_key, task.to_pickled_vec());
"Rewrote the queue row from its task."
} else {
batch.clear(due_key);
"Removed a queue row with no task."
};
match server.store().write(batch.build_all()).await {
Ok(_) => {
repaired = true;
trc::event!(
TaskManager(TaskManagerEvent::TaskIgnored),
Id = id,
Due = trc::Value::Timestamp(due),
Reason = action,
);
}
Err(err) if err.matches(trc::EventType::Store(trc::StoreEvent::AssertValueFailed)) => {
// The task went away meanwhile; the next scan looks again
}
Err(err) => {
trc::error!(
err.id(id)
.details("Failed to repair an unreadable queue row.")
.caused_by(trc::location!())
);
}
}
}
repaired
}
/// inbuxa: the task type a report reschedule meant, when a queue row holds
/// an internal report's object type (the value that reschedule wrote).
fn rescheduled_report_type(value: &[u8]) -> Option<TaskType> {
let id = u16::from_be_bytes(value.get(..2)?.try_into().ok()?);
match ObjectType::from_id(id)? {
ObjectType::DmarcInternalReport => Some(TaskType::DmarcReport),
ObjectType::TlsInternalReport => Some(TaskType::TlsReport),
_ => None,
}
}
+4 -2
View File
@@ -47,8 +47,10 @@ impl SpawnQueueManager for IpcReceivers {
// inbuxa: upstream started these only when the node's role included
// outboundMta at boot, so turning the role on later did nothing and
// turning it off left them delivering until a restart. They now run
// on every node and follow the role live (see Queue::start and the
// report scheduler). This also drains the queue channel on nodes
// on every node: the queue follows the role live (see Queue::start),
// and the report scheduler records DMARC and TLS results on every
// node, whatever its role (see reporting/scheduler.rs). This also
// drains the queue channel on nodes
// without the role, where every queued message's refresh used to sit
// in a channel nobody read until it filled and queueing blocked.
if !core.storage.registry.is_recovery_mode() {
+43 -23
View File
@@ -2,9 +2,12 @@
* SPDX-FileCopyrightText: 2020 Stalwart Labs LLC <[email protected]>
*
* SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-SEL
*
* Modified by Coffey Labs in 2026 for INBUXA.
*/
use super::AggregateTimestamp;
use super::shared::{MAX_WRITE_RETRIES, Revisioned, write_retry_pause};
use crate::{
core::Session,
queue::RecipientDomain,
@@ -349,29 +352,43 @@ impl DmarcReporting for Server {
let object_id = ObjectType::DmarcInternalReport.to_id();
let key = ValueClass::Registry(RegistryClass::Item { object_id, item_id });
let Some(report) = self
.store()
.get_value::<DmarcInternalReport>(ValueKey::from(key.clone()))
.await
.caused_by(trc::location!())?
else {
return Ok(());
};
// Delete report. inbuxa: only the version read here, so a record
// another node appends meanwhile is sent with it rather than lost
let mut attempt = 0;
let report = loop {
let Some(Revisioned {
revision,
value: report,
}) = self
.store()
.get_value::<Revisioned<DmarcInternalReport>>(ValueKey::from(key.clone()))
.await
.caused_by(trc::location!())?
else {
return Ok(());
};
// Delete report
let mut batch = BatchBuilder::new();
batch.clear(key).clear(RegistryClass::PrimaryKey {
object_id: object_id.into(),
index_id: Property::Domain.to_id(),
key: KeySerializer::new(report.domain.len() + U64_LEN)
.write(&report.domain)
.write(report.policy_identifier)
.finalize(),
});
self.store()
.write(batch.build_all())
.await
.caused_by(trc::location!())?;
let mut batch = BatchBuilder::new();
batch
.assert_value(key.clone(), AssertValue::Hash(revision))
.clear(key.clone())
.clear(RegistryClass::PrimaryKey {
object_id: object_id.into(),
index_id: Property::Domain.to_id(),
key: KeySerializer::new(report.domain.len() + U64_LEN)
.write(&report.domain)
.write(report.policy_identifier)
.finalize(),
});
match self.store().write(batch.build_all()).await {
Ok(_) => break report,
Err(err) if err.is_assertion_failure() && attempt < MAX_WRITE_RETRIES => {
attempt += 1;
write_retry_pause(attempt).await;
}
Err(err) => return Err(err.caused_by(trc::location!())),
}
};
let span_id = self.inner.data.span_id_gen.generate();
let event_from = report.report.date_range_begin.timestamp() as u64;
@@ -676,8 +693,11 @@ impl DmarcReporting for Server {
break;
}
Err(err) => {
if err.is_assertion_failure() && rety_count < 3 {
// inbuxa: another node appended first; try again
// after a short pause
if err.is_assertion_failure() && rety_count < MAX_WRITE_RETRIES {
rety_count += 1;
write_retry_pause(rety_count).await;
continue;
}
trc::error!(
+29 -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 registry::{
@@ -40,35 +42,49 @@ pub trait InternalReportIndex: ObjectImpl {
fn primary_key(&self) -> ValueClass;
/// Moves the report's delivery, and its queued task, to `at`.
///
/// inbuxa: the new queue row carries the task's type, as
/// `schedule_task_with_id` writes it, and the task row gets the new due
/// too. `queued` is the task as stored: its due, not the report's
/// `deliverAt`, is the queue row that exists (they differ once the task
/// has been retried).
fn reschedule_ops(
&mut self,
batch: &mut BatchBuilder,
item_id: u64,
revision: u64,
at: UTCDateTime,
queued: Option<&Task>,
) {
let current_deliver_at = self.deliver_at();
let current_due = current_deliver_at.timestamp() as u64;
let queued_due = queued.map_or(current_due, |task| task.due_timestamp());
let new_due = at.timestamp() as u64;
if current_deliver_at != at {
if current_deliver_at != at || queued_due != new_due {
let object = Self::OBJECT;
let object_id = object.to_id();
let key = ValueClass::Registry(RegistryClass::Item { object_id, item_id });
self.set_deliver_at(at);
batch
.assert_value(key.clone(), AssertValue::Hash(revision))
.clear(ValueClass::TaskQueue(TaskQueueClass::Due {
batch.assert_value(key.clone(), AssertValue::Hash(revision));
if queued_due != new_due {
batch.clear(ValueClass::TaskQueue(TaskQueueClass::Due {
id: item_id,
due: current_deliver_at.timestamp() as u64,
}))
.set(
ValueClass::TaskQueue(TaskQueueClass::Due {
id: item_id,
due: at.timestamp() as u64,
}),
object_id.serialize(),
)
due: queued_due,
}));
}
// A row an earlier reschedule left at the report's deliverAt
if current_due != new_due && current_due != queued_due {
batch.clear(ValueClass::TaskQueue(TaskQueueClass::Due {
id: item_id,
due: current_due,
}));
}
batch
.schedule_task_with_id(item_id, self.task(item_id))
.set(key, self.to_pickled_vec());
}
}
+3
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 common::config::smtp::report::AggregateFrequency;
@@ -15,6 +17,7 @@ pub mod inbound;
pub mod index;
pub mod scheduler;
pub mod send;
pub mod shared; // inbuxa: reports written by every node
pub mod spf;
pub mod tls;
+11 -8
View File
@@ -20,14 +20,17 @@ impl SpawnReport for mpsc::Receiver<ReportingEvent> {
tokio::spawn(async move {
while let Some(event) = self.recv().await {
let server = inner.build_server();
// inbuxa: reports are the outbound MTA's business, as at
// boot, but the role is read per event so a change applies
// without a restart. Events that arrive while the role is
// off are dropped, as they were on a node started without it
if !matches!(event, ReportingEvent::Stop) && !server.core.network.roles.outbound_mta
{
continue;
}
// inbuxa: every node records what it received, whatever its
// role. An aggregate report covers all of a domain's mail,
// whichever node took it, and recording is a store write
// that nodes already share: the report's primary key is
// versioned, so concurrent appends from several nodes retry
// rather than overwrite. Only building and sending the
// report (the DmarcReport and TlsReport tasks) belongs to
// the outbound MTA; the task manager keeps those to nodes
// with that role. Upstream ran this only on outbound MTA
// nodes, so mail received anywhere else never reached a
// report.
match event {
ReportingEvent::Dmarc(event) => server.schedule_dmarc(event).await,
ReportingEvent::Tls(event) => server.schedule_tls(event).await,
+45
View File
@@ -0,0 +1,45 @@
/*
* SPDX-FileCopyrightText: 2026 Coffey Labs
*
* SPDX-License-Identifier: AGPL-3.0-only
*/
//! inbuxa: internal DMARC and TLS reports are shared by every node. Any node
//! that receives mail appends to them, so several nodes can write one report
//! at once, and the node that sends it may do so while another is appending.
//! Appends already guard the report's versioned primary key and retry when
//! another writer got there first; these helpers give those retries room and
//! let the sender delete exactly the report it read.
use rand::RngExt;
use std::time::Duration;
use store::{Deserialize, xxhash_rust::xxh3::xxh3_64};
/// How many times a report write that lost to another writer is retried.
/// Upstream retried three times, when only outbound MTA nodes wrote.
pub(crate) const MAX_WRITE_RETRIES: u32 = 10;
/// A short random pause, longer on each attempt, before retrying a report
/// write that lost to another node, so the writers spread out instead of
/// colliding again.
pub(crate) async fn write_retry_pause(attempt: u32) {
let ms = rand::rng().random_range(5..=25u64) * u64::from(attempt.max(1));
tokio::time::sleep(Duration::from_millis(ms)).await;
}
/// A stored value with the hash of the bytes it was read from, for
/// `AssertValue::Hash`: a write asserting it fails if anyone changed the
/// value since.
pub(crate) struct Revisioned<T> {
pub revision: u64,
pub value: T,
}
impl<T: Deserialize> Deserialize for Revisioned<T> {
fn deserialize(bytes: &[u8]) -> trc::Result<Self> {
Ok(Revisioned {
revision: xxh3_64(bytes),
value: T::deserialize(bytes)?,
})
}
}
+40 -22
View File
@@ -2,9 +2,12 @@
* SPDX-FileCopyrightText: 2020 Stalwart Labs LLC <[email protected]>
*
* SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-SEL
*
* Modified by Coffey Labs in 2026 for INBUXA.
*/
use super::AggregateTimestamp;
use super::shared::{MAX_WRITE_RETRIES, Revisioned, write_retry_pause};
use crate::{
queue::RecipientDomain,
reporting::{index::InternalReportIndex, send::MtaReportSend},
@@ -70,28 +73,40 @@ impl TlsReporting for Server {
let object_id = ObjectType::TlsInternalReport.to_id();
let key = ValueClass::Registry(RegistryClass::Item { object_id, item_id });
let Some(report) = self
.store()
.get_value::<TlsInternalReport>(ValueKey::from(key.clone()))
.await
.caused_by(trc::location!())?
else {
return Ok(());
};
// Delete report. inbuxa: only the version read here, so a result
// another node appends meanwhile is sent with it rather than lost
let mut attempt = 0;
let report = loop {
let Some(Revisioned {
revision,
value: report,
}) = self
.store()
.get_value::<Revisioned<TlsInternalReport>>(ValueKey::from(key.clone()))
.await
.caused_by(trc::location!())?
else {
return Ok(());
};
// Delete report
let mut batch = BatchBuilder::new();
batch.clear(key).clear(RegistryClass::PrimaryKey {
object_id: object_id.into(),
index_id: Property::Domain.to_id(),
key: report.domain.as_bytes().to_vec(),
});
self.core
.storage
.data
.write(batch.build_all())
.await
.caused_by(trc::location!())?;
let mut batch = BatchBuilder::new();
batch
.assert_value(key.clone(), AssertValue::Hash(revision))
.clear(key.clone())
.clear(RegistryClass::PrimaryKey {
object_id: object_id.into(),
index_id: Property::Domain.to_id(),
key: report.domain.as_bytes().to_vec(),
});
match self.core.storage.data.write(batch.build_all()).await {
Ok(_) => break report,
Err(err) if err.is_assertion_failure() && attempt < MAX_WRITE_RETRIES => {
attempt += 1;
write_retry_pause(attempt).await;
}
Err(err) => return Err(err.caused_by(trc::location!())),
}
};
let domain_name = report.domain.as_str();
let event_from = report.report.date_range_start.timestamp() as u64;
@@ -477,8 +492,11 @@ impl TlsReporting for Server {
break;
}
Err(err) => {
if err.is_assertion_failure() && rety_count < 3 {
// inbuxa: another node appended first; try again
// after a short pause
if err.is_assertion_failure() && rety_count < MAX_WRITE_RETRIES {
rety_count += 1;
write_retry_pause(rety_count).await;
continue;
}
trc::error!(
+97 -1
View File
@@ -51,7 +51,9 @@ impl PostgresStore {
let all_fields = index.all_fields();
let fields = document.fields;
// inbuxa: keyword text (addresses, contact fields, ...) is split into
// words before it reaches the text parser, see keyword_terms().
// words before it reaches the text parser, see keyword_terms();
// language text gets the words inside its URLs, host names and
// file names added, see url_terms().
let keywords = primary_keys
.iter()
.chain(all_fields)
@@ -60,6 +62,9 @@ impl PostgresStore {
value,
language: Language::None,
}) if field.is_text() => Some(keyword_terms(value)),
Some(SearchValue::Text { value, .. }) if field.is_text() => {
url_terms(value)
}
_ => None,
})
.collect::<Vec<_>>();
@@ -290,14 +295,36 @@ impl PostgresStore {
continue;
}
} else {
// inbuxa: a query word written as a URL, host,
// file or hyphenated word also matches as its word
// parts, which url_terms() indexes
let parts = match value {
SearchValue::Text { value, .. } => query_url_terms(value),
_ => None,
};
let parts_pos = value_pos + 1;
let _ = write!(query, "@@ ({method}('{config}', ${value_pos})");
if parts.is_some() {
let _ = write!(query, " || {method}('{config}', ${parts_pos})");
}
for fallback in [PG_FALLBACK_LANG, PG_UNSTEMMED_LANG] {
if fallback != config && self.ts_configs.contains(fallback) {
let _ =
write!(query, " || {method}('{fallback}', ${value_pos})");
if parts.is_some() {
let _ = write!(
query,
" || {method}('{fallback}', ${parts_pos})"
);
}
}
}
query.push(')');
values.push(SqlParam::Ref(value));
if let Some(parts) = parts {
values.push(SqlParam::Owned(parts));
}
continue;
}
values.push(SqlParam::Ref(value));
} else if let SearchValue::KeyValues(kv) = value {
@@ -391,6 +418,75 @@ pub(crate) fn keyword_terms(value: &str) -> String {
terms
}
// inbuxa: in language text (subject, body, attachments) PostgreSQL's parser
// keeps a URL, a host name, a path or a file name as tokens of its own:
// "https://x.example/shipping-support/" gives a url, a host and a url_path,
// "invoice-2024.pdf" a file, so a body search for "shipping" or "invoice"
// missed messages where the word appears only there, while the built-in index
// splits them into words. The text is indexed as it was, followed by the word
// parts of each such token (SpaceTokenizer, as keyword_terms() splits), so
// they go through the same configuration and stemming as the words around
// them. On sample mail the text vector grows by about 15% for a newsletter
// full of tracking links and 30% for a short order notice with three links.
// Plain words, and words that only carry punctuation ("end.", "(see"),
// add nothing; hyphenated words are already split by the parser. Returns None
// when there is nothing to add, so most text is indexed exactly as before.
/// Characters that join the parts of a URL, host, path, address or file name.
const URL_SEPARATORS: [char; 13] = [
'/', '.', '@', ':', '?', '=', '&', '#', '_', '%', '+', '~', '\\',
];
pub(crate) fn url_terms(value: &str) -> Option<String> {
let mut terms = String::new();
// Each word is added once: a phrase search still finds the first URL it
// is in, and a newsletter's hundred tracking links don't add a hundred
// positions for "utm" and "campaign"
let mut seen = std::collections::HashSet::new();
for token in value.split(|c: char| {
c.is_whitespace() || matches!(c, '<' | '>' | '"' | '(' | ')' | '[' | ']' | '{' | '}')
}) {
let token = token.trim_matches(|c: char| !c.is_alphanumeric());
if token.contains(URL_SEPARATORS) {
for word in SpaceTokenizer::new(token, MAX_TOKEN_LENGTH) {
if !seen.insert(word.clone()) {
continue;
}
if terms.is_empty() {
terms.reserve(value.len() + 64);
terms.push_str(value);
terms.push('\n');
} else {
terms.push(' ');
}
terms.push_str(&word);
}
}
}
(!terms.is_empty()).then_some(terms)
}
/// The query side of url_terms(): each query word that is a URL, host, file
/// name or hyphenated word replaced by its word parts, or None when there is
/// none. It is searched in addition to the query as written, so documents
/// indexed before url_terms() still match as they did.
pub(crate) fn query_url_terms(value: &str) -> Option<String> {
let mut terms = String::with_capacity(value.len());
let mut changed = false;
for token in value.split_whitespace() {
let word = token.trim_matches(|c: char| !c.is_alphanumeric());
if !terms.is_empty() {
terms.push(' ');
}
if word.contains(URL_SEPARATORS) || word.contains('-') {
changed = true;
terms.push_str(&keyword_terms(word));
} else {
terms.push_str(token);
}
}
changed.then_some(terms)
}
pub(super) enum SqlParam<'x> {
Ref(&'x (dyn ToSql + Sync)),
Owned(String),
+21 -3
View File
@@ -245,9 +245,27 @@ impl Collector {
Update::RegisterReceiver { receiver } => {
self.receivers.push(receiver);
}
Update::RegisterSubscriber { subscriber } => {
ACTIVE_SUBSCRIBERS.lock().push(subscriber.id.clone());
self.subscribers.push(subscriber);
Update::RegisterSubscriber { mut subscriber } => {
// inbuxa: a subscriber registered under the id of a
// running one replaces it (a tracer whose settings
// changed). Every event collected so far went to the old
// one, every later event goes to the new one: the old
// one's batch is sent first (anything its full channel
// can't take moves over, rather than being dropped), and
// dropping it closes its channel, so its task writes
// what is queued and ends.
if let Some(old) = self.subscribers.iter_mut().find(|s| s.id == subscriber.id) {
let _ = old.send_batch();
if !old.batch.is_empty() {
let mut batch = std::mem::take(&mut old.batch);
batch.append(&mut subscriber.batch);
subscriber.batch = batch;
}
*old = subscriber;
} else {
ACTIVE_SUBSCRIBERS.lock().push(subscriber.id.clone());
self.subscribers.push(subscriber);
}
}
Update::UnregisterSubscriber { id } => {
ACTIVE_SUBSCRIBERS.lock().retain(|s| s != &id);
+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 std::sync::Arc;
@@ -105,6 +107,9 @@ impl SubscriberBuilder {
self
}
/// Registers the subscriber with the collector. inbuxa: one registered
/// under the id of a running subscriber replaces it, handing over at an
/// event boundary; the old one's channel then closes.
pub fn register(self) -> (mpsc::Sender<EventBatch>, mpsc::Receiver<EventBatch>) {
let (tx, rx) = mpsc::channel(8192);
+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.3"
"2026.9.25"
};
}
+317
View File
@@ -0,0 +1,317 @@
/*
* SPDX-FileCopyrightText: 2026 Coffey Labs
*
* SPDX-License-Identifier: AGPL-3.0-only
*/
//! DMARC results recorded on a node without outboundMta reach the aggregate
//! report, which a node with outboundMta builds and sends. Before, a front
//! node's results were dropped (or, before live roles, left in a channel
//! nobody read), so the report covered only the mail the outbound nodes
//! received. Also checks that nodes appending to one report at once lose
//! nothing. Needs a store the nodes can share (STORE=PostgreSql or MySql).
use crate::{smtp::inbound::TestMessage, utils::server::TestServerBuilder};
use common::{Server, config::smtp::report::AggregateFrequency, ipc::DmarcEvent};
use mail_auth::{
common::parse::TxtRecordParser,
dmarc::Dmarc,
report::{ActionDisposition, DmarcResult, Record, Report},
};
use registry::{
schema::{
enums::ClusterTaskType,
prelude::{ObjectType, Property},
structs::{
ClusterListenerGroup, ClusterRole, ClusterTaskGroup, ClusterTaskGroupProperties,
DmarcInternalReport, DmarcReportSettings, Expression, Task, TaskDmarcReport,
TaskStatus,
},
},
types::{EnumImpl, map::Map},
};
use smtp::reporting::{dmarc::DmarcReporting, send::MtaReportSend};
use std::{
collections::BTreeSet,
net::IpAddr,
sync::Arc,
time::{Duration, Instant},
};
use store::{
ValueKey,
registry::{RegistryFilter, RegistryFilterValue, RegistryQuery},
write::{BatchBuilder, RegistryClass, TaskQueueClass, ValueClass, now},
};
use types::id::Id;
const FRONT_ROLE: &str = "front_reports_front";
const MTA_ROLE: &str = "front_reports_mta";
const DOMAIN: &str = "front-reports.example";
#[tokio::test(flavor = "multi_thread")]
pub async fn front_node_report_tests() {
if matches!(
std::env::var("STORE").as_deref(),
Ok("RocksDb" | "Sqlite") | Err(_)
) {
println!("Skipping front node report tests: they need a store the nodes can share.");
return;
}
println!(
"Running front node report tests on {}...",
std::env::var("STORE").unwrap_or_default()
);
// A front role without outboundMta, an MTA role with it
let seed = TestServerBuilder::new("front_reports_seed").await;
seed.insert_object(role(FRONT_ROLE, &[ClusterTaskType::PushNotifications]))
.await;
seed.insert_object(role(MTA_ROLE, &[ClusterTaskType::OutboundMta]))
.await;
seed.insert_object(DmarcReportSettings {
aggregate_max_report_size: Expression {
else_: "1048576".into(),
..Default::default()
},
..Default::default()
})
.await;
let seed = seed.disable_services().build().await;
// The front node receives mail from two sources: the events its SMTP
// sessions hand the report scheduler
let front = TestServerBuilder::new_with_role(
"front_reports_front",
"front.front-reports.example".into(),
Some(FRONT_ROLE.into()),
false,
)
.await
.build_with_opts(false)
.await;
let front_server = front.server.clone();
assert!(!front_server.core.network.roles.outbound_mta);
for ip in ["192.0.2.1", "192.0.2.2"] {
front_server.schedule_report(event(ip)).await;
}
// Both are recorded in the shared report. Upstream, and main after live
// roles, left the front node's results out
let report_id = wait_for_report(&front_server, 2).await;
// Make the report due now. The front node leaves it alone: building and
// sending it is the outbound MTA's
move_task(&front_server, report_id, TaskStatus::now()).await;
tokio::time::sleep(Duration::from_secs(3)).await;
front_server.notify_task_queue();
tokio::time::sleep(Duration::from_secs(2)).await;
assert!(
task_exists(&front_server, report_id).await,
"the front node ran the report task"
);
// Several writers append to the report at once, from both nodes: none
// of their records is lost. The report waits in the future meanwhile, or
// the MTA node would send it as soon as it starts
move_task(
&front_server,
report_id,
TaskStatus::at(now() as i64 + 3600),
)
.await;
let mut mta = TestServerBuilder::new_with_role(
"front_reports_mta",
"mta.front-reports.example".into(),
Some(MTA_ROLE.into()),
false,
)
.await
.capture_queue()
.build_with_opts(false)
.await;
let mta_server = mta.server.clone();
assert!(mta_server.core.network.roles.outbound_mta);
let concurrent: Vec<String> = (10..18).map(|n| format!("192.0.2.{n}")).collect();
let mut handles = Vec::new();
for (n, ip) in concurrent.iter().enumerate() {
let server = if n % 2 == 0 {
front_server.clone()
} else {
mta_server.clone()
};
let ip = ip.clone();
handles.push(tokio::spawn(async move {
server.schedule_dmarc(Box::new(event(&ip))).await;
}));
}
for handle in handles {
handle.await.unwrap();
}
// Due again, the MTA node sends the report with every record in it
move_task(&mta_server, report_id, TaskStatus::now()).await;
let message = mta.expect_message().await;
let report =
Report::parse_rfc5322(message.read_message(&mta).await.as_bytes(), usize::MAX).unwrap();
assert_eq!(report.domain(), DOMAIN);
let sent: BTreeSet<IpAddr> = report
.records()
.iter()
.map(|r| r.source_ip().unwrap())
.collect();
let expected: BTreeSet<IpAddr> = ["192.0.2.1", "192.0.2.2"]
.into_iter()
.map(String::from)
.chain(concurrent)
.map(|ip| ip.parse().unwrap())
.collect();
assert_eq!(sent, expected);
wait_for(Duration::from_secs(20), "report task to finish", || async {
!task_exists(&mta_server, report_id).await
})
.await;
assert!(reports(&mta_server).await.is_empty());
if seed.is_reset() {
seed.temp_dir.delete();
front.temp_dir.delete();
mta.temp_dir.delete();
}
}
fn role(name: &str, tasks: &[ClusterTaskType]) -> ClusterRole {
ClusterRole {
name: name.into(),
description: None,
listeners: ClusterListenerGroup::EnableAll,
tasks: ClusterTaskGroup::EnableSome(ClusterTaskGroupProperties {
task_types: Map::new(tasks.to_vec()),
}),
}
}
fn event(ip: &str) -> DmarcEvent {
DmarcEvent {
domain: DOMAIN.to_string(),
report_record: Record::new()
.with_source_ip(ip.parse().unwrap())
.with_action_disposition(ActionDisposition::Pass)
.with_dmarc_dkim_result(DmarcResult::Pass)
.with_dmarc_spf_result(DmarcResult::Pass)
.with_envelope_from("sender.example")
.with_header_from("sender.example"),
dmarc_record: Arc::new(
Dmarc::parse(format!("v=DMARC1; p=reject; rua=mailto:reports@{DOMAIN}").as_bytes())
.unwrap(),
),
interval: AggregateFrequency::Daily,
span_id: 0,
}
}
async fn reports(server: &Server) -> Vec<(u64, DmarcInternalReport)> {
let ids = server
.registry()
.query::<Vec<Id>>(RegistryQuery::new(ObjectType::DmarcInternalReport).filter(
RegistryFilter::greater_than(
Property::Domain,
RegistryFilterValue::Bytes(vec![]),
true,
),
))
.await
.unwrap();
let mut reports = Vec::new();
for id in ids {
if let Some(report) = server
.store()
.get_value::<DmarcInternalReport>(ValueKey::from(ValueClass::Registry(
RegistryClass::Item {
object_id: ObjectType::DmarcInternalReport.to_id(),
item_id: id.id(),
},
)))
.await
.unwrap()
{
reports.push((id.id(), report));
}
}
reports
}
/// Waits for the report for `DOMAIN` to hold `records` records; returns its id.
async fn wait_for_report(server: &Server, records: usize) -> u64 {
let started = Instant::now();
loop {
let found = reports(server)
.await
.into_iter()
.find(|(_, report)| report.domain == DOMAIN);
if let Some((id, report)) = &found
&& report.report.records.len() == records
{
return *id;
}
assert!(
started.elapsed() < Duration::from_secs(10),
"no report with {records} records for {DOMAIN}: {found:?}"
);
tokio::time::sleep(Duration::from_millis(200)).await;
}
}
/// Reschedules the report's task.
async fn move_task(server: &Server, id: u64, status: TaskStatus) {
let task = server
.store()
.get_value::<Task>(ValueKey::from(ValueClass::TaskQueue(
TaskQueueClass::Task { id },
)))
.await
.unwrap()
.expect("report task missing");
let mut batch = BatchBuilder::new();
batch
.clear(ValueClass::TaskQueue(TaskQueueClass::Due {
id,
due: task.due_timestamp(),
}))
.schedule_task_with_id(
id,
Task::DmarcReport(TaskDmarcReport {
report_id: id.into(),
status,
}),
);
server.store().write(batch.build_all()).await.unwrap();
server.notify_task_queue();
}
async fn task_exists(server: &Server, id: u64) -> bool {
server
.store()
.get_value::<Task>(ValueKey::from(ValueClass::TaskQueue(
TaskQueueClass::Task { id },
)))
.await
.unwrap()
.is_some()
}
async fn wait_for<F, Fut>(within: Duration, what: &str, mut check: F)
where
F: FnMut() -> Fut,
Fut: Future<Output = bool>,
{
let started = Instant::now();
while !check().await {
assert!(
started.elapsed() < within,
"still waiting for the {what} after {:?}",
started.elapsed()
);
tokio::time::sleep(Duration::from_millis(250)).await;
}
}
+1
View File
@@ -7,6 +7,7 @@
*/
pub mod broadcast;
pub mod front_reports; // inbuxa: every node records DMARC and TLS results
pub mod live_roles; // inbuxa: role edits apply without a restart
#[cfg(feature = "nats")]
pub mod coordinator; // inbuxa: coordinator reconnects
+3
View File
@@ -2,9 +2,12 @@
* SPDX-FileCopyrightText: 2020 Stalwart Labs LLC <[email protected]>
*
* SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-SEL
*
* Modified by Coffey Labs in 2026 for INBUXA.
*/
pub mod analyze;
pub mod dmarc;
pub mod reschedule; // inbuxa: report reschedules and unreadable queue rows
pub mod scheduler;
pub mod tls;
+370
View File
@@ -0,0 +1,370 @@
/*
* SPDX-FileCopyrightText: 2026 Coffey Labs
*
* SPDX-License-Identifier: AGPL-3.0-only
*/
//! Rescheduling an internal DMARC or TLS report over JMAP moves its task: the
//! task runs at the new time, x:Task/get shows the new due, and tasks due
//! after it still run. A task queue row whose type can't be read is logged
//! and repaired rather than stopping every task due after it, including the
//! rows an earlier reschedule wrote with the report's object type.
use crate::utils::server::{TestServer, TestServerBuilder};
use common::{
Server,
config::smtp::report::AggregateFrequency,
ipc::{DmarcEvent, PolicyType, TlsEvent},
};
use mail_auth::{
common::parse::TxtRecordParser,
dmarc::Dmarc,
mta_sts::TlsRpt,
report::{ActionDisposition, DmarcResult, Record},
};
use registry::{
schema::{
enums::{TaskStoreMaintenanceType, TaskType},
prelude::{ObjectType, Property},
structs::{
DmarcInternalReport, DmarcReportSettings, Expression, Task, TaskStatus,
TaskStoreMaintenance, TlsInternalReport, TlsReportSettings,
},
},
types::{EnumImpl, ObjectImpl, datetime::UTCDateTime},
};
use serde_json::json;
use smtp::reporting::{index::InternalReportIndex, send::MtaReportSend};
use std::{
sync::Arc,
time::{Duration, Instant},
};
use store::{
SerializeInfallible, ValueKey,
write::{BatchBuilder, RegistryClass, TaskQueueClass, ValueClass, now},
};
use types::id::Id;
use utils::snowflake::SnowflakeIdGenerator;
#[tokio::test(flavor = "multi_thread")]
#[serial_test::serial]
async fn report_reschedule() {
let mut test = TestServerBuilder::new("smtp_report_reschedule")
.await
.with_http_listener(19057)
.await
.capture_queue()
.build()
.await;
let admin = test.account("admin");
admin
.registry_create_object(TlsReportSettings {
max_report_size: Expression {
else_: "1024".into(),
..Default::default()
},
..Default::default()
})
.await;
admin
.registry_create_object(DmarcReportSettings {
aggregate_max_report_size: Expression {
else_: "1024".into(),
..Default::default()
},
..Default::default()
})
.await;
admin.reload_settings().await;
test.reload_core();
test.expect_reload_settings().await;
let admin = test.account("admin");
// A daily DMARC and TLS report, due a day from now
schedule_dmarc(&test, "foobar.org").await;
schedule_tls(&test, "foobar.org").await;
let dmarc_id = wait_for_report::<DmarcInternalReport>(&test, "foobar.org").await;
let tls_id = wait_for_report::<TlsInternalReport>(&test, "foobar.org").await;
// Reschedule both to a few seconds from now, with a task due after them
let at = now() + 3;
let later = marker_task(&test.server, at + 3).await;
for (object, id, task_type) in [
(
ObjectType::DmarcInternalReport,
dmarc_id,
TaskType::DmarcReport,
),
(ObjectType::TlsInternalReport, tls_id, TaskType::TlsReport),
] {
admin
.registry_update_object(
object,
id,
json!({
Property::DeliverAt: UTCDateTime::from_timestamp(at as i64),
}),
)
.await;
// x:Task/get shows the new due, and the queue row carries the task's
// type. Upstream wrote the report's object type there and left the
// task at its old due
let task = admin.registry_get::<Task>(id).await;
assert_eq!(task.object_type(), task_type);
assert_eq!(
task.due_timestamp(),
at,
"{object:?} task due not moved: {task:?}"
);
assert_eq!(
queue_row(&test.server, id.id(), at).await,
Some(task_type.to_id().serialize()),
"{object:?} queue row"
);
}
// Both reports go out at the new time, and the later task still runs
wait_until_run(&test.server, &[dmarc_id.id(), tls_id.id(), later]).await;
assert!(now() >= at, "the reports went out before their new time");
assert!(
admin
.registry_get_all::<DmarcInternalReport>()
.await
.is_empty()
);
assert!(
admin
.registry_get_all::<TlsInternalReport>()
.await
.is_empty()
);
// Rows an earlier reschedule may have left in a store: one with the
// report's object type and the task left at its old due, and one that
// is unreadable and has no task behind it. Neither may hold back a task
// due after them.
schedule_dmarc(&test, "foobar.net").await;
let dmarc_id = wait_for_report::<DmarcInternalReport>(&test, "foobar.net").await;
let at = now() + 2;
let old_due = old_style_reschedule(&test.server, dmarc_id.id(), at).await;
let orphan = SnowflakeIdGenerator::global_id().unwrap();
let mut batch = BatchBuilder::new();
batch.set(
ValueClass::TaskQueue(TaskQueueClass::Due {
id: orphan,
due: at,
}),
vec![0xff, 0xff],
);
test.server.store().write(batch.build_all()).await.unwrap();
let later = marker_task(&test.server, at + 2).await;
wait_until_run(&test.server, &[dmarc_id.id(), later]).await;
assert!(
admin
.registry_get_all::<DmarcInternalReport>()
.await
.is_empty()
);
assert_eq!(queue_row(&test.server, orphan, at).await, None);
assert_eq!(queue_row(&test.server, dmarc_id.id(), at).await, None);
assert_eq!(queue_row(&test.server, dmarc_id.id(), old_due).await, None);
// x:Task/query by type skips an unreadable row rather than failing
let mut batch = BatchBuilder::new();
let due = now() + 3600;
batch.set(
ValueClass::TaskQueue(TaskQueueClass::Due { id: orphan, due }),
vec![0xff, 0xff],
);
test.server.store().write(batch.build_all()).await.unwrap();
admin
.registry_query_ids(
ObjectType::Task,
vec![(Property::Type, TaskType::DmarcReport.as_str())],
Vec::<&str>::new(),
)
.await;
let mut batch = BatchBuilder::new();
batch.clear(ValueClass::TaskQueue(TaskQueueClass::Due {
id: orphan,
due,
}));
test.server.store().write(batch.build_all()).await.unwrap();
if test.is_reset() {
test.temp_dir.delete();
}
}
async fn schedule_dmarc(test: &TestServer, domain: &str) {
test.server
.schedule_report(DmarcEvent {
domain: domain.to_string(),
report_record: Record::new()
.with_source_ip("192.168.1.2".parse().unwrap())
.with_action_disposition(ActionDisposition::Pass)
.with_dmarc_dkim_result(DmarcResult::Pass)
.with_dmarc_spf_result(DmarcResult::Fail)
.with_envelope_from("[email protected]")
.with_envelope_to("[email protected]")
.with_header_from("[email protected]"),
dmarc_record: Arc::new(
Dmarc::parse(format!("v=DMARC1; p=reject; rua=mailto:reports@{domain}").as_bytes())
.unwrap(),
),
interval: AggregateFrequency::Daily,
span_id: 0,
})
.await;
}
async fn schedule_tls(test: &TestServer, domain: &str) {
test.server
.schedule_report(TlsEvent {
domain: domain.to_string(),
policy: PolicyType::None,
failure: None,
tls_record: Arc::new(
TlsRpt::parse(format!("v=TLSRPTv1;rua=mailto:reports@{domain}").as_bytes())
.unwrap(),
),
interval: AggregateFrequency::Daily,
span_id: 0,
})
.await;
}
trait ReportDomain: ObjectImpl {
fn report_domain(&self) -> &str;
}
impl ReportDomain for DmarcInternalReport {
fn report_domain(&self) -> &str {
&self.domain
}
}
impl ReportDomain for TlsInternalReport {
fn report_domain(&self) -> &str {
&self.domain
}
}
async fn wait_for_report<T: ReportDomain>(test: &TestServer, domain: &str) -> Id {
let admin = test.account("admin");
for _ in 0..100 {
if let Some((id, _)) = admin
.registry_get_all::<T>()
.await
.into_iter()
.find(|(_, report)| report.report_domain() == domain)
{
return id;
}
tokio::time::sleep(Duration::from_millis(100)).await;
}
panic!("No {} for {domain}", T::OBJECT.as_str());
}
/// A task that succeeds when it runs, due at `due`.
async fn marker_task(server: &Server, due: u64) -> u64 {
let id = SnowflakeIdGenerator::global_id().unwrap();
let mut batch = BatchBuilder::new();
batch.schedule_task_with_id(
id,
Task::StoreMaintenance(TaskStoreMaintenance {
maintenance_type: TaskStoreMaintenanceType::RemoveLockDav,
shard_index: Some(0),
status: TaskStatus::at(due as i64),
}),
);
server.store().write(batch.build_all()).await.unwrap();
server.notify_task_queue();
id
}
/// What the reschedule before this fix wrote: the report's object type in
/// the new queue row, and the task row left at its old due. Returns that
/// old due.
async fn old_style_reschedule(server: &Server, item_id: u64, at: u64) -> u64 {
let object_id = ObjectType::DmarcInternalReport.to_id();
let key = ValueClass::Registry(RegistryClass::Item { object_id, item_id });
let mut report = server
.store()
.get_value::<DmarcInternalReport>(ValueKey::from(key.clone()))
.await
.unwrap()
.unwrap();
let old_due = report.deliver_at().timestamp() as u64;
report.set_deliver_at(UTCDateTime::from_timestamp(at as i64));
let mut batch = BatchBuilder::new();
batch
.clear(ValueClass::TaskQueue(TaskQueueClass::Due {
id: item_id,
due: old_due,
}))
.set(
ValueClass::TaskQueue(TaskQueueClass::Due {
id: item_id,
due: at,
}),
object_id.serialize(),
)
.set(key, report.to_pickled_vec());
server.store().write(batch.build_all()).await.unwrap();
server.notify_task_queue();
old_due
}
struct RawValue(Vec<u8>);
impl store::Deserialize for RawValue {
fn deserialize(bytes: &[u8]) -> trc::Result<Self> {
Ok(RawValue(bytes.to_vec()))
}
}
async fn queue_row(server: &Server, id: u64, due: u64) -> Option<Vec<u8>> {
server
.store()
.get_value::<RawValue>(ValueKey::from(ValueClass::TaskQueue(TaskQueueClass::Due {
id,
due,
})))
.await
.unwrap()
.map(|raw| raw.0)
}
async fn task_exists(server: &Server, id: u64) -> bool {
server
.store()
.get_value::<Task>(ValueKey::from(ValueClass::TaskQueue(
TaskQueueClass::Task { id },
)))
.await
.unwrap()
.is_some()
}
async fn wait_until_run(server: &Server, ids: &[u64]) {
let started = Instant::now();
loop {
let mut pending = Vec::new();
for id in ids {
if task_exists(server, *id).await {
pending.push(*id);
}
}
if pending.is_empty() {
return;
}
if started.elapsed() > Duration::from_secs(30) {
panic!("tasks {pending:?} never ran");
}
tokio::time::sleep(Duration::from_millis(200)).await;
}
}
+81
View File
@@ -133,6 +133,11 @@ pub async fn test(test: &TestServer) {
println!("Running address search tests...");
test_address_search(store.clone()).await;
// inbuxa: words inside URLs, host names and file names in body text
// are found on every backend
println!("Running URL word search tests...");
test_url_word_search(store.clone()).await;
// Large document insert test
println!("Running large document insert tests...");
let mut large_text = String::with_capacity(20 * 1024 * 1024);
@@ -1129,3 +1134,79 @@ async fn test_address_search(store: SearchStore) {
.await
.unwrap();
}
async fn test_url_word_search(store: SearchStore) {
const ACCOUNT_ID: u32 = 8;
let bodies = [
"Track your parcel here: https://x.example/shipping-support/ and reply.",
"Reset it at https://mail.example.com/login/?password=reset&user=jane now.",
"Attached is invoice-2024.pdf for your records.",
"Shipping was fast, thanks again.",
"Nothing to see at www.example.org/about-us, really.",
];
let mut documents = Vec::new();
let mut mask = RoaringBitmap::new();
for (document_id, body) in bodies.iter().enumerate() {
let mut document = IndexDocument::new(SearchIndex::Email)
.with_account_id(ACCOUNT_ID)
.with_document_id(document_id as u32);
document.index_text(EmailSearchField::Body, body, Language::English);
document.index_unsigned(EmailSearchField::ReceivedAt, document_id as u64);
documents.push(document);
mask.insert(document_id as u32);
}
store.index(documents).await.unwrap();
if let SearchStore::ElasticSearch(store) = &store {
store.refresh_index(SearchIndex::Email).await.unwrap();
}
for (text, expected) in [
// only inside a URL path, a query string or a file name
("shipping", vec![0u32, 3]),
("support", vec![0]),
("password", vec![1]),
("login", vec![1]),
("jane", vec![1]),
("invoice", vec![2]),
("pdf", vec![2]),
("2024", vec![2]),
// host names
("example", vec![0, 1, 4]),
("mail", vec![1]),
// written as they appear
("https://x.example/shipping-support/", vec![0]),
("shipping-support", vec![0]),
("invoice-2024.pdf", vec![2]),
("mail.example.com", vec![1]),
// plain words are unaffected
("parcel", vec![0]),
("records", vec![2]),
("thanks", vec![3]),
// no match
("billing", vec![]),
("example.net", vec![]),
] {
let ids = store
.query_account(
SearchQuery::new(SearchIndex::Email)
.with_filters(vec![
SearchFilter::eq(SearchField::AccountId, ACCOUNT_ID),
SearchFilter::has_english_text(EmailSearchField::Body, text),
])
.with_comparator(SearchComparator::ascending(EmailSearchField::ReceivedAt))
.with_mask(mask.clone()),
)
.await
.unwrap();
assert_eq!(ids, expected, "Body {text:?}");
}
store
.unindex(
SearchQuery::new(SearchIndex::Email)
.with_filter(SearchFilter::eq(SearchField::AccountId, ACCOUNT_ID)),
)
.await
.unwrap();
}
+53 -1
View File
@@ -116,12 +116,64 @@ async fn test_write_applies(test: &TestServer) {
for name in &names {
assert!(has_schedule(test, name), "{name} missing");
}
// A burst of separate requests shares a reload or two: each arrives
// tens of milliseconds after the last, so none overlaps a running
// reload, and the reload waits for writes to settle instead
let reloads = test.server.inner.data.settings_reload.reloads();
let started = std::time::Instant::now();
let burst = (0..10)
.map(|i| format!("autoreload-burst-{i}"))
.collect::<Vec<_>>();
let mut writes = Vec::new();
for name in &burst {
writes.push(admin.registry_create([MtaDeliverySchedule {
name: name.clone(),
queue_id,
..Default::default()
}]));
}
for response in futures::future::join_all(writes).await {
assert_applied(&response);
schedule_ids.push(response.created_id(0));
}
let burst_reloads = test.server.inner.data.settings_reload.reloads() - reloads;
println!(
"10 concurrent writes: {burst_reloads} reload(s), {} ms",
started.elapsed().as_millis()
);
assert!(
(1..=2).contains(&burst_reloads),
"{burst_reloads} reloads for 10 concurrent writes"
);
for name in &burst {
assert!(has_schedule(test, name), "{name} missing");
}
// A single write still reloads promptly
let reloads = test.server.inner.data.settings_reload.reloads();
let started = std::time::Instant::now();
let response = admin
.registry_create([MtaDeliverySchedule {
name: "autoreload-single".into(),
queue_id,
..Default::default()
}])
.await;
assert_applied(&response);
schedule_ids.push(response.created_id(0));
println!("1 write: {} ms", started.elapsed().as_millis());
assert_eq!(
test.server.inner.data.settings_reload.reloads() - reloads,
1
);
assert!(has_schedule(test, "autoreload-single"));
// Several objects in one request: one reload
let response = admin
.registry_destroy(ObjectType::MtaDeliverySchedule, schedule_ids.iter())
.await;
assert_applied(&response);
for name in &names {
for name in names.iter().chain(&burst) {
assert!(!has_schedule(test, name), "{name} still present");
}
+1
View File
@@ -24,6 +24,7 @@ pub mod quota;
pub mod reload; // inbuxa: reloads and build errors
pub mod security;
pub mod task;
pub mod tracer_reload; // inbuxa: tracers start over when their settings change
pub mod tenant;
pub mod undelete;
+204
View File
@@ -0,0 +1,204 @@
/*
* SPDX-FileCopyrightText: 2026 Coffey Labs
*
* SPDX-License-Identifier: AGPL-3.0-only
*/
// inbuxa: a tracer whose own settings change is started over by the reload
// that follows the write: a Log tracer moved to another directory writes
// there from then on, and no event is lost or written twice on the way.
use crate::utils::{
jmap::JmapResponse,
server::{TestServer, TestServerBuilder},
};
use registry::{
schema::{
enums::{EventPolicy, LogRotateFrequency, TracingLevel},
prelude::ObjectType,
structs::{Expression, MtaStageAuth, Tracer, TracerLog},
},
types::map::Map,
};
use serde_json::json;
use std::{
path::{Path, PathBuf},
time::{Duration, Instant},
};
use trc::{EventType, ServerEvent};
const PREFIX: &str = "tracer-reload";
#[tokio::test(flavor = "multi_thread")]
pub async fn tracer_reload_tests() {
let mut test = TestServerBuilder::new("tracer_reload_tests")
.await
.with_default_listeners()
.await
.with_object(MtaStageAuth {
require: Expression {
else_: "false".to_string(),
..Default::default()
},
..Default::default()
})
.await
.build()
.await;
let admin = test
.create_user_account(
"admin",
"[email protected]",
"these_pretzels_are_making_me_thirsty",
&[],
"Admin",
)
.await;
test.account("admin")
.assign_roles_to_account(admin.id(), &["user", "system"])
.await;
test.insert_account(admin);
test_log_tracer_moves(&test).await;
if test.is_reset() {
test.temp_dir.delete();
}
}
async fn test_log_tracer_moves(test: &TestServer) {
println!("Running Log tracer path change...");
let admin = test.account("[email protected]");
let old_dir = test.temp_dir.path.join("tracer-old");
let new_dir = test.temp_dir.path.join("tracer-new");
for dir in [&old_dir, &new_dir] {
let _ = std::fs::remove_dir_all(dir);
std::fs::create_dir_all(dir).unwrap();
}
let old_file = old_dir.join(PREFIX);
let new_file = new_dir.join(PREFIX);
// A Log tracer for one event type, written to the old directory
let response = admin
.registry_create([Tracer::Log(TracerLog {
path: old_dir.to_string_lossy().into_owned(),
prefix: PREFIX.into(),
rotate: LogRotateFrequency::Never,
ansi: false,
multiline: false,
enable: true,
level: TracingLevel::Trace,
lossy: false,
events: Map::new(vec![EventType::Server(ServerEvent::Licensing)]),
events_policy: EventPolicy::Include,
})])
.await;
assert_applied(&response);
let tracer_id = response.created_id(0);
emit("marker-before");
wait_for(&old_file, "marker-before").await;
// Events keep coming while the path changes
let stream = tokio::spawn(async {
for i in 0..2000u32 {
emit(&format!("seq-{i:05}-end"));
if i % 50 == 0 {
tokio::time::sleep(Duration::from_millis(1)).await;
}
}
});
tokio::time::sleep(Duration::from_millis(5)).await;
let response = admin
.registry_update(
ObjectType::Tracer,
[(tracer_id, json!({"path": new_dir.to_string_lossy()}))],
)
.await;
assert_applied(&response);
stream.await.unwrap();
// Once the reload has run, events go to the new file only
tokio::time::sleep(Duration::from_millis(200)).await;
emit("marker-after");
wait_for(&new_file, "marker-after").await;
wait_for(&new_file, "seq-01999-end").await;
tokio::time::sleep(Duration::from_millis(200)).await;
let old = read(&old_file);
let new = read(&new_file);
assert!(!old.contains("marker-after"), "old file still written to");
assert!(!new.contains("marker-before"));
// Every event written once, in one file or the other
let old_seq = count_seq(&old);
let new_seq = count_seq(&new);
println!(
"{} events in the old file, {} in the new one",
old_seq.iter().filter(|c| **c > 0).count(),
new_seq.iter().filter(|c| **c > 0).count()
);
for i in 0..2000 {
assert_eq!(
old_seq[i] + new_seq[i],
1,
"seq-{i:05} written {} + {} times",
old_seq[i],
new_seq[i]
);
}
assert!(
new_seq.iter().any(|c| *c > 0),
"no event of the stream reached the new file"
);
// Removing the tracer stops it
let response = admin
.registry_destroy(ObjectType::Tracer, [tracer_id])
.await;
assert_applied(&response);
tokio::time::sleep(Duration::from_millis(200)).await;
emit("marker-removed");
tokio::time::sleep(Duration::from_millis(300)).await;
assert!(!read(&new_file).contains("marker-removed"));
assert!(!read(&old_file).contains("marker-removed"));
}
fn emit(marker: &str) {
trc::event!(Server(ServerEvent::Licensing), Details = marker.to_string());
}
fn read(path: &Path) -> String {
std::fs::read_to_string(path).unwrap_or_default()
}
fn count_seq(text: &str) -> Vec<u32> {
let mut counts = vec![0u32; 2000];
for part in text.split("seq-").skip(1) {
if let Some(n) = part.get(..5).and_then(|n| n.parse::<usize>().ok())
&& part[5..].starts_with("-end")
{
counts[n] += 1;
}
}
counts
}
async fn wait_for(path: &PathBuf, marker: &str) {
let started = Instant::now();
while !read(path).contains(marker) {
assert!(
started.elapsed() < Duration::from_secs(10),
"{marker} not in {}",
path.display()
);
tokio::time::sleep(Duration::from_millis(20)).await;
}
}
fn assert_applied(response: &JmapResponse) {
assert_eq!(
response.pointer("/methodResponses/0/1/x:settingsReload"),
Some(&json!({"applied": true})),
"{response:?}"
);
}