Eight conflicted files resolved, plus the lock file and the schema:
- crates/services/src/task_manager/spam_classifier.rs: upstream's rules
update now replaces existing rules, DNSBL servers, lookups and file
extensions, keeping only whether each is on. Taken, with one difference:
an object an admin edited is kept as it is. Every object an update writes
is fingerprinted (content without `enable`, SHA-256, stored under
SUBSPACE_INBUXA "Sf"), and only one that still matches is replaced.
Scores are never replaced, as upstream has it. The AU-1.10 summary record
now names what was added, replaced and kept, and the bundled rules are
marked applied only when the update fully succeeded, so a failure runs
again on the next start. The marker becomes "3.0.2+2", which runs the
update once on upgrade to fingerprint every rule still as bundled.
- crates/common/src/network/autoconfig/autodiscover.rs: upstream's rewrite
(implicit TLS first, labeled SSL), with the per-protocol switches (LP-7,
LP-14a) passed in as a filter.
- crates/store/src/backend/mysql/{search,write}.rs: upstream's chunked
deletes (no unbounded first DELETE, stop on a short chunk, halve the
chunk on the new chunk-too-large errors) inside the fork's query timeout.
- crates/smtp/src/lib.rs: the fork's queue spawn kept. It already fixed the
stall upstream fixes here (a node without outboundMta stops accepting
mail at about 1024 queued messages), and follows role changes live.
- crates/jmap/src/registry/mapping/bootstrap.rs: the log path stays
/var/log/inbuxa/; upstream's PowerDNS mapping taken.
- crates/main/Cargo.toml: the AGPL-only license kept, version 0.16.24.
- tests/src/jmap/principal/get.rs: the fork's capabilities kept.
- resources/schema/schema.json.gz: merged as JSON; upstream relabeled the
vendor Sieve extensions "(Stalwart)", kept as "(vnd.inbuxa)".
- Cargo.lock: upstream's, with the fork's crates added by Cargo.
Also:
- tests/src/smtp/inbound/spam_rules_kept.rs: an edited rule survives an
update, an unedited one is updated, rules from before fingerprints are
handled, and the audit summary says so. Upstream's own spam_rules test
passes unchanged.
- tests/src/smtp/reporting/reschedule.rs moves to port 19058; upstream's
new spam_rules test took 19057.
- tools/fork/renames.py renames the "(Stalwart)" labels and the default
log path, so neither conflicts again.
- tools/fork/notice-check.py compares against the newest snapshot in the
checked-out history instead of the upstream branch head, so moving the
branch no longer fails other open pull requests.
- tests/src/directory/issuer.rs (since v0.16.23) stays out, and is on the
build check's known list: it tests issuer-based directory routing, which
the fork doesn't have (DIR-2).
- Strip report: docs/fork/strip-reports/v0.16.24.{md,json}.
210 lines
8.3 KiB
Rust
210 lines
8.3 KiB
Rust
/*
|
|
* 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::{Event, PURGE_EVERY, SEND_TIMEOUT, push::spawn_push_manager};
|
|
use crate::state_manager::IpcSubscriber;
|
|
use common::{
|
|
Inner,
|
|
ipc::{BroadcastEvent, PushEvent},
|
|
};
|
|
use std::{sync::Arc, time::Instant};
|
|
use store::ahash::AHashMap;
|
|
use tokio::sync::mpsc::{self, error::TrySendError};
|
|
use trc::ServerEvent;
|
|
|
|
#[derive(Default)]
|
|
struct Subscriber {
|
|
ipc: Vec<IpcSubscriber>,
|
|
is_push: bool,
|
|
}
|
|
|
|
#[allow(clippy::unwrap_or_default)]
|
|
pub fn spawn_push_router(inner: Arc<Inner>, mut change_rx: mpsc::Receiver<PushEvent>) {
|
|
let push_tx = spawn_push_manager(inner.clone());
|
|
|
|
tokio::spawn(async move {
|
|
let mut subscribers: AHashMap<u32, Subscriber> = AHashMap::default();
|
|
let mut last_purge = Instant::now();
|
|
|
|
while let Some(event) = change_rx.recv().await {
|
|
let mut purge_needed = last_purge.elapsed() >= PURGE_EVERY;
|
|
|
|
match event {
|
|
PushEvent::Stop => {
|
|
if push_tx.send(Event::Reset).await.is_err() {
|
|
trc::event!(
|
|
Server(ServerEvent::ThreadError),
|
|
Details = "Error sending push reset.",
|
|
CausedBy = trc::location!()
|
|
);
|
|
}
|
|
break;
|
|
}
|
|
|
|
PushEvent::Subscribe {
|
|
account_ids,
|
|
types,
|
|
tx,
|
|
} => {
|
|
let owner = account_ids.first().copied().unwrap_or(u32::MAX);
|
|
for account_id in account_ids {
|
|
subscribers
|
|
.entry(account_id)
|
|
.or_default()
|
|
.ipc
|
|
.push(IpcSubscriber {
|
|
types,
|
|
tx: tx.clone(),
|
|
owner,
|
|
});
|
|
}
|
|
}
|
|
|
|
// inbuxa: SCIM-52: dropping every sender closes the session's
|
|
// channel, which ends it
|
|
PushEvent::Revoke { account_id } => {
|
|
for subscriber_list in subscribers.values_mut() {
|
|
subscriber_list
|
|
.ipc
|
|
.retain(|subscriber| subscriber.owner != account_id);
|
|
}
|
|
purge_needed = true;
|
|
}
|
|
|
|
PushEvent::PushServerRegister { activate, expired } => {
|
|
for account_id in activate {
|
|
subscribers.entry(account_id).or_default().is_push = true;
|
|
}
|
|
|
|
for account_id in expired {
|
|
let mut remove_account = false;
|
|
if let Some(subscriber_list) = subscribers.get_mut(&account_id) {
|
|
subscriber_list.is_push = false;
|
|
remove_account = subscriber_list.ipc.is_empty();
|
|
}
|
|
if remove_account {
|
|
subscribers.remove(&account_id);
|
|
}
|
|
}
|
|
}
|
|
|
|
PushEvent::Publish {
|
|
notification,
|
|
broadcast,
|
|
} => {
|
|
// Publish event to cluster
|
|
if broadcast
|
|
&& let Some(broadcast_tx) = inner.ipc.broadcast_tx.as_ref()
|
|
&& broadcast_tx
|
|
.send(BroadcastEvent::PushNotification(notification.clone()))
|
|
.await
|
|
.is_err()
|
|
{
|
|
trc::event!(
|
|
Server(trc::ServerEvent::ThreadError),
|
|
Details = "Error sending broadcast event.",
|
|
CausedBy = trc::location!()
|
|
);
|
|
}
|
|
|
|
let account_id = notification.account_id();
|
|
if let Some(subscribers) = subscribers.get(&account_id) {
|
|
for subscriber in &subscribers.ipc {
|
|
if let Some(notification) = notification.filter_types(&subscriber.types)
|
|
{
|
|
match subscriber.tx.try_send(notification) {
|
|
Ok(()) => {}
|
|
Err(TrySendError::Full(notification)) => {
|
|
let subscriber_tx = subscriber.tx.clone();
|
|
|
|
tokio::spawn(async move {
|
|
// Timeout after 500ms in case there is a blocked client
|
|
if subscriber_tx
|
|
.send_timeout(notification, SEND_TIMEOUT)
|
|
.await
|
|
.is_err()
|
|
{
|
|
trc::event!(
|
|
Server(ServerEvent::ThreadError),
|
|
Details =
|
|
"Error sending state change to subscriber.",
|
|
CausedBy = trc::location!()
|
|
);
|
|
}
|
|
});
|
|
}
|
|
Err(TrySendError::Closed(_)) => {
|
|
purge_needed = true;
|
|
}
|
|
}
|
|
}
|
|
}
|
|
|
|
if subscribers.is_push
|
|
&& push_tx.send(Event::Push { notification }).await.is_err()
|
|
{
|
|
trc::event!(
|
|
Server(ServerEvent::ThreadError),
|
|
Details = "Error sending push updates.",
|
|
CausedBy = trc::location!()
|
|
);
|
|
}
|
|
}
|
|
}
|
|
|
|
PushEvent::PushServerUpdate {
|
|
account_id,
|
|
broadcast,
|
|
} => {
|
|
// Publish event to cluster
|
|
if broadcast
|
|
&& let Some(broadcast_tx) = inner.ipc.broadcast_tx.as_ref()
|
|
&& broadcast_tx
|
|
.send(BroadcastEvent::PushServerUpdate(account_id))
|
|
.await
|
|
.is_err()
|
|
{
|
|
trc::event!(
|
|
Server(trc::ServerEvent::ThreadError),
|
|
Details = "Error sending broadcast event.",
|
|
CausedBy = trc::location!()
|
|
);
|
|
}
|
|
|
|
// Notify push manager
|
|
if push_tx.send(Event::Update { account_id }).await.is_err() {
|
|
trc::event!(
|
|
Server(ServerEvent::ThreadError),
|
|
Details = "Error sending push updates.",
|
|
CausedBy = trc::location!()
|
|
);
|
|
}
|
|
}
|
|
}
|
|
|
|
if purge_needed {
|
|
let mut remove_account_ids = Vec::new();
|
|
|
|
for (account_id, subscribers) in &mut subscribers {
|
|
subscribers.ipc.retain(|subscriber| subscriber.is_valid());
|
|
|
|
if subscribers.ipc.is_empty() && !subscribers.is_push {
|
|
remove_account_ids.push(*account_id);
|
|
}
|
|
}
|
|
|
|
for remove_account_id in remove_account_ids {
|
|
subscribers.remove(&remove_account_id);
|
|
}
|
|
|
|
last_purge = Instant::now();
|
|
}
|
|
}
|
|
});
|
|
}
|