Merge upstream v0.16.25

Brings in the stripped v0.16.25 snapshot (72f8ddd), a bug-fix release:
DKIM rotation (keys retired before their successor is published, keys
made under manual DNS never rotating, manual DNS activating an
unpublished key), IMAP answering failed logins with an untagged NO,
DNSBL scoring only the first return code and caching "not listed" for
24 hours, Pyzor digesting empty input, queue quotas with an empty match
never enforced, RocksDB's info log growing without limit, MaskedEmail
creation with several domains, and Autodiscover answering other schemas
with the Outlook settings.

Conflicts:
- crates/main/Cargo.toml: inbuxa's name and AGPL-only license, version
  0.16.25.
- SECURITY.md and .github/PULL_REQUEST_TEMPLATE.md: inbuxa's own.
- Cargo.lock: upstream's, re-resolved against inbuxa's manifests.
- crates/common/src/network/autoconfig/autodiscover.rs: upstream now
  parses the request into a struct, so the legacy-protocol switch
  (LP-7) reads the address from request.email; ActiveSync and other
  schemas get upstream's error 601 untouched.
- resources/schema: unchanged upstream apart from two labels the rename
  pass now covers, which main already had.

AGENTS.md, new upstream, is left out: it is about contributing to
upstream, which doesn't apply to this repository.
This commit is contained in:
jcoffey-dev committed 2026-10-05 21:20:26 -07:00
commit db4135e481
60 files changed
+1574 -488

No files matched your search

