Merge upstream v0.16.24

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}.
This commit is contained in:
2026-09-28 06:30:20 -07:00
94 changed files with 4206 additions and 1065 deletions
+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "common"
version = "0.16.23"
version = "0.16.24"
edition = "2024"
build = "build.rs"
+3 -3
View File
@@ -535,8 +535,8 @@ async fn insert_safe_defaults(bp: &mut Bootstrap) -> trc::Result<()> {
// inbuxa: rules are always to hand, since a copy ships with the server
// (spam_rules). They load on first boot, and again when the bundled
// version differs from the one last loaded, which only adds what's
// missing: new tags and rules, never a changed score.
// rules differ from the ones last loaded: new tags and rules, fixes to
// rules nobody edited, never a changed score or an admin's edit.
let rules_url = super::spam_rules::rules_url(
bp.registry
.object::<SpamSettings>(Id::singleton())
@@ -547,7 +547,7 @@ async fn insert_safe_defaults(bp: &mut Bootstrap) -> trc::Result<()> {
&& super::spam_rules::applied_version(&bp.data_store)
.await?
.as_deref()
!= Some(super::spam_rules::BUNDLED_SPAM_RULES_VERSION);
!= Some(super::spam_rules::BUNDLED_SPAM_RULES_APPLIED);
if bp.registry.count_object(ObjectType::SpamRule).await? == 0 || bundled_is_new {
let mut batch = BatchBuilder::new();
batch.schedule_task(Task::SpamFilterMaintenance(TaskSpamFilterMaintenance {
+61 -8
View File
@@ -12,11 +12,16 @@
//! and license) and uses it whenever no other source is configured. The rules
//! URL remains an operator override (`https://` or `file://`).
//!
//! Loading rules only ever adds what's missing, never changes an existing rule
//! or score. They load on first boot, and again whenever the bundled version
//! differs from the one last applied, so an upgrade brings new tags (the AI
//! classifier's `LLM_*` scores, say) to an install that already had rules.
//! Loading rules adds what's missing and brings an existing rule up to date,
//! but never touches one an admin edited: every object an update writes is
//! fingerprinted, and one that no longer matches its fingerprint is kept as
//! it is. Tags (scores) are never replaced. Switching a rule on or off isn't
//! an edit, and is kept either way. They load on first boot, and again
//! whenever the bundled rules differ from the ones last applied, so an
//! upgrade brings new tags (the AI classifier's `LLM_*` scores, say) and
//! fixed rules to an install that already had rules.
use registry::{schema::prelude::ObjectType, types::EnumImpl};
use std::io::Read;
use store::{
SUBSPACE_INBUXA, Store, ValueKey,
@@ -27,13 +32,17 @@ use trc::AddContext;
/// The version of spam-filter the embedded rules come from.
pub const BUNDLED_SPAM_RULES_VERSION: &str = "3.0.2";
/// What's recorded once the bundled rules are loaded: their version, then the
/// fork's own generation of the update, so a change to how an update applies
/// runs it once more. Generation 2 fingerprints (upstream v0.16.24).
pub const BUNDLED_SPAM_RULES_APPLIED: &str = "3.0.2+2";
static BUNDLED_SPAM_RULES: &[u8] =
include_bytes!("../../../../resources/spam-filter/spam-filter-rules.json.gz");
/// Upstream's default rules source, the value every install created before
/// the rules were bundled has saved. Read only to treat it as unset.
const LEGACY_DEFAULT_URL: &str =
"https://github.com/stalwartlabs/spam-filter/releases/latest/download/spam-filter-rules.json.gz";
const LEGACY_DEFAULT_URL: &str = "https://github.com/stalwartlabs/spam-filter/releases/latest/download/spam-filter-rules.json.gz";
/// The URL to fetch rules from, or `None` for the bundled rules. An empty
/// setting and upstream's old default both mean the bundled rules.
@@ -57,14 +66,49 @@ fn applied_key() -> ValueClass {
})
}
/// The bundled version last loaded into the registry, if any.
fn fingerprint_key(object: ObjectType, id: u64) -> ValueClass {
let mut key = b"Sf".to_vec();
key.extend_from_slice(object.as_str().as_bytes());
key.push(0);
key.extend_from_slice(&id.to_be_bytes());
ValueClass::Any(AnyClass {
subspace: SUBSPACE_INBUXA,
key,
})
}
/// The fingerprint of what a rules update last wrote to this object, if one
/// did.
pub async fn fingerprint(data: &Store, object: ObjectType, id: u64) -> trc::Result<Option<String>> {
data.get_value::<String>(ValueKey::from(fingerprint_key(object, id)))
.await
.caused_by(trc::location!())
}
/// Records the fingerprint of what a rules update wrote to this object.
pub async fn set_fingerprint(
data: &Store,
object: ObjectType,
id: u64,
fingerprint: &str,
) -> trc::Result<()> {
let mut batch = BatchBuilder::new();
batch.set(fingerprint_key(object, id), fingerprint.as_bytes().to_vec());
data.write(batch.build_all())
.await
.caused_by(trc::location!())
.map(|_| ())
}
/// The bundled rules last loaded into the registry, if any
/// ([`BUNDLED_SPAM_RULES_APPLIED`]'s form).
pub async fn applied_version(data: &Store) -> trc::Result<Option<String>> {
data.get_value::<String>(ValueKey::from(applied_key()))
.await
.caused_by(trc::location!())
}
/// Records that the bundled rules of this version have been loaded.
/// Records that the bundled rules have been loaded.
pub async fn set_applied_version(data: &Store, version: &str) -> trc::Result<()> {
let mut batch = BatchBuilder::new();
batch.set(applied_key(), version.as_bytes().to_vec());
@@ -90,6 +134,15 @@ mod tests {
);
}
#[test]
fn applied_marker_names_the_bundled_version() {
assert!(
BUNDLED_SPAM_RULES_APPLIED
.strip_prefix(BUNDLED_SPAM_RULES_VERSION)
.is_some_and(|generation| generation.starts_with('+'))
);
}
#[test]
fn bundled_rules_parse_and_score_the_ai_tags() {
let rules: serde_json::Value = serde_json::from_slice(&bundled_rules().unwrap()).unwrap();
@@ -10,8 +10,9 @@ use crate::{Server, manager::application::Resource};
use quick_xml::Reader;
use quick_xml::XmlVersion;
use quick_xml::events::Event;
use registry::schema::enums::ServiceProtocol;
use registry::schema::{enums::ServiceProtocol, structs::Service};
use std::fmt::Write;
use utils::map::vec_map::VecMap;
impl Server {
pub async fn handle_autodiscover_request(
@@ -26,89 +27,103 @@ impl Server {
.details("Failed to parse autodiscover request")
.ctx(trc::Key::Reason, err)
})?;
let default_host = &self.core.network.server_name;
// Build XML response
let mut config = String::with_capacity(1024);
let _ = writeln!(&mut config, "<?xml version=\"1.0\" encoding=\"UTF-8\"?>");
let _ = writeln!(
&mut config,
"<Autodiscover xmlns=\"http://schemas.microsoft.com/exchange/autodiscover/responseschema/2006\">"
);
let _ = writeln!(
&mut config,
"\t<Response xmlns=\"http://schemas.microsoft.com/exchange/autodiscover/outlook/responseschema/2006a\">"
);
let _ = writeln!(&mut config, "\t\t<User>");
let _ = writeln!(
&mut config,
"\t\t\t<DisplayName>{emailaddress}</DisplayName>"
);
let _ = writeln!(
&mut config,
"\t\t\t<AutoDiscoverSMTPAddress>{emailaddress}</AutoDiscoverSMTPAddress>"
);
// DeploymentId is a required field of User but we are not a MS Exchange server so use a random value
let _ = writeln!(
&mut config,
"\t\t\t<DeploymentId>644560b8-a1ce-429c-8ace-23395843f701</DeploymentId>"
);
let _ = writeln!(&mut config, "\t\t</User>");
let _ = writeln!(&mut config, "\t\t<Account>");
let _ = writeln!(&mut config, "\t\t\t<AccountType>email</AccountType>");
let _ = writeln!(&mut config, "\t\t\t<Action>settings</Action>");
// inbuxa: legacy-protocols LP-7, LP-14a
let legacy_off = match emailaddress.rsplit_once('@') {
Some((_, domain)) => self.legacy_off_for(domain).await?,
None => self.legacy_off_for("").await?,
};
for (protocol, service) in &self.core.network.info.services {
if legacy_off.service(protocol) {
continue;
}
let (protocol, ports) = match protocol {
ServiceProtocol::Imap => ("IMAP", [143, 993]),
ServiceProtocol::Pop3 => ("POP3", [110, 995]),
ServiceProtocol::Smtp => ("SMTP", [587, 465]),
_ => continue,
};
for (is_tls, port) in ports.into_iter().enumerate() {
if is_tls == 1 || service.cleartext {
let server_name = service.hostname.as_deref().unwrap_or(default_host);
let _ = writeln!(&mut config, "\t\t\t<Protocol>");
let _ = writeln!(&mut config, "\t\t\t\t<Type>{protocol}</Type>",);
let _ = writeln!(&mut config, "\t\t\t\t<Server>{server_name}</Server>");
let _ = writeln!(&mut config, "\t\t\t\t<Port>{port}</Port>");
let _ = writeln!(&mut config, "\t\t\t\t<LoginName>{emailaddress}</LoginName>");
let _ = writeln!(&mut config, "\t\t\t\t<AuthRequired>on</AuthRequired>");
let _ = writeln!(&mut config, "\t\t\t\t<DirectoryPort>0</DirectoryPort>");
let _ = writeln!(&mut config, "\t\t\t\t<ReferralPort>0</ReferralPort>");
let _ = writeln!(
&mut config,
"\t\t\t\t<SSL>{}</SSL>",
if is_tls == 1 { "on" } else { "off" }
);
if is_tls == 1 {
let _ = writeln!(&mut config, "\t\t\t\t<Encryption>TLS</Encryption>");
}
let _ = writeln!(&mut config, "\t\t\t\t<SPA>off</SPA>");
let _ = writeln!(&mut config, "\t\t\t</Protocol>");
}
}
}
let _ = writeln!(&mut config, "\t\t</Account>");
let _ = writeln!(&mut config, "\t</Response>");
let _ = writeln!(&mut config, "</Autodiscover>");
Ok(Resource::new(
"application/xml; charset=utf-8",
config.into_bytes(),
build_autodiscover_response(
&emailaddress,
&self.core.network.server_name,
&self.core.network.info.services,
|protocol| legacy_off.service(protocol),
)
.into_bytes(),
))
}
}
fn build_autodiscover_response(
emailaddress: &str,
default_host: &str,
services: &VecMap<ServiceProtocol, Service>,
switched_off: impl Fn(&ServiceProtocol) -> bool,
) -> String {
// Build XML response
let mut config = String::with_capacity(1024);
let _ = writeln!(&mut config, "<?xml version=\"1.0\" encoding=\"UTF-8\"?>");
let _ = writeln!(
&mut config,
"<Autodiscover xmlns=\"http://schemas.microsoft.com/exchange/autodiscover/responseschema/2006\">"
);
let _ = writeln!(
&mut config,
"\t<Response xmlns=\"http://schemas.microsoft.com/exchange/autodiscover/outlook/responseschema/2006a\">"
);
let _ = writeln!(&mut config, "\t\t<User>");
let _ = writeln!(
&mut config,
"\t\t\t<DisplayName>{emailaddress}</DisplayName>"
);
let _ = writeln!(
&mut config,
"\t\t\t<AutoDiscoverSMTPAddress>{emailaddress}</AutoDiscoverSMTPAddress>"
);
// DeploymentId is a required field of User but we are not a MS Exchange server so use a random value
let _ = writeln!(
&mut config,
"\t\t\t<DeploymentId>644560b8-a1ce-429c-8ace-23395843f701</DeploymentId>"
);
let _ = writeln!(&mut config, "\t\t</User>");
let _ = writeln!(&mut config, "\t\t<Account>");
let _ = writeln!(&mut config, "\t\t\t<AccountType>email</AccountType>");
let _ = writeln!(&mut config, "\t\t\t<Action>settings</Action>");
for (protocol, service) in services {
if switched_off(protocol) {
continue;
}
let (protocol, ports) = match protocol {
ServiceProtocol::Imap => ("IMAP", [(993, true), (143, false)]),
ServiceProtocol::Pop3 => ("POP3", [(995, true), (110, false)]),
ServiceProtocol::Smtp => ("SMTP", [(465, true), (587, false)]),
_ => continue,
};
// Implicit TLS is listed first so that it is preferred (RFC 8314)
for (port, is_tls) in ports {
if is_tls || service.cleartext {
let server_name = service.hostname.as_deref().unwrap_or(default_host);
let _ = writeln!(&mut config, "\t\t\t<Protocol>");
let _ = writeln!(&mut config, "\t\t\t\t<Type>{protocol}</Type>",);
let _ = writeln!(&mut config, "\t\t\t\t<Server>{server_name}</Server>");
let _ = writeln!(&mut config, "\t\t\t\t<Port>{port}</Port>");
let _ = writeln!(&mut config, "\t\t\t\t<LoginName>{emailaddress}</LoginName>");
let _ = writeln!(&mut config, "\t\t\t\t<AuthRequired>on</AuthRequired>");
let _ = writeln!(&mut config, "\t\t\t\t<DirectoryPort>0</DirectoryPort>");
let _ = writeln!(&mut config, "\t\t\t\t<ReferralPort>0</ReferralPort>");
let (ssl, encryption) = if is_tls {
("on", "SSL")
} else {
("off", "TLS")
};
let _ = writeln!(&mut config, "\t\t\t\t<SSL>{ssl}</SSL>");
let _ = writeln!(&mut config, "\t\t\t\t<Encryption>{encryption}</Encryption>");
let _ = writeln!(&mut config, "\t\t\t\t<SPA>off</SPA>");
let _ = writeln!(&mut config, "\t\t\t</Protocol>");
}
}
}
let _ = writeln!(&mut config, "\t\t</Account>");
let _ = writeln!(&mut config, "\t</Response>");
let _ = writeln!(&mut config, "</Autodiscover>");
config
}
fn parse_autodiscover_request(bytes: &[u8]) -> Result<String, String> {
if bytes.is_empty() {
return Err("Empty request body".to_string());
@@ -211,4 +226,79 @@ mod tests {
"[email protected]"
);
}
#[test]
fn autodiscover_encryption() {
use registry::schema::{enums::ServiceProtocol, structs::Service};
use utils::map::vec_map::VecMap;
fn tag<'x>(block: &'x str, name: &str) -> &'x str {
block
.split_once(&format!("<{name}>"))
.and_then(|(_, rest)| rest.split_once(&format!("</{name}>")))
.map(|(value, _)| value)
.unwrap()
}
for (cleartext, expected) in [
(
false,
vec![
("IMAP", "993", "on", "SSL"),
("POP3", "995", "on", "SSL"),
("SMTP", "465", "on", "SSL"),
],
),
(
true,
vec![
("IMAP", "993", "on", "SSL"),
("IMAP", "143", "off", "TLS"),
("POP3", "995", "on", "SSL"),
("POP3", "110", "off", "TLS"),
("SMTP", "465", "on", "SSL"),
("SMTP", "587", "off", "TLS"),
],
),
] {
let services: VecMap<ServiceProtocol, Service> = [
ServiceProtocol::Imap,
ServiceProtocol::Pop3,
ServiceProtocol::Smtp,
ServiceProtocol::Jmap,
]
.into_iter()
.map(|protocol| {
(
protocol,
Service {
hostname: None,
cleartext,
},
)
})
.collect();
let response = super::build_autodiscover_response(
"[email protected]",
"mail.example.com",
&services,
|_| false,
);
assert_eq!(
response
.split("<Protocol>")
.skip(1)
.map(|block| (
tag(block, "Type"),
tag(block, "Port"),
tag(block, "SSL"),
tag(block, "Encryption"),
))
.collect::<Vec<_>>(),
expected,
"cleartext: {cleartext}"
);
}
}
}
+14
View File
@@ -961,6 +961,20 @@ impl DnsUpdater {
)
.map_err(|err| format!("Failed to build DNS updater: {}", err))?,
}),
DnsServer::PowerDns(server) => Ok(DnsUpdater {
polling_interval: server.polling_interval.into_inner(),
propagation_timeout: server.propagation_timeout.into_inner(),
propagation_delay: server.propagation_delay.map(|d| d.into_inner()),
ttl: server.ttl.into_inner(),
core,
updater: dns_update::DnsUpdater::new_pdns(
server.api_key.secret().await?,
server.endpoint,
server.server_id,
server.timeout.into_inner().into(),
)
.map_err(|err| format!("Failed to build DNS updater: {}", err))?,
}),
DnsServer::Safedns(server) => Ok(DnsUpdater {
polling_interval: server.polling_interval.into_inner(),
propagation_timeout: server.propagation_timeout.into_inner(),
+102 -36
View File
@@ -6,33 +6,79 @@
* Modified by Coffey Labs in 2026 for INBUXA.
*/
use ahash::AHashMap;
use base64::{Engine, engine::general_purpose::URL_SAFE_NO_PAD};
use p256::{
SecretKey,
ecdsa::{Signature, SigningKey, signature::Signer},
pkcs8::{DecodePrivateKey, PrivateKeyInfo, der::SecretDocument},
};
use parking_lot::Mutex;
use reqwest::{Url, header::HeaderValue};
use std::sync::Arc;
const VAPID_TOKEN_TTL: u64 = 12 * 60 * 60;
const VAPID_TOKEN_REFRESH: u64 = VAPID_TOKEN_TTL / 2;
#[derive(Clone)]
pub struct Vapid {
key: VapidKey,
contact: Option<String>,
tokens: Arc<Mutex<AHashMap<String, VapidToken>>>,
}
struct VapidToken {
authorization: HeaderValue,
issued_at: u64,
}
impl Vapid {
pub fn new(key: VapidKey, contact: Option<String>) -> Self {
Self { key, contact }
Self {
key,
contact,
tokens: Arc::default(),
}
}
pub fn public_key(&self) -> &str {
self.key.public_key()
}
pub fn authorization(&self, endpoint: &str, now: u64) -> Option<String> {
self.key
.authorization(endpoint, self.contact.as_deref(), now)
pub fn authorization(&self, endpoint: &str, now: u64) -> Option<HeaderValue> {
let prefix = endpoint_prefix(endpoint)?;
if let Some(token) = self
.tokens
.lock()
.get(prefix)
.filter(|token| token.is_fresh(now))
{
return Some(token.authorization.clone());
}
let authorization = HeaderValue::try_from(self.key.authorization(
endpoint,
self.contact.as_deref(),
now,
)?)
.ok()?;
let mut tokens = self.tokens.lock();
tokens.retain(|_, token| token.is_fresh(now));
tokens.insert(
prefix.to_string(),
VapidToken {
authorization: authorization.clone(),
issued_at: now,
},
);
Some(authorization)
}
}
impl VapidToken {
fn is_fresh(&self, now: u64) -> bool {
now.checked_sub(self.issued_at)
.is_some_and(|age| age < VAPID_TOKEN_REFRESH)
}
}
@@ -105,41 +151,15 @@ impl VapidKey {
}
}
fn endpoint_origin(url: &str) -> Option<String> {
fn endpoint_prefix(url: &str) -> Option<&str> {
let (scheme, rest) = url.split_once("://")?;
let scheme = scheme.to_ascii_lowercase();
let authority = rest.split(['/', '?', '#']).next()?;
let authority = authority
.rsplit_once('@')
.map(|(_, host)| host)
.unwrap_or(authority);
if authority.is_empty() {
return None;
}
url.get(..scheme.len() + "://".len() + authority.len())
}
let (host, port) = if let Some(rest) = authority.strip_prefix('[') {
let (addr, tail) = rest.split_once(']')?;
(
format!("[{}]", addr.to_ascii_lowercase()),
tail.strip_prefix(':').filter(|port| !port.is_empty()),
)
} else if let Some((host, port)) = authority.rsplit_once(':') {
(
host.to_ascii_lowercase(),
Some(port).filter(|p| !p.is_empty()),
)
} else {
(authority.to_ascii_lowercase(), None)
};
match port {
Some(port)
if !((scheme == "https" && port == "443") || (scheme == "http" && port == "80")) =>
{
Some(format!("{scheme}://{host}:{port}"))
}
_ => Some(format!("{scheme}://{host}")),
}
fn endpoint_origin(url: &str) -> Option<String> {
let origin = Url::parse(url).ok()?.origin();
origin.is_tuple().then(|| origin.ascii_serialization())
}
pub fn normalize_contact(contact: &str) -> Option<String> {
@@ -206,7 +226,12 @@ mod tests {
endpoint_origin("http://[2001:DB8::1]:80/p").unwrap(),
"http://[2001:db8::1]"
);
assert_eq!(
endpoint_origin("https://attacker.example\\@fcm.googleapis.com/fcm/send/x").unwrap(),
"https://attacker.example"
);
assert!(endpoint_origin("not-a-url").is_none());
assert!(endpoint_origin("mailto:[email protected]").is_none());
}
#[test]
@@ -336,6 +361,47 @@ B4yDfR2rGOd2H6Kv3fQNHPj9Nu5Tks8QYMLzrX8ONCNoFnNUQl9S0r0QS6phVqD0
}
}
#[test]
fn authorization_is_reused_per_endpoint_prefix() {
let vapid = Vapid::new(test_key(), None);
let now = 1_700_000_000;
let token = vapid
.authorization("https://push.example.com/push/a", now)
.unwrap();
assert_eq!(
vapid
.authorization("https://push.example.com/push/b?x=1", now + 60)
.unwrap(),
token
);
assert_ne!(
vapid
.authorization("https://other.example.com/push/a", now)
.unwrap(),
token
);
assert_ne!(
vapid
.authorization("https://push.example.com/push/a", now - 1)
.unwrap(),
token
);
let refreshed = vapid
.authorization("https://push.example.com/push/a", now + VAPID_TOKEN_REFRESH)
.unwrap();
assert_ne!(refreshed, token);
assert_eq!(
vapid
.authorization(
"https://push.example.com/push/c",
now + VAPID_TOKEN_REFRESH + 1
)
.unwrap(),
refreshed
);
}
#[test]
fn authorization_omits_subject_when_no_contact() {
let key = test_key();
+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "coordinator"
version = "0.16.23"
version = "0.16.24"
edition = "2024"
[dependencies]
+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "dav-proto"
version = "0.16.23"
version = "0.16.24"
edition = "2024"
[dependencies]
+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "dav"
version = "0.16.23"
version = "0.16.24"
edition = "2024"
[dependencies]
+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "directory"
version = "0.16.23"
version = "0.16.24"
edition = "2024"
[dependencies]
+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "email"
version = "0.16.23"
version = "0.16.24"
edition = "2024"
[dependencies]
+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "groupware"
version = "0.16.23"
version = "0.16.24"
edition = "2024"
[dependencies]
+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "http_proto"
version = "0.16.23"
version = "0.16.24"
edition = "2024"
[dependencies]
+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "http"
version = "0.16.23"
version = "0.16.24"
edition = "2024"
[dependencies]
+103 -70
View File
@@ -12,7 +12,7 @@ use common::{
},
};
use hyper::body::{Bytes, Frame};
use mail_auth::{IpLookupStrategy, mta_sts::TlsRpt};
use mail_auth::{DnssecStatus, IpLookupStrategy, mta_sts::TlsRpt};
use serde::{Deserialize, Serialize};
use smtp::outbound::{
client::{SmtpClient, StartTlsResult},
@@ -382,81 +382,25 @@ async fn delivery_diagnose(
}
}
// Fetch TLSA record
tx.send(DeliveryStage::TlsaLookupStart).await?;
let now = Instant::now();
let dane_policy = match server.tlsa_lookup(format!("_25._tcp.{hostname}.")).await {
Ok(TlsaResult::Secure(tlsa)) if tlsa.has_end_entities => {
tx.send(DeliveryStage::TlsaLookupSuccess {
record: tlsa.as_ref().clone(),
elapsed: now.elapsed_ms(),
})
.await?;
Some(tlsa)
}
Ok(TlsaResult::Secure(_)) => {
tx.send(DeliveryStage::TlsaLookupError {
elapsed: now.elapsed_ms(),
reason: "TLSA record does not have end entities".to_string(),
})
.await?;
None
}
Ok(TlsaResult::Bogus) => {
tx.send(DeliveryStage::TlsaLookupError {
elapsed: now.elapsed_ms(),
reason: "Bogus TLSA record".to_string(),
})
.await?;
continue 'outer;
}
Ok(TlsaResult::Missing) => {
tx.send(DeliveryStage::TlsaNotFound {
elapsed: now.elapsed_ms(),
reason: "No TLSA DNSSEC records found".to_string(),
})
.await?;
None
}
Err(err) => {
if matches!(
&err,
mail_auth::Error::Dns(mail_auth::DnsError::RecordNotFound(_))
) {
tx.send(DeliveryStage::TlsaNotFound {
elapsed: now.elapsed_ms(),
reason: "No TLSA records found for MX".to_string(),
})
.await?;
None
} else {
tx.send(DeliveryStage::TlsaLookupError {
elapsed: now.elapsed_ms(),
reason: err.to_string(),
})
.await?;
continue 'outer;
}
}
};
tx.send(DeliveryStage::IpLookupStart).await?;
let now = Instant::now();
let remote_ips = match host.fqdn_hostname() {
let validate_addresses = server.core.smtp.resolvers.dnssec_available
&& host.dnssec_status() == DnssecStatus::Secure;
let (remote_ips, addresses_dnssec_status) = match host.fqdn_hostname() {
HostOrIp::Host(hostname) => {
match server
.ip_lookup(&hostname, IpLookupStrategy::Ipv4thenIpv6, usize::MAX, false)
.ip_lookup(
&hostname,
IpLookupStrategy::Ipv4thenIpv6,
usize::MAX,
validate_addresses,
)
.await
{
Ok((remote_ips, _)) if !remote_ips.is_empty() => remote_ips,
Ok((remote_ips, dnssec_status)) if !remote_ips.is_empty() => {
(remote_ips, dnssec_status)
}
Ok(_) => {
tx.send(DeliveryStage::IpLookupError {
reason: "No IP addresses found for host".to_string(),
@@ -475,7 +419,7 @@ async fn delivery_diagnose(
}
}
}
HostOrIp::Ip(ip) => vec![ip],
HostOrIp::Ip(ip) => (vec![ip], DnssecStatus::Indeterminate),
};
tx.send(DeliveryStage::IpLookupSuccess {
@@ -484,6 +428,95 @@ async fn delivery_diagnose(
})
.await?;
// Fetch TLSA record
tx.send(DeliveryStage::TlsaLookupStart).await?;
let now = Instant::now();
let dane_policy = match host.dane_status(addresses_dnssec_status) {
(DnssecStatus::Secure, _) => {
match server.tlsa_lookup(format!("_25._tcp.{hostname}.")).await {
Ok(TlsaResult::Secure(tlsa)) if tlsa.has_end_entities => {
tx.send(DeliveryStage::TlsaLookupSuccess {
record: tlsa.as_ref().clone(),
elapsed: now.elapsed_ms(),
})
.await?;
Some(tlsa)
}
Ok(TlsaResult::Secure(_)) => {
tx.send(DeliveryStage::TlsaLookupError {
elapsed: now.elapsed_ms(),
reason: "TLSA record does not have end entities".to_string(),
})
.await?;
None
}
Ok(TlsaResult::Bogus) => {
tx.send(DeliveryStage::TlsaLookupError {
elapsed: now.elapsed_ms(),
reason: "Bogus TLSA record".to_string(),
})
.await?;
continue 'outer;
}
Ok(TlsaResult::Missing) => {
tx.send(DeliveryStage::TlsaNotFound {
elapsed: now.elapsed_ms(),
reason: "No TLSA DNSSEC records found".to_string(),
})
.await?;
None
}
Err(err) => {
if matches!(
&err,
mail_auth::Error::Dns(mail_auth::DnsError::RecordNotFound(_))
) {
tx.send(DeliveryStage::TlsaNotFound {
elapsed: now.elapsed_ms(),
reason: "No TLSA records found for MX".to_string(),
})
.await?;
None
} else {
tx.send(DeliveryStage::TlsaLookupError {
elapsed: now.elapsed_ms(),
reason: err.to_string(),
})
.await?;
continue 'outer;
}
}
}
}
(DnssecStatus::Bogus, dnssec_entity) => {
tx.send(DeliveryStage::TlsaLookupError {
elapsed: now.elapsed_ms(),
reason: format!("Bogus {dnssec_entity} records were found"),
})
.await?;
continue 'outer;
}
(_, dnssec_entity) => {
tx.send(DeliveryStage::TlsaNotFound {
elapsed: now.elapsed_ms(),
reason: format!(
"{dnssec_entity} records are not DNSSEC signed, DANE does not apply"
),
})
.await?;
None
}
};
for remote_ip in remote_ips {
// Start connection
tx.send(DeliveryStage::ConnectionStart { remote_ip })
+3 -1
View File
@@ -36,7 +36,7 @@ use hyper::{
server::conn::http1,
service::service_fn,
};
use hyper_util::rt::TokioIo;
use hyper_util::rt::{TokioIo, TokioTimer};
use jmap::{
api::{
ToJmapHttpResponse, event_source::EventSourceHandler, request::RequestHandler,
@@ -690,6 +690,7 @@ async fn handle_session<T: SessionStream>(inner: Arc<Inner>, session: SessionDat
let is_tls = session.stream.is_tls();
if let Err(http_err) = http1::Builder::new()
.timer(TokioTimer::new())
.keep_alive(true)
.serve_connection(
TokioIo::new(session.stream),
@@ -875,6 +876,7 @@ async fn handle_session<T: SessionStream>(inner: Arc<Inner>, session: SessionDat
)
.with_upgrades()
.await
&& !http_err.is_timeout()
{
if http_err.is_parse() {
let server = inner.build_server();
+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "imap_proto"
version = "0.16.23"
version = "0.16.24"
edition = "2024"
[dependencies]
+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "imap"
version = "0.16.23"
version = "0.16.24"
edition = "2024"
[dependencies]
+249 -158
View File
@@ -9,6 +9,7 @@ use crate::{
core::{MailboxId, SelectedMailbox, Session, SessionData},
spawn_op,
};
use ahash::AHashMap;
use common::{ipc::PushNotification, network::SessionStream, storage::index::ObjectIndexBuilder};
use email::{
cache::{MessageCacheFetch, email::MessageCacheAccess},
@@ -22,8 +23,13 @@ use email::{
use imap_proto::{
Command, ResponseCode, StatusResponse, protocol::copy_move::Arguments, receiver::Request,
};
use rand::RngExt;
use registry::schema::enums::Permission;
use std::{sync::Arc, time::Instant};
use std::{
ops::RangeInclusive,
sync::Arc,
time::{Duration, Instant},
};
use store::{
ValueKey,
roaring::RoaringBitmap,
@@ -36,6 +42,9 @@ use types::{
type_state::{DataType, StateChange},
};
const MAX_MOVE_RETRIES: u32 = 3;
const MOVE_RETRY_BACKOFF_MS: RangeInclusive<u64> = 1..=15;
impl<T: SessionStream> Session<T> {
pub async fn handle_copy_move(
&mut self,
@@ -236,148 +245,180 @@ impl<T: SessionStream> SessionData<T> {
// Mailboxes are in the same account
let account_id = src_mailbox.id.account_id;
let dest_mailbox_id = UidMailbox::new_unassigned(dest_mailbox_id);
let mut batch = BatchBuilder::new();
let mut written_uids = AHashMap::with_capacity(ids.len());
let mut retries = 0;
for (id, imap_id) in ids {
// Obtain mailbox tags
let data_ = if let Some(result) = self
.get_message_data(account_id, id)
.await
.imap_ctx(&arguments.tag, trc::location!())?
{
result
} else {
continue;
};
loop {
let mut batch = BatchBuilder::new();
copied_ids.clear();
did_move = false;
// Deserialize
let data = data_
.to_unarchived::<MessageData>()
.imap_ctx(&arguments.tag, trc::location!())?;
for (&id, imap_id) in &ids {
// Obtain mailbox tags
let data_ = if let Some(result) = self
.get_message_data(account_id, id)
.await
.imap_ctx(&arguments.tag, trc::location!())?
{
result
} else {
continue;
};
// Make sure the message still belongs to this mailbox
if !data
.inner
.mailboxes
.iter()
.any(|mailbox| mailbox.mailbox_id == src_mailbox.id.mailbox_id)
{
continue;
}
// Deserialize
let data = data_
.to_unarchived::<MessageData>()
.imap_ctx(&arguments.tag, trc::location!())?;
// If the message is already in the destination mailbox, skip it.
if let Some(mailbox) = data
.inner
.mailboxes
.iter()
.find(|mailbox| mailbox.mailbox_id == dest_mailbox_id.mailbox_id)
{
copied_ids.push((imap_id.uid, mailbox.uid.to_native()));
if is_move {
let mut new_data = data.inner.to_builder();
new_data.remove_mailbox(src_mailbox.id.mailbox_id);
batch
.with_account_id(account_id)
.with_collection(Collection::Email)
.with_document(id)
.custom(
ObjectIndexBuilder::new()
.with_current(data)
.with_changes(new_data.seal()),
)
.imap_ctx(&arguments.tag, trc::location!())?
.log_vanished_item(
VanishedCollection::Email,
(src_mailbox.id.mailbox_id, imap_id.uid),
)
.commit_point();
did_move = true;
// Make sure the message still belongs to this mailbox
if !data
.inner
.mailboxes
.iter()
.any(|mailbox| mailbox.mailbox_id == src_mailbox.id.mailbox_id)
{
// Moved by a chunk of a previous attempt
if let Some(&uid) = written_uids.get(&id)
&& data.inner.message_uid(dest_mailbox_id.mailbox_id) == Some(uid)
{
copied_ids.push((imap_id.uid, uid));
did_move = true;
}
continue;
}
continue;
}
// If the message is already in the destination mailbox, skip it.
if let Some(mailbox) = data
.inner
.mailboxes
.iter()
.find(|mailbox| mailbox.mailbox_id == dest_mailbox_id.mailbox_id)
{
let uid = mailbox.uid.to_native();
copied_ids.push((imap_id.uid, uid));
// Prepare changes
let mut new_data = data.inner.to_builder();
if is_move {
let mut new_data = data.inner.to_builder();
new_data.remove_mailbox(src_mailbox.id.mailbox_id);
batch
.with_account_id(account_id)
.with_collection(Collection::Email)
.with_document(id)
.custom(
ObjectIndexBuilder::new()
.with_current(data)
.with_changes(new_data.seal()),
)
.imap_ctx(&arguments.tag, trc::location!())?
.log_vanished_item(
VanishedCollection::Email,
(src_mailbox.id.mailbox_id, imap_id.uid),
)
.commit_point();
written_uids.insert(id, uid);
did_move = true;
}
// Add destination folder
new_data.add_mailbox(dest_mailbox_id);
if is_move {
new_data.remove_mailbox(src_mailbox.id.mailbox_id);
}
continue;
}
// Assign IMAP UIDs
let ids = self
.server
.assign_email_ids(
account_id,
new_data
.mailboxes
.iter()
.filter(|m| m.uid == 0)
.map(|m| m.mailbox_id),
false,
)
.await
.caused_by(trc::location!())?;
// Prepare changes
let mut new_data = data.inner.to_builder();
for (uid_mailbox, uid) in new_data
.mailboxes
.iter_mut()
.filter(|m| m.uid == 0)
.zip(ids)
{
copied_ids.push((imap_id.uid, uid));
uid_mailbox.uid = uid;
}
// Add destination folder
new_data.add_mailbox(dest_mailbox_id);
if is_move {
new_data.remove_mailbox(src_mailbox.id.mailbox_id);
}
// Prepare write batch
batch
.with_account_id(account_id)
.with_collection(Collection::Email)
.with_document(id)
.custom(
ObjectIndexBuilder::new()
.with_current(data)
.with_changes(new_data.seal()),
)
.imap_ctx(&arguments.tag, trc::location!())?;
if is_move {
batch.log_vanished_item(
VanishedCollection::Email,
(src_mailbox.id.mailbox_id, imap_id.uid),
);
}
// Add message to training queue
if dest_mailbox_id.mailbox_id == JUNK_ID {
self.server
.add_account_spam_sample(&mut batch, account_id, id, true, self.session_id)
.await
.imap_ctx(&arguments.tag, trc::location!())?;
} else if src_mailbox.id.mailbox_id == JUNK_ID
&& dest_mailbox_id.mailbox_id != TRASH_ID
{
self.server
.add_account_spam_sample(&mut batch, account_id, id, false, self.session_id)
// Assign IMAP UIDs
let ids = self
.server
.assign_email_ids(
account_id,
new_data
.mailboxes
.iter()
.filter(|m| m.uid == 0)
.map(|m| m.mailbox_id),
false,
)
.await
.caused_by(trc::location!())?;
for (uid_mailbox, uid) in new_data
.mailboxes
.iter_mut()
.filter(|m| m.uid == 0)
.zip(ids)
{
copied_ids.push((imap_id.uid, uid));
written_uids.insert(id, uid);
uid_mailbox.uid = uid;
}
// Prepare write batch
batch
.with_account_id(account_id)
.with_collection(Collection::Email)
.with_document(id)
.custom(
ObjectIndexBuilder::new()
.with_current(data)
.with_changes(new_data.seal()),
)
.imap_ctx(&arguments.tag, trc::location!())?;
if is_move {
batch.log_vanished_item(
VanishedCollection::Email,
(src_mailbox.id.mailbox_id, imap_id.uid),
);
}
// Add message to training queue
if dest_mailbox_id.mailbox_id == JUNK_ID {
self.server
.add_account_spam_sample(
&mut batch,
account_id,
id,
true,
self.session_id,
)
.await
.imap_ctx(&arguments.tag, trc::location!())?;
} else if src_mailbox.id.mailbox_id == JUNK_ID
&& dest_mailbox_id.mailbox_id != TRASH_ID
{
self.server
.add_account_spam_sample(
&mut batch,
account_id,
id,
false,
self.session_id,
)
.await
.imap_ctx(&arguments.tag, trc::location!())?;
}
batch.commit_point();
// Update changelog
if is_move {
did_move = true;
}
}
batch.commit_point();
// Update changelog
if is_move {
did_move = true;
// Write changes
match self.server.commit_batch(batch).await {
Ok(_) => break,
Err(err) => {
retry_after_conflict(err, &mut retries, &arguments.tag, trc::location!())
.await?
}
}
}
// Write changes
self.server
.commit_batch(batch)
.await
.imap_ctx(&arguments.tag, trc::location!())?;
} else {
// Obtain quota for target account
let src_account_id = src_mailbox.id.account_id;
@@ -400,7 +441,7 @@ impl<T: SessionStream> SessionData<T> {
let mut train_batch = BatchBuilder::new();
let mut did_train = false;
train_batch.with_account_id(src_account_id);
for (id, imap_id) in ids {
'next_message: for (id, imap_id) in ids {
match self
.server
.copy_message(
@@ -448,22 +489,27 @@ impl<T: SessionStream> SessionData<T> {
{
copied_ids.push((imap_id.uid, uid));
} else {
let data_ = if let Some(data_) = self
.get_message_data(dest_account_id, existing_id)
.await
.imap_ctx(&arguments.tag, trc::location!())?
{
data_
} else {
continue;
};
let data = data_
.to_unarchived::<MessageData>()
.imap_ctx(&arguments.tag, trc::location!())?;
let mut retries = 0;
loop {
let data_ = if let Some(data_) = self
.get_message_data(dest_account_id, existing_id)
.await
.imap_ctx(&arguments.tag, trc::location!())?
{
data_
} else {
continue 'next_message;
};
let data = data_
.to_unarchived::<MessageData>()
.imap_ctx(&arguments.tag, trc::location!())?;
if let Some(uid) = data.inner.message_uid(dest_mailbox_id) {
copied_ids.push((imap_id.uid, uid));
break;
}
if let Some(uid) = data.inner.message_uid(dest_mailbox_id) {
copied_ids.push((imap_id.uid, uid));
} else {
let mut new_data = data.inner.to_builder();
new_data.add_mailbox(UidMailbox::new_unassigned(dest_mailbox_id));
@@ -504,15 +550,27 @@ impl<T: SessionStream> SessionData<T> {
)
.imap_ctx(&arguments.tag, trc::location!())?;
dest_change_id = self
match self
.server
.commit_batch(batch)
.await
.and_then(|ids| ids.last_change_id(dest_account_id))
.imap_ctx(&arguments.tag, trc::location!())?
.into();
copied_ids.push((imap_id.uid, assigned_uid));
{
Ok(change_id) => {
dest_change_id = change_id.into();
copied_ids.push((imap_id.uid, assigned_uid));
break;
}
Err(err) => {
retry_after_conflict(
err,
&mut retries,
&arguments.tag,
trc::location!(),
)
.await?
}
}
}
}
}
@@ -554,21 +612,33 @@ impl<T: SessionStream> SessionData<T> {
// Untag or delete emails
if !destroy_ids.is_empty() {
let mut batch = BatchBuilder::new();
self.email_untag_or_delete(
src_account_id,
src_mailbox.id.mailbox_id,
&destroy_ids,
&mut batch,
)
.await
.imap_ctx(&arguments.tag, trc::location!())?;
let mut retries = 0;
self.server
.commit_batch(batch)
loop {
let mut batch = BatchBuilder::new();
self.email_untag_or_delete(
src_account_id,
src_mailbox.id.mailbox_id,
&destroy_ids,
&mut batch,
)
.await
.imap_ctx(&arguments.tag, trc::location!())?;
match self.server.commit_batch(batch).await {
Ok(_) => break,
Err(err) => {
retry_after_conflict(
err,
&mut retries,
&arguments.tag,
trc::location!(),
)
.await?
}
}
}
did_move = true;
}
@@ -731,3 +801,24 @@ impl<T: SessionStream> SessionData<T> {
}
}
}
async fn retry_after_conflict(
err: trc::Error,
retries: &mut u32,
tag: &str,
location: &'static str,
) -> trc::Result<()> {
if !err.is_assertion_failure() {
Err(err).imap_ctx(tag, location)
} else if *retries < MAX_MOVE_RETRIES {
*retries += 1;
let backoff = rand::rng().random_range(MOVE_RETRY_BACKOFF_MS);
tokio::time::sleep(Duration::from_millis(backoff)).await;
Ok(())
} else {
Err(trc::ImapEvent::Error
.into_err()
.details("Some messages were modified by another process.")
.id(tag.to_string()))
}
}
+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "jmap_proto"
version = "0.16.23"
version = "0.16.24"
edition = "2024"
[dependencies]
+1 -2
View File
@@ -12,7 +12,6 @@ use crate::{
email::{EmailProperty, EmailValue},
},
request::{
MaybeInvalid,
deserialize::{DeserializeArguments, deserialize_request},
reference::{MaybeIdReference, MaybeResultReference, ResultReference},
},
@@ -33,7 +32,7 @@ pub struct ImportEmailRequest {
#[derive(Debug, Clone, Default)]
pub struct ImportEmail {
pub blob_id: MaybeInvalid<BlobId>,
pub blob_id: MaybeIdReference<BlobId>,
pub mailbox_ids: MaybeResultReference<Vec<MaybeIdReference<Id>>>,
pub keywords: Vec<Keyword>,
pub received_at: Option<UTCDate>,
@@ -388,6 +388,10 @@ impl ResolveReference for ImportEmailRequest {
fn resolve_references(&mut self, response: &Response<'_>) -> trc::Result<()> {
// Resolve email mailbox references
for email in self.emails.values_mut() {
if let MaybeIdReference::Reference(ir) = &email.blob_id {
email.blob_id = MaybeIdReference::Id(response.eval_blob_id_reference(ir)?);
}
match &mut email.mailbox_ids {
MaybeResultReference::Reference(reference) => {
email.mailbox_ids = MaybeResultReference::Value(
@@ -105,6 +105,12 @@ impl<V: Default> Default for MaybeResultReference<V> {
}
}
impl<V: FromStr> Default for MaybeIdReference<V> {
fn default() -> Self {
MaybeIdReference::Invalid(String::new())
}
}
impl<T: Default> MaybeResultReference<T> {
pub fn unwrap(self) -> T {
match self {
+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "jmap"
version = "0.16.23"
version = "0.16.24"
edition = "2024"
[dependencies]
+2 -2
View File
@@ -18,7 +18,7 @@ use jmap_proto::{
error::set::{SetError, SetErrorType},
method::import::{ImportEmailRequest, ImportEmailResponse},
object::email::EmailProperty,
request::MaybeInvalid,
request::reference::MaybeIdReference,
types::state::State,
};
use mail_parser::{HeaderName, MessageParser};
@@ -128,7 +128,7 @@ impl EmailImport for Server {
}
}
let MaybeInvalid::Value(blob_id) = email.blob_id else {
let MaybeIdReference::Id(blob_id) = email.blob_id else {
response.not_created.append(
id,
SetError::invalid_properties()
+5 -2
View File
@@ -854,8 +854,11 @@ impl EmailSet for Server {
new_data.set_mailboxes(
ids.into_expanded_boolean_set()
.filter_map(|id| {
UidMailbox::new_unassigned(
id.try_into_property()?.try_into_id()?.document_id(),
let mailbox_id =
id.try_into_property()?.try_into_id()?.document_id();
UidMailbox::new(
mailbox_id,
data.inner.message_uid(mailbox_id).unwrap_or(0),
)
.into()
})
@@ -641,6 +641,7 @@ fn map_dns_server(dns_server: &DnsServerBootstrap) -> Option<registry::schema::s
DnsServerBootstrap::Ns1(inner) => DnsServer::Ns1(inner.clone()).into(),
DnsServerBootstrap::OracleCloud(inner) => DnsServer::OracleCloud(inner.clone()).into(),
DnsServerBootstrap::Plesk(inner) => DnsServer::Plesk(inner.clone()).into(),
DnsServerBootstrap::PowerDns(inner) => DnsServer::PowerDns(inner.clone()).into(),
DnsServerBootstrap::Safedns(inner) => DnsServer::Safedns(inner.clone()).into(),
DnsServerBootstrap::Scaleway(inner) => DnsServer::Scaleway(inner.clone()).into(),
DnsServerBootstrap::TencentCloud(inner) => DnsServer::TencentCloud(inner.clone()).into(),
+1 -1
View File
@@ -7,7 +7,7 @@ keywords = ["imap", "jmap", "smtp", "email", "mail", "webdav", "server"]
categories = ["email"]
# Upstream offers AGPL-3.0-only OR LicenseRef-SEL; inbuxa takes the AGPL only.
license = "AGPL-3.0-only"
version = "0.16.23"
version = "0.16.24"
edition = "2024"
[[bin]]
+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "managesieve"
version = "0.16.23"
version = "0.16.24"
edition = "2024"
[dependencies]
+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "migration"
version = "0.16.23"
version = "0.16.24"
edition = "2024"
[dependencies]
+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "nlp"
version = "0.16.23"
version = "0.16.24"
edition = "2024"
[dependencies]
+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "pop3"
version = "0.16.23"
version = "0.16.24"
edition = "2024"
[dependencies]
+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "registry"
version = "0.16.23"
version = "0.16.24"
edition = "2024"
[dependencies]
+2
View File
@@ -620,6 +620,7 @@ pub enum DnsServerBootstrapType {
Vultr = 68,
WebSupport = 69,
YandexCloud = 70,
PowerDns = 71,
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Hash)]
@@ -696,6 +697,7 @@ pub enum DnsServerType {
Vultr = 67,
WebSupport = 68,
YandexCloud = 69,
PowerDns = 70,
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Hash)]
+8 -2
View File
@@ -3023,6 +3023,7 @@ impl EnumImpl for DnsServerBootstrapType {
b"Vultr" => DnsServerBootstrapType::Vultr,
b"WebSupport" => DnsServerBootstrapType::WebSupport,
b"YandexCloud" => DnsServerBootstrapType::YandexCloud,
b"PowerDns" => DnsServerBootstrapType::PowerDns,
}
.copied()
}
@@ -3100,6 +3101,7 @@ impl EnumImpl for DnsServerBootstrapType {
DnsServerBootstrapType::Vultr => "Vultr",
DnsServerBootstrapType::WebSupport => "WebSupport",
DnsServerBootstrapType::YandexCloud => "YandexCloud",
DnsServerBootstrapType::PowerDns => "PowerDns",
}
}
@@ -3180,11 +3182,12 @@ impl EnumImpl for DnsServerBootstrapType {
68 => Some(DnsServerBootstrapType::Vultr),
69 => Some(DnsServerBootstrapType::WebSupport),
70 => Some(DnsServerBootstrapType::YandexCloud),
71 => Some(DnsServerBootstrapType::PowerDns),
_ => None,
}
}
const COUNT: usize = 71;
const COUNT: usize = 72;
}
impl serde::Serialize for DnsServerBootstrapType {
@@ -3281,6 +3284,7 @@ impl EnumImpl for DnsServerType {
b"Vultr" => DnsServerType::Vultr,
b"WebSupport" => DnsServerType::WebSupport,
b"YandexCloud" => DnsServerType::YandexCloud,
b"PowerDns" => DnsServerType::PowerDns,
}
.copied()
}
@@ -3357,6 +3361,7 @@ impl EnumImpl for DnsServerType {
DnsServerType::Vultr => "Vultr",
DnsServerType::WebSupport => "WebSupport",
DnsServerType::YandexCloud => "YandexCloud",
DnsServerType::PowerDns => "PowerDns",
}
}
@@ -3436,11 +3441,12 @@ impl EnumImpl for DnsServerType {
67 => Some(DnsServerType::Vultr),
68 => Some(DnsServerType::WebSupport),
69 => Some(DnsServerType::YandexCloud),
70 => Some(DnsServerType::PowerDns),
_ => None,
}
}
const COUNT: usize = 70;
const COUNT: usize = 71;
}
impl serde::Serialize for DnsServerType {
+1
View File
@@ -1042,6 +1042,7 @@ pub enum Property {
SentinelUsername = 914,
Separator = 97,
ServerHostname = 121,
ServerId = 934,
Servers = 308,
ServiceAccountJson = 316,
ServiceName = 913,
@@ -1195,6 +1195,7 @@ impl EnumImpl for Property {
b"sentinelUsername" => Property::SentinelUsername,
b"separator" => Property::Separator,
b"serverHostname" => Property::ServerHostname,
b"serverId" => Property::ServerId,
b"servers" => Property::Servers,
b"serviceAccountJson" => Property::ServiceAccountJson,
b"serviceName" => Property::ServiceName,
@@ -2134,6 +2135,7 @@ impl EnumImpl for Property {
Property::SentinelUsername => "sentinelUsername",
Property::Separator => "separator",
Property::ServerHostname => "serverHostname",
Property::ServerId => "serverId",
Property::Servers => "servers",
Property::ServiceAccountJson => "serviceAccountJson",
Property::ServiceName => "serviceName",
@@ -3077,6 +3079,7 @@ impl EnumImpl for Property {
914 => Some(Property::SentinelUsername),
97 => Some(Property::Separator),
121 => Some(Property::ServerHostname),
934 => Some(Property::ServerId),
308 => Some(Property::Servers),
316 => Some(Property::ServiceAccountJson),
913 => Some(Property::ServiceName),
@@ -3227,7 +3230,7 @@ impl EnumImpl for Property {
}
}
const COUNT: usize = 934;
const COUNT: usize = 935;
}
impl serde::Serialize for Property {
@@ -4493,6 +4496,7 @@ impl ObjectInner {
ObjectInner::DnsServer(DnsServer::Vultr(obj)) => obj.member_tenant_id,
ObjectInner::DnsServer(DnsServer::WebSupport(obj)) => obj.member_tenant_id,
ObjectInner::DnsServer(DnsServer::YandexCloud(obj)) => obj.member_tenant_id,
ObjectInner::DnsServer(DnsServer::PowerDns(obj)) => obj.member_tenant_id,
ObjectInner::Domain(obj) => obj.member_tenant_id,
ObjectInner::MailingList(obj) => obj.member_tenant_id,
ObjectInner::OAuthClient(obj) => obj.member_tenant_id,
@@ -4595,6 +4599,7 @@ impl ObjectInner {
ObjectInner::DnsServer(DnsServer::Vultr(obj)) => obj.member_tenant_id = Some(id),
ObjectInner::DnsServer(DnsServer::WebSupport(obj)) => obj.member_tenant_id = Some(id),
ObjectInner::DnsServer(DnsServer::YandexCloud(obj)) => obj.member_tenant_id = Some(id),
ObjectInner::DnsServer(DnsServer::PowerDns(obj)) => obj.member_tenant_id = Some(id),
ObjectInner::Domain(obj) => obj.member_tenant_id = Some(id),
ObjectInner::MailingList(obj) => obj.member_tenant_id = Some(id),
ObjectInner::OAuthClient(obj) => obj.member_tenant_id = Some(id),
+27
View File
@@ -1457,6 +1457,7 @@ pub enum DnsServer {
Vultr(DnsServerCloud),
WebSupport(DnsServerWebSupport),
YandexCloud(DnsServerYandexCloud),
PowerDns(DnsServerPowerDns),
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
@@ -1672,6 +1673,7 @@ pub enum DnsServerBootstrap {
Vultr(DnsServerCloud),
WebSupport(DnsServerWebSupport),
YandexCloud(DnsServerYandexCloud),
PowerDns(DnsServerPowerDns),
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
@@ -2435,6 +2437,31 @@ pub struct DnsServerPorkbun {
pub propagation_delay: Option<Duration>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(default)]
pub struct DnsServerPowerDns {
#[serde(rename = "apiKey")]
pub api_key: SecretKey,
#[serde(rename = "endpoint")]
pub endpoint: Option<String>,
#[serde(rename = "serverId")]
pub server_id: Option<String>,
#[serde(rename = "description")]
pub description: String,
#[serde(rename = "memberTenantId")]
pub member_tenant_id: Option<Id>,
#[serde(rename = "timeout")]
pub timeout: Duration,
#[serde(rename = "ttl")]
pub ttl: Duration,
#[serde(rename = "pollingInterval")]
pub polling_interval: Duration,
#[serde(rename = "propagationTimeout")]
pub propagation_timeout: Duration,
#[serde(rename = "propagationDelay")]
pub propagation_delay: Option<Duration>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(default)]
pub struct DnsServerRoute53 {
+179
View File
@@ -10559,6 +10559,7 @@ impl ObjectImpl for DnsServer {
DnsServer::Vultr(inner) => inner.validate(errors),
DnsServer::WebSupport(inner) => inner.validate(errors),
DnsServer::YandexCloud(inner) => inner.validate(errors),
DnsServer::PowerDns(inner) => inner.validate(errors),
}
}
@@ -10772,6 +10773,9 @@ impl ObjectImpl for DnsServer {
DnsServer::YandexCloud(object) => {
object.index(i);
}
DnsServer::PowerDns(object) => {
object.index(i);
}
}
}
}
@@ -11064,6 +11068,10 @@ impl Pickle for DnsServer {
69u16.pickle(out);
inner.pickle(out);
}
DnsServer::PowerDns(inner) => {
70u16.pickle(out);
inner.pickle(out);
}
}
}
@@ -11139,6 +11147,7 @@ impl Pickle for DnsServer {
67 => Pickle::unpickle(stream).map(DnsServer::Vultr),
68 => Pickle::unpickle(stream).map(DnsServer::WebSupport),
69 => Pickle::unpickle(stream).map(DnsServer::YandexCloud),
70 => Pickle::unpickle(stream).map(DnsServer::PowerDns),
_ => None,
}
}
@@ -11635,6 +11644,13 @@ impl IntoValue for DnsServer {
.insert_unchecked(Property::Type, JmapValue::Str("YandexCloud".into()));
obj
}
DnsServer::PowerDns(obj) => {
let mut obj = obj.into_value();
obj.as_object_mut()
.unwrap()
.insert_unchecked(Property::Type, JmapValue::Str("PowerDns".into()));
obj
}
}
}
}
@@ -11719,6 +11735,7 @@ impl RegistryJsonPatch for DnsServer {
DnsServerType::Vultr => *self = DnsServer::Vultr(Default::default()),
DnsServerType::WebSupport => *self = DnsServer::WebSupport(Default::default()),
DnsServerType::YandexCloud => *self = DnsServer::YandexCloud(Default::default()),
DnsServerType::PowerDns => *self = DnsServer::PowerDns(Default::default()),
}
}
match self {
@@ -11792,6 +11809,7 @@ impl RegistryJsonPatch for DnsServer {
DnsServer::Vultr(inner) => inner.patch(pointer, value),
DnsServer::WebSupport(inner) => inner.patch(pointer, value),
DnsServer::YandexCloud(inner) => inner.patch(pointer, value),
DnsServer::PowerDns(inner) => inner.patch(pointer, value),
}
}
}
@@ -11869,6 +11887,7 @@ impl DnsServer {
DnsServer::Vultr(_) => DnsServerType::Vultr,
DnsServer::WebSupport(_) => DnsServerType::WebSupport,
DnsServer::YandexCloud(_) => DnsServerType::YandexCloud,
DnsServer::PowerDns(_) => DnsServerType::PowerDns,
}
}
}
@@ -12704,6 +12723,7 @@ impl DnsServerBootstrap {
DnsServerBootstrap::Vultr(inner) => inner.validate(errors),
DnsServerBootstrap::WebSupport(inner) => inner.validate(errors),
DnsServerBootstrap::YandexCloud(inner) => inner.validate(errors),
DnsServerBootstrap::PowerDns(inner) => inner.validate(errors),
}
}
}
@@ -12999,6 +13019,10 @@ impl Pickle for DnsServerBootstrap {
70u16.pickle(out);
inner.pickle(out);
}
DnsServerBootstrap::PowerDns(inner) => {
71u16.pickle(out);
inner.pickle(out);
}
}
}
@@ -13075,6 +13099,7 @@ impl Pickle for DnsServerBootstrap {
68 => Pickle::unpickle(stream).map(DnsServerBootstrap::Vultr),
69 => Pickle::unpickle(stream).map(DnsServerBootstrap::WebSupport),
70 => Pickle::unpickle(stream).map(DnsServerBootstrap::YandexCloud),
71 => Pickle::unpickle(stream).map(DnsServerBootstrap::PowerDns),
_ => None,
}
}
@@ -13576,6 +13601,13 @@ impl IntoValue for DnsServerBootstrap {
.insert_unchecked(Property::Type, JmapValue::Str("YandexCloud".into()));
obj
}
DnsServerBootstrap::PowerDns(obj) => {
let mut obj = obj.into_value();
obj.as_object_mut()
.unwrap()
.insert_unchecked(Property::Type, JmapValue::Str("PowerDns".into()));
obj
}
}
}
}
@@ -13793,6 +13825,9 @@ impl RegistryJsonPatch for DnsServerBootstrap {
DnsServerBootstrapType::YandexCloud => {
*self = DnsServerBootstrap::YandexCloud(Default::default())
}
DnsServerBootstrapType::PowerDns => {
*self = DnsServerBootstrap::PowerDns(Default::default())
}
}
}
match self {
@@ -13867,6 +13902,7 @@ impl RegistryJsonPatch for DnsServerBootstrap {
DnsServerBootstrap::Vultr(inner) => inner.patch(pointer, value),
DnsServerBootstrap::WebSupport(inner) => inner.patch(pointer, value),
DnsServerBootstrap::YandexCloud(inner) => inner.patch(pointer, value),
DnsServerBootstrap::PowerDns(inner) => inner.patch(pointer, value),
}
}
}
@@ -13945,6 +13981,7 @@ impl DnsServerBootstrap {
DnsServerBootstrap::Vultr(_) => DnsServerBootstrapType::Vultr,
DnsServerBootstrap::WebSupport(_) => DnsServerBootstrapType::WebSupport,
DnsServerBootstrap::YandexCloud(_) => DnsServerBootstrapType::YandexCloud,
DnsServerBootstrap::PowerDns(_) => DnsServerBootstrapType::PowerDns,
}
}
}
@@ -18228,6 +18265,148 @@ impl RegistryJsonPropertyPatch for DnsServerPorkbun {
}
}
impl DnsServerPowerDns {
fn validate(&self, errors: &mut Vec<ValidationError>) -> bool {
let neb = errors.len();
let value = &self.api_key;
value.validate(errors);
if let Some(value) = &self.endpoint {
if value.is_empty() {
errors.push(ValidationError::required(Property::Endpoint));
}
}
if let Some(value) = &self.server_id {
if value.is_empty() {
errors.push(ValidationError::required(Property::ServerId));
}
}
let value = &self.description;
if value.is_empty() {
errors.push(ValidationError::required(Property::Description));
}
if let Some(value) = &self.member_tenant_id {
if !value.is_valid() {
errors.push(ValidationError::required(Property::MemberTenantId));
}
}
errors.len() == neb
}
fn index<'x>(&'x self, i: &mut IndexBuilder<'x>) {
i.foreign_key(ObjectType::Tenant, self.member_tenant_id, None);
if let Some(value) = &self.member_tenant_id {
i.search(Property::MemberTenantId, value);
}
}
}
impl Pickle for DnsServerPowerDns {
fn pickle(&self, out: &mut Vec<u8>) {
self.api_key.pickle(out);
self.endpoint.pickle(out);
self.server_id.pickle(out);
self.description.pickle(out);
self.member_tenant_id.pickle(out);
self.timeout.pickle(out);
self.ttl.pickle(out);
self.polling_interval.pickle(out);
self.propagation_timeout.pickle(out);
self.propagation_delay.pickle(out);
}
fn unpickle(stream: &mut crate::pickle::PickledStream<'_>) -> Option<Self> {
let mut this = Self::default();
this.api_key = Pickle::unpickle(stream)?;
this.endpoint = Pickle::unpickle(stream)?;
this.server_id = Pickle::unpickle(stream)?;
this.description = Pickle::unpickle(stream)?;
this.member_tenant_id = Pickle::unpickle(stream)?;
this.timeout = Pickle::unpickle(stream)?;
this.ttl = Pickle::unpickle(stream)?;
this.polling_interval = Pickle::unpickle(stream)?;
this.propagation_timeout = Pickle::unpickle(stream)?;
this.propagation_delay = Pickle::unpickle(stream)?;
Some(this)
}
}
impl Default for DnsServerPowerDns {
fn default() -> Self {
Self {
api_key: Default::default(),
endpoint: Default::default(),
server_id: Default::default(),
description: Default::default(),
member_tenant_id: Default::default(),
timeout: Duration::from_millis(30000),
ttl: Duration::from_millis(300000),
polling_interval: Duration::from_millis(15000),
propagation_timeout: Duration::from_millis(60000),
propagation_delay: Default::default(),
}
}
}
impl IntoValue for DnsServerPowerDns {
fn into_value(self) -> JmapValue<'static> {
let mut map = jmap_tools::Map::with_capacity(12);
map.insert_unchecked(Property::ApiKey, self.api_key.into_value());
map.insert_unchecked(Property::Endpoint, self.endpoint.into_value());
map.insert_unchecked(Property::ServerId, self.server_id.into_value());
map.insert_unchecked(Property::Description, self.description.into_value());
map.insert_unchecked(Property::MemberTenantId, self.member_tenant_id.into_value());
map.insert_unchecked(Property::Timeout, self.timeout.into_value());
map.insert_unchecked(Property::Ttl, self.ttl.into_value());
map.insert_unchecked(
Property::PollingInterval,
self.polling_interval.into_value(),
);
map.insert_unchecked(
Property::PropagationTimeout,
self.propagation_timeout.into_value(),
);
map.insert_unchecked(
Property::PropagationDelay,
self.propagation_delay.into_value(),
);
JmapValue::Object(map)
}
}
impl RegistryJsonPropertyPatch for DnsServerPowerDns {
fn patch_property<'x>(
&mut self,
mut pointer: JsonPointerPatch<'_>,
value: JmapValue<'x>,
) -> PatchResult<'x> {
match pointer.next_property() {
Some(Property::ApiKey) => self.api_key.patch(pointer, value),
Some(Property::Endpoint) => self
.endpoint
.patch(pointer.with_validators(&[StringValidator::Trim]), value),
Some(Property::ServerId) => self
.server_id
.patch(pointer.with_validators(&[StringValidator::Trim]), value),
Some(Property::Description) => self
.description
.patch(pointer.with_validators(&[StringValidator::Trim]), value),
Some(Property::MemberTenantId) => self
.member_tenant_id
.patch(pointer.assert_can_set_tenant()?, value),
Some(Property::Timeout) => self.timeout.patch(pointer, value),
Some(Property::Ttl) => self.ttl.patch(pointer, value),
Some(Property::PollingInterval) => self.polling_interval.patch(pointer, value),
Some(Property::PropagationTimeout) => self.propagation_timeout.patch(pointer, value),
Some(Property::PropagationDelay) => self.propagation_delay.patch(pointer, value),
Some(Property::Type) => Ok(MaybeUnpatched::Unpatched {
property: Property::Type,
value,
}),
_ => Err(PatchError::new(pointer, "Invalid property")),
}
}
}
impl DnsServerRoute53 {
fn validate(&self, errors: &mut Vec<ValidationError>) -> bool {
let neb = errors.len();
+1
View File
@@ -14,6 +14,7 @@ pub mod dkim;
pub mod http;
pub mod report;
pub mod secret;
pub mod spam;
pub mod task;
impl Roles {
+59
View File
@@ -0,0 +1,59 @@
/*
* SPDX-FileCopyrightText: 2020 Stalwart Labs LLC <hello@stalw.art>
*
* SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-SEL
*/
use crate::schema::prelude::{SpamDnsblServer, SpamRule};
impl SpamRule {
pub fn enable(&self) -> bool {
match self {
SpamRule::Any(rule) => rule.enable,
SpamRule::Url(rule) => rule.enable,
SpamRule::Domain(rule) => rule.enable,
SpamRule::Email(rule) => rule.enable,
SpamRule::Ip(rule) => rule.enable,
SpamRule::Header(rule) => rule.enable,
SpamRule::Body(rule) => rule.enable,
}
}
pub fn set_enable(&mut self, enable: bool) {
match self {
SpamRule::Any(rule) => rule.enable = enable,
SpamRule::Url(rule) => rule.enable = enable,
SpamRule::Domain(rule) => rule.enable = enable,
SpamRule::Email(rule) => rule.enable = enable,
SpamRule::Ip(rule) => rule.enable = enable,
SpamRule::Header(rule) => rule.enable = enable,
SpamRule::Body(rule) => rule.enable = enable,
}
}
}
impl SpamDnsblServer {
pub fn enable(&self) -> bool {
match self {
SpamDnsblServer::Any(server) => server.enable,
SpamDnsblServer::Url(server) => server.enable,
SpamDnsblServer::Domain(server) => server.enable,
SpamDnsblServer::Email(server) => server.enable,
SpamDnsblServer::Ip(server) => server.enable,
SpamDnsblServer::Header(server) => server.enable,
SpamDnsblServer::Body(server) => server.enable,
}
}
pub fn set_enable(&mut self, enable: bool) {
match self {
SpamDnsblServer::Any(server) => server.enable = enable,
SpamDnsblServer::Url(server) => server.enable = enable,
SpamDnsblServer::Domain(server) => server.enable = enable,
SpamDnsblServer::Email(server) => server.enable = enable,
SpamDnsblServer::Ip(server) => server.enable = enable,
SpamDnsblServer::Header(server) => server.enable = enable,
SpamDnsblServer::Body(server) => server.enable = enable,
}
}
}
+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "scim-proto"
version = "0.16.23"
version = "0.16.24"
edition = "2024"
[dependencies]
+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "scim"
version = "0.16.23"
version = "0.16.24"
edition = "2024"
[dependencies]
+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "services"
version = "0.16.23"
version = "0.16.24"
edition = "2024"
[dependencies]
+71 -57
View File
@@ -9,7 +9,7 @@ use super::{
ece::{ECE_WEBPUSH_MAX_PLAINTEXT_SIZE, WEBPUSH_MAX_BODY_SIZE, ece_encrypt},
email_push::build_email_push_object,
};
use crate::state_manager::PushRegistration;
use crate::state_manager::{PushBatch, PushRegistration};
use calcard::jscalendar::JSCalendarDateTime;
use common::{Server, ipc::PushNotification, network::webpush::Vapid};
use email::push::{PushSubscription, Urgency};
@@ -29,7 +29,7 @@ use std::time::{Duration, Instant};
use store::write::now;
use tokio::sync::mpsc;
use trc::PushSubscriptionEvent;
use types::{id::Id, type_state::DataType};
use types::id::Id;
use utils::map::vec_map::VecMap;
const MAX_ERROR_RESPONSE_LEN: usize = 1024;
@@ -48,34 +48,29 @@ impl PushRegistration {
pub fn send(
&mut self,
id: Id,
push_client: &Client,
push_tx: mpsc::Sender<Event>,
push_timeout: Duration,
server: Server,
) {
let subscription = self.server.clone();
let push_client = self.client.clone();
let notifications = std::mem::take(&mut self.notifications);
let push_client = push_client.clone();
let batch = std::mem::take(&mut self.pending);
self.in_flight = true;
self.last_request = Instant::now();
tokio::spawn(async move {
let mut changed: VecMap<Id, VecMap<DataType, State>> = VecMap::new();
let vapid = server.core.jmap.vapid.as_ref();
let mut email_pushes: VecMap<Id, EmailPushObject> = VecMap::new();
let mut failed_state_change = false;
let mut failed = PushBatch::default();
let mut failed_email_pushes = Vec::new();
let mut failed_calendar_alerts = Vec::new();
for notification in &notifications {
for notification in &batch.notifications {
match notification {
PushNotification::StateChange(state_change) => {
for type_state in state_change.types {
changed
.get_mut_or_insert(state_change.account_id.into())
.set(type_state, State::Exact(state_change.change_id));
}
}
PushNotification::StateChange(_) => {}
PushNotification::CalendarAlert(calendar_alert) => {
let payload = PushObject::CalendarAlert {
account_id: calendar_alert.account_id.into(),
@@ -86,12 +81,12 @@ impl PushRegistration {
}),
alert_id: calendar_alert.alert_id.clone(),
};
if !http_request(
if !post_object(
&push_client,
&subscription,
serde_json::to_string(&payload).unwrap().into_bytes(),
&payload,
push_timeout,
server.core.jmap.vapid.as_ref(),
vapid,
Urgency::Normal,
)
.await
@@ -155,18 +150,23 @@ impl PushRegistration {
}
}
if !changed.is_empty() {
failed_state_change = !http_request(
if !batch.state_changes.is_empty() {
let payload = PushObject::StateChange {
changed: batch.state_changes,
};
if !post_object(
&push_client,
&subscription,
serde_json::to_string(&PushObject::StateChange { changed })
.unwrap()
.into_bytes(),
&payload,
push_timeout,
server.core.jmap.vapid.as_ref(),
vapid,
Urgency::Normal,
)
.await;
.await
&& let PushObject::StateChange { changed } = payload
{
failed.state_changes = changed;
}
}
for (account_id, email_push) in email_pushes {
@@ -180,12 +180,12 @@ impl PushRegistration {
state: email_push.change_id.map(State::Exact),
};
if !http_request(
if !post_object(
&push_client,
&subscription,
serde_json::to_string(&payload).unwrap().into_bytes(),
&payload,
push_timeout,
server.core.jmap.vapid.as_ref(),
vapid,
email_push.urgency,
)
.await
@@ -194,44 +194,26 @@ impl PushRegistration {
}
}
let result = if !failed_state_change
let result = if failed.state_changes.is_empty()
&& failed_email_pushes.is_empty()
&& failed_calendar_alerts.is_empty()
{
Event::DeliverySuccess { id }
} else {
let mut failed_notifications = Vec::with_capacity(
failed_state_change as usize
+ failed_email_pushes.len()
+ failed_calendar_alerts.len(),
);
for notification in notifications {
match &notification {
PushNotification::StateChange(_) => {
if failed_state_change {
failed_notifications.push(notification);
}
}
failed.notifications = batch
.notifications
.into_iter()
.filter(|notification| match notification {
PushNotification::StateChange(_) => false,
PushNotification::EmailPush(email_push) => {
if failed_email_pushes.contains(&email_push.account_id) {
failed_notifications.push(notification);
}
failed_email_pushes.contains(&email_push.account_id)
}
PushNotification::CalendarAlert(calendar_alert) => {
if failed_calendar_alerts
.contains(&(calendar_alert.account_id, calendar_alert.event_id))
{
failed_notifications.push(notification);
}
}
}
}
PushNotification::CalendarAlert(calendar_alert) => failed_calendar_alerts
.contains(&(calendar_alert.account_id, calendar_alert.event_id)),
})
.collect();
Event::DeliveryFailure {
id,
notifications: failed_notifications,
}
Event::DeliveryFailure { id, failed }
};
push_tx.send(result).await.ok();
@@ -239,6 +221,38 @@ impl PushRegistration {
}
}
async fn post_object(
push_client: &Client,
subscription: &PushSubscription,
object: &PushObject,
push_timeout: Duration,
vapid: Option<&Vapid>,
urgency: Urgency,
) -> bool {
match serde_json::to_vec(object) {
Ok(body) => {
http_request(
push_client,
subscription,
body,
push_timeout,
vapid,
urgency,
)
.await
}
Err(err) => {
trc::event!(
PushSubscription(PushSubscriptionEvent::Error),
Details = "Failed to serialize push object",
Url = subscription.url.to_string(),
Reason = err.to_string()
);
true
}
}
}
pub(crate) fn build_push_client() -> Client {
utils::http::http_client_builder(cfg!(feature = "test_mode"))
.redirect(Policy::custom(|attempt| match attempt.previous().last() {
+26 -22
View File
@@ -14,7 +14,7 @@ use common::{
};
use std::{sync::Arc, time::Instant};
use store::ahash::AHashMap;
use tokio::sync::mpsc;
use tokio::sync::mpsc::{self, error::TrySendError};
use trc::ServerEvent;
#[derive(Default)]
@@ -99,7 +99,7 @@ pub fn spawn_push_router(inner: Arc<Inner>, mut change_rx: mpsc::Receiver<PushEv
} => {
// Publish event to cluster
if broadcast
&& let Some(broadcast_tx) = &inner.ipc.broadcast_tx.clone()
&& let Some(broadcast_tx) = inner.ipc.broadcast_tx.as_ref()
&& broadcast_tx
.send(BroadcastEvent::PushNotification(notification.clone()))
.await
@@ -117,26 +117,30 @@ pub fn spawn_push_router(inner: Arc<Inner>, mut change_rx: mpsc::Receiver<PushEv
for subscriber in &subscribers.ipc {
if let Some(notification) = notification.filter_types(&subscriber.types)
{
if subscriber.is_valid() {
let subscriber_tx = subscriber.tx.clone();
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!()
);
}
});
} else {
purge_needed = true;
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;
}
}
}
}
@@ -159,7 +163,7 @@ pub fn spawn_push_router(inner: Arc<Inner>, mut change_rx: mpsc::Receiver<PushEv
} => {
// Publish event to cluster
if broadcast
&& let Some(broadcast_tx) = &inner.ipc.broadcast_tx.clone()
&& let Some(broadcast_tx) = inner.ipc.broadcast_tx.as_ref()
&& broadcast_tx
.send(BroadcastEvent::PushServerUpdate(account_id))
.await
+140 -17
View File
@@ -14,14 +14,14 @@ pub mod push;
use common::ipc::PushNotification;
use email::push::PushSubscription;
use reqwest::Client;
use jmap_proto::types::state::State;
use std::{
sync::Arc,
time::{Duration, Instant},
};
use tokio::sync::mpsc;
use types::{id::Id, type_state::DataType};
use utils::map::bitmap::Bitmap;
use utils::map::{bitmap::Bitmap, vec_map::VecMap};
const PURGE_EVERY: Duration = Duration::from_secs(3600);
const SEND_TIMEOUT: Duration = Duration::from_millis(500);
@@ -40,26 +40,22 @@ pub struct PushRegistration {
member_account_ids: Vec<u32>,
num_attempts: u32,
last_request: Instant,
notifications: Vec<PushNotification>,
pending: PushBatch,
in_flight: bool,
client: Client,
}
#[derive(Debug, Default)]
pub struct PushBatch {
state_changes: VecMap<Id, VecMap<DataType, State>>,
notifications: Vec<PushNotification>,
}
#[derive(Debug)]
pub enum Event {
Push {
notification: PushNotification,
},
Update {
account_id: u32,
},
DeliverySuccess {
id: Id,
},
DeliveryFailure {
id: Id,
notifications: Vec<PushNotification>,
},
Push { notification: PushNotification },
Update { account_id: u32 },
DeliverySuccess { id: Id },
DeliveryFailure { id: Id, failed: PushBatch },
Reset,
}
@@ -68,3 +64,130 @@ impl IpcSubscriber {
!self.tx.is_closed()
}
}
impl PushBatch {
pub fn push(&mut self, notification: PushNotification) {
match notification {
PushNotification::StateChange(state_change) => {
if !state_change.types.is_empty() {
let states = self
.state_changes
.get_mut_or_insert(Id::from(state_change.account_id));
for data_type in state_change.types {
merge_state(states, data_type, state_change.change_id);
}
}
}
notification => self.notifications.push(notification),
}
}
pub fn merge_failed(&mut self, failed: PushBatch) {
for (account_id, failed_states) in failed.state_changes {
let states = self.state_changes.get_mut_or_insert(account_id);
for (data_type, state) in failed_states {
if let State::Exact(change_id) = state {
merge_state(states, data_type, change_id);
}
}
}
if !failed.notifications.is_empty() {
let mut notifications = failed.notifications;
notifications.append(&mut self.notifications);
self.notifications = notifications;
}
}
pub fn is_empty(&self) -> bool {
self.state_changes.is_empty() && self.notifications.is_empty()
}
pub fn clear(&mut self) {
self.state_changes.clear();
self.notifications.clear();
}
}
fn merge_state(states: &mut VecMap<DataType, State>, data_type: DataType, change_id: u64) {
match states.get_mut(&data_type) {
Some(State::Exact(current)) if *current >= change_id => {}
Some(state) => *state = State::Exact(change_id),
None => states.append(data_type, State::Exact(change_id)),
}
}
#[cfg(test)]
mod tests {
use super::PushBatch;
use common::ipc::{EmailPush, PushNotification};
use jmap_proto::types::state::State;
use types::{
id::Id,
type_state::{DataType, StateChange},
};
use utils::map::bitmap::Bitmap;
fn state_change<const N: usize>(change_id: u64, types: [DataType; N]) -> PushNotification {
PushNotification::StateChange(StateChange {
account_id: 1,
change_id,
types: Bitmap::from_iter(types),
})
}
fn email_push(email_id: u32) -> PushNotification {
PushNotification::EmailPush(EmailPush {
account_id: 1,
email_id,
change_id: email_id.into(),
})
}
#[test]
fn batch_keeps_newest_state_per_type() {
let mut batch = PushBatch::default();
batch.push(state_change(5, [DataType::Email]));
batch.push(state_change(7, [DataType::Email, DataType::Mailbox]));
batch.push(state_change(6, [DataType::Mailbox]));
let mut failed = PushBatch::default();
failed.push(state_change(
4,
[DataType::Email, DataType::Mailbox, DataType::Thread],
));
batch.merge_failed(failed);
let states = batch.state_changes.get(&Id::from(1u32)).unwrap();
assert_eq!(states.len(), 3);
assert_eq!(states.get(&DataType::Email), Some(&State::Exact(7)));
assert_eq!(states.get(&DataType::Mailbox), Some(&State::Exact(7)));
assert_eq!(states.get(&DataType::Thread), Some(&State::Exact(4)));
assert!(!batch.is_empty());
batch.clear();
assert!(batch.is_empty());
}
#[test]
fn failed_notifications_are_retried_first() {
let mut batch = PushBatch::default();
batch.push(email_push(3));
let mut failed = PushBatch::default();
failed.push(email_push(1));
failed.push(email_push(2));
batch.merge_failed(failed);
let email_ids = batch
.notifications
.iter()
.filter_map(|notification| match notification {
PushNotification::EmailPush(email_push) => Some(email_push.email_id),
_ => None,
})
.collect::<Vec<_>>();
assert_eq!(email_ids, [1, 2, 3]);
assert!(batch.state_changes.is_empty());
}
}
+157 -89
View File
@@ -8,13 +8,14 @@ use super::{
Event,
http::{build_push_client, http_request},
};
use crate::state_manager::PushRegistration;
use crate::state_manager::{PushBatch, PushRegistration};
use common::{
BuildServer, IPC_CHANNEL_BUFFER, Inner, LONG_1Y_SLUMBER, Server,
auth::BuildAccessToken,
ipc::{PushEvent, PushNotification},
};
use email::push::{PushSubscription, PushSubscriptions, Urgency};
use reqwest::Client;
use std::{
collections::hash_map::Entry,
sync::Arc,
@@ -36,7 +37,10 @@ pub fn spawn_push_manager(inner: Arc<Inner>) -> mpsc::Sender<Event> {
tokio::spawn(async move {
let mut push_servers: AHashMap<Id, PushRegistration> = AHashMap::default();
let mut account_push_ids: AHashMap<u32, AHashSet<Id>> = AHashMap::default();
let mut last_verify: AHashMap<u32, Instant> = AHashMap::default();
let mut last_verify: AHashMap<u32, (Instant, u32)> = AHashMap::default();
let mut pending_verify: AHashMap<u32, (Instant, Arc<PushSubscription>)> =
AHashMap::default();
let mut next_verify: Option<Instant> = None;
let mut last_retry = Instant::now();
let mut retry_timeout = LONG_1Y_SLUMBER;
let mut retry_ids = AHashSet::default();
@@ -90,10 +94,9 @@ pub fn spawn_push_manager(inner: Arc<Inner>) -> mpsc::Sender<Event> {
last_request: Instant::now()
- (server.core.jmap.push_throttle
+ Duration::from_millis(1)),
notifications: Vec::new(),
pending: PushBatch::default(),
server: subscription.clone(),
in_flight: false,
client: push_client.clone(),
},
);
}
@@ -129,8 +132,39 @@ pub fn spawn_push_manager(inner: Arc<Inner>) -> mpsc::Sender<Event> {
}
loop {
if let Some(verify_due) = next_verify {
let current_instant = Instant::now();
if verify_due <= current_instant {
let server = inner.build_server();
let push_timeout = server.core.jmap.push_timeout;
let current_time = now();
next_verify = None;
pending_verify.retain(|account_id, (verify_due, subscription)| {
if *verify_due > current_instant {
next_verify =
Some(next_verify.map_or(*verify_due, |next| next.min(*verify_due)));
true
} else {
if subscription.expires > current_time {
last_verify.insert(*account_id, (current_instant, subscription.id));
send_verification(
&push_client,
subscription.clone(),
&server,
push_timeout,
);
}
false
}
});
}
}
// Wait for the next event or timeout
let event_or_timeout = tokio::time::timeout(retry_timeout, push_rx.recv()).await;
let wait_timeout = next_verify.map_or(retry_timeout, |verify_due| {
retry_timeout.min(verify_due.saturating_duration_since(Instant::now()))
});
let event_or_timeout = tokio::time::timeout(wait_timeout, push_rx.recv()).await;
// Load settings
let server = inner.build_server();
@@ -193,10 +227,9 @@ pub fn spawn_push_manager(inner: Arc<Inner>) -> mpsc::Sender<Event> {
num_attempts: 0,
last_request: Instant::now()
- (push_throttle + Duration::from_millis(1)),
notifications: Vec::new(),
pending: PushBatch::default(),
server: subscription.clone(),
in_flight: false,
client: push_client.clone(),
});
}
}
@@ -213,52 +246,56 @@ pub fn spawn_push_manager(inner: Arc<Inner>) -> mpsc::Sender<Event> {
#[cfg(feature = "test_mode")]
if subscription.url.contains("skip_checks") {
last_verify.insert(
account_id,
current_time - (push_verify_timeout + Duration::from_millis(1)),
);
last_verify.remove(&account_id);
}
if last_verify
match last_verify
.get(&account_id)
.map(|last_verify| {
current_time - *last_verify > push_verify_timeout
.map(|(verified_at, verified_id)| {
(*verified_at + push_verify_timeout, *verified_id)
})
.unwrap_or(true)
.filter(|(verify_due, _)| *verify_due >= current_time)
{
let core = server.core.clone();
let push_client = push_client.clone();
tokio::spawn(async move {
http_request(
None => {
last_verify.retain(|_, (verified_at, _)| {
current_time.duration_since(*verified_at)
<= push_verify_timeout
});
last_verify.insert(account_id, (current_time, subscription.id));
pending_verify.remove(&account_id);
send_verification(
&push_client,
&subscription,
format!(
concat!(
"{{\"@type\":\"PushVerification\",",
"\"pushSubscriptionId\":\"{}\",",
"\"verificationCode\":\"{}\"}}"
),
Id::from(subscription.id),
subscription.verification_code
)
.into_bytes(),
subscription,
&server,
push_timeout,
core.jmap.vapid.as_ref(),
Urgency::Normal,
)
.await;
});
last_verify.insert(account_id, current_time);
} else {
trc::event!(
PushSubscription(PushSubscriptionEvent::Error),
Details = "Failed to verify push subscription",
Url = subscription.url.clone(),
AccountId = account_id,
Reason = "Too many requests"
);
);
}
Some((_, verified_id)) if verified_id == subscription.id => {
trc::event!(
PushSubscription(PushSubscriptionEvent::Error),
Details = "Failed to verify push subscription",
Url = subscription.url.clone(),
AccountId = account_id,
Reason = "Too many requests"
);
pending_verify.remove(&account_id);
}
Some((verify_due, _)) => {
trc::event!(
PushSubscription(PushSubscriptionEvent::Error),
Details = "Push subscription verification deferred",
Url = subscription.url.clone(),
AccountId = account_id,
Reason = "Too many requests"
);
next_verify = Some(
next_verify.map_or(verify_due, |next| next.min(verify_due)),
);
pending_verify.insert(account_id, (verify_due, subscription));
}
}
} else {
pending_verify.remove(&account_id);
}
// Update subscriptions
@@ -342,7 +379,7 @@ pub fn spawn_push_manager(inner: Arc<Inner>) -> mpsc::Sender<Event> {
);
}
subscription.notifications.push(notification);
subscription.pending.push(notification);
let last_request = subscription.last_request.elapsed();
if !subscription.in_flight
@@ -354,6 +391,7 @@ pub fn spawn_push_manager(inner: Arc<Inner>) -> mpsc::Sender<Event> {
{
subscription.send(
*id,
&push_client,
push_tx.clone(),
push_timeout,
server.clone(),
@@ -402,19 +440,25 @@ pub fn spawn_push_manager(inner: Arc<Inner>) -> mpsc::Sender<Event> {
Event::Reset => {
push_servers.clear();
account_push_ids.clear();
pending_verify.clear();
next_verify = None;
}
Event::DeliverySuccess { id } => {
if let Some(subscription) = push_servers.get_mut(&id) {
subscription.num_attempts = 0;
subscription.in_flight = false;
retry_ids.remove(&id);
if subscription.pending.is_empty() {
retry_ids.remove(&id);
} else {
retry_ids.insert(id);
}
}
}
Event::DeliveryFailure { id, notifications } => {
Event::DeliveryFailure { id, failed } => {
if let Some(subscription) = push_servers.get_mut(&id) {
subscription.last_request = Instant::now();
subscription.num_attempts += 1;
subscription.notifications.extend(notifications);
subscription.pending.merge_failed(failed);
subscription.in_flight = false;
retry_ids.insert(id);
}
@@ -430,52 +474,46 @@ pub fn spawn_push_manager(inner: Arc<Inner>) -> mpsc::Sender<Event> {
let last_retry_elapsed = last_retry.elapsed();
if last_retry_elapsed >= push_retry_interval {
let mut remove_ids = Vec::with_capacity(retry_ids.len());
retry_ids.retain(|retry_id| {
let Some(subscription) = push_servers.get_mut(retry_id) else {
return false;
};
let last_request = subscription.last_request.elapsed();
let is_due = !subscription.in_flight
&& ((subscription.num_attempts == 0 && last_request >= push_throttle)
|| (subscription.num_attempts > 0
&& last_request >= push_attempt_interval));
if !is_due {
return true;
}
for retry_id in &retry_ids {
if let Some(subscription) = push_servers.get_mut(retry_id) {
let last_request = subscription.last_request.elapsed();
if !subscription.in_flight
&& ((subscription.num_attempts == 0
&& last_request >= push_throttle)
|| (subscription.num_attempts > 0
&& last_request >= push_attempt_interval))
{
if subscription.num_attempts < push_attempts_max {
subscription.send(
*retry_id,
push_tx.clone(),
push_timeout,
server.clone(),
);
} else {
trc::event!(
PushSubscription(PushSubscriptionEvent::Error),
Details = "Failed to deliver push subscription",
Url = subscription.server.url.clone(),
Reason = "Too many failed attempts"
);
subscription.notifications.clear();
subscription.num_attempts = 0;
}
remove_ids.push(*retry_id);
}
if subscription.num_attempts < push_attempts_max {
subscription.send(
*retry_id,
&push_client,
push_tx.clone(),
push_timeout,
server.clone(),
);
} else {
remove_ids.push(*retry_id);
}
}
trc::event!(
PushSubscription(PushSubscriptionEvent::Error),
Details = "Failed to deliver push subscription",
Url = subscription.server.url.clone(),
Reason = "Too many failed attempts"
);
if remove_ids.len() < retry_ids.len() {
for remove_id in remove_ids {
retry_ids.remove(&remove_id);
subscription.pending.clear();
subscription.num_attempts = 0;
}
false
});
if retry_ids.is_empty() {
LONG_1Y_SLUMBER
} else {
last_retry = Instant::now();
push_retry_interval
} else {
retry_ids.clear();
LONG_1Y_SLUMBER
}
} else {
push_retry_interval - last_retry_elapsed
@@ -489,6 +527,36 @@ pub fn spawn_push_manager(inner: Arc<Inner>) -> mpsc::Sender<Event> {
push_tx_
}
fn send_verification(
push_client: &Client,
subscription: Arc<PushSubscription>,
server: &Server,
push_timeout: Duration,
) {
let core = server.core.clone();
let push_client = push_client.clone();
tokio::spawn(async move {
http_request(
&push_client,
&subscription,
format!(
concat!(
"{{\"@type\":\"PushVerification\",",
"\"pushSubscriptionId\":\"{}\",",
"\"verificationCode\":\"{}\"}}"
),
Id::from(subscription.id),
subscription.verification_code
)
.into_bytes(),
push_timeout,
core.jmap.vapid.as_ref(),
Urgency::Normal,
)
.await;
});
}
async fn load_push_subscriptions(
server: &Server,
account_id: u32,
@@ -15,13 +15,13 @@ use common::{
use registry::{
schema::{
enums::TaskSpamFilterMaintenanceType,
prelude::ObjectType,
prelude::{Object, ObjectInner, ObjectType},
structs::{
HttpLookup, MemoryLookupKey, SpamDnsblServer, SpamFileExtension, SpamRule, SpamTag,
TaskSpamFilterMaintenance,
},
},
types::EnumImpl,
types::{EnumImpl, ObjectImpl},
};
use spam_filter::modules::classifier::SpamClassifier;
use std::time::{Duration, Instant};
@@ -100,13 +100,90 @@ struct Rules {
file_exts: Vec<SpamFileExtension>,
}
#[derive(Default)]
trait UpstreamObject: ObjectImpl + PartialEq + From<Object> + Into<ObjectInner> {
fn replacement_for(self, _local: &Self) -> Option<Self> {
Some(self)
}
// inbuxa: the object as its fingerprint sees it. Switching a rule on or
// off isn't an edit, so `enable` is left out.
fn without_enable(self) -> Self {
self
}
}
impl UpstreamObject for SpamRule {
fn replacement_for(mut self, local: &Self) -> Option<Self> {
self.set_enable(local.enable());
Some(self)
}
fn without_enable(mut self) -> Self {
self.set_enable(true);
self
}
}
impl UpstreamObject for SpamDnsblServer {
fn replacement_for(mut self, local: &Self) -> Option<Self> {
self.set_enable(local.enable());
Some(self)
}
fn without_enable(mut self) -> Self {
self.set_enable(true);
self
}
}
impl UpstreamObject for HttpLookup {
fn replacement_for(mut self, local: &Self) -> Option<Self> {
self.enable = local.enable;
Some(self)
}
fn without_enable(mut self) -> Self {
self.enable = true;
self
}
}
impl UpstreamObject for SpamTag {
fn replacement_for(self, _local: &Self) -> Option<Self> {
None
}
}
impl UpstreamObject for MemoryLookupKey {}
impl UpstreamObject for SpamFileExtension {}
struct RuleUpdateResult {
success: usize,
already_exists: usize,
object_type: ObjectType,
added: usize,
updated: usize,
unchanged: usize,
kept: usize,
failed: usize,
}
impl RuleUpdateResult {
fn new(object_type: ObjectType) -> Self {
RuleUpdateResult {
object_type,
added: 0,
updated: 0,
unchanged: 0,
kept: 0,
failed: 0,
}
}
fn has_changes(&self) -> bool {
self.added + self.updated > 0
}
}
async fn update_spam_rules(server: &Server) -> trc::Result<TaskResult> {
let started = Instant::now();
let bundled = server.core.spam.spam_rules_url.is_none();
@@ -121,171 +198,67 @@ async fn update_spam_rules(server: &Server) -> trc::Result<TaskResult> {
}
};
let registry = server.registry();
let mut stats: AHashMap<ObjectType, RuleUpdateResult> = AHashMap::new();
let settings = [
apply_upstream(server, rules.rules).await?,
apply_upstream(server, rules.dnsbls).await?,
apply_upstream(server, rules.tags).await?,
apply_upstream(server, rules.file_exts).await?,
];
let lookups = [
apply_upstream(server, rules.http_lookups).await?,
apply_upstream(server, rules.key_lookups).await?,
];
let mut reload_settings = false;
let mut reload_lookups = false;
for rule in rules.rules {
match registry.write(RegistryWrite::insert(&rule.into())).await? {
RegistryWriteResult::Success(_) => {
stats.entry(ObjectType::SpamRule).or_default().success += 1;
reload_settings = true;
}
RegistryWriteResult::PrimaryKeyConflict { .. } => {
stats
.entry(ObjectType::SpamRule)
.or_default()
.already_exists += 1;
}
_ => {
stats.entry(ObjectType::SpamRule).or_default().failed += 1;
}
let mut reload_errors = Vec::new();
for object in [
settings
.iter()
.any(RuleUpdateResult::has_changes)
.then_some(ObjectType::SpamRule),
lookups
.iter()
.any(RuleUpdateResult::has_changes)
.then_some(ObjectType::MemoryLookupKey),
]
.into_iter()
.flatten()
{
if let Err(reason) = reload_and_broadcast(server, object).await {
reload_errors.push(reason);
}
}
for dnsbl in rules.dnsbls {
match registry.write(RegistryWrite::insert(&dnsbl.into())).await? {
RegistryWriteResult::Success(_) => {
stats
.entry(ObjectType::SpamDnsblServer)
.or_default()
.success += 1;
reload_settings = true;
}
RegistryWriteResult::PrimaryKeyConflict { .. } => {
stats
.entry(ObjectType::SpamDnsblServer)
.or_default()
.already_exists += 1;
}
_ => {
stats.entry(ObjectType::SpamDnsblServer).or_default().failed += 1;
}
}
}
for tag in rules.tags {
match registry.write(RegistryWrite::insert(&tag.into())).await? {
RegistryWriteResult::Success(_) => {
stats.entry(ObjectType::SpamTag).or_default().success += 1;
reload_settings = true;
}
RegistryWriteResult::PrimaryKeyConflict { .. } => {
stats.entry(ObjectType::SpamTag).or_default().already_exists += 1;
}
_ => {
stats.entry(ObjectType::SpamTag).or_default().failed += 1;
}
}
}
for lookup in rules.http_lookups {
match registry
.write(RegistryWrite::insert(&lookup.into()))
.await?
{
RegistryWriteResult::Success(_) => {
stats.entry(ObjectType::HttpLookup).or_default().success += 1;
reload_lookups = true;
}
RegistryWriteResult::PrimaryKeyConflict { .. } => {
stats
.entry(ObjectType::HttpLookup)
.or_default()
.already_exists += 1;
}
_ => {
stats.entry(ObjectType::HttpLookup).or_default().failed += 1;
}
}
}
for key_lookup in rules.key_lookups {
match registry
.write(RegistryWrite::insert(&key_lookup.into()))
.await?
{
RegistryWriteResult::Success(_) => {
stats
.entry(ObjectType::MemoryLookupKey)
.or_default()
.success += 1;
reload_lookups = true;
}
RegistryWriteResult::PrimaryKeyConflict { .. } => {
stats
.entry(ObjectType::MemoryLookupKey)
.or_default()
.already_exists += 1;
}
_ => {
stats.entry(ObjectType::MemoryLookupKey).or_default().failed += 1;
}
}
}
for ext in rules.file_exts {
match registry.write(RegistryWrite::insert(&ext.into())).await? {
RegistryWriteResult::Success(_) => {
stats
.entry(ObjectType::SpamFileExtension)
.or_default()
.success += 1;
reload_settings = true;
}
RegistryWriteResult::PrimaryKeyConflict { .. } => {
stats
.entry(ObjectType::SpamFileExtension)
.or_default()
.already_exists += 1;
}
_ => {
stats
.entry(ObjectType::SpamFileExtension)
.or_default()
.failed += 1;
}
}
}
if reload_settings {
if let Err(err) =
Box::pin(server.reload_registry(RegistryChange::Reload(ObjectType::SpamRule))).await
{
trc::error!(err.details("Failed to reload registry after updating spam rules"));
}
server
.cluster_broadcast(BroadcastEvent::RegistryChange(RegistryChange::Reload(
ObjectType::SpamRule,
)))
.await;
}
if reload_lookups {
if let Err(err) =
Box::pin(server.reload_registry(RegistryChange::Reload(ObjectType::MemoryLookupKey)))
.await
{
trc::error!(err.details("Failed to reload registry after updating spam rules"));
}
server
.cluster_broadcast(BroadcastEvent::RegistryChange(RegistryChange::Reload(
ObjectType::MemoryLookupKey,
)))
.await;
}
// inbuxa: AU-1.10: what the update added, as one audit record
let added = stats
let failed: usize = settings
.iter()
.filter(|(_, result)| result.success > 0)
.map(|(object_type, result)| format!("{} {}", result.success, object_type.as_str()))
.collect::<Vec<_>>();
if !added.is_empty() {
let mut added = added;
added.sort();
.chain(&lookups)
.map(|result| result.failed)
.sum();
// inbuxa: AU-1.10: what the update changed, as one audit record
let summary = |count: fn(&RuleUpdateResult) -> usize| {
let mut parts = settings
.iter()
.chain(&lookups)
.filter(|result| count(result) > 0)
.map(|result| format!("{} {}", count(result), result.object_type.as_str()))
.collect::<Vec<_>>();
parts.sort();
parts.join(", ")
};
let details = [
("added", summary(|result| result.added)),
("replaced", summary(|result| result.updated)),
("kept as edited locally", summary(|result| result.kept)),
("failed", summary(|result| result.failed)),
]
.into_iter()
.filter(|(_, part)| !part.is_empty())
.map(|(what, part)| format!("{what} {part}"))
.collect::<Vec<_>>();
if settings
.iter()
.chain(&lookups)
.any(RuleUpdateResult::has_changes)
{
server
.audit_note(inbuxa_features::audit::Record {
at: std::time::SystemTime::now()
@@ -301,7 +274,7 @@ async fn update_spam_rules(server: &Server) -> trc::Result<TaskResult> {
..Default::default()
},
changes: vec![],
details: Some(format!("Rules update added {}", added.join(", "))),
details: Some(format!("Rules update {}", details.join("; "))),
reason: None,
outcome: inbuxa_features::audit::Outcome::success(),
})
@@ -310,27 +283,173 @@ async fn update_spam_rules(server: &Server) -> trc::Result<TaskResult> {
trc::event!(
Spam(SpamEvent::RulesUpdated),
Details = stats
Details = settings
.into_iter()
.map(|(object_type, result)| {
.chain(lookups)
.map(|result| {
Value::Array(vec![
Value::String(object_type.as_str().into()),
Value::from(result.success),
Value::from(result.already_exists),
Value::String(result.object_type.as_str().into()),
Value::from(result.added),
Value::from(result.updated),
Value::from(result.unchanged),
Value::from(result.failed),
Value::from(result.kept),
])
})
.collect::<Vec<_>>(),
Elapsed = started.elapsed(),
);
// inbuxa: so the next start knows these bundled rules are in
if bundled {
spam_rules::set_applied_version(server.store(), spam_rules::BUNDLED_SPAM_RULES_VERSION)
.await?;
if !reload_errors.is_empty() {
Ok(TaskResult::permanent(format!(
"Spam rules were stored but not activated ({}); fix the logged errors and run Reload settings",
reload_errors.join("; ")
)))
} else if failed > 0 {
Ok(TaskResult::permanent(format!(
"{failed} spam filter objects failed to import or update"
)))
} else {
// inbuxa: so the next start knows these bundled rules are in. Only
// once they all are: a failed update runs again on the next start.
if bundled {
spam_rules::set_applied_version(server.store(), spam_rules::BUNDLED_SPAM_RULES_APPLIED)
.await?;
}
Ok(TaskResult::Success(vec![]))
}
}
async fn reload_and_broadcast(server: &Server, object: ObjectType) -> Result<(), String> {
match Box::pin(server.reload_registry(RegistryChange::Reload(object))).await {
Ok(result) => {
result.log();
if result.has_errors() {
return Err(format!("{} configuration errors", result.errors.len()));
}
server
.cluster_broadcast(BroadcastEvent::RegistryChange(RegistryChange::Reload(
object,
)))
.await;
Ok(())
}
Err(err) => {
let reason = err.to_string();
trc::error!(err.details("Failed to reload registry after updating spam rules"));
Err(reason)
}
}
}
async fn apply_upstream<T: UpstreamObject>(
server: &Server,
objects: Vec<T>,
) -> trc::Result<RuleUpdateResult> {
let registry = server.registry();
let mut result = RuleUpdateResult::new(T::OBJECT);
for upstream in objects {
// inbuxa: every object the update writes is fingerprinted, and an
// existing object is replaced only while it still matches: one an
// admin edited is kept as it is.
let upstream_print = fingerprint(&upstream);
let written = Object::from(upstream.clone());
let existing_id = match registry.write(RegistryWrite::insert(&written)).await? {
RegistryWriteResult::Success(id) => {
spam_rules::set_fingerprint(server.store(), T::OBJECT, id.id(), &upstream_print)
.await?;
result.added += 1;
continue;
}
RegistryWriteResult::PrimaryKeyConflict { existing_id, .. }
if existing_id.object() == T::OBJECT =>
{
existing_id
}
RegistryWriteResult::PrimaryKeyConflict { .. } => {
result.unchanged += 1;
continue;
}
_ => {
result.failed += 1;
continue;
}
};
let Some(local) = registry.get(existing_id).await? else {
result.failed += 1;
continue;
};
let revision = local.revision;
let local = T::from(local);
let local_print = fingerprint(&local);
let written_print =
spam_rules::fingerprint(server.store(), T::OBJECT, existing_id.id().id()).await?;
if local_print == upstream_print {
// The same as upstream's. An install from before fingerprints
// gets one here, so the next release can replace it.
if written_print.as_deref() != Some(upstream_print.as_str()) {
spam_rules::set_fingerprint(
server.store(),
T::OBJECT,
existing_id.id().id(),
&upstream_print,
)
.await?;
}
result.unchanged += 1;
continue;
}
if written_print.as_deref() != Some(local_print.as_str()) {
// Changed since the update wrote it, or never written by one.
result.kept += 1;
continue;
}
let Some(replacement) = upstream
.replacement_for(&local)
.filter(|replacement| replacement != &local)
else {
result.unchanged += 1;
continue;
};
let replacement = Object::from(replacement);
let local = Object::with_revision(local.into(), revision);
match registry
.write(RegistryWrite::update(
existing_id.id(),
&replacement,
&local,
))
.await?
{
RegistryWriteResult::Success(_) => {
spam_rules::set_fingerprint(
server.store(),
T::OBJECT,
existing_id.id().id(),
&upstream_print,
)
.await?;
result.updated += 1
}
_ => result.failed += 1,
}
}
Ok(TaskResult::Success(vec![]))
Ok(result)
}
/// inbuxa: a digest of an object's content, `enable` aside, in hex.
fn fingerprint<T: UpstreamObject>(object: &T) -> String {
use sha2::Digest;
sha2::Sha256::digest(serde_json::to_vec(&object.clone().without_enable()).unwrap_or_default())
.iter()
.map(|byte| format!("{byte:02x}"))
.collect()
}
async fn fetch_spam_rules(server: &Server) -> Result<Rules, RuleUpdateError> {
@@ -347,13 +466,12 @@ async fn fetch_spam_rules(server: &Server) -> Result<Rules, RuleUpdateError> {
reason,
}),
};
let rules_json: AHashMap<String, Vec<serde_json::Value>> =
bytes.and_then(|bytes| {
serde_json::from_slice(&bytes).map_err(|err| RuleUpdateError {
typ: TaskFailureType::Permanent,
reason: format!("Failed to parse spam rules JSON: {err}"),
})
})?;
let rules_json: AHashMap<String, Vec<serde_json::Value>> = bytes.and_then(|bytes| {
serde_json::from_slice(&bytes).map_err(|err| RuleUpdateError {
typ: TaskFailureType::Permanent,
reason: format!("Failed to parse spam rules JSON: {err}"),
})
})?;
let mut rules = Rules::default();
for (object_type, values) in rules_json {
+1 -1
View File
@@ -6,7 +6,7 @@ homepage = "https://inbuxa.org"
keywords = ["smtp", "email", "mail", "server"]
categories = ["email"]
license = "AGPL-3.0-only OR LicenseRef-SEL"
version = "0.16.23"
version = "0.16.24"
edition = "2024"
[dependencies]
+6 -1
View File
@@ -696,7 +696,12 @@ impl<T: SessionStream> Session<T> {
.map(|a| a.as_str())
.unwrap_or_default(),
)
.with_message(parsed_message);
.with_message(
edited_message
.as_deref()
.and_then(|message| MessageParser::new().parse(message))
.unwrap_or(parsed_message),
);
let modifications = match self.run_script(script_id, script.clone(), params).await {
ScriptResult::Accept { modifications } => modifications,
+1 -1
View File
@@ -132,7 +132,7 @@ impl<T: SessionStream> Session<T> {
name: "X-Quarantine".into(),
value: "true".into(),
});
FilterResponse::accept()
continue;
}
};
+42 -10
View File
@@ -531,11 +531,23 @@ impl QueuedMessage {
};
// Obtain remote hosts list
let mx_unvalidated = mx_config.is_some() && !tls_strategy.try_dane();
let mx_list;
if let Some(mx_config) = mx_config {
// Lookup MX
let time = Instant::now();
mx_list = match server.mx_lookup(domain).await {
let mx_lookup = if mx_unvalidated {
server
.core
.smtp
.resolvers
.dns
.mx_lookup(domain, Some(&server.inner.cache.dns_mx))
.await
} else {
server.mx_lookup(domain).await
};
mx_list = match mx_lookup {
Ok(mx) => mx,
Err(mail_auth::Error::Dns(mail_auth::DnsError::RecordNotFound(_))) => {
trc::event!(
@@ -674,6 +686,33 @@ impl QueuedMessage {
message.span_id,
);
let validated_host;
let remote_host = if mx_unvalidated && tls_strategy.try_dane() {
let time = Instant::now();
let dnssec_status = match server.mx_lookup(domain).await {
Ok(mx) => mx.dnssec_status,
Err(mail_auth::Error::Dns(mail_auth::DnsError::RecordNotFound(_))) => {
DnssecStatus::Indeterminate
}
Err(err) => {
trc::event!(
Delivery(DeliveryEvent::MxLookupFailed),
SpanId = message.span_id,
Domain = domain.to_string(),
CausedBy = trc::Error::from(err.clone()),
Elapsed = time.elapsed(),
);
last_status = Status::from_mail_auth_error(domain, err);
continue 'next_host;
}
};
validated_host = remote_host.with_dnssec_status(dnssec_status);
&validated_host
} else {
remote_host
};
// Obtain source and remote IPs
let time = Instant::now();
let validate_addresses = server.core.smtp.resolvers.dnssec_available
@@ -722,15 +761,8 @@ impl QueuedMessage {
let time = Instant::now();
let strict = tls_strategy.is_dane_required();
let (dnssec_status, dnssec_entity) = match remote_host.dnssec_status() {
DnssecStatus::Secure => match addresses_dnssec_status {
status @ (DnssecStatus::Insecure | DnssecStatus::Bogus) => {
(status, "A/AAAA")
}
_ => (DnssecStatus::Secure, "MX"),
},
status => (status, "MX"),
};
let (dnssec_status, dnssec_entity) =
remote_host.dane_status(addresses_dnssec_status);
match dnssec_status {
DnssecStatus::Secure => {
+28 -1
View File
@@ -350,12 +350,39 @@ impl NextHop<'_> {
}
}
fn dnssec_status(&self) -> DnssecStatus {
pub fn dnssec_status(&self) -> DnssecStatus {
match self {
NextHop::MX { dnssec_status, .. } => *dnssec_status,
NextHop::Relay(_) => DnssecStatus::Indeterminate,
}
}
fn with_dnssec_status(&self, dnssec_status: DnssecStatus) -> Self {
match self {
NextHop::MX {
is_implicit,
host,
config,
..
} => NextHop::MX {
is_implicit: *is_implicit,
host,
config,
dnssec_status,
},
NextHop::Relay(relay) => NextHop::Relay(relay),
}
}
pub fn dane_status(&self, addresses: DnssecStatus) -> (DnssecStatus, &'static str) {
match self.dnssec_status() {
DnssecStatus::Secure => match addresses {
status @ (DnssecStatus::Insecure | DnssecStatus::Bogus) => (status, "A/AAAA"),
_ => (DnssecStatus::Secure, "MX"),
},
status => (status, "MX"),
}
}
}
impl DeliveryResult {
+5
View File
@@ -67,6 +67,10 @@ impl SpawnQueue for mpsc::Receiver<QueueEvent> {
Queue::new(core, self).start().await;
});
}
fn discard(mut self) {
tokio::spawn(async move { while self.recv().await.is_some() {} });
}
}
const BACK_PRESSURE_WARN_INTERVAL: Duration = Duration::from_secs(60);
@@ -510,6 +514,7 @@ impl Recipient {
pub trait SpawnQueue {
fn spawn(self, core: Arc<Inner>);
fn discard(self);
}
impl QueueStats {
+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "spam-filter"
version = "0.16.23"
version = "0.16.24"
edition = "2024"
[dependencies]
+28 -2
View File
@@ -203,7 +203,7 @@ impl SpamFilterAnalyzeUrl for Server {
if !ctx.result.has_tag("URL_REDIRECTOR_NESTED") {
let mut redirect_count = 1;
let mut url_redirect = Cow::Borrowed(url.element.url.as_str());
let mut url_redirect = url.element.request_url();
while redirect_count <= 3 {
match http_get_header(
@@ -224,7 +224,8 @@ impl SpamFilterAnalyzeUrl for Server {
)
.await
{
url_redirect = Cow::Owned(location.url);
url_redirect =
Cow::Owned(location.request_url().into_owned());
redirect_count += 1;
continue;
} else {
@@ -487,6 +488,15 @@ impl<'x> UrlParts<'x> {
.is_some_and(|url| url.host.fqdn.starts_with("www."))
}
pub fn request_url(&self) -> Cow<'_, str> {
let url = self.url_original.trim();
if self.has_scheme {
Cow::Borrowed(url)
} else {
Cow::Owned([HTTPS_SCHEME, url].concat())
}
}
fn parse(url: &str) -> Option<UrlParsed> {
url.parse::<Uri>().ok().and_then(|parts| {
parts
@@ -505,3 +515,19 @@ impl<'x> UrlParts<'x> {
}
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn request_url_keeps_case() {
let url = UrlParts::new(" https://Bit.ly/3AbCdEf ");
assert_eq!(url.url, "https://bit.ly/3abcdef");
assert_eq!(url.request_url(), "https://Bit.ly/3AbCdEf");
let url = UrlParts::no_scheme("Bit.ly/3AbCdEf");
assert_eq!(url.url, "https://bit.ly/3abcdef");
assert_eq!(url.request_url(), "https://Bit.ly/3AbCdEf");
}
}
+1 -3
View File
@@ -259,9 +259,7 @@ impl SpamClassifier for Server {
remove_entries = true;
}
if trainer.last_id == 0 {
trainer.last_id = id;
}
trainer.last_id = trainer.last_id.max(id);
Ok(true)
},
+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "store"
version = "0.16.23"
version = "0.16.24"
edition = "2024"
[dependencies]
+142 -25
View File
@@ -4,21 +4,19 @@
* SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-SEL
*/
use super::{HttpStore, HttpStoreConfig};
use crate::{Value, backend::http::HttpStoreFormat, write::now};
use ahash::AHashMap;
use compact_str::ToCompactString;
use rand::seq::IndexedRandom;
use std::{
borrow::Cow,
io::{BufRead, BufReader},
sync::{Arc, atomic::Ordering},
time::Instant,
};
use ahash::AHashMap;
use compact_str::ToCompactString;
use rand::seq::IndexedRandom;
use utils::HttpLimitResponse;
use crate::{Value, backend::http::HttpStoreFormat, write::now};
use super::HttpStore;
const BROWSER_USER_AGENTS: [&str; 5] = [
"Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/120.0.0.0 Safari/537.36",
"Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Edge/120.0.0.0 Safari/537.36",
@@ -36,7 +34,7 @@ pub(crate) trait HttpStoreGet {
impl HttpStoreGet for Arc<HttpStore> {
fn get(&self, key: &str) -> Option<Value<'static>> {
self.refresh();
self.entries.load().get(key).cloned()
self.entries.load().get(lookup_key(key).as_ref()).cloned()
}
fn contains(&self, key: &str) -> bool {
@@ -59,7 +57,7 @@ impl HttpStoreGet for Arc<HttpStore> {
}
self.refresh();
self.entries.load().contains_key(key)
self.entries.load().contains_key(lookup_key(key).as_ref())
}
fn refresh(&self) {
@@ -142,9 +140,10 @@ impl HttpStore {
Box::new(&bytes[..])
};
let mut entries = AHashMap::new();
for (pos, line) in BufReader::new(reader).lines().enumerate() {
let line_ = line.map_err(|err| {
let entries = self
.config
.parse_entries(BufReader::new(reader))
.map_err(|err| {
trc::StoreEvent::HttpStoreError
.into_err()
.reason(err)
@@ -153,11 +152,31 @@ impl HttpStore {
.details("Failed to read line")
})?;
match &self.config.format {
trc::event!(
Store(trc::StoreEvent::HttpStoreFetch),
Url = self.config.url.to_compact_string(),
Total = entries.len(),
Elapsed = time.elapsed(),
);
Ok(entries)
}
}
impl HttpStoreConfig {
fn parse_entries(
&self,
reader: impl BufRead,
) -> std::io::Result<AHashMap<String, Value<'static>>> {
let mut entries = AHashMap::new();
for (pos, line) in reader.lines().enumerate() {
let line_ = line?;
match &self.format {
HttpStoreFormat::List => {
let line = line_.trim();
if !line.is_empty() {
entries.insert(line.to_string(), Value::Integer(1));
entries.insert(lookup_key(line).into_owned(), Value::Integer(1));
}
}
HttpStoreFormat::Csv {
@@ -188,12 +207,12 @@ impl HttpStore {
}
} else if col_num == *index_key {
entry_key.push(ch);
if entry_key.len() > self.config.max_entry_size {
if entry_key.len() > self.max_entry_size {
break;
}
} else if index_value.is_some_and(|v| col_num == v) {
entry_value.push(ch);
if entry_value.len() > self.config.max_entry_size {
if entry_value.len() > self.max_entry_size {
break;
}
}
@@ -209,24 +228,122 @@ impl HttpStore {
} else {
Value::Integer(1)
};
let entry_key = match lookup_key(&entry_key) {
Cow::Owned(key) => key,
Cow::Borrowed(_) => entry_key,
};
entries.insert(entry_key, entry_value);
}
}
_ => (),
}
if entries.len() == self.config.max_entries {
if entries.len() == self.max_entries {
break;
}
}
trc::event!(
Store(trc::StoreEvent::HttpStoreFetch),
Url = self.config.url.to_compact_string(),
Total = entries.len(),
Elapsed = time.elapsed(),
);
Ok(entries)
}
}
fn lookup_key(key: &str) -> Cow<'_, str> {
if key.bytes().any(|b| !b.is_ascii() || b.is_ascii_uppercase()) {
Cow::Owned(key.to_lowercase())
} else {
Cow::Borrowed(key)
}
}
#[cfg(test)]
mod tests {
use super::*;
use arc_swap::ArcSwap;
use reqwest::Client;
use std::{
sync::atomic::{AtomicBool, AtomicU64},
time::Duration,
};
fn http_store(format: HttpStoreFormat, feed: &str) -> Arc<HttpStore> {
let config = HttpStoreConfig {
id: "test".into(),
url: "https://lists.example.org/feed".into(),
retry: 0,
refresh: 0,
timeout: Duration::from_secs(1),
gzipped: false,
max_size: 1024 * 1024,
max_entries: 100,
max_entry_size: 512,
format,
};
let entries = config
.parse_entries(feed.as_bytes())
.expect("feed is readable");
Arc::new(HttpStore {
entries: ArcSwap::from_pointee(entries),
expires: AtomicU64::new(u64::MAX),
in_flight: AtomicBool::new(false),
config,
client: Client::new(),
})
}
#[test]
fn list_keys_ignore_case() {
let store = http_store(
HttpStoreFormat::List,
"https://phish.example.org/Account/Verify?Token=AbC123\n\
https://PHISH.example.net/lower\n",
);
assert!(store.contains("https://phish.example.org/account/verify?token=abc123"));
assert!(store.contains("https://phish.example.org/Account/Verify?Token=AbC123"));
assert!(store.contains("https://phish.example.net/lower"));
assert!(store.contains("HTTPS://PHISH.EXAMPLE.NET/LOWER"));
assert!(!store.contains("https://phish.example.org/account/verify"));
}
#[test]
fn csv_keys_ignore_case() {
let store = http_store(
HttpStoreFormat::Csv {
index_key: 1,
index_value: None,
separator: ',',
skip_first: true,
},
"phish_id,url,phish_detail_url\n\
1,\"https://phish.example.org/Login.PHP?Id=Xy\",https://phishtank.example/1\n",
);
assert!(store.contains("https://phish.example.org/login.php?id=xy"));
assert!(store.contains("https://phish.example.org/Login.PHP?Id=Xy"));
assert!(!store.contains("phish_id"));
assert!(!store.contains("url"));
}
#[test]
fn csv_values_keep_case() {
let store = http_store(
HttpStoreFormat::Csv {
index_key: 0,
index_value: Some(1),
separator: ',',
skip_first: false,
},
"Example.ORG,Some Value\n",
);
assert_eq!(
store.get("example.org"),
Some(Value::Text("Some Value".into()))
);
assert_eq!(
store.get("EXAMPLE.org"),
Some(Value::Text("Some Value".into()))
);
}
}
+14
View File
@@ -99,7 +99,10 @@ pub(crate) fn into_error(err: impl Display) -> trc::Error {
trc::StoreEvent::MysqlError.reason(err)
}
const ER_UNKNOWN_ERROR: u16 = 1105;
const ER_TRANS_CACHE_FULL: u16 = 1197;
const ER_LOCK_WAIT_TIMEOUT: u16 = 1205;
const ER_LOCK_TABLE_FULL: u16 = 1206;
const ER_STATEMENT_TIMEOUT: u16 = 1969;
const ER_QUERY_TIMEOUT: u16 = 3024;
@@ -116,6 +119,17 @@ pub(crate) fn is_timeout_error(err: &mysql_async::Error) -> bool {
)
}
#[inline(always)]
pub(crate) fn is_chunk_too_large_error(err: &mysql_async::Error) -> bool {
is_timeout_error(err)
|| matches!(err, mysql_async::Error::Server(err)
if matches!(
err.code,
ER_UNKNOWN_ERROR | ER_TRANS_CACHE_FULL | ER_LOCK_TABLE_FULL
)
)
}
impl SearchIndex {
pub fn mysql_table(&self) -> &'static str {
match self {
+5 -12
View File
@@ -11,7 +11,7 @@ use crate::{
MAX_TOKEN_LENGTH,
mysql::{
DELETE_CHUNK_SIZE, MIN_DELETE_CHUNK_SIZE, MysqlSearchField, MysqlStore, bounded,
into_error, is_timeout_error,
into_error, is_chunk_too_large_error,
},
},
search::{
@@ -123,14 +123,6 @@ impl MysqlStore {
let mut conn = self.conn().await?;
let limit = self.timeouts.maintenance;
let result = tokio::time::timeout(limit, async {
let s = conn.prep(&query).await.map_err(into_error)?;
match conn.exec_drop(s, params.clone()).await {
Ok(_) => return Ok(conn.affected_rows()),
Err(err) if is_timeout_error(&err) => (),
Err(err) => return Err(into_error(err)),
}
let mut chunk_size = DELETE_CHUNK_SIZE;
let mut deleted = 0;
@@ -144,13 +136,14 @@ impl MysqlStore {
match conn.exec_drop(&s, params.clone()).await {
Ok(_) => {
let affected = conn.affected_rows();
if affected == 0 {
deleted += affected;
if affected < chunk_size as u64 {
return Ok(deleted);
}
deleted += affected;
}
Err(err)
if is_timeout_error(&err) && chunk_size > MIN_DELETE_CHUNK_SIZE =>
if is_chunk_too_large_error(&err)
&& chunk_size > MIN_DELETE_CHUNK_SIZE =>
{
chunk_size = (chunk_size / 2).max(MIN_DELETE_CHUNK_SIZE);
break;
+17 -12
View File
@@ -7,7 +7,8 @@
*/
use super::{
DELETE_CHUNK_SIZE, MIN_DELETE_CHUNK_SIZE, MysqlStore, bounded, into_error, is_timeout_error,
DELETE_CHUNK_SIZE, MIN_DELETE_CHUNK_SIZE, MysqlStore, bounded, into_error,
is_chunk_too_large_error,
};
use crate::{
IndexKey, Key, LogKey, SUBSPACE_COUNTER, SUBSPACE_IN_MEMORY_COUNTER, SUBSPACE_QUOTA,
@@ -416,12 +417,6 @@ impl MysqlStore {
.await
.map_err(into_error)?;
match conn.exec_drop(&delete, (&from, &to)).await {
Ok(_) => return Ok(()),
Err(err) if is_timeout_error(&err) => (),
Err(err) => return Err(into_error(err)),
}
let mut chunk_size = DELETE_CHUNK_SIZE;
loop {
@@ -438,7 +433,10 @@ impl MysqlStore {
.await
{
Ok(next) => next,
Err(err) if is_timeout_error(&err) && chunk_size > MIN_DELETE_CHUNK_SIZE => {
Err(err)
if is_chunk_too_large_error(&err)
&& chunk_size > MIN_DELETE_CHUNK_SIZE =>
{
chunk_size = (chunk_size / 2).max(MIN_DELETE_CHUNK_SIZE);
break;
}
@@ -450,7 +448,10 @@ impl MysqlStore {
.await
{
Ok(_) => (),
Err(err) if is_timeout_error(&err) && chunk_size > MIN_DELETE_CHUNK_SIZE => {
Err(err)
if is_chunk_too_large_error(&err)
&& chunk_size > MIN_DELETE_CHUNK_SIZE =>
{
chunk_size = (chunk_size / 2).max(MIN_DELETE_CHUNK_SIZE);
break;
}
@@ -477,7 +478,7 @@ async fn purge_table(conn: &mut Conn, table: char) -> trc::Result<()> {
match conn.exec_drop(&s, ()).await {
Ok(_) => return Ok(()),
Err(err) if is_timeout_error(&err) => (),
Err(err) if is_chunk_too_large_error(&err) => (),
Err(err) => return Err(into_error(err)),
}
@@ -505,7 +506,9 @@ async fn purge_table(conn: &mut Conn, table: char) -> trc::Result<()> {
loop {
let next = match conn.exec_first::<Vec<u8>, _, _>(&boundary, (&from,)).await {
Ok(next) => next,
Err(err) if is_timeout_error(&err) && chunk_size > MIN_DELETE_CHUNK_SIZE => {
Err(err)
if is_chunk_too_large_error(&err) && chunk_size > MIN_DELETE_CHUNK_SIZE =>
{
chunk_size = (chunk_size / 2).max(MIN_DELETE_CHUNK_SIZE);
break;
}
@@ -519,7 +522,9 @@ async fn purge_table(conn: &mut Conn, table: char) -> trc::Result<()> {
match result {
Ok(_) => (),
Err(err) if is_timeout_error(&err) && chunk_size > MIN_DELETE_CHUNK_SIZE => {
Err(err)
if is_chunk_too_large_error(&err) && chunk_size > MIN_DELETE_CHUNK_SIZE =>
{
chunk_size = (chunk_size / 2).max(MIN_DELETE_CHUNK_SIZE);
break;
}
+6 -14
View File
@@ -9,16 +9,7 @@
use super::{RedisPool, RedisStore, into_error};
use crate::{Deserialize, write::now};
use deadpool::managed::{Manager, Object, Pool};
use redis::{AsyncCommands, RedisError, RedisResult, RetryMethod, Script};
use std::sync::LazyLock;
static INCR_EXPIRE: LazyLock<Script> = LazyLock::new(|| {
Script::new(
"redis.call('INCRBY', KEYS[1], ARGV[1])
redis.call('EXPIRE', KEYS[1], ARGV[2])
return redis.call('GET', KEYS[1])",
)
});
use redis::{AsyncCommands, RedisError, RedisResult, RetryMethod};
impl RedisStore {
pub async fn key_set(&self, key: &[u8], value: &[u8], expires: Option<u64>) -> trc::Result<()> {
@@ -48,19 +39,19 @@ impl RedisStore {
match &self.pool {
RedisPool::Single(pool) => {
with_conn(pool, async |conn| {
Self::key_incr_(conn, key, value, expires).await
self.key_incr_(conn, key, value, expires).await
})
.await
}
RedisPool::Cluster(pool) => {
with_conn(pool, async |conn| {
Self::key_incr_(conn, key, value, expires).await
self.key_incr_(conn, key, value, expires).await
})
.await
}
RedisPool::Sentinel(pool) => {
with_conn(pool, async |conn| {
Self::key_incr_(conn, key, value, expires).await
self.key_incr_(conn, key, value, expires).await
})
.await
}
@@ -219,13 +210,14 @@ impl RedisStore {
}
async fn key_incr_(
&self,
conn: &mut impl AsyncCommands,
key: &[u8],
value: i64,
expires: Option<u64>,
) -> RedisResult<i64> {
if let Some(expires) = expires {
INCR_EXPIRE
self.incr_expire
.key(key)
.arg(value)
.arg(expires as i64)
+22 -10
View File
@@ -10,7 +10,7 @@ use deadpool::{
managed::{Manager, Pool},
};
use redis::{
Client, ConnectionAddr, IntoConnectionInfo, ProtocolVersion, TlsMode,
Client, ConnectionAddr, IntoConnectionInfo, ProtocolVersion, Script, TlsMode,
cluster::{ClusterClient, ClusterClientBuilder},
cluster_read_routing::RandomReplicaStrategy,
sentinel::{SentinelClient, SentinelClientBuilder, SentinelServerType},
@@ -27,6 +27,7 @@ pub mod pool;
#[derive(Debug)]
pub struct RedisStore {
pub pool: RedisPool,
incr_expire: Script,
}
pub struct RedisConnectionManager {
@@ -51,9 +52,20 @@ pub enum RedisPool {
}
impl RedisStore {
fn new(pool: RedisPool) -> Self {
RedisStore {
pool,
incr_expire: Script::new(
"redis.call('INCRBY', KEYS[1], ARGV[1])
redis.call('EXPIRE', KEYS[1], ARGV[2])
return redis.call('GET', KEYS[1])",
),
}
}
pub async fn open_single(config: structs::RedisStore) -> Result<InMemoryStore, String> {
Ok(InMemoryStore::Redis(Arc::new(RedisStore {
pool: RedisPool::Single(build_pool(
Ok(InMemoryStore::Redis(Arc::new(RedisStore::new(
RedisPool::Single(build_pool(
RedisConnectionManager {
client: Client::open(config.url)
.map_err(|err| format!("Failed to open Redis client: {err:?}"))?,
@@ -64,7 +76,7 @@ impl RedisStore {
config.pool_timeout_wait,
config.pool_timeout_recycle,
)?),
})))
))))
}
pub async fn open_cluster(config: structs::RedisClusterStore) -> Result<InMemoryStore, String> {
@@ -95,8 +107,8 @@ impl RedisStore {
.build()
.map_err(|err| format!("Failed to open Redis client: {err:?}"))?;
Ok(InMemoryStore::Redis(Arc::new(RedisStore {
pool: RedisPool::Cluster(build_pool(
Ok(InMemoryStore::Redis(Arc::new(RedisStore::new(
RedisPool::Cluster(build_pool(
RedisClusterConnectionManager {
client,
timeout: config.timeout.into_inner(),
@@ -106,7 +118,7 @@ impl RedisStore {
config.pool_timeout_wait,
config.pool_timeout_recycle,
)?),
})))
))))
}
pub async fn open_sentinel(
@@ -167,8 +179,8 @@ impl RedisStore {
.build()
.map_err(|err| format!("Failed to open Redis Sentinel client: {err:?}"))?;
Ok(InMemoryStore::Redis(Arc::new(RedisStore {
pool: RedisPool::Sentinel(build_pool(
Ok(InMemoryStore::Redis(Arc::new(RedisStore::new(
RedisPool::Sentinel(build_pool(
RedisSentinelConnectionManager {
client: tokio::sync::Mutex::new(client),
timeout: config.timeout.into_inner(),
@@ -178,7 +190,7 @@ impl RedisStore {
config.pool_timeout_wait,
config.pool_timeout_recycle,
)?),
})))
))))
}
}
+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "trc"
version = "0.16.23"
version = "0.16.24"
edition = "2024"
[dependencies]
+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "event_macro"
version = "0.16.23"
version = "0.16.24"
edition = "2024"
[lib]
+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "types"
version = "0.16.23"
version = "0.16.24"
edition = "2024"
[dependencies]
+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "utils"
version = "0.16.23"
version = "0.16.24"
edition = "2024"
[dependencies]
+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "proc_macros"
version = "0.16.23"
version = "0.16.24"
edition = "2024"
[lib]