+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "spam-filter"
version = "0.16.24"
version = "0.16.25"
edition = "2024"
[dependencies]
+91 -65
View File
@@ -37,7 +37,7 @@ use std::{
hash::{Hash, RandomState},
sync::Arc,
};
use store::ahash::AHashSet;
use store::ahash::AHashMap;
use store::rand::seq::SliceRandom;
use store::write::{BlobLink, RegistryClass, now};
use store::{
@@ -92,7 +92,14 @@ struct TrainingTask {
sample: TrainingSample,
is_spam: bool,
is_replay: bool,
remove: Option<u64>,
remove: SampleRemoval,
}
#[derive(Debug, Clone, Copy)]
enum SampleRemoval {
Keep,
Item,
ItemAndLink { until: u64 },
}
#[derive(rkyv::Archive, rkyv::Deserialize, rkyv::Serialize, Debug)]
@@ -199,7 +206,7 @@ impl SpamClassifier for Server {
object_id,
item_id: u64::MAX,
}));
let mut seen_samples = AHashSet::new();
let mut seen_samples = AHashMap::new();
let mut spam_count = 0;
let mut ham_count = 0;
self.store()
@@ -220,43 +227,57 @@ impl SpamClassifier for Server {
.unwrap_or(u32::MAX),
};
if seen_samples.insert(sample.clone()) {
// Add to reservoir
if !do_remove {
trainer.reservoir.update_reservoir(
&sample,
match seen_samples.entry(sample.clone()) {
Entry::Vacant(entry) => {
entry.insert((!do_remove).then_some(until));
// Add to reservoir
if !do_remove {
trainer.reservoir.update_reservoir(
&sample,
is_spam,
config.reservoir_capacity,
);
} else {
trainer.reservoir.update_counts(is_spam);
}
samples.push(TrainingTask {
id,
sample,
is_spam,
config.reservoir_capacity,
);
} else {
trainer.reservoir.update_counts(is_spam);
is_replay: false,
remove: if do_remove {
SampleRemoval::ItemAndLink { until }
} else {
SampleRemoval::Keep
},
});
remove_entries |= do_remove;
// Update trainer stats
if is_spam {
spam_count += 1;
} else {
ham_count += 1;
}
}
samples.push(TrainingTask {
id,
sample,
is_spam,
is_replay: false,
remove: do_remove.then_some(until),
});
remove_entries |= do_remove;
// Update trainer stats
if is_spam {
spam_count += 1;
} else {
ham_count += 1;
Entry::Occupied(entry) => {
let remove = if *entry.get() == Some(until) {
SampleRemoval::Item
} else {
SampleRemoval::ItemAndLink { until }
};
duplicate_samples.push(TrainingTask {
id,
sample,
is_spam,
is_replay: false,
remove,
});
remove_entries = true;
}
} else {
duplicate_samples.push(TrainingTask {
id,
sample,
is_spam,
is_replay: false,
remove: Some(until),
});
remove_entries = true;
}
trainer.last_id = trainer.last_id.max(id);
@@ -315,7 +336,7 @@ impl SpamClassifier for Server {
sample: sample.clone(),
is_spam: false,
is_replay: true,
remove: None,
remove: SampleRemoval::Keep,
}),
);
} else if ham_count > spam_count {
@@ -329,7 +350,7 @@ impl SpamClassifier for Server {
sample: sample.clone(),
is_spam: true,
is_replay: true,
remove: None,
remove: SampleRemoval::Keep,
}),
);
}
@@ -884,33 +905,38 @@ async fn delete_samples(
let object_id = ObjectType::SpamTrainingSample.to_id();
let mut batch = BatchBuilder::new();
for sample in samples.into_iter().chain(duplicate_samples) {
if let Some(until) = sample.remove {
batch
.with_account_id(sample.sample.account_id)
.clear(BlobOp::Link {
hash: sample.sample.hash,
to: BlobLink::Temporary { until },
})
.clear(ValueClass::Registry(RegistryClass::Item {
object_id,
item_id: sample.id,
}))
.clear(ValueClass::Registry(RegistryClass::Index {
index_id: Property::AccountId.to_id(),
object_id,
item_id: sample.id,
key: (sample.sample.account_id as u64).serialize(),
}));
let until = match sample.remove {
SampleRemoval::Keep => continue,
SampleRemoval::Item => None,
SampleRemoval::ItemAndLink { until } => Some(until),
};
batch.with_account_id(sample.sample.account_id);
if let Some(until) = until {
batch.clear(BlobOp::Link {
hash: sample.sample.hash,
to: BlobLink::Temporary { until },
});
}
batch
.clear(ValueClass::Registry(RegistryClass::Item {
object_id,
item_id: sample.id,
}))
.clear(ValueClass::Registry(RegistryClass::Index {
index_id: Property::AccountId.to_id(),
object_id,
item_id: sample.id,
key: (sample.sample.account_id as u64).serialize(),
}));
if batch.is_large_batch() {
server
.store()
.write(batch.build_all())
.await
.caused_by(trc::location!())?;
batch = BatchBuilder::new();
batch.with_account_id(sample.sample.account_id);
}
if batch.is_large_batch() {
server
.store()
.write(batch.build_all())
.await
.caused_by(trc::location!())?;
batch = BatchBuilder::new();
batch.with_account_id(sample.sample.account_id);
}
}
if !batch.is_empty() {
+143 -95
View File
@@ -11,7 +11,14 @@ use common::{
config::mailstore::spamfilter::{DnsBlServer, Element, IpResolver, Location},
expr::functions::ResolveVariable,
};
use mail_auth::{Error, common::resolver::ToFqdn};
use mail_auth::common::resolver::ToFqdn;
#[cfg(not(feature = "test_mode"))]
use mail_auth::hickory_resolver::{
net::{DnsError, NetError},
proto::rr::{Name, RData},
};
#[cfg(feature = "test_mode")]
use mail_auth::{DnsError, Error};
use std::{
net::Ipv4Addr,
sync::Arc,
@@ -19,6 +26,18 @@ use std::{
};
use trc::SpamEvent;
const MAX_NEGATIVE_TTL: u32 = 3600;
enum DnsblAnswer {
Listed {
ips: Vec<Ipv4Addr>,
expires: Instant,
},
NotListed {
expires: Option<Instant>,
},
}
pub(crate) async fn check_dnsbl(
server: &Server,
ctx: &mut SpamFilterContext<'_>,
@@ -49,7 +68,7 @@ pub(crate) async fn check_dnsbl(
for dnsbl in &server.core.spam.dnsbl.servers {
if dnsbl.scope == scope
&& checks < max_checks
&& let Some(tag) = is_dnsbl(
&& let Some(codes) = dnsbl_codes(
server,
dnsbl,
SpamFilterResolver::new(ctx, resolver, location),
@@ -58,7 +77,19 @@ pub(crate) async fn check_dnsbl(
)
.await
{
ctx.result.add_tag(tag);
for code in codes.iter() {
let tag = server
.eval_if::<String, _>(
&dnsbl.tags,
&SpamFilterResolver::new(ctx, code, location),
ctx.input.span_id,
)
.await;
if let Some(tag) = tag {
ctx.result.add_tag(tag);
}
}
}
}
@@ -71,13 +102,13 @@ pub(crate) async fn check_dnsbl(
}
}
async fn is_dnsbl(
async fn dnsbl_codes(
server: &Server,
config: &DnsBlServer,
resolver: SpamFilterResolver<'_, impl ResolveVariable>,
element: Element,
checks: &mut usize,
) -> Option<String> {
) -> Option<Arc<[IpResolver]>> {
let time = Instant::now();
let zone = server
.eval_if::<String, _>(&config.zone, &resolver, resolver.ctx.input.span_id)
@@ -92,105 +123,122 @@ async fn is_dnsbl(
{
None
} else {
server
.eval_if(
&config.tags,
&SpamFilterResolver::new(
resolver.ctx,
&IpResolver::new(
format!("127.0.{}.{}", parts[1], parts[0]).parse().unwrap(),
),
resolver.location,
),
resolver.ctx.input.span_id,
)
.await
Some(Arc::from([IpResolver::new(
format!("127.0.{}.{}", parts[1], parts[0]).parse().unwrap(),
)]))
};
}
}
let result = match server.inner.cache.dns_rbl.get(zone.as_str()) {
Some(Some(result)) => result,
Some(None) => return None,
None => {
*checks += 1;
if let Some(codes) = server.inner.cache.dns_rbl.get(zone.as_str()) {
return codes;
}
match server
.core
.smtp
.resolvers
.dns
.ipv4_lookup_raw(zone.to_fqdn().as_ref())
.await
{
Ok(result) => {
trc::event!(
Spam(SpamEvent::Dnsbl),
Hostname = zone.clone(),
Result = result
.entry
.iter()
.map(|ip| trc::Value::from(ip.to_string()))
.collect::<Vec<_>>(),
Details = element.as_str(),
Elapsed = time.elapsed()
);
*checks += 1;
let entry = Arc::new(IpResolver::new(
result
.entry
.iter()
.copied()
.next()
.unwrap_or(Ipv4Addr::BROADCAST)
.into(),
));
match resolve_zone(server, zone.to_fqdn().as_ref()).await {
Ok(DnsblAnswer::Listed { ips, expires }) => {
trc::event!(
Spam(SpamEvent::Dnsbl),
Hostname = zone.clone(),
Result = ips
.iter()
.map(|ip| trc::Value::from(ip.to_string()))
.collect::<Vec<_>>(),
Details = element.as_str(),
Elapsed = time.elapsed()
);
server.inner.cache.dns_rbl.insert_with_expiry(
zone.into(),
Some(entry.clone()),
result.expires,
);
let codes: Arc<[IpResolver]> = ips
.into_iter()
.map(|ip| IpResolver::new(ip.into()))
.collect();
entry
}
Err(Error::Dns(mail_auth::DnsError::RecordNotFound(_))) => {
trc::event!(
Spam(SpamEvent::Dnsbl),
Hostname = zone.clone(),
Result = trc::Value::None,
Details = element.as_str(),
Elapsed = time.elapsed()
);
server.inner.cache.dns_rbl.insert_with_expiry(
zone.into(),
Some(codes.clone()),
expires,
);
server.inner.cache.dns_rbl.insert(
zone.into(),
None,
Duration::from_secs(86400),
);
return None;
}
Err(err) => {
trc::event!(
Spam(SpamEvent::DnsblError),
Hostname = zone,
Elapsed = time.elapsed(),
Details = element.as_str(),
CausedBy = err.to_string()
);
return None;
}
}
Some(codes)
}
};
Ok(DnsblAnswer::NotListed { expires }) => {
trc::event!(
Spam(SpamEvent::Dnsbl),
Hostname = zone.clone(),
Result = trc::Value::None,
Details = element.as_str(),
Elapsed = time.elapsed()
);
server
.eval_if(
&config.tags,
&SpamFilterResolver::new(resolver.ctx, result.as_ref(), resolver.location),
resolver.ctx.input.span_id,
)
.await
if let Some(expires) = expires {
server
.inner
.cache
.dns_rbl
.insert_with_expiry(zone.into(), None, expires);
}
None
}
Err(err) => {
trc::event!(
Spam(SpamEvent::DnsblError),
Hostname = zone,
Elapsed = time.elapsed(),
Details = element.as_str(),
CausedBy = err
);
None
}
}
}
#[cfg(not(feature = "test_mode"))]
async fn resolve_zone(server: &Server, zone: &str) -> Result<DnsblAnswer, String> {
let name = Name::from_str_relaxed(zone).map_err(|err| err.to_string())?;
match server.core.smtp.resolvers.dns.0.ipv4_lookup(name).await {
Ok(lookup) => {
let expires = lookup.valid_until();
let ips = lookup
.answers()
.iter()
.filter_map(|record| match &record.data {
RData::A(a) => Some(a.0),
_ => None,
})
.collect::<Vec<_>>();
Ok(if !ips.is_empty() {
DnsblAnswer::Listed { ips, expires }
} else {
DnsblAnswer::NotListed {
expires: Some(expires),
}
})
}
Err(NetError::Dns(DnsError::NoRecordsFound(no_records))) => Ok(DnsblAnswer::NotListed {
expires: no_records
.negative_ttl
.filter(|ttl| *ttl > 0)
.map(|ttl| Instant::now() + Duration::from_secs(ttl.min(MAX_NEGATIVE_TTL).into())),
}),
Err(err) => Err(err.to_string()),
}
}
#[cfg(feature = "test_mode")]
async fn resolve_zone(server: &Server, zone: &str) -> Result<DnsblAnswer, String> {
match server.core.smtp.resolvers.dns.ipv4_lookup_raw(zone).await {
Ok(result) => Ok(DnsblAnswer::Listed {
ips: result.entry.to_vec(),
expires: result.expires,
}),
Err(Error::Dns(DnsError::RecordNotFound(_))) => Ok(DnsblAnswer::NotListed {
expires: Some(Instant::now() + Duration::from_secs(MAX_NEGATIVE_TTL.into())),
}),
Err(err) => Err(err.to_string()),
}
}
+39 -35
View File
@@ -33,17 +33,9 @@ pub(crate) async fn pyzor_check(
message: &Message<'_>,
config: &PyzorConfig,
) -> trc::Result<Option<PyzorResponse>> {
// Make sure there is at least one text part
if !message
.parts
.iter()
.any(|p| matches!(p.body, PartType::Text(_) | PartType::Html(_)))
{
let Some(request) = message.pyzor_check_message() else {
return Ok(None);
}
// Hash message
let request = message.pyzor_check_message();
};
// Send message to address. inbuxa: in tests, a fixed table answers
// instead of a public server (test_response).
@@ -167,15 +159,15 @@ impl PyzorWrite for Sha1 {
}
trait PyzorDigest<W: PyzorWrite> {
fn pyzor_digest(&self, writer: W) -> W;
fn pyzor_digest(&self, writer: W) -> Option<W>;
}
pub trait PyzorCheck {
fn pyzor_check_message(&self) -> String;
fn pyzor_check_message(&self) -> Option<String>;
}
impl<W: PyzorWrite> PyzorDigest<W> for Message<'_> {
fn pyzor_digest(&self, writer: W) -> W {
fn pyzor_digest(&self, writer: W) -> Option<W> {
let parts = self
.parts
.iter()
@@ -191,7 +183,7 @@ impl<W: PyzorWrite> PyzorDigest<W> for Message<'_> {
}
impl PyzorCheck for Message<'_> {
fn pyzor_check_message(&self) -> String {
fn pyzor_check_message(&self) -> Option<String> {
let time = SystemTime::now()
.duration_since(SystemTime::UNIX_EPOCH)
.map_or(0, |d| d.as_secs());
@@ -204,9 +196,9 @@ impl PyzorCheck for Message<'_> {
}
}
fn pyzor_create_message(message: &Message<'_>, time: u64, thread: u16) -> String {
fn pyzor_create_message(message: &Message<'_>, time: u64, thread: u16) -> Option<String> {
// Hash message
let hash = message.pyzor_digest(Sha1::new()).finalize().hex_encode();
let hash = message.pyzor_digest(Sha1::new())?.finalize().hex_encode();
// Hash key
let mut hash_key = Sha1::new();
hash_key.update("anonymous:".as_bytes());
@@ -226,10 +218,10 @@ fn pyzor_create_message(message: &Message<'_>, time: u64, thread: u16) -> String
sig.update(format!(":{time}:{hash_key}"));
let sig = sig.finalize().hex_encode();
format!("{message}\nSig: {sig}\n")
Some(format!("{message}\nSig: {sig}\n"))
}
fn pyzor_digest<'x, I, W>(mut writer: W, lines: I) -> W
fn pyzor_digest<'x, I, W>(mut writer: W, lines: I) -> Option<W>
where
I: Iterator<Item = &'x str>,
W: PyzorWrite,
@@ -291,6 +283,10 @@ where
}
}
if result.is_empty() {
return None;
}
if result.len() > ATOMIC_NUM_LINES {
for (offset, length) in DIGEST_SPEC {
for i in 0..*length {
@@ -305,7 +301,7 @@ where
}
}
writer
Some(writer)
}
fn html_to_text(input: &str) -> String {
@@ -489,7 +485,8 @@ mod test {
&MessageParser::new().parse(HTML_TEXT_STYLE_SCRIPT).unwrap(),
1697468672,
49005,
);
)
.unwrap();
assert_eq!(
message,
@@ -522,10 +519,9 @@ mod test {
"http://spammer.com/special-offers?buy=now",
] {
assert_eq!(
String::from_utf8(pyzor_digest(
Vec::new(),
format!("Test {strip_me} Test2").lines(),
))
String::from_utf8(
pyzor_digest(Vec::new(), format!("Test {strip_me} Test2").lines()).unwrap()
)
.unwrap(),
"TestTest2"
);
@@ -533,20 +529,26 @@ mod test {
// Test short lines
assert_eq!(
String::from_utf8(pyzor_digest(
Vec::new(),
concat!("This line is included\n", "not this\n", "This also").lines(),
))
String::from_utf8(
pyzor_digest(
Vec::new(),
concat!("This line is included\n", "not this\n", "This also").lines(),
)
.unwrap()
)
.unwrap(),
"ThislineisincludedThisalso"
);
// Test atomic
assert_eq!(
String::from_utf8(pyzor_digest(
Vec::new(),
"All this message\nShould be included\nIn the digest".lines(),
))
String::from_utf8(
pyzor_digest(
Vec::new(),
"All this message\nShould be included\nIn the digest".lines(),
)
.unwrap()
)
.unwrap(),
"AllthismessageShouldbeincludedInthedigest"
);
@@ -561,7 +563,7 @@ mod test {
expected += format!("Line{i}testtesttest").as_str();
}
assert_eq!(
String::from_utf8(pyzor_digest(Vec::new(), text.lines(),)).unwrap(),
String::from_utf8(pyzor_digest(Vec::new(), text.lines()).unwrap()).unwrap(),
expected
);
@@ -593,7 +595,8 @@ mod test {
MessageParser::new()
.parse(input)
.unwrap()
.pyzor_digest(Vec::new(),)
.pyzor_digest(Vec::new())
.unwrap()
)
.unwrap(),
expected,
@@ -606,7 +609,8 @@ mod test {
MessageParser::new()
.parse(HTML_TEXT_STYLE_SCRIPT)
.unwrap()
.pyzor_digest(Sha1::new(),)
.pyzor_digest(Sha1::new())
.unwrap()
.finalize()
.hex_encode(),
"b2c27325a034c581df0c9ef37e4a0d63208a3e7e",