diff --git a/CHANGELOG.md b/CHANGELOG.md index 6cf7732..9cc23c6 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -2,6 +2,39 @@ All notable changes to this project will be documented in this file. This project adheres to [Semantic Versioning](http://semver.org/). +## [0.16.23] - 2026-09-21 + +If you are upgrading from v0.16.x, replace the binary (or run `docker pull`). If you are upgrading from v0.15.x and below, please read the [upgrading documentation](https://github.com/stalwartlabs/stalwart/blob/main/UPGRADING/v0_16.md) for more information on how to upgrade from previous versions. + +## Added +- Expressions: `bit_and` function. + +## Changed + +## Fixed +- MTA: + - A mailing list whose recipients include another mailing list is accepted at `RCPT TO` and then rejected at local delivery with `550 5.5.0 Mailbox not found`. + - DMARC aggregate reports carry two `spf` elements per record and the `version` element of a DMARC aggregate report is written as `1` instead of `1.0`. + - DSNs generated for an alias rewrite or a list expansion emit a doubled `addr-type` in `Original-Recipient` (`rfc822;rfc822;user@example.org`). + - DSNs that cannot be written to the store are discarded, the recipients are flagged as notified and the original message is removed from the queue, losing both the bounce and the message. +- POP3: + - `TOP msg n` counts the `n` lines from the first byte of the message instead of from the first byte of the body. + - A message whose very first line begins with `.` is not byte-stuffed. +- Spam filter: Moving or copying a message from one account into another creates no training sample, so the classifier never learns from it. +- Sieve: `envelope "orcpt"` yields the bare address for an `ORCPT` supplied over SMTP. It now carries the `addr-type` prefix in every case, as required by RFC 6009. +- ACME: The `_acme-challenge` TXT records published for a DNS-01 authorization are never removed. +- DNS: The DNSSEC resolver queries a single nameserver at a time, working around a `hickory-resolver` race that cancels the TCP retry when two nameservers return a truncated response in parallel. +- Troubleshoot tool: + - MX records are resolved through the DNSSEC-validating resolver, matching the resolver used by the delivery path. + - A TLSA lookup that fails or returns bogus records stops the delivery attempt for that host, instead of continuing without DANE. +- OIDC: Bearer tokens that carry no `email`, `preferred_username` or `upn` claim are always authenticated against the default directory. +- Meilisearch: A confirmation timeout is treated as a failed write even when `failOnTimeout` is disabled, so an index whose batches take longer than `pollInterval` x `maxRetries` never completes an indexing task and resubmits the same batch indefinitely. +- WebUI: A failed update no longer takes an `Application` offline. +- FoundationDB: The cached read version is invalidated when any broadcast is received from another node. +- Redis: + - On a cluster, the rate limiter and the blob upload quota issue `INCR` and `EXPIRE` as a `MULTI`/`EXEC` transaction, whose `MOVED` redirects collapse into a single `EXECABORT` that never refreshes the slot map. + - A connection that fails because it is addressing the wrong server is returned to the pool and reused, since the recycle check only issues `PING`. + ## [0.16.22] - 2026-09-13 If you are upgrading from v0.16.x, replace the binary (or run `docker pull`). If you are upgrading from v0.15.x and below, please read the [upgrading documentation](https://github.com/stalwartlabs/stalwart/blob/main/UPGRADING/v0_16.md) for more information on how to upgrade from previous versions. diff --git a/Cargo.lock b/Cargo.lock index 75f8177..ed5d81c 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -234,7 +234,7 @@ dependencies = [ "proc-macro2", "quote", "syn 2.0.119", - "synstructure", + "synstructure 0.13.2", ] [[package]] @@ -277,9 +277,9 @@ dependencies = [ [[package]] name = "async-compression" -version = "0.4.46" +version = "0.4.48" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "4f10dafd0c8d2e51ae9a748805777613ed0bbe17bf586b76c8311f45c020a32f" +checksum = "fb61aea1a7def73ee7c350a184f0e70b32c182344e2e75bf70c9b621b83417fd" dependencies = [ "compression-codecs", "compression-core", @@ -310,7 +310,7 @@ dependencies = [ "memchr", "pin-project", "portable-atomic", - "rand 0.10.2", + "rand 0.10.3", "regex", "rustls-native-certs", "rustls-pki-types", @@ -369,7 +369,7 @@ checksum = "82f6aeea286b8eb4dd3431a1be1b59d290ace00f5bfd8e2a159bc2a05e2c1667" dependencies = [ "proc-macro2", "quote", - "syn 3.0.5", + "syn 3.0.6", ] [[package]] @@ -874,7 +874,7 @@ dependencies = [ "log", "num", "pin-project-lite", - "rand 0.10.2", + "rand 0.10.3", "rustls", "rustls-native-certs", "rustls-pki-types", @@ -984,7 +984,7 @@ checksum = "46d07918caa9eeaaf06b7873925c53a61daac173539b4f7715090745e44e4e69" dependencies = [ "proc-macro2", "quote", - "syn 3.0.5", + "syn 3.0.6", ] [[package]] @@ -1110,9 +1110,9 @@ dependencies = [ [[package]] name = "cc" -version = "1.4.6" +version = "1.4.7" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "a3eb0f42d6c360dc3f8a821f6bf2fdea7f72bfd36b3076eb0e6d1e9e0752fff4" +checksum = "54413ede23c2daf518f35156dfde027feb2374004d63bd497f983c8db9c0e313" dependencies = [ "find-msvc-tools", "jobserver", @@ -1160,9 +1160,9 @@ dependencies = [ [[package]] name = "cfg-if" -version = "1.0.4" +version = "1.0.5" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "9330f8b2ff13f34540b44e946ef35111825727b38d33286ef986142615121801" +checksum = "4e7648175b45a9a48536d676f68d918270699102aa8dab5496df06904c914600" [[package]] name = "cfg_aliases" @@ -1302,7 +1302,7 @@ dependencies = [ [[package]] name = "common" -version = "0.16.22" +version = "0.16.23" dependencies = [ "aes-gcm-siv", "ahash", @@ -1402,9 +1402,9 @@ dependencies = [ [[package]] name = "compression-codecs" -version = "0.4.41" +version = "0.4.43" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "58a6d0db8759036a783bc7c3f7a07f8cef3bf9470eb1db3bc86e8bcd1c5d0fe8" +checksum = "bef16c47ba2797aa6a909cc37d39911f3a6743811fe7408ac0b0cc0276b656e9" dependencies = [ "compression-core", "flate2", @@ -1487,7 +1487,7 @@ checksum = "3d52eff69cd5e647efe296129160853a42795992097e8af39800e1060caeea9b" [[package]] name = "coordinator" -version = "0.16.22" +version = "0.16.23" dependencies = [ "async-nats", "futures", @@ -1849,7 +1849,7 @@ dependencies = [ "proc-macro2", "quote", "strsim", - "syn 3.0.5", + "syn 3.0.6", ] [[package]] @@ -1882,7 +1882,7 @@ checksum = "2ac7135c3ef02b2f7833bbeb1be5ba7f966dcde8a87c6b87f65a778d71a02785" dependencies = [ "darling_core 0.24.1", "quote", - "syn 3.0.5", + "syn 3.0.6", ] [[package]] @@ -1899,7 +1899,7 @@ checksum = "4583a4551df46e2792f82ceeac45e850d2e2d5debba0b91f102385cda5b11f06" [[package]] name = "dav" -version = "0.16.22" +version = "0.16.23" dependencies = [ "calcard", "chrono", @@ -1922,7 +1922,7 @@ dependencies = [ [[package]] name = "dav-proto" -version = "0.16.22" +version = "0.16.23" dependencies = [ "calcard", "chrono", @@ -2135,7 +2135,7 @@ dependencies = [ [[package]] name = "directory" -version = "0.16.22" +version = "0.16.23" dependencies = [ "ahash", "argon2 0.6.0", @@ -2192,7 +2192,7 @@ checksum = "c6232dd377dcc64799954cbd3a9bb882e9cdc1308ccd87b1c098f1fb2eaf82a8" dependencies = [ "proc-macro2", "quote", - "syn 3.0.5", + "syn 3.0.6", ] [[package]] @@ -2376,7 +2376,7 @@ dependencies = [ [[package]] name = "email" -version = "0.16.22" +version = "0.16.23" dependencies = [ "aes 0.9.3", "aes-gcm 0.11.1", @@ -2485,10 +2485,10 @@ dependencies = [ [[package]] name = "event_macro" -version = "0.16.22" +version = "0.16.23" dependencies = [ "quote", - "syn 3.0.5", + "syn 3.0.6", ] [[package]] @@ -2571,7 +2571,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "ee93edf3c501f0035bbeffeccfed0b79e14c311f12195ec0e661e114a0f60da4" dependencies = [ "portable-atomic", - "rand 0.10.2", + "rand 0.10.3", "web-time", ] @@ -2594,9 +2594,9 @@ checksum = "28dea519a9695b9977216879a3ebfddf92f1c08c05d984f8996aecd6ecdc811d" [[package]] name = "find-msvc-tools" -version = "0.1.12" +version = "0.1.13" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "3e0f1c7c3a72c66fd80abe965175f7523475c0489a87d3ff9d6e8c87d87a9d2d" +checksum = "ef25905e51abafe4dcea6c15fec58c57b601cdbd0ee53d22ea1d3016c587d39b" [[package]] name = "fixed_decimal" @@ -2710,7 +2710,7 @@ dependencies = [ "foundationdb-sys", "foundationdb-tuple", "futures", - "rand 0.10.2", + "rand 0.10.3", "serde", "serde_bytes", "serde_json", @@ -2842,7 +2842,7 @@ checksum = "9fb9654ba8355388abeb8dcb4fc62f511300867002afc858860463bdd9fe0c44" dependencies = [ "proc-macro2", "quote", - "syn 3.0.5", + "syn 3.0.6", ] [[package]] @@ -3013,7 +3013,7 @@ dependencies = [ [[package]] name = "groupware" -version = "0.16.22" +version = "0.16.23" dependencies = [ "ahash", "calcard", @@ -3169,7 +3169,7 @@ dependencies = [ "jni", "lru-cache", "parking_lot", - "rand 0.10.2", + "rand 0.10.3", "rustls", "rustls-pki-types", "rustls-platform-verifier", @@ -3196,7 +3196,7 @@ dependencies = [ "jni", "once_cell", "prefix-trie", - "rand 0.10.2", + "rand 0.10.3", "ring", "rustls-pki-types", "thiserror 2.0.20", @@ -3223,7 +3223,7 @@ dependencies = [ "ndk-context", "once_cell", "parking_lot", - "rand 0.10.2", + "rand 0.10.3", "resolv-conf", "rustls", "smallvec", @@ -3302,7 +3302,7 @@ dependencies = [ [[package]] name = "http" -version = "0.16.22" +version = "0.16.23" dependencies = [ "async-stream", "base64 0.23.1", @@ -3398,7 +3398,7 @@ dependencies = [ [[package]] name = "http_proto" -version = "0.16.22" +version = "0.16.23" dependencies = [ "common", "compact_str", @@ -3488,9 +3488,9 @@ dependencies = [ [[package]] name = "hyper-rustls" -version = "0.27.9" +version = "0.27.10" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "33ca68d021ef39cf6463ab54c1d0f5daf03377b70561305bb89a8f83aab66e0f" +checksum = "dfa8e654703247911e29c23fbeaa261834bd9bb74efba2f9acddc37bfb127f53" dependencies = [ "http 1.5.0", "hyper", @@ -3533,7 +3533,7 @@ dependencies = [ "libc", "percent-encoding", "pin-project-lite", - "socket2 0.5.10", + "socket2 0.6.5", "tokio", "tower-service", "tracing", @@ -3884,7 +3884,7 @@ checksum = "65b27460c2c92b037f3f94c538ed9a3342f3fdf923606781629ccb35f82d042a" [[package]] name = "imap" -version = "0.16.22" +version = "0.16.23" dependencies = [ "ahash", "common", @@ -3897,7 +3897,7 @@ dependencies = [ "md5", "nlp", "parking_lot", - "rand 0.10.2", + "rand 0.10.3", "registry", "store", "tokio", @@ -3909,7 +3909,7 @@ dependencies = [ [[package]] name = "imap_proto" -version = "0.16.22" +version = "0.16.23" dependencies = [ "ahash", "base64 0.23.1", @@ -3924,7 +3924,7 @@ dependencies = [ [[package]] name = "inbuxa" -version = "0.16.22" +version = "0.16.23" dependencies = [ "common", "coordinator", @@ -3932,7 +3932,7 @@ dependencies = [ "directory", "email", "groupware", - "http 0.16.22", + "http 0.16.23", "http_proto", "imap", "jmap", @@ -4134,25 +4134,24 @@ checksum = "4d3667095d64c3ecffc96463a21157b04bf3e252f6e8d5750b20c02e33c194e3" [[package]] name = "jieba-macros" -version = "0.10.3" +version = "0.10.4" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "34904340bc65749a9e9a02fcc7f3368e675427c18447b9bbe02df52c15c9a36a" +checksum = "455f837e9d0255b68a712200db247c68fdad4941b72471b76bfa61c3b0c1f79f" dependencies = [ "phf_codegen", ] [[package]] name = "jieba-rs" -version = "0.10.3" +version = "0.10.4" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "bb5bdea4dc241d589e179f39d2a778f31490f3370aa2f626223dbd930ebc5c9d" +checksum = "b6a8bbb0f77ee810f0689a30b7cec56b875751ef4ec2e74fd995613dc52b3ae1" dependencies = [ "bytecount", "cedarwood", "include-flate", "jieba-macros", "phf 0.13.1", - "regex", "rustc-hash", ] @@ -4212,7 +4211,7 @@ dependencies = [ [[package]] name = "jmap" -version = "0.16.22" +version = "0.16.23" dependencies = [ "async-stream", "base64 0.23.1", @@ -4236,7 +4235,7 @@ dependencies = [ "mail-parser", "nlp", "p256", - "rand 0.10.2", + "rand 0.10.3", "registry", "reqwest 0.13.5", "rkyv", @@ -4294,7 +4293,7 @@ dependencies = [ [[package]] name = "jmap_proto" -version = "0.16.22" +version = "0.16.23" dependencies = [ "ahash", "calcard", @@ -4699,9 +4698,9 @@ dependencies = [ [[package]] name = "lru-slab" -version = "0.1.2" +version = "0.1.3" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "112b39cec0b298b6c1999fee3e31427f74f676e4cb9879ed1a121b43661a4154" +checksum = "4050469837a6ff301cd14c1f8f24f88549e6d548f24f64e2148eb0f72cebc51f" [[package]] name = "lz4-sys" @@ -4742,9 +4741,9 @@ dependencies = [ [[package]] name = "mail-auth" -version = "0.13.2" +version = "0.13.3" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "e11f19d98aac923fc5b7ee30c3509733a013ef546a226acb959b9202f5ca58f0" +checksum = "8505122ba86e1f4adeb664196c1e787c3f29bb6e7c128e4a366d47d209911440" dependencies = [ "aws-lc-rs", "flate2", @@ -4757,7 +4756,7 @@ dependencies = [ "mail-parser", "memchr", "quick-xml 0.42.0", - "rand 0.10.2", + "rand 0.10.3", "rkyv", "rsa", "rustls-pki-types", @@ -4800,7 +4799,7 @@ dependencies = [ [[package]] name = "managesieve" -version = "0.16.22" +version = "0.16.23" dependencies = [ "common", "compact_str", @@ -4935,7 +4934,7 @@ checksum = "c797b9d6bb23aab2fc369c65f871be49214f5c759af65bde26ffaaa2b646b492" [[package]] name = "migration" -version = "0.16.22" +version = "0.16.23" dependencies = [ "common", "email", @@ -5066,7 +5065,7 @@ dependencies = [ "proc-macro2", "quote", "rustversion", - "syn 3.0.5", + "syn 3.0.6", ] [[package]] @@ -5131,7 +5130,7 @@ dependencies = [ "lru", "mysql_common", "percent-encoding", - "rand 0.10.2", + "rand 0.10.3", "rustls", "serde", "socket2 0.6.5", @@ -5206,14 +5205,14 @@ dependencies = [ [[package]] name = "nlp" -version = "0.16.22" +version = "0.16.23" dependencies = [ "ahash", "hashify", "jieba-rs", "maplit", "psl", - "rand 0.10.2", + "rand 0.10.3", "rkyv", "rust-stemmers", "serde", @@ -6038,7 +6037,7 @@ dependencies = [ [[package]] name = "pop3" -version = "0.16.22" +version = "0.16.23" dependencies = [ "common", "directory", @@ -6082,7 +6081,7 @@ dependencies = [ "hmac 0.13.0", "md-5 0.11.0", "memchr", - "rand 0.10.2", + "rand 0.10.3", "sha2 0.11.0", "stringprep", ] @@ -6119,9 +6118,9 @@ checksum = "439ee305def115ba05938db6eb1644ff94165c5ab5e9420d1c1bcedbba909391" [[package]] name = "ppmd-rust" -version = "1.4.1" +version = "1.5.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "9e9219bcb9d7aca6b2f63c83cf100cf78bcd619ac46e6ecbd0dd90869a39345d" +checksum = "196a7c80b9a7652aba7cc070827516c2abe4ccdf53d128e1944003cf5726cff1" [[package]] name = "ppv-lite86" @@ -6206,7 +6205,7 @@ dependencies = [ "proc-macro-error-attr3", "proc-macro2", "quote", - "syn 3.0.5", + "syn 3.0.6", ] [[package]] @@ -6260,7 +6259,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "b570b25f7617e43d59005d0990ccb79e950a423952cea19671b7a876da390adf" dependencies = [ "anyhow", - "itertools 0.13.0", + "itertools 0.14.0", "proc-macro2", "quote", "syn 2.0.119", @@ -6287,9 +6286,9 @@ dependencies = [ [[package]] name = "psl" -version = "2.1.232" +version = "2.1.235" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "62834e308cc83aea5e30cd8c80b8aa82cdb104a3240c7f210d4f68d46e29f308" +checksum = "8319b56ff38ca0522b4e1e40bfa2b5de7f62dc89fc1e9033eac365551ec58e0e" dependencies = [ "psl-types", ] @@ -6317,7 +6316,7 @@ checksum = "1c8d9ca532f185d5d4db7a7c9d51420b452168ea1c2b913953281bd6fe1fcbd0" dependencies = [ "proc-macro2", "quote", - "syn 3.0.5", + "syn 3.0.6", ] [[package]] @@ -6388,9 +6387,9 @@ dependencies = [ [[package]] name = "quinn" -version = "0.11.11" +version = "0.11.12" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "0c1a41e437b6bbd489372cd4971de128e85c855f56c57f283d20ff016cf7c0a8" +checksum = "4051e23e9185c255a7e33ef59cdbca87a22d359052eecd22fc6b901fb37d9d11" dependencies = [ "bytes", "cfg_aliases", @@ -6399,7 +6398,7 @@ dependencies = [ "quinn-udp", "rustc-hash", "rustls", - "socket2 0.5.10", + "socket2 0.6.5", "thiserror 2.0.20", "tokio", "tracing", @@ -6408,16 +6407,16 @@ dependencies = [ [[package]] name = "quinn-proto" -version = "0.11.17" +version = "0.11.18" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "04759210543be93709136e28212294a659ef5001836ff4eab4d663e4529bba83" +checksum = "a9746dbde176634f4f2f1faf2404e30a31b2bc1e9cafb5329c95d8177a18c9fc" dependencies = [ "aws-lc-rs", "bytes", "fastbloom", "getrandom 0.4.3", "lru-slab", - "rand 0.10.2", + "rand 0.10.3", "rand_pcg", "ring", "rustc-hash", @@ -6440,7 +6439,7 @@ dependencies = [ "cfg_aliases", "libc", "once_cell", - "socket2 0.5.10", + "socket2 0.6.5", "tracing", "windows-sys 0.61.2", ] @@ -6534,9 +6533,9 @@ dependencies = [ [[package]] name = "rand" -version = "0.10.2" +version = "0.10.3" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "c7f5fa3a058cd35567ef9bfa5e75732bee0f9e4c55fa90477bef2dfcdbc4be80" +checksum = "65c9fb96cbc91e3478eaae79a69fcd3f1ae4ad052e471fe6732fff548984b4af" dependencies = [ "chacha20", "getrandom 0.4.3", @@ -6776,7 +6775,7 @@ dependencies = [ "num-bigint 0.5.1", "percent-encoding", "pin-project-lite", - "rand 0.10.2", + "rand 0.10.3", "rustls", "rustls-native-certs", "ryu", @@ -6800,11 +6799,10 @@ dependencies = [ [[package]] name = "redox_users" -version = "0.5.2" +version = "0.5.3" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "a4e608c6638b9c18977b00b475ac1f28d14e84b27d8d42f70e0bf1e3dec127ac" +checksum = "60dc65c0ff1a7ae1294b0c67b9f14baf70b644404010370171787bfac1038fc0" dependencies = [ - "getrandom 0.2.17", "libredox", "thiserror 2.0.20", ] @@ -6826,7 +6824,7 @@ checksum = "92ecd8964f8453721699a1ed72037b0db49ce2f5a5138486ee89bed6f67cdf3a" dependencies = [ "proc-macro2", "quote", - "syn 3.0.5", + "syn 3.0.6", ] [[package]] @@ -6860,7 +6858,7 @@ checksum = "d6f6ff9a378485b298a5286656da665ba74413d36db0979633275d2e708145d4" [[package]] name = "registry" -version = "0.16.22" +version = "0.16.23" dependencies = [ "ahash", "hashify", @@ -7044,7 +7042,7 @@ checksum = "1c25ef604ac7dd839d44d64648952ea23c97866f124ff671b0ed2cf3ad9bb06e" dependencies = [ "proc-macro2", "quote", - "syn 3.0.5", + "syn 3.0.6", ] [[package]] @@ -7225,9 +7223,9 @@ dependencies = [ [[package]] name = "rustix" -version = "1.1.4" +version = "1.1.5" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "b6fe4565b9518b83ef4f91bb47ce29620ca828bd32cb7e408f0062e9930ba190" +checksum = "891efababe418670775f199f0d233d84843c227a0949a883ce15b37c78d6629d" dependencies = [ "bitflags 2.13.2", "errno", @@ -7412,12 +7410,12 @@ dependencies = [ "proc-macro2", "quote", "serde_derive_internals", - "syn 3.0.5", + "syn 3.0.6", ] [[package]] name = "scim" -version = "0.16.22" +version = "0.16.23" dependencies = [ "ahash", "base64 0.23.1", @@ -7443,7 +7441,7 @@ dependencies = [ [[package]] name = "scim-proto" -version = "0.16.22" +version = "0.16.23" dependencies = [ "hashify", "serde", @@ -7580,7 +7578,7 @@ dependencies = [ "sha2 0.10.9", "sha3 0.10.9", "slh-dsa", - "thiserror 1.0.69", + "thiserror 2.0.20", "twofish", "typenum", "x25519-dalek", @@ -7624,7 +7622,7 @@ checksum = "e7a5d71263a5a7d47b41f6b3f06ba276f10cc18b0931f1799f710578e2309348" dependencies = [ "proc-macro2", "quote", - "syn 3.0.5", + "syn 3.0.6", ] [[package]] @@ -7635,7 +7633,7 @@ checksum = "f852137cce035d6a4df67ccce505ff6b3e9fd3a10e3e52b24dc71e650bb1a9bd" dependencies = [ "proc-macro2", "quote", - "syn 3.0.5", + "syn 3.0.6", ] [[package]] @@ -7671,7 +7669,7 @@ checksum = "8d3b1629de253c70a0508c3899572da79ca359fdab27c7920ff00406df418906" dependencies = [ "proc-macro2", "quote", - "syn 3.0.5", + "syn 3.0.6", ] [[package]] @@ -7716,7 +7714,7 @@ dependencies = [ "darling 0.24.1", "proc-macro2", "quote", - "syn 3.0.5", + "syn 3.0.6", ] [[package]] @@ -7764,12 +7762,12 @@ checksum = "a22144e767da4ddd8416dbf383700542ffd8a5dc493dfecedfe1fe3ad03c98ae" dependencies = [ "proc-macro2", "quote", - "syn 3.0.5", + "syn 3.0.6", ] [[package]] name = "services" -version = "0.16.22" +version = "0.16.23" dependencies = [ "aes-gcm 0.11.1", "aho-corasick", @@ -8084,7 +8082,7 @@ checksum = "ba467056f1b547ed52077911161fc86985becbc60e8e1857c8a144dab0def891" [[package]] name = "smtp" -version = "0.16.22" +version = "0.16.23" dependencies = [ "ahash", "base64 0.23.1", @@ -8099,7 +8097,7 @@ dependencies = [ "mail-builder 1.0.0", "mail-parser", "parking_lot", - "rand 0.10.2", + "rand 0.10.3", "registry", "reqwest 0.13.5", "rkyv", @@ -8175,7 +8173,7 @@ dependencies = [ [[package]] name = "spam-filter" -version = "0.16.22" +version = "0.16.23" dependencies = [ "common", "compact_str", @@ -8295,7 +8293,7 @@ checksum = "a2eb9349b6444b326872e140eb1cf5e7c522154d69e7a0ffb0fb81c06b37543f" [[package]] name = "store" -version = "0.16.22" +version = "0.16.23" dependencies = [ "ahash", "arc-swap", @@ -8321,7 +8319,7 @@ dependencies = [ "parking_lot", "r2d2", "radsort", - "rand 0.10.2", + "rand 0.10.3", "rayon", "redis", "registry", @@ -8417,9 +8415,9 @@ dependencies = [ [[package]] name = "syn" -version = "3.0.5" +version = "3.0.6" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "12df2e0110f65b775f769bb17ef989067a1d931b2eb822bd4346631eeada89f9" +checksum = "8593e8e72159ed2257d083c7a454a85cbf854f37a0966d8d483aff8c8a3ebcee" dependencies = [ "proc-macro2", "quote", @@ -8446,6 +8444,17 @@ dependencies = [ "syn 2.0.119", ] +[[package]] +name = "synstructure" +version = "0.14.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "901704edd0dfe137f1987838ee4f259e4e063c31371bdb423f7ae38ec6f77f02" +dependencies = [ + "proc-macro2", + "quote", + "syn 3.0.6", +] + [[package]] name = "sysinfo" version = "0.37.2" @@ -8544,7 +8553,7 @@ dependencies = [ [[package]] name = "tests" -version = "0.16.22" +version = "0.16.23" dependencies = [ "ahash", "aws-lc-rs", @@ -8566,7 +8575,7 @@ dependencies = [ "form_urlencoded", "futures", "groupware", - "http 0.16.22", + "http 0.16.23", "http_proto", "hyper", "hyper-util", @@ -8652,7 +8661,7 @@ checksum = "bc04cd3e1236dd4a98afca4569f2deb3f120e5422a4023be2cb683f8486292af" dependencies = [ "proc-macro2", "quote", - "syn 3.0.5", + "syn 3.0.6", ] [[package]] @@ -8790,7 +8799,7 @@ checksum = "78773a2a397f451582ce068015985c33193cf6dea8b74d2a639fe457b2f07b0e" dependencies = [ "proc-macro2", "quote", - "syn 3.0.5", + "syn 3.0.6", ] [[package]] @@ -8812,7 +8821,7 @@ dependencies = [ "pin-project-lite", "postgres-protocol", "postgres-types", - "rand 0.10.2", + "rand 0.10.3", "socket2 0.6.5", "tokio", "tokio-util", @@ -8997,7 +9006,7 @@ dependencies = [ "constant_time_eq", "hmac 0.13.0", "percent-encoding", - "rand 0.10.2", + "rand 0.10.3", "serde", "sha1 0.11.0", "sha2 0.11.0", @@ -9136,7 +9145,7 @@ dependencies = [ [[package]] name = "trc" -version = "0.16.22" +version = "0.16.23" dependencies = [ "ahash", "base64 0.23.1", @@ -9195,7 +9204,7 @@ dependencies = [ "http 1.5.0", "httparse", "log", - "rand 0.10.2", + "rand 0.10.3", "sha1 0.11.0", "thiserror 2.0.20", ] @@ -9245,7 +9254,7 @@ checksum = "b6f5e870be6c3b371b77fe0ee0bafb859fa4964b4404c27de1d380043c4dda20" [[package]] name = "types" -version = "0.16.22" +version = "0.16.23" dependencies = [ "blake3", "compact_str", @@ -9297,9 +9306,9 @@ checksum = "0b993bddc193ae5bd0d623b49ec06ac3e9312875fdae725a975c51db1cc1677f" [[package]] name = "unicode-ident" -version = "1.0.24" +version = "1.0.26" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "e6e4313cd5fcd3dad5cafa179702e2b244f760991f45397d14d4ebf38247da75" +checksum = "d245f478577f809a851594d02313b640fb437e0bb33866753cff937863096954" [[package]] name = "unicode-normalization" @@ -9414,7 +9423,7 @@ checksum = "b6c140620e7ffbb22c2dee59cafe6084a59b5ffc27a8859a5f0d494b5d52b6be" [[package]] name = "utils" -version = "0.16.22" +version = "0.16.23" dependencies = [ "ahash", "arcstr", @@ -9616,7 +9625,7 @@ dependencies = [ "bumpalo", "proc-macro2", "quote", - "syn 3.0.5", + "syn 3.0.6", "wasm-bindgen-shared", ] @@ -10125,14 +10134,14 @@ dependencies = [ [[package]] name = "yoke-derive" -version = "0.8.2" +version = "0.8.3" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "de844c262c8848816172cef550288e7dc6c7b7814b4ee56b3e1553f275f1858e" +checksum = "33811428bee40dbceb6d545e95754741d17a6aef9a4849f0fd62e2ba4f412a78" dependencies = [ "proc-macro2", "quote", - "syn 2.0.119", - "synstructure", + "syn 3.0.6", + "synstructure 0.14.0", ] [[package]] @@ -10655,14 +10664,14 @@ dependencies = [ [[package]] name = "zerofrom-derive" -version = "0.1.7" +version = "0.1.8" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "11532158c46691caf0f2593ea8358fed6bbf68a0315e80aae9bd41fbade684a1" +checksum = "f75b4683f6c7f45248d4d64056a24298c6281e0993356d7d1b4a1a962ef10d4a" dependencies = [ "proc-macro2", "quote", - "syn 2.0.119", - "synstructure", + "syn 3.0.6", + "synstructure 0.14.0", ] [[package]] @@ -10717,7 +10726,7 @@ checksum = "34df6fc39dbd26ddc9c10e6a2984476e13acce22e64e4487636ef494369225da" dependencies = [ "proc-macro2", "quote", - "syn 3.0.5", + "syn 3.0.6", ] [[package]] @@ -10749,9 +10758,9 @@ dependencies = [ [[package]] name = "zlib-rs" -version = "0.6.7" +version = "0.6.8" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "34b31d188d9d685a4f9c7b46d6e36631b07058d2cfe190267adce54dc230bf12" +checksum = "b268e58e7c693d7c271f93ffc4ba3b380412554231c85bf61ca7af91042a4112" [[package]] name = "zmij" diff --git a/crates/common/Cargo.toml b/crates/common/Cargo.toml index 99dae45..37bd7d1 100644 --- a/crates/common/Cargo.toml +++ b/crates/common/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "common" -version = "0.16.22" +version = "0.16.23" edition = "2024" build = "build.rs" diff --git a/crates/common/src/auth/authentication.rs b/crates/common/src/auth/authentication.rs index bcec345..fc7f354 100644 --- a/crates/common/src/auth/authentication.rs +++ b/crates/common/src/auth/authentication.rs @@ -11,7 +11,7 @@ use crate::{ auth::{ AccessToken, AuthRequest, DomainCache, credential::{ApiKey, AppPassword}, - oauth::GrantType, + oauth::{GrantType, token::TOKEN_HEADER}, }, }; use base64::{Engine, engine::general_purpose}; @@ -23,7 +23,8 @@ use registry::schema::{ enums::Permission, structs::{self, Credential}, }; -use std::{net::IpAddr, sync::Arc}; +use serde::Deserialize; +use std::{borrow::Cow, net::IpAddr, sync::Arc}; use store::write::now; use trc::AddContext; @@ -321,19 +322,12 @@ impl Server { // Obtain external directory, if any. When no username is supplied // (e.g. HTTP bearer auth), peek at the JWT claims to find the // user's domain so per-domain OIDC directories are reachable. - let directory = if let Some(username) = username.as_deref().map(UsernameParts::new) - { - if let Some(domain_name) = username.auth_as().domain() { - self.get_directory_for_domain(domain_name).await? - } else if let Some(domain_name) = extract_jwt_domain(token) { - self.get_directory_for_domain(&domain_name).await? - } else { - self.get_default_directory() - } - } else if let Some(domain_name) = extract_jwt_domain(token) { - self.get_directory_for_domain(&domain_name).await? - } else { - self.get_default_directory() + let directory = match username.as_deref().map(UsernameParts::new) { + Some(username) => match username.auth_as().domain() { + Some(domain_name) => self.get_directory_for_domain(domain_name).await?, + None => self.get_directory_for_token(token).await?, + }, + None => self.get_directory_for_token(token).await?, }; // Try external directory authentication first if supported, then fallback to internal OAuth. @@ -563,6 +557,29 @@ impl Server { }) } + async fn get_directory_for_token(&self, token: &str) -> trc::Result>> { + let Some(payload) = JwtClaims::decode_payload(token) else { + return Ok(self.get_default_directory()); + }; + let Some(claims) = JwtClaims::parse(&payload) else { + return Ok(self.get_default_directory()); + }; + + match (claims.domain(), claims.iss.as_deref()) { + (Some(domain_name), _) => self.get_directory_for_domain(domain_name).await, + (None, Some(issuer)) => Ok(self + .get_directory_for_issuer(issuer) + .or_else(|| self.get_default_directory())), + (None, None) => Ok(self.get_default_directory()), + } + } + + /// inbuxa: DIR-2: a token naming no address gets the server default, so + /// no directory is chosen by issuer. + fn get_directory_for_issuer(&self, _issuer: &str) -> Option<&Arc> { + None + } + /// inbuxa: DIR-1, DIR-5: as above, for a domain already read. A /// `directoryId` naming no directory the server built is unavailable, /// never the internal directory. @@ -622,25 +639,50 @@ pub fn unavailable_directory() -> &'static Arc { }) } -fn extract_jwt_domain(token: &str) -> Option { - let mut parts = token.split('.'); - let _header = parts.next()?; - let payload = parts.next()?; - let _signature = parts.next()?; - if parts.next().is_some() { - return None; - } - let payload_bytes = general_purpose::URL_SAFE_NO_PAD.decode(payload).ok()?; - let claims: serde_json::Value = serde_json::from_slice(&payload_bytes).ok()?; - for claim in ["email", "preferred_username", "upn"] { - if let Some(val) = claims.get(claim).and_then(|v| v.as_str()) - && let Some((_, domain)) = val.rsplit_once('@') - && !domain.is_empty() - { - return Some(domain.to_ascii_lowercase()); +#[derive(Deserialize)] +struct JwtClaims<'x> { + #[serde(borrow, default)] + iss: Option>, + #[serde(borrow, default)] + email: Option>, + #[serde(borrow, default)] + preferred_username: Option>, + #[serde(borrow, default)] + upn: Option>, +} + +impl<'x> JwtClaims<'x> { + fn decode_payload(token: &str) -> Option> { + if token.starts_with(TOKEN_HEADER) { + return None; } + + let mut parts = token.split('.'); + let _header = parts.next()?; + let payload = parts.next()?; + let _signature = parts.next()?; + if parts.next().is_some() { + return None; + } + + general_purpose::URL_SAFE_NO_PAD.decode(payload).ok() + } + + fn parse(payload: &'x [u8]) -> Option { + serde_json::from_slice(payload).ok() + } + + fn domain(&self) -> Option<&str> { + [&self.email, &self.preferred_username, &self.upn] + .into_iter() + .flatten() + .find_map(|claim| { + claim + .rsplit_once('@') + .map(|(_, domain)| domain) + .filter(|domain| !domain.is_empty()) + }) } - None } impl UsernameParts { @@ -738,3 +780,76 @@ impl AuthRequest { } } } + +#[cfg(test)] +mod tests { + use super::*; + + fn jwt(payload: &str) -> String { + format!( + "eyJhbGciOiJSUzI1NiJ9.{}.c2lnbmF0dXJl", + general_purpose::URL_SAFE_NO_PAD.encode(payload) + ) + } + + fn hints(token: &str) -> Option<(Option, Option)> { + let payload = JwtClaims::decode_payload(token)?; + let claims = JwtClaims::parse(&payload)?; + + Some(( + claims.domain().map(str::to_string), + claims.iss.as_deref().map(str::to_string), + )) + } + + #[test] + fn jwt_claims_are_extracted() { + for (payload, domain, issuer) in [ + ( + r#"{"iss":"https://idp.example.org","email":"John@Example.ORG"}"#, + Some("Example.ORG"), + Some("https://idp.example.org"), + ), + ( + r#"{"preferred_username":"jane@example.net","upn":"jane@example.com"}"#, + Some("example.net"), + None, + ), + ( + r#"{"email":"broken@","upn":"jane@example.com"}"#, + Some("example.com"), + None, + ), + ( + r#"{"iss":"https://idp.example.org","sub":"5db2d1b6","aud":["a","b"],"scope":"openid"}"#, + None, + Some("https://idp.example.org"), + ), + (r#"{"sub":"5db2d1b6"}"#, None, None), + (r#"{"email":"jane@example.net"}"#, Some("example.net"), None), + ] { + assert_eq!( + hints(&jwt(payload)), + Some((domain.map(str::to_string), issuer.map(str::to_string))), + "Unexpected claims for {payload}" + ); + } + } + + #[test] + fn non_jwt_tokens_are_ignored() { + for token in [ + "sw1.eyJhbGciOiJSUzI1NiJ9.eyJpc3MiOiJodHRwczovL2lkcC5leGFtcGxlLm9yZyJ9", + "sw1.eyJhbGciOiJSUzI1NiJ9", + "opaque-token", + "one.two", + "one.two.three.four", + "", + ] { + assert!( + JwtClaims::decode_payload(token).is_none(), + "Token {token:?} was parsed as a JWT" + ); + } + } +} diff --git a/crates/common/src/auth/oauth/token.rs b/crates/common/src/auth/oauth/token.rs index cc4ef4f..044c4a7 100644 --- a/crates/common/src/auth/oauth/token.rs +++ b/crates/common/src/auth/oauth/token.rs @@ -17,7 +17,7 @@ pub const FAILED_TO_DECODE_TOKEN: &str = concat!( "the Authentication object." ); -const TOKEN_HEADER: &str = "sw1."; +pub(crate) const TOKEN_HEADER: &str = "sw1."; const TOKEN_KEY_CONTEXT: &str = "stalwart-oauth-token-sw1"; const OAUTH_EPOCH: u64 = 946684800; // Jan 1, 2000 diff --git a/crates/common/src/config/smtp/resolver.rs b/crates/common/src/config/smtp/resolver.rs index cedf560..a23a51f 100644 --- a/crates/common/src/config/smtp/resolver.rs +++ b/crates/common/src/config/smtp/resolver.rs @@ -214,6 +214,7 @@ impl Resolvers { let config_dnssec = resolver_config.clone(); let mut opts_dnssec = opts.clone(); opts_dnssec.validate = true; + opts_dnssec.num_concurrent_reqs = 1; let dnssec = DnssecResolver { resolver: TokioResolver::builder_with_config( @@ -343,6 +344,7 @@ impl Default for Resolvers { let config_dnssec = config.clone(); let mut opts_dnssec = opts.clone(); opts_dnssec.validate = true; + opts_dnssec.num_concurrent_reqs = 1; Self { dns: MessageAuthenticator::new(config, opts).expect("Failed to build DNS resolver"), diff --git a/crates/common/src/expr/functions/misc.rs b/crates/common/src/expr/functions/misc.rs index 924a625..a40407e 100644 --- a/crates/common/src/expr/functions/misc.rs +++ b/crates/common/src/expr/functions/misc.rs @@ -23,6 +23,13 @@ pub(crate) fn fn_is_number(v: Vec) -> Variable { matches!(&v[0], Variable::Integer(_) | Variable::Float(_)).into() } +pub(crate) fn fn_bit_and(v: Vec) -> Variable { + match (v[0].to_integer(), v[1].to_integer()) { + (Some(lhs), Some(rhs)) => Variable::Integer(lhs & rhs), + _ => Variable::Integer(0), + } +} + pub(crate) fn fn_is_ip_addr(v: Vec) -> Variable { v[0].to_string() .as_str() diff --git a/crates/common/src/expr/functions/mod.rs b/crates/common/src/expr/functions/mod.rs index 78cc03e..0c4a13c 100644 --- a/crates/common/src/expr/functions/mod.rs +++ b/crates/common/src/expr/functions/mod.rs @@ -46,6 +46,7 @@ pub(crate) const FUNCTIONS: &[(&str, fn(Vec) -> Variable, u32)] = &[ ("email_part", email::fn_email_part, 2), ("is_empty", misc::fn_is_empty, 1), ("is_number", misc::fn_is_number, 1), + ("bit_and", misc::fn_bit_and, 2), ("is_ip_addr", misc::fn_is_ip_addr, 1), ("is_ipv4_addr", misc::fn_is_ipv4_addr, 1), ("is_ipv6_addr", misc::fn_is_ipv6_addr, 1), diff --git a/crates/common/src/manager/application.rs b/crates/common/src/manager/application.rs index a5f5370..b3eb6a3 100644 --- a/crates/common/src/manager/application.rs +++ b/crates/common/src/manager/application.rs @@ -2,6 +2,8 @@ * SPDX-FileCopyrightText: 2020 Stalwart Labs LLC * * SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-SEL + * + * Modified by Coffey Labs in 2026 for INBUXA. */ use crate::{Server, manager::fetch_resource}; @@ -11,8 +13,11 @@ use registry::schema::{enums::CompressionAlgo, structs::Application}; use std::{ borrow::Cow, io::{self, Cursor, Read}, - path::PathBuf, - sync::Arc, + path::{Path, PathBuf}, + sync::{ + Arc, + atomic::{AtomicU64, Ordering}, + }, time::Duration, }; use store::{ @@ -36,16 +41,18 @@ enum IndexEdit<'x> { pub struct WebApplications { applications: ArcSwap>, routes: ArcSwap>>, + generation: AtomicU64, } pub struct AppRoutes { resources: AHashMap>, oauth_client_id_meta: Option, + _bundle_dir: TempDir, } #[derive(Clone)] pub struct WebApplicationManager { - bundle_path: TempDir, + base_path: PathBuf, prefixes: Vec, description: String, url: String, @@ -79,6 +86,7 @@ impl WebApplications { Self { applications: ArcSwap::new(Arc::new(Vec::new())), routes: ArcSwap::new(Arc::new(AHashMap::new())), + generation: AtomicU64::new(0), } } @@ -128,48 +136,55 @@ impl WebApplications { } pub async fn unpack_all(&self, server: &Server, update: bool) { - let mut routes = AHashMap::new(); + let previous = self.routes.load_full(); + let sweep_orphans = previous.is_empty(); + let mut routes = AHashMap::with_capacity(previous.len()); + for app in self.applications.load().as_ref() { - if update && let Err(err) = app.delete(server).await { - trc::event!( - Resource(trc::ResourceEvent::Error), - Reason = err, - Url = app.url.clone(), - Details = format!( - "Failed to delete application bundle for prefixes: {}", - app.prefixes.join(", ") - ) - ); - } - match app.unpack(server).await { - Ok(resources) => { - let app_routes = Arc::new(AppRoutes { - resources, - oauth_client_id_meta: app - .oauth_client_id - .as_deref() - .map(oauth_client_id_meta), - }); + match app + .unpack(server, self.next_generation(), update, sweep_orphans) + .await + { + Ok(app_routes) => { + let app_routes = Arc::new(app_routes); for prefix in &app.prefixes { routes.insert(prefix.clone(), app_routes.clone()); } } Err(err) => { + let mut is_retained = false; + for prefix in &app.prefixes { + if let Some(app_routes) = previous.get(prefix) { + routes.insert(prefix.clone(), app_routes.clone()); + is_retained = true; + } + } + trc::event!( Resource(trc::ResourceEvent::Error), Reason = err, Url = app.url.clone(), Details = format!( - "Failed to unpack application for prefixes: {}", - app.prefixes.join(", ") + "Failed to unpack application for prefixes: {}, {}", + app.prefixes.join(", "), + if is_retained { + "the previously unpacked bundle remains in service" + } else { + "no bundle is available to serve" + } ) ); } } } + self.routes.store(Arc::new(routes)); } + + fn next_generation(&self) -> u64 { + self.generation.fetch_add(1, Ordering::Relaxed) + } } impl WebApplicationManager { @@ -182,7 +197,7 @@ impl WebApplicationManager { .join(app.id.id().to_string()); Self { - bundle_path: TempDir::new(base_path), + base_path, blob_key: BlobHash::generate(format!("{}{}", APP_BLOB_PREFIX, app.id.id()).as_bytes()), url: app.object.resource_url, description: app.object.description, @@ -202,82 +217,43 @@ impl WebApplicationManager { } } - async fn unpack(&self, server: &Server) -> trc::Result>> { - // Delete any existing bundles - self.bundle_path.clean().await.map_err(unpack_error)?; - - // Obtain application bundle - let bundle = if let Some(bundle) = server - .blob_store() - .get_blob(self.blob_key.as_slice(), 0..usize::MAX) - .await? - { - bundle + async fn unpack( + &self, + server: &Server, + generation: u64, + force_refresh: bool, + sweep_orphans: bool, + ) -> trc::Result { + let cached = if force_refresh { + None } else { - // Fetch app bundle - let resource = fetch_resource(&self.url, None, Duration::from_secs(60), MAX_APP_SIZE) - .await - .map_err(|err| { - trc::ResourceEvent::Error - .caused_by(trc::location!()) - .ctx(Key::Url, self.url.clone()) - .reason(err) - .details("Failed to fetch application bundle") - })?; - - // Store in blob store for future use server .blob_store() - .put_blob(self.blob_key.as_slice(), &resource, CompressionAlgo::None) - .await - .caused_by(trc::location!())?; - - // Schedule expiration - let mut batch = BatchBuilder::new(); - batch - .set( - BlobOp::Link { - hash: self.blob_key.clone(), - to: BlobLink::Temporary { - until: now() + self.expiry, - }, - }, - vec![], - ) - .set( - BlobOp::Commit { - hash: self.blob_key.clone(), - }, - Vec::new(), - ); - server - .store() - .write(batch.build_all()) - .await - .caused_by(trc::location!())?; - - trc::event!( - Resource(trc::ResourceEvent::ApplicationUpdated), - Url = self.url.clone(), - Details = self.description.clone(), - ); - - resource + .get_blob(self.blob_key.as_slice(), 0..usize::MAX) + .await? + }; + let is_cached = cached.is_some(); + let bundle = match cached { + Some(bundle) => bundle, + None => self.fetch().await?, }; + let staging = TempDir::new(self.base_path.join(format!("{:x}-{generation:x}", now()))); + staging.create().await.map_err(unpack_error)?; + let url = self.url.clone(); - let bundle_path = self.bundle_path.path.clone(); - let routes = tokio::task::spawn_blocking(move || -> trc::Result<_> { - let mut bundle = zip::ZipArchive::new(Cursor::new(bundle)).map_err(|err| { + let bundle_path = staging.path.clone(); + let (resources, bundle) = tokio::task::spawn_blocking(move || -> trc::Result<_> { + let mut archive = zip::ZipArchive::new(Cursor::new(bundle)).map_err(|err| { trc::ResourceEvent::Error .caused_by(trc::location!()) .reason(err) .ctx(Key::Url, url.clone()) .details("Failed to decompress application bundle") })?; - let mut routes = AHashMap::new(); - for i in 0..bundle.len() { - let mut file = bundle.by_index(i).map_err(|err| { + let mut resources = AHashMap::with_capacity(archive.len()); + for i in 0..archive.len() { + let mut file = archive.by_index(i).map_err(|err| { trc::ResourceEvent::Error .caused_by(trc::location!()) .reason(err) @@ -315,9 +291,9 @@ impl WebApplicationManager { contents: path, }; - routes.insert(file_name, resource); + resources.insert(file_name, resource); } - Ok(routes) + Ok((resources, archive.into_inner().into_inner())) }) .await .map_err(|err| { @@ -327,21 +303,81 @@ impl WebApplicationManager { .details("Bundle unpack task panicked") })??; + if !is_cached && let Err(err) = self.cache(server, &bundle).await { + trc::event!( + Resource(trc::ResourceEvent::Error), + Reason = err, + Url = self.url.clone(), + Details = "Failed to cache application bundle, it will be downloaded again" + ); + } + + if sweep_orphans { + remove_siblings(&self.base_path, &staging.path).await; + } + trc::event!( Resource(trc::ResourceEvent::ApplicationUnpacked), Url = self.url.clone(), - Path = self.bundle_path.path.to_string_lossy().into_owned(), + Path = staging.path.to_string_lossy().into_owned(), ); - Ok(routes) + Ok(AppRoutes { + resources, + oauth_client_id_meta: self.oauth_client_id.as_deref().map(oauth_client_id_meta), + _bundle_dir: staging, + }) } - async fn delete(&self, server: &Server) -> trc::Result<()> { + async fn fetch(&self) -> trc::Result> { + fetch_resource(&self.url, None, Duration::from_secs(60), MAX_APP_SIZE) + .await + .map_err(|err| { + trc::ResourceEvent::Error + .caused_by(trc::location!()) + .ctx(Key::Url, self.url.clone()) + .reason(err) + .details("Failed to fetch application bundle") + }) + } + + async fn cache(&self, server: &Server, bundle: &[u8]) -> trc::Result<()> { server .blob_store() - .delete_blob(self.blob_key.as_slice()) + .put_blob(self.blob_key.as_slice(), bundle, CompressionAlgo::None) .await - .map(|_| ()) + .caused_by(trc::location!())?; + + let mut batch = BatchBuilder::new(); + batch + .set( + BlobOp::Link { + hash: self.blob_key.clone(), + to: BlobLink::Temporary { + until: now() + self.expiry, + }, + }, + vec![], + ) + .set( + BlobOp::Commit { + hash: self.blob_key.clone(), + }, + Vec::new(), + ); + server + .store() + .write(batch.build_all()) + .await + .caused_by(trc::location!())?; + + trc::event!( + Resource(trc::ResourceEvent::ApplicationUpdated), + Url = self.url.clone(), + Details = self.description.clone(), + ); + + Ok(()) } pub async fn delete_bundle(server: &Server, app_id: Id) -> trc::Result<()> { @@ -361,7 +397,6 @@ impl Resource> { } } -#[derive(Clone)] pub struct TempDir { pub path: PathBuf, } @@ -371,11 +406,36 @@ impl TempDir { TempDir { path } } - pub async fn clean(&self) -> io::Result<()> { + pub async fn create(&self) -> io::Result<()> { if tokio::fs::metadata(&self.path).await.is_ok() { let _ = tokio::fs::remove_dir_all(&self.path).await; } - tokio::fs::create_dir(&self.path).await + tokio::fs::create_dir_all(&self.path).await + } +} + +impl Drop for TempDir { + fn drop(&mut self) { + let _ = std::fs::remove_dir_all(&self.path); + } +} + +async fn remove_siblings(base_path: &Path, keep: &Path) { + let Ok(mut entries) = tokio::fs::read_dir(base_path).await else { + return; + }; + + while let Ok(Some(entry)) = entries.next_entry().await { + let path = entry.path(); + if path == keep { + continue; + } + + if matches!(entry.file_type().await, Ok(file_type) if file_type.is_dir()) { + let _ = tokio::fs::remove_dir_all(&path).await; + } else { + let _ = tokio::fs::remove_file(&path).await; + } } } @@ -385,12 +445,6 @@ fn unpack_error(err: std::io::Error) -> trc::Error { .details("Failed to unpack application bundle") } -impl Drop for TempDir { - fn drop(&mut self) { - let _ = std::fs::remove_dir_all(&self.path); - } -} - impl Default for WebApplications { fn default() -> Self { Self::new() @@ -521,9 +575,9 @@ mod tests { ); } - async fn fixture(name: &str, client_id: Option<&str>) -> (WebApplications, TempDir) { + async fn fixture(name: &str, client_id: Option<&str>) -> WebApplications { let dir = TempDir::new(std::env::temp_dir().join(format!("inbuxa-app-{name}"))); - dir.clean().await.unwrap(); + dir.create().await.unwrap(); tokio::fs::write(dir.path.join("index.html"), INDEX) .await .unwrap(); @@ -544,6 +598,7 @@ mod tests { let routes = Arc::new(AppRoutes { resources, oauth_client_id_meta: client_id.map(oauth_client_id_meta), + _bundle_dir: dir, }); let mut map = AHashMap::new(); @@ -553,7 +608,7 @@ mod tests { let apps = WebApplications::new(); apps.routes.store(Arc::new(map)); - (apps, dir) + apps } async fn serve_html(apps: &WebApplications, prefix: &str, path: &str) -> String { @@ -565,7 +620,7 @@ mod tests { #[tokio::test] async fn serving_index_injects_the_prefix_and_client_id() { - let (apps, _dir) = fixture("serve-configured", Some("pocket-id-client")).await; + let apps = fixture("serve-configured", Some("pocket-id-client")).await; let html = serve_html(&apps, "admin", "index.html").await; assert!(html.contains(""), "{html}"); @@ -584,7 +639,7 @@ mod tests { #[tokio::test] async fn unknown_paths_fall_back_to_a_rewritten_index() { - let (apps, _dir) = fixture("serve-fallback", Some("pocket-id-client")).await; + let apps = fixture("serve-fallback", Some("pocket-id-client")).await; let html = serve_html(&apps, "admin", "settings/directory").await; assert!(html.contains(""), "{html}"); @@ -596,7 +651,7 @@ mod tests { #[tokio::test] async fn assets_and_unknown_prefixes_are_untouched() { - let (apps, _dir) = fixture("serve-assets", Some("pocket-id-client")).await; + let apps = fixture("serve-assets", Some("pocket-id-client")).await; let served = apps.serve("admin", "app.js").await.unwrap().unwrap(); assert_eq!(served.resource.contents, b"export const x = 1;\n"); @@ -608,7 +663,7 @@ mod tests { #[tokio::test] async fn serving_index_without_a_client_id_keeps_the_placeholder() { - let (apps, _dir) = fixture("serve-unconfigured", None).await; + let apps = fixture("serve-unconfigured", None).await; let html = serve_html(&apps, "admin", "index.html").await; assert!(html.contains(""), "{html}"); @@ -624,4 +679,65 @@ mod tests { assert_eq!(rewrite_index(bundle, "admin", None), bundle.as_bytes()); } + #[tokio::test] + async fn missing_parent_directories_are_created() { + let base = std::env::temp_dir().join("inbuxa-app-nested"); + let _ = tokio::fs::remove_dir_all(&base).await; + + let dir = TempDir::new(base.join("webui").join("0")); + dir.create().await.unwrap(); + + assert!(tokio::fs::metadata(&dir.path).await.is_ok()); + + drop(dir); + let _ = tokio::fs::remove_dir_all(&base).await; + } + + #[tokio::test] + async fn dropping_the_routes_removes_the_bundle_directory() { + let apps = fixture("drop-guard", None).await; + let path = apps + .routes + .load() + .get("admin") + .unwrap() + ._bundle_dir + .path + .clone(); + + assert!(tokio::fs::metadata(&path).await.is_ok()); + + apps.routes.store(Arc::new(AHashMap::new())); + + assert!(tokio::fs::metadata(&path).await.is_err()); + } + + #[tokio::test] + async fn sweeping_orphans_spares_the_current_generation() { + let base = std::env::temp_dir().join("inbuxa-app-sweep"); + let _ = tokio::fs::remove_dir_all(&base).await; + + let current = TempDir::new(base.join("1")); + current.create().await.unwrap(); + let orphan = base.join("0"); + tokio::fs::create_dir_all(&orphan).await.unwrap(); + let stray = base.join("webui.zip"); + tokio::fs::write(&stray, b"not a bundle").await.unwrap(); + + remove_siblings(&base, ¤t.path).await; + + assert!(tokio::fs::metadata(¤t.path).await.is_ok()); + assert!(tokio::fs::metadata(&orphan).await.is_err()); + assert!(tokio::fs::metadata(&stray).await.is_err()); + + drop(current); + let _ = tokio::fs::remove_dir_all(&base).await; + } + + #[test] + fn generations_never_repeat() { + let apps = WebApplications::new(); + + assert_ne!(apps.next_generation(), apps.next_generation()); + } } diff --git a/crates/common/src/network/acme/order.rs b/crates/common/src/network/acme/order.rs index 259e9fc..5ed1fed 100644 --- a/crates/common/src/network/acme/order.rs +++ b/crates/common/src/network/acme/order.rs @@ -98,7 +98,38 @@ impl AcmeRequestBuilder { reuse_key_pem: Option, dns_parameters: Option, ) -> AcmeResult { - let mut params = CertificateParams::new(domains.clone()).map_err(|err| { + let mut published = BTreeSet::new(); + let result = self + .run_order( + server, + &domains, + reuse_key_pem, + dns_parameters.as_ref(), + &mut published, + ) + .await; + + if let Some(dns_parameters) = &dns_parameters { + for (zone, challenge_name) in published { + let _ = dns_parameters + .updater + .delete_rrset(&zone, &challenge_name, dns_update::DnsRecordType::TXT) + .await; + } + } + + result + } + + async fn run_order( + &self, + server: &Server, + domains: &[String], + reuse_key_pem: Option, + dns_parameters: Option<&AcmeDnsParameters>, + published: &mut BTreeSet<(String, String)>, + ) -> AcmeResult { + let mut params = CertificateParams::new(domains.to_vec()).map_err(|err| { AcmeError::Crypto(format!("Failed to create certificate params: {}", err)) })?; params.distinguished_name = DistinguishedName::new(); @@ -110,7 +141,7 @@ impl AcmeRequestBuilder { AcmeError::Crypto(format!("Failed to generate key pair: {}", err)) })?, }; - let response = self.new_order(domains.clone()).await?; + let response = self.new_order(domains.to_vec()).await?; let order_url = response.location; let mut order = response.body; let mut retry_after = None; @@ -119,7 +150,7 @@ impl AcmeRequestBuilder { Acme(AcmeEvent::OrderStart), Url = self.directory.new_order.to_string(), Details = order_url.to_string(), - Hostname = domains.as_slice(), + Hostname = domains, Type = self.challenge.as_str(), ); @@ -128,19 +159,20 @@ impl AcmeRequestBuilder { OrderStatus::Pending => { if matches!(self.challenge, ChallengeType::Dns01) { for url in &order.authorizations { - self.authorize(server, url, dns_parameters.as_ref()).await?; + self.authorize(server, url, dns_parameters, Some(published)) + .await?; } } else { let auth_futures = order .authorizations .iter() - .map(|url| self.authorize(server, url, dns_parameters.as_ref())); + .map(|url| self.authorize(server, url, dns_parameters, None)); try_join_all(auth_futures).await?; } trc::event!( Acme(AcmeEvent::AuthCompleted), Url = self.directory.new_order.to_string(), - Hostname = domains.as_slice(), + Hostname = domains, ); let response = self.order(&order_url).await?; order = response.body; @@ -151,7 +183,7 @@ impl AcmeRequestBuilder { trc::event!( Acme(AcmeEvent::OrderProcessing), Url = self.directory.new_order.to_string(), - Hostname = domains.as_slice(), + Hostname = domains, Total = i, ); @@ -179,7 +211,7 @@ impl AcmeRequestBuilder { trc::event!( Acme(AcmeEvent::OrderReady), Url = self.directory.new_order.to_string(), - Hostname = domains.as_slice(), + Hostname = domains, ); let csr = params.serialize_request(&key_pair).map_err(|err| { @@ -192,10 +224,10 @@ impl AcmeRequestBuilder { trc::event!( Acme(AcmeEvent::OrderValid), Url = self.directory.new_order.to_string(), - Hostname = domains.as_slice(), + Hostname = domains, ); - let certificate = self.select_certificate(&domains, certificate).await?; + let certificate = self.select_certificate(domains, certificate).await?; return Ok(PemCert { certificate, @@ -213,7 +245,7 @@ impl AcmeRequestBuilder { Acme(AcmeEvent::OrderInvalid), Url = self.directory.new_order.to_string(), Details = order_url.to_string(), - Hostname = domains.as_slice(), + Hostname = domains, Reason = reason.clone(), ); @@ -228,6 +260,7 @@ impl AcmeRequestBuilder { server: &Server, url: &String, dns_parameters: Option<&AcmeDnsParameters>, + published: Option<&mut BTreeSet<(String, String)>>, ) -> AcmeResult<()> { let response = self .auth(url) @@ -289,7 +322,12 @@ impl AcmeRequestBuilder { .await?; } ChallengeType::Dns01 => { - let dns_parameters = dns_parameters.unwrap(); + let Some(dns_parameters) = dns_parameters else { + return Err(AcmeError::Invalid( + "DNS-01 challenge requested but a DNS provider was not configured" + .to_string(), + )); + }; let domain = domain.strip_prefix("*.").unwrap_or(&domain); let zone = dns_parameters @@ -310,6 +348,11 @@ impl AcmeRequestBuilder { ) .await .map_err(AcmeError::Dns)?; + + if let Some(published) = published { + published.insert((zone.to_string(), challenge_name.clone())); + } + dns_parameters .updater .wait_for_txt_propagation(&challenge_name, zone, &proof) diff --git a/crates/common/src/network/dns/update.rs b/crates/common/src/network/dns/update.rs index 394ba85..0d88d52 100644 --- a/crates/common/src/network/dns/update.rs +++ b/crates/common/src/network/dns/update.rs @@ -1150,6 +1150,36 @@ impl DnsUpdater { Ok(()) } + pub async fn delete_rrset( + &self, + origin: &str, + name: &str, + record_type: DnsRecordType, + ) -> Result<(), String> { + if let Err(err) = self + .updater + .set_rrset( + name, + record_type, + self.ttl.as_secs() as u32, + Vec::new(), + origin, + ) + .await + { + trc::event!( + Dns(DnsEvent::RecordDeletionFailed), + Hostname = name.to_string(), + Details = origin.to_string(), + Type = record_type.as_str(), + Reason = err.to_string(), + ); + return Err(format!("Failed to delete DNS RRSet: {}", err)); + } + + Ok(()) + } + pub async fn add_to_rrset( &self, origin: &str, diff --git a/crates/common/src/network/mta.rs b/crates/common/src/network/mta.rs index 1d99dc4..696d7a9 100644 --- a/crates/common/src/network/mta.rs +++ b/crates/common/src/network/mta.rs @@ -23,6 +23,7 @@ use crate::{ manager::SPAM_CLASSIFIER_KEY, network::RcptResolution, }; +use ahash::AHashSet; use directory::Recipient; use mail_auth::IpLookupStrategy; use registry::schema::enums::ExpressionVariable; @@ -37,6 +38,8 @@ use store::{ write::{AlignedBytes, Archive, QueueClass, ValueClass}, }; use trc::{AddContext, SpamEvent}; +use types::id::Id; +use utils::DomainPart; impl Server { pub async fn rcpt_resolve( @@ -163,7 +166,10 @@ impl Server { } EmailCache::MailingList(id) => { if let Some(list) = self.try_list(id).await? { - return Ok(RcptResolution::Expand(list.recipients.clone())); + return Ok(RcptResolution::Expand( + self.expand_nested_lists(id, list.recipients.clone()) + .await?, + )); } else { self.inner .cache @@ -195,6 +201,56 @@ impl Server { } } + async fn expand_nested_lists( + &self, + list_id: u32, + recipients: Arc<[Box]>, + ) -> trc::Result]>> { + let mut has_nested = false; + for member in recipients.iter() { + if let Some(EmailCache::MailingList(_)) = self.rcpt_id_from_email(member).await? { + has_nested = true; + break; + } + } + if !has_nested { + return Ok(recipients); + } + + let mut expanded = Vec::with_capacity(recipients.len()); + let mut seen: AHashSet> = AHashSet::with_capacity(recipients.len()); + let mut visited = AHashSet::from_iter([list_id]); + let mut pending: Vec]>> = Vec::new(); + let mut members = recipients; + + loop { + for member in members.iter() { + if let Some(EmailCache::MailingList(nested_id)) = + self.rcpt_id_from_email(member).await? + { + if !visited.insert(nested_id) { + continue; + } + if let Some(nested) = self.try_list(nested_id).await? { + pending.push(nested.recipients.clone()); + continue; + } + } + + if seen.insert(member.to_canonical_address().into()) { + expanded.push(member.clone()); + } + } + + let Some(next) = pending.pop() else { + break; + }; + members = next; + } + + Ok(expanded.into()) + } + pub async fn get_dkim_signers( &self, domain: &str, diff --git a/crates/coordinator/Cargo.toml b/crates/coordinator/Cargo.toml index d1e161b..ea63a67 100644 --- a/crates/coordinator/Cargo.toml +++ b/crates/coordinator/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "coordinator" -version = "0.16.22" +version = "0.16.23" edition = "2024" [dependencies] diff --git a/crates/dav-proto/Cargo.toml b/crates/dav-proto/Cargo.toml index 2388a3e..469a56a 100644 --- a/crates/dav-proto/Cargo.toml +++ b/crates/dav-proto/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "dav-proto" -version = "0.16.22" +version = "0.16.23" edition = "2024" [dependencies] diff --git a/crates/dav/Cargo.toml b/crates/dav/Cargo.toml index a80ae27..6c7879c 100644 --- a/crates/dav/Cargo.toml +++ b/crates/dav/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "dav" -version = "0.16.22" +version = "0.16.23" edition = "2024" [dependencies] diff --git a/crates/directory/Cargo.toml b/crates/directory/Cargo.toml index 3a019b3..1352552 100644 --- a/crates/directory/Cargo.toml +++ b/crates/directory/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "directory" -version = "0.16.22" +version = "0.16.23" edition = "2024" [dependencies] diff --git a/crates/email/Cargo.toml b/crates/email/Cargo.toml index 890fb75..12cc2ae 100644 --- a/crates/email/Cargo.toml +++ b/crates/email/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "email" -version = "0.16.22" +version = "0.16.23" edition = "2024" [dependencies] diff --git a/crates/email/src/message/delivery.rs b/crates/email/src/message/delivery.rs index 456f6f7..08c9444 100644 --- a/crates/email/src/message/delivery.rs +++ b/crates/email/src/message/delivery.rs @@ -22,6 +22,8 @@ use std::{borrow::Cow, future::Future}; use store::ahash::AHashMap; use types::blob_hash::BlobHash; +pub const ORCPT_ADDR_TYPE: &str = "rfc822;"; + #[derive(Debug)] pub struct IngestMessage { pub sender_address: String, @@ -40,6 +42,12 @@ pub struct IngestRecipient { } impl IngestRecipient { + pub fn orcpt_parameter(&self) -> Option { + self.orcpt + .as_deref() + .map(|orcpt| format!("{ORCPT_ADDR_TYPE}{orcpt}")) + } + pub fn is_spam(&self) -> bool { self.spam_percentage .is_some_and(|percentage| percentage >= 50) diff --git a/crates/email/src/sieve/ingest.rs b/crates/email/src/sieve/ingest.rs index ada7e5a..ef13c36 100644 --- a/crates/email/src/sieve/ingest.rs +++ b/crates/email/src/sieve/ingest.rs @@ -126,6 +126,7 @@ impl SieveScriptIngest for Server { .caused_by(trc::location!())?; // Create Sieve instance + let orcpt = envelope_to.orcpt_parameter(); let mut instance = self.core.sieve.untrusted_runtime.filter_parsed(message); // Set account name and email @@ -141,7 +142,7 @@ impl SieveScriptIngest for Server { // Set envelope instance.set_envelope(Envelope::From, envelope_from); instance.set_envelope(Envelope::To, envelope_to.address.as_str()); - if let Some(orcpt) = &envelope_to.orcpt { + if let Some(orcpt) = &orcpt { instance.set_envelope(Envelope::Orcpt, orcpt.as_str()); } instance.set_spam_status(spam_status(envelope_to.spam_percentage)); diff --git a/crates/groupware/Cargo.toml b/crates/groupware/Cargo.toml index 932f1f8..3b4cf57 100644 --- a/crates/groupware/Cargo.toml +++ b/crates/groupware/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "groupware" -version = "0.16.22" +version = "0.16.23" edition = "2024" [dependencies] diff --git a/crates/http-proto/Cargo.toml b/crates/http-proto/Cargo.toml index e0f575b..51bb7f1 100644 --- a/crates/http-proto/Cargo.toml +++ b/crates/http-proto/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "http_proto" -version = "0.16.22" +version = "0.16.23" edition = "2024" [dependencies] diff --git a/crates/http/Cargo.toml b/crates/http/Cargo.toml index ee84554..ddf9497 100644 --- a/crates/http/Cargo.toml +++ b/crates/http/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "http" -version = "0.16.22" +version = "0.16.23" edition = "2024" [dependencies] diff --git a/crates/http/src/api/diagnose.rs b/crates/http/src/api/diagnose.rs index 9c2e7dd..e6c757f 100644 --- a/crates/http/src/api/diagnose.rs +++ b/crates/http/src/api/diagnose.rs @@ -230,14 +230,7 @@ async fn delivery_diagnose( // Lookup MX let now = Instant::now(); - let mxs = match server - .core - .smtp - .resolvers - .dns - .mx_lookup(&domain, Some(&server.inner.cache.dns_mx)) - .await - { + let mxs = match server.mx_lookup(domain.as_str()).await { Ok(mxs) => mxs, Err(err) => { tx.send(DeliveryStage::MxLookupError { @@ -419,7 +412,7 @@ async fn delivery_diagnose( }) .await?; - None + continue 'outer; } Ok(TlsaResult::Missing) => { tx.send(DeliveryStage::TlsaNotFound { @@ -440,14 +433,17 @@ async fn delivery_diagnose( 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; } - None } }; diff --git a/crates/imap-proto/Cargo.toml b/crates/imap-proto/Cargo.toml index 7d40a5e..363a4d9 100644 --- a/crates/imap-proto/Cargo.toml +++ b/crates/imap-proto/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "imap_proto" -version = "0.16.22" +version = "0.16.23" edition = "2024" [dependencies] diff --git a/crates/imap/Cargo.toml b/crates/imap/Cargo.toml index 668dfbe..9916808 100644 --- a/crates/imap/Cargo.toml +++ b/crates/imap/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "imap" -version = "0.16.22" +version = "0.16.23" edition = "2024" [dependencies] diff --git a/crates/imap/src/op/copy_move.rs b/crates/imap/src/op/copy_move.rs index 97de5c2..7cfd562 100644 --- a/crates/imap/src/op/copy_move.rs +++ b/crates/imap/src/op/copy_move.rs @@ -390,6 +390,16 @@ impl SessionData { .await .imap_ctx(&arguments.tag, trc::location!())?; let mut dest_cache = None; + let train_spam = if dest_mailbox_id == JUNK_ID { + Some(true) + } else if src_mailbox.id.mailbox_id == JUNK_ID && dest_mailbox_id != TRASH_ID { + Some(false) + } else { + None + }; + 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 { match self .server @@ -515,11 +525,33 @@ impl SessionData { } }; + if let Some(is_spam) = train_spam { + self.server + .add_account_spam_sample( + &mut train_batch, + src_account_id, + id, + is_spam, + self.session_id, + ) + .await + .imap_ctx(&arguments.tag, trc::location!())?; + train_batch.commit_point(); + did_train = true; + } + if is_move { destroy_ids.insert(id); } } + if did_train { + self.server + .commit_batch(train_batch) + .await + .imap_ctx(&arguments.tag, trc::location!())?; + } + // Untag or delete emails if !destroy_ids.is_empty() { let mut batch = BatchBuilder::new(); diff --git a/crates/jmap-proto/Cargo.toml b/crates/jmap-proto/Cargo.toml index 8041c3a..0382fab 100644 --- a/crates/jmap-proto/Cargo.toml +++ b/crates/jmap-proto/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "jmap_proto" -version = "0.16.22" +version = "0.16.23" edition = "2024" [dependencies] diff --git a/crates/jmap/Cargo.toml b/crates/jmap/Cargo.toml index a55745e..67154d6 100644 --- a/crates/jmap/Cargo.toml +++ b/crates/jmap/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "jmap" -version = "0.16.22" +version = "0.16.23" edition = "2024" [dependencies] diff --git a/crates/jmap/src/email/copy.rs b/crates/jmap/src/email/copy.rs index dfb7779..96b98ff 100644 --- a/crates/jmap/src/email/copy.rs +++ b/crates/jmap/src/email/copy.rs @@ -11,7 +11,11 @@ use crate::{ use common::{Server, auth::AccessToken}; use email::{ cache::{MessageCacheFetch, email::MessageCacheAccess, mailbox::MailboxCacheAccess}, - message::copy::{CopyMessageError, EmailCopy}, + mailbox::JUNK_ID, + message::{ + copy::{CopyMessageError, EmailCopy}, + ingest::EmailIngest, + }, }; use http_proto::HttpSessionData; use jmap_proto::{ @@ -29,6 +33,7 @@ use jmap_proto::{ }; use jmap_tools::{Key, Value}; use std::future::Future; +use store::write::BatchBuilder; use trc::AddContext; use types::acl::Acl; use utils::map::vec_map::VecMap; @@ -87,6 +92,9 @@ impl JmapEmailCopy for Server { }; let on_success_delete = request.on_success_destroy_original.unwrap_or(false); let mut destroy_ids = Vec::new(); + let mut train_batch = BatchBuilder::new(); + let mut did_train = false; + train_batch.with_account_id(from_account_id); 'create: for (id, create) in request.create.into_valid() { let mut from_message_id = None; @@ -208,6 +216,7 @@ impl JmapEmailCopy for Server { } // Add response + let train_spam = mailboxes.contains(&JUNK_ID); match self .copy_message( from_account_id, @@ -221,6 +230,20 @@ impl JmapEmailCopy for Server { .await? { Ok(email) => { + if train_spam { + self.add_account_spam_sample( + &mut train_batch, + from_account_id, + from_message_id.document_id(), + true, + session.session_id, + ) + .await + .caused_by(trc::location!())?; + train_batch.commit_point(); + did_train = true; + } + response .created .append(id, ingested_into_object(email).into()); @@ -245,6 +268,12 @@ impl JmapEmailCopy for Server { } } + if did_train { + self.commit_batch(train_batch) + .await + .caused_by(trc::location!())?; + } + // Update state if !response.created.is_empty() { response.new_state = self.get_cached_messages(account_id).await?.get_state(false); diff --git a/crates/main/Cargo.toml b/crates/main/Cargo.toml index ed5f73c..78a6c86 100644 --- a/crates/main/Cargo.toml +++ b/crates/main/Cargo.toml @@ -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.22" +version = "0.16.23" edition = "2024" [[bin]] diff --git a/crates/main/src/test_data.rs b/crates/main/src/test_data.rs index cc7d91c..d2f73a1 100644 --- a/crates/main/src/test_data.rs +++ b/crates/main/src/test_data.rs @@ -57,7 +57,7 @@ pub async fn insert_test_data(server: &Server) { server.inner.data.queue_id_gen.generate(), QueueName::default(), ); - assert!(qm.save_changes(server, None).await); + assert!(qm.save_changes(server, None, None).await); } for report in sample_tls_internal_reports() { @@ -163,7 +163,7 @@ fn sample_queued_messages(blob_hashes: Vec) -> Vec { }, }), flags: RCPT_DSN_SENT, - orcpt: Some("rfc822;bob@example.org".into()), + orcpt: Some("bob@example.org".into()), }, ], received_from_ip: std::net::IpAddr::V4(Ipv4Addr::new(192, 168, 1, 10)), diff --git a/crates/managesieve/Cargo.toml b/crates/managesieve/Cargo.toml index 945be75..1eaec11 100644 --- a/crates/managesieve/Cargo.toml +++ b/crates/managesieve/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "managesieve" -version = "0.16.22" +version = "0.16.23" edition = "2024" [dependencies] diff --git a/crates/migration/Cargo.toml b/crates/migration/Cargo.toml index d6b2186..75c0928 100644 --- a/crates/migration/Cargo.toml +++ b/crates/migration/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "migration" -version = "0.16.22" +version = "0.16.23" edition = "2024" [dependencies] diff --git a/crates/nlp/Cargo.toml b/crates/nlp/Cargo.toml index 1056faf..379e51c 100644 --- a/crates/nlp/Cargo.toml +++ b/crates/nlp/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "nlp" -version = "0.16.22" +version = "0.16.23" edition = "2024" [dependencies] diff --git a/crates/pop3/Cargo.toml b/crates/pop3/Cargo.toml index b344b3e..d850c69 100644 --- a/crates/pop3/Cargo.toml +++ b/crates/pop3/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "pop3" -version = "0.16.22" +version = "0.16.23" edition = "2024" [dependencies] diff --git a/crates/pop3/src/op/fetch.rs b/crates/pop3/src/op/fetch.rs index fedc41a..f1f576d 100644 --- a/crates/pop3/src/op/fetch.rs +++ b/crates/pop3/src/op/fetch.rs @@ -64,14 +64,8 @@ impl Session { ) .get_full_range(); - self.write_bytes( - Response::Message:: { - bytes, - lines: lines.unwrap_or(0), - } - .serialize(), - ) - .await + self.write_bytes(Response::Message:: { bytes, lines }.serialize()) + .await } else { Err(trc::Pop3Event::Error .into_err() diff --git a/crates/pop3/src/protocol/response.rs b/crates/pop3/src/protocol/response.rs index 8c58ea7..052c588 100644 --- a/crates/pop3/src/protocol/response.rs +++ b/crates/pop3/src/protocol/response.rs @@ -16,7 +16,7 @@ pub enum Response<'x, T> { List(Vec), Message { bytes: SliceRange<'x>, - lines: u32, + lines: Option, }, Capability { mechanisms: Vec, @@ -54,40 +54,65 @@ impl<'x, T: Display> Response<'x, T> { buf } Response::Message { bytes, lines } => { - let mut buf = Vec::with_capacity(bytes.len() + 10); - buf.extend_from_slice(b"+OK "); - buf.extend_from_slice(bytes.len().to_string().as_bytes()); - buf.extend_from_slice(b" octets\r\n"); - - let mut line_count = 0; - let mut last_byte = 0; + let lines = *lines; + let mut message = Vec::with_capacity(bytes.len() + 16); + let mut octets = 0; + let mut last_byte = b'\n'; + let mut in_headers = lines.is_some(); + let mut is_blank_line = true; + let mut body_lines = 0; // Transparency procedure for &byte in bytes.into_iter() { // POP3 requires that lines end with CRLF, do this check to ensure that if byte == b'\n' && last_byte != b'\r' { - buf.push(b'\r'); + message.push(b'\r'); + octets += 1; } if byte == b'.' && last_byte == b'\n' { - buf.push(b'.'); + message.push(b'.'); } - buf.push(byte); + message.push(byte); + octets += 1; last_byte = byte; - if *lines > 0 && byte == b'\n' { - line_count += 1; - if line_count == *lines { - break; + match byte { + b'\n' => { + if in_headers { + in_headers = !is_blank_line; + } else { + body_lines += 1; + } + if !in_headers && lines.is_some_and(|lines| body_lines >= lines) { + break; + } + is_blank_line = true; + } + b'\r' => {} + _ => { + is_blank_line = false; } } } if last_byte != b'\n' { - buf.extend_from_slice(b"\r\n"); + message.extend_from_slice(b"\r\n"); + octets += 2; } - buf.extend_from_slice(b".\r\n"); + if in_headers { + message.extend_from_slice(b"\r\n"); + octets += 2; + } + + message.extend_from_slice(b".\r\n"); + + let mut buf = Vec::with_capacity(message.len() + 24); + buf.extend_from_slice(b"+OK "); + buf.extend_from_slice(octets.to_string().as_bytes()); + buf.extend_from_slice(b" octets\r\n"); + buf.extend_from_slice(&message); buf } Response::Capability { mechanisms, stls } => { @@ -208,9 +233,51 @@ mod tests { ( Response::Message { bytes: SliceRange::Split(b"Subject: test\r\n\r\n.\r\n", b"test.\r\n.test\r\na"), - lines: 0, + lines: None, }, - "+OK 35 octets\r\nSubject: test\r\n\r\n..\r\ntest.\r\n..test\r\na\r\n.\r\n", + "+OK 37 octets\r\nSubject: test\r\n\r\n..\r\ntest.\r\n..test\r\na\r\n.\r\n", + ), + ( + Response::Message { + bytes: SliceRange::Split(b"Subject: test\r\n\r\n.\r\n", b"test.\r\n.test\r\na"), + lines: Some(0), + }, + "+OK 17 octets\r\nSubject: test\r\n\r\n.\r\n", + ), + ( + Response::Message { + bytes: SliceRange::Split(b"Subject: test\r\n\r\n.\r\n", b"test.\r\n.test\r\na"), + lines: Some(2), + }, + "+OK 27 octets\r\nSubject: test\r\n\r\n..\r\ntest.\r\n.\r\n", + ), + ( + Response::Message { + bytes: SliceRange::Split(b"Subject: test\r\n\r\n.\r\n", b"test.\r\n.test\r\na"), + lines: Some(100), + }, + "+OK 37 octets\r\nSubject: test\r\n\r\n..\r\ntest.\r\n..test\r\na\r\n.\r\n", + ), + ( + Response::Message { + bytes: SliceRange::Single(b"Subject: test\n\nbody\n"), + lines: None, + }, + "+OK 23 octets\r\nSubject: test\r\n\r\nbody\r\n.\r\n", + ), + ( + Response::Message { + bytes: SliceRange::Single(b"Subject: test\n\n.leading dot\n"), + lines: Some(1), + }, + "+OK 31 octets\r\nSubject: test\r\n\r\n..leading dot\r\n.\r\n", + ), + ( + Response::Message { + bytes: SliceRange::Single(b".dot\r\nSubject: test\r\n"), + lines: Some(3), + }, + "+OK 23 octets\r\n..dot\r\nSubject: test\r\n\r\n.\r\n", ), ] { assert_eq!(expected, String::from_utf8(cmd.serialize()).unwrap()); diff --git a/crates/registry/Cargo.toml b/crates/registry/Cargo.toml index ef5c7dc..4adb8f3 100644 --- a/crates/registry/Cargo.toml +++ b/crates/registry/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "registry" -version = "0.16.22" +version = "0.16.23" edition = "2024" [dependencies] diff --git a/crates/scim-proto/Cargo.toml b/crates/scim-proto/Cargo.toml index 42b037c..099ee79 100644 --- a/crates/scim-proto/Cargo.toml +++ b/crates/scim-proto/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "scim-proto" -version = "0.16.22" +version = "0.16.23" edition = "2024" [dependencies] diff --git a/crates/scim/Cargo.toml b/crates/scim/Cargo.toml index cfd84e3..0636477 100644 --- a/crates/scim/Cargo.toml +++ b/crates/scim/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "scim" -version = "0.16.22" +version = "0.16.23" edition = "2024" [dependencies] diff --git a/crates/services/Cargo.toml b/crates/services/Cargo.toml index b2374e2..8e98b97 100644 --- a/crates/services/Cargo.toml +++ b/crates/services/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "services" -version = "0.16.22" +version = "0.16.23" edition = "2024" [dependencies] diff --git a/crates/services/src/broadcast/subscriber.rs b/crates/services/src/broadcast/subscriber.rs index f8b832d..f509e93 100644 --- a/crates/services/src/broadcast/subscriber.rs +++ b/crates/services/src/broadcast/subscriber.rs @@ -96,6 +96,13 @@ pub fn spawn_broadcast_subscriber(inner: Arc, mut shutdown_rx: watch::Rec } }; + inner + .shared_core + .load() + .storage + .data + .invalidate_read_snapshot(); + loop { match batch.next_event() { Ok(Some(event)) => { @@ -174,9 +181,7 @@ pub fn spawn_broadcast_subscriber(inner: Arc, mut shutdown_rx: watch::Rec .await; } BroadcastEvent::QueueRefresh => { - let core = inner.shared_core.load_full(); - if core.network.roles.outbound_mta { - core.storage.data.invalidate_read_snapshot(); + if inner.shared_core.load().network.roles.outbound_mta { let _ = inner .ipc .queue_tx @@ -185,9 +190,7 @@ pub fn spawn_broadcast_subscriber(inner: Arc, mut shutdown_rx: watch::Rec } } BroadcastEvent::RegistryChange(change) => { - let server = inner.build_server(); - server.store().invalidate_read_snapshot(); - match Box::pin(server.reload_registry(change)).await { + match Box::pin(inner.build_server().reload_registry(change)).await { Ok(result) => { result.log(); } diff --git a/crates/services/src/task_manager/destroy_account.rs b/crates/services/src/task_manager/destroy_account.rs index a616922..ee88695 100644 --- a/crates/services/src/task_manager/destroy_account.rs +++ b/crates/services/src/task_manager/destroy_account.rs @@ -6,7 +6,7 @@ * Modified by Coffey Labs in 2026 for INBUXA. */ -use crate::task_manager::TaskResult; +use crate::task_manager::{TaskResult, deferred_retry_time}; use common::Server; use email::{message::metadata::MessageMetadata, sieve::SieveScript}; use groupware::file::FileNode; @@ -41,7 +41,7 @@ impl DestroyAccountTask for Server { match destroy_account(self, task).await { Ok(result) => result, Err(err) => { - let result = TaskResult::temporary(err.to_string()); + let result = TaskResult::deferred(deferred_retry_time(&err), err.to_string()); trc::error!( err.account_id(task.account_id.document_id()) .details("Failed to destroy account") diff --git a/crates/services/src/task_manager/index.rs b/crates/services/src/task_manager/index.rs index 41c35dc..8942980 100644 --- a/crates/services/src/task_manager/index.rs +++ b/crates/services/src/task_manager/index.rs @@ -6,7 +6,7 @@ * Modified by Coffey Labs in 2026 for INBUXA. */ -use crate::task_manager::{Task, TaskDetails, TaskFailureType, TaskResult}; +use crate::task_manager::{Task, TaskDetails, TaskFailureType, TaskResult, deferred_retry_time}; use common::Server; use email::{ cache::MessageCacheFetch, @@ -274,7 +274,7 @@ impl SearchIndexTask for Server { ); for r in results.iter_mut() { if r.task_type == TaskType::Insert && r.result.is_success() { - r.result = search_store_failure(retry_at, "Failed to index documents"); + r.result = TaskResult::deferred(retry_at, "Failed to index documents"); } } return results; @@ -331,7 +331,7 @@ impl SearchIndexTask for Server { for r in results.iter_mut() { if r.task_type == TaskType::Delete && r.result.is_success() { r.result = - search_store_failure(retry_at, "Failed to delete documents from index"); + TaskResult::deferred(retry_at, "Failed to delete documents from index"); } } return results; @@ -445,22 +445,6 @@ pub(crate) async fn reindex_account(server: &Server, account_id: u32) -> trc::Re Ok(()) } -fn deferred_retry_time(err: &trc::Error) -> Option { - err.value(trc::Key::NextRetry) - .and_then(|value| value.to_uint()) -} - -fn search_store_failure(retry_at: Option, message: &'static str) -> TaskResult { - match retry_at { - Some(retry_at) => TaskResult::Failure { - typ: TaskFailureType::Retry(retry_at), - message: message.into(), - max_attempts: None, - }, - None => TaskResult::temporary(message), - } -} - fn attempt_number(status: &TaskStatus) -> u64 { match status { TaskStatus::Pending(_) => 0, diff --git a/crates/services/src/task_manager/manager.rs b/crates/services/src/task_manager/manager.rs index 23cb05b..180aeac 100644 --- a/crates/services/src/task_manager/manager.rs +++ b/crates/services/src/task_manager/manager.rs @@ -621,6 +621,7 @@ pub fn perpetual_retry_time(typ: TaskType, attempt: u64) -> Option { | TaskType::DkimManagement | TaskType::IndexDocument | TaskType::UnindexDocument + | TaskType::DestroyAccount ) .then(|| { now().saturating_add( diff --git a/crates/services/src/task_manager/mod.rs b/crates/services/src/task_manager/mod.rs index 3575b74..d2631f5 100644 --- a/crates/services/src/task_manager/mod.rs +++ b/crates/services/src/task_manager/mod.rs @@ -137,4 +137,20 @@ impl TaskResult { max_attempts: None, } } + + pub fn deferred(retry_at: Option, message: impl Into) -> Self { + match retry_at { + Some(retry_at) => TaskResult::Failure { + typ: TaskFailureType::Retry(retry_at), + message: message.into(), + max_attempts: None, + }, + None => TaskResult::temporary(message), + } + } +} + +pub(crate) fn deferred_retry_time(err: &trc::Error) -> Option { + err.value(trc::Key::NextRetry) + .and_then(|value| value.to_uint()) } diff --git a/crates/smtp/Cargo.toml b/crates/smtp/Cargo.toml index 1581c7a..bd19341 100644 --- a/crates/smtp/Cargo.toml +++ b/crates/smtp/Cargo.toml @@ -7,7 +7,7 @@ homepage = "https://stalw.art/smtp" keywords = ["smtp", "email", "mail", "server"] categories = ["email"] license = "AGPL-3.0-only OR LicenseRef-SEL" -version = "0.16.22" +version = "0.16.23" edition = "2024" [dependencies] diff --git a/crates/smtp/src/core/mod.rs b/crates/smtp/src/core/mod.rs index 63b301a..485b3c3 100644 --- a/crates/smtp/src/core/mod.rs +++ b/crates/smtp/src/core/mod.rs @@ -11,6 +11,7 @@ use common::{ config::smtp::auth::VerifyStrategy, network::{ServerInstance, asn::AsnGeoLookupResult}, }; +use email::message::delivery::ORCPT_ADDR_TYPE; use mail_auth::{IprevOutput, SpfOutput}; use smtp_proto::request::receiver::{ BdatReceiver, DataReceiver, DummyDataReceiver, DummyLineReceiver, LineReceiver, RequestReceiver, @@ -306,10 +307,13 @@ impl SessionAddress { } } - pub fn report_address(&self) -> &str { + pub fn orig_address(&self) -> &str { + self.dsn_info.as_deref().unwrap_or(&self.address_lcase) + } + + pub fn orcpt_parameter(&self) -> Option { self.dsn_info - .as_ref() - .and_then(|v| v.strip_prefix("rfc822;")) - .unwrap_or(&self.address_lcase) + .as_deref() + .map(|orcpt| format!("{ORCPT_ADDR_TYPE}{}", orcpt.to_lowercase())) } } diff --git a/crates/smtp/src/inbound/data.rs b/crates/smtp/src/inbound/data.rs index a8d50b3..a3ee5ab 100644 --- a/crates/smtp/src/inbound/data.rs +++ b/crates/smtp/src/inbound/data.rs @@ -443,7 +443,7 @@ impl Session { if !rc.analysis.forward { self.data .rcpt_to - .retain(|rcpt| !rc.analysis.is_report_address(rcpt.report_address())); + .retain(|rcpt| !rc.analysis.is_report_address(rcpt.orig_address())); } if self.data.rcpt_to.is_empty() { diff --git a/crates/smtp/src/inbound/rcpt.rs b/crates/smtp/src/inbound/rcpt.rs index 32a290b..1eb25e0 100644 --- a/crates/smtp/src/inbound/rcpt.rs +++ b/crates/smtp/src/inbound/rcpt.rs @@ -202,8 +202,8 @@ impl Session { let mut new_addr = SessionAddress::new(address); if !self.data.rcpt_to.contains(&new_addr) { - new_addr.dsn_info = format!("rfc822;{}", orig_addr.address_lcase).into(); new_addr.flags = orig_addr.flags; + new_addr.dsn_info = orig_addr.address_lcase.into(); self.data.rcpt_to.push(new_addr); } else { trc::event!( @@ -353,7 +353,6 @@ impl Session { // Expand list if let Some(members) = rcpt_members { let list_addr = self.data.rcpt_to.pop().unwrap(); - let orcpt = format!("rfc822;{}", list_addr.address_lcase); for member in members.as_ref() { let member_lcase = member.to_lowercase(); let is_local = match self @@ -399,7 +398,7 @@ impl Session { if !self.data.rcpt_to.contains(&member_addr) && member_addr.address_lcase != list_addr.address_lcase { - member_addr.dsn_info = orcpt.clone().into(); + member_addr.dsn_info = list_addr.address_lcase.clone().into(); member_addr.flags = list_addr.flags; self.data.rcpt_to.push(member_addr); } diff --git a/crates/smtp/src/inbound/spam.rs b/crates/smtp/src/inbound/spam.rs index 0ebb670..94e85e5 100644 --- a/crates/smtp/src/inbound/spam.rs +++ b/crates/smtp/src/inbound/spam.rs @@ -89,17 +89,7 @@ impl Session { .iter() .map(|r| r.address_lcase.as_str()) .collect(), - env_rcpt_orig_to: self - .data - .rcpt_to - .iter() - .map(|r| { - r.dsn_info - .as_deref() - .and_then(|info| info.strip_prefix("rfc822;")) - .unwrap_or(r.address_lcase.as_str()) - }) - .collect(), + env_rcpt_orig_to: self.data.rcpt_to.iter().map(|r| r.orig_address()).collect(), is_test: false, is_train: false, } diff --git a/crates/smtp/src/outbound/delivery.rs b/crates/smtp/src/outbound/delivery.rs index cd83494..4c9c0eb 100644 --- a/crates/smtp/src/outbound/delivery.rs +++ b/crates/smtp/src/outbound/delivery.rs @@ -15,8 +15,8 @@ use crate::outbound::lookup::{DnsLookup, SourceIp}; use crate::outbound::mta_sts::lookup::MtaStsLookup; use crate::outbound::mta_sts::verify::VerifyPolicy; use crate::outbound::{client::StartTlsResult, dane::verify::TlsaVerify}; -use crate::queue::dsn::SendDsn; -use crate::queue::spool::SmtpSpool; +use crate::queue::dsn::{DsnStatus, SendDsn}; +use crate::queue::spool::{DSN_RETRY, SmtpSpool}; use crate::queue::throttle::IsAllowed; use crate::queue::{ Error, FROM_REPORT, HostResponse, MessageWrapper, Metadata, QueueEnvelope, QueuedMessage, @@ -155,7 +155,7 @@ impl QueuedMessage { let span_id = message.span_id; // Send any due Delivery Status Notifications - server.send_dsn(&mut message).await; + let dsn_status = server.send_dsn(&mut message).await; match has_pending_delivery { PendingDelivery::Yes(true) @@ -163,21 +163,27 @@ impl QueuedMessage { .message .next_delivery_event(self.queue_name.into()) .is_some_and(|due| due <= now()) => {} - PendingDelivery::No => { + PendingDelivery::No if dsn_status == DsnStatus::Completed => { trc::event!( Delivery(DeliveryEvent::Completed), SpanId = span_id, Elapsed = trc::Value::Duration((now() - message.message.created) * 1000) ); - // All message recipients expired, do not re-queue. (DSN has been already sent) + // All message recipients expired, do not re-queue. message.remove(&server, self.due.into()).await; return QueueEventStatus::Completed; } + PendingDelivery::No => { + message + .save_changes(&server, self.due.into(), Some(now() + DSN_RETRY)) + .await; + return QueueEventStatus::Deferred; + } _ => { // Re-queue the message if its not yet due for delivery - message.save_changes(&server, self.due.into()).await; + message.save_changes(&server, self.due.into(), None).await; return QueueEventStatus::Deferred; } } @@ -208,7 +214,7 @@ impl QueuedMessage { } } - message.save_changes(&server, self.due.into()).await; + message.save_changes(&server, self.due.into(), None).await; return QueueEventStatus::Deferred; } @@ -1485,7 +1491,7 @@ impl QueuedMessage { } // Send Delivery Status Notifications - server.send_dsn(&mut message).await; + let dsn_status = server.send_dsn(&mut message).await; // Notify queue manager if message.message.next_event(None).is_some() { @@ -1501,7 +1507,13 @@ impl QueuedMessage { ); // Save changes to disk - message.save_changes(&server, self.due.into()).await; + message.save_changes(&server, self.due.into(), None).await; + + QueueEventStatus::Deferred + } else if dsn_status == DsnStatus::Deferred { + message + .save_changes(&server, self.due.into(), Some(now() + DSN_RETRY)) + .await; QueueEventStatus::Deferred } else { diff --git a/crates/smtp/src/outbound/local.rs b/crates/smtp/src/outbound/local.rs index 8a89092..bd38fc6 100644 --- a/crates/smtp/src/outbound/local.rs +++ b/crates/smtp/src/outbound/local.rs @@ -120,7 +120,7 @@ impl MessageWrapper { ) .await; - message + let _ = message .queue( QueueParams::new(&autogenerated.message, self.span_id, server) .with_dkim_signers(dkim_signers) diff --git a/crates/smtp/src/queue/dsn.rs b/crates/smtp/src/queue/dsn.rs index f83c8c6..c08271b 100644 --- a/crates/smtp/src/queue/dsn.rs +++ b/crates/smtp/src/queue/dsn.rs @@ -10,9 +10,10 @@ use super::{ Recipient, Status, }; use crate::inbound::dkim::DkimSign; -use crate::queue::spool::QueueParams; +use crate::queue::spool::{DSN_RETRY, QueueParams}; use crate::queue::{MessageWrapper, UnexpectedResponse}; use common::Server; +use email::message::delivery::ORCPT_ADDR_TYPE; use mail_builder::MessageBuilder; use mail_builder::headers::HeaderType; use mail_builder::headers::content_type::ContentType; @@ -25,16 +26,24 @@ use std::fmt::Write; use std::future::Future; use store::write::now; +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum DsnStatus { + Completed, + Deferred, +} + pub trait SendDsn: Sync + Send { - fn send_dsn(&self, message: &mut MessageWrapper) -> impl Future + Send; + fn send_dsn(&self, message: &mut MessageWrapper) -> impl Future + Send; fn log_dsn(&self, message: &MessageWrapper) -> impl Future + Send; } impl SendDsn for Server { - async fn send_dsn(&self, message: &mut MessageWrapper) { + async fn send_dsn(&self, message: &mut MessageWrapper) -> DsnStatus { // Send DSN events self.log_dsn(message).await; + let mut status = DsnStatus::Completed; + if !message.message.return_path.is_empty() { // Build DSN if let Some(dsn) = message.build_dsn(self).await { @@ -51,12 +60,19 @@ impl SendDsn for Server { message.span_id, ) .await; - dsn_message + if dsn_message .queue( QueueParams::new(&dsn, message.span_id, self) .with_dkim_signers(dkim_signers), ) - .await; + .await + { + message.mark_dsn_sent(); + } else { + status = DsnStatus::Deferred; + } + } else { + message.mark_dsn_sent(); } } else { // Handle double bounce @@ -64,7 +80,9 @@ impl SendDsn for Server { } // Update next DSN notify times - message.update_next_dsn(self).await; + message.update_next_dsn(self, status).await; + + status } async fn log_dsn(&self, message: &MessageWrapper) { @@ -132,7 +150,7 @@ impl SendDsn for Server { const MAX_HEADER_SIZE: usize = 4096; impl MessageWrapper { - pub async fn build_dsn(&mut self, server: &Server) -> Option> { + pub async fn build_dsn(&self, server: &Server) -> Option> { let config = &server.core.smtp.queue; let now = now(); @@ -141,13 +159,12 @@ impl MessageWrapper { let mut txt_failed = String::new(); let mut dsn = String::new(); - for rcpt in &mut self.message.recipients { + for rcpt in &self.message.recipients { if rcpt.has_flag(RCPT_DSN_SENT | RCPT_NOTIFY_NEVER) { continue; } match &rcpt.status { Status::Completed(response) => { - rcpt.flags |= RCPT_DSN_SENT; if !rcpt.has_flag(RCPT_NOTIFY_SUCCESS) { continue; } @@ -164,7 +181,6 @@ impl MessageWrapper { response.write_dsn_text(&rcpt.address, &mut txt_delay); } Status::PermanentFailure(response) => { - rcpt.flags |= RCPT_DSN_SENT; if !rcpt.has_flag(RCPT_NOTIFY_FAILURE) { continue; } @@ -357,7 +373,7 @@ impl MessageWrapper { .into() } - pub async fn update_next_dsn(&mut self, server: &Server) { + pub async fn update_next_dsn(&mut self, server: &Server, status: DsnStatus) { let now = now(); let mut notify_changes = Vec::new(); for (rcpt_idx, rcpt) in self.message.recipients.iter().enumerate() { @@ -366,6 +382,11 @@ impl MessageWrapper { Status::TemporaryFailure(_) | Status::Scheduled ) && rcpt.notify.due <= now { + if status == DsnStatus::Deferred { + notify_changes.push((rcpt_idx, 0, now + DSN_RETRY)); + continue; + } + let envelope = QueueEnvelope::new(&self.message, rcpt); let queue_id = server @@ -391,6 +412,19 @@ impl MessageWrapper { } } + fn mark_dsn_sent(&mut self) { + for rcpt in &mut self.message.recipients { + if !rcpt.has_flag(RCPT_DSN_SENT | RCPT_NOTIFY_NEVER) + && matches!( + rcpt.status, + Status::Completed(_) | Status::PermanentFailure(_) + ) + { + rcpt.flags |= RCPT_DSN_SENT; + } + } + } + fn handle_double_bounce(&mut self) { let mut is_double_bounce = Vec::with_capacity(0); let now = now(); @@ -523,7 +557,7 @@ impl Message { impl Recipient { fn write_dsn(&self, dsn: &mut String) { if let Some(orcpt) = &self.orcpt { - let _ = write!(dsn, "Original-Recipient: rfc822;{orcpt}\r\n"); + let _ = write!(dsn, "Original-Recipient: {ORCPT_ADDR_TYPE}{orcpt}\r\n"); } let _ = write!(dsn, "Final-Recipient: rfc822;{}\r\n", self.address); } diff --git a/crates/smtp/src/queue/spool.rs b/crates/smtp/src/queue/spool.rs index 14a756f..c3e80ce 100644 --- a/crates/smtp/src/queue/spool.rs +++ b/crates/smtp/src/queue/spool.rs @@ -45,6 +45,7 @@ use utils::DomainPart; pub const LOCK_EXPIRY: u64 = 10 * 60; // 10 minutes pub const QUEUE_REFRESH: u64 = 5 * 60; // 5 minutes +pub const DSN_RETRY: u64 = 5 * 60; // 5 minutes pub(crate) const INFINITE_LOCK: u64 = 60 * 60 * 24 * 365; // 1 year const CANDIDATE_OVERSCAN: usize = 4; const MAX_PREALLOCATED_CANDIDATES: usize = 1024; @@ -370,6 +371,7 @@ pub(crate) struct QueueParams<'x, 'y> { } impl MessageWrapper { + #[must_use] pub(crate) async fn queue<'x, 'y>(mut self, mut params: QueueParams<'x, 'y>) -> bool { // Add DKIM signatures let dkim_headers = if params.dkim_signers.is_some() { @@ -669,7 +671,12 @@ impl MessageWrapper { recipient.queue = queue.virtual_queue; } - pub async fn save_changes(mut self, server: &Server, prev_event: Option) -> bool { + pub async fn save_changes( + mut self, + server: &Server, + prev_event: Option, + retry_at: Option, + ) -> bool { // Release quota for completed deliveries let mut batch = BatchBuilder::new(); self.release_quota(&mut batch); @@ -684,7 +691,12 @@ impl MessageWrapper { }, ))); } - for (queue_name, due) in self.message.next_events() { + let mut next_events = self.message.next_events(); + if let Some(retry_at) = retry_at { + let due = next_events.entry(self.queue_name).or_insert(retry_at); + *due = std::cmp::min(*due, retry_at); + } + for (queue_name, due) in next_events { batch.set( ValueClass::Queue(QueueClass::MessageEvent(store::write::QueueEvent { due, diff --git a/crates/smtp/src/reporting/dmarc.rs b/crates/smtp/src/reporting/dmarc.rs index bee93b2..57f1272 100644 --- a/crates/smtp/src/reporting/dmarc.rs +++ b/crates/smtp/src/reporting/dmarc.rs @@ -316,9 +316,6 @@ impl Session { if let Some(dkim2_output) = dkim2_output { report_record = report_record.with_dkim2_output(dkim2_output); } - if let Some(spf_ehlo) = &self.data.spf_ehlo { - report_record = report_record.with_spf_output(spf_ehlo, SPFDomainScope::Helo); - } if let Some(spf_mail_from) = &self.data.spf_mail_from { report_record = report_record.with_spf_output(spf_mail_from, SPFDomainScope::MailFrom); } diff --git a/crates/smtp/src/reporting/inbound.rs b/crates/smtp/src/reporting/inbound.rs index bba4f0e..82f1e5b 100644 --- a/crates/smtp/src/reporting/inbound.rs +++ b/crates/smtp/src/reporting/inbound.rs @@ -41,7 +41,7 @@ impl Session { self.data .rcpt_to .iter() - .any(|addr| analysis.is_report_address(addr.report_address())) + .any(|addr| analysis.is_report_address(addr.orig_address())) } } diff --git a/crates/smtp/src/reporting/send.rs b/crates/smtp/src/reporting/send.rs index ad2de84..7bde376 100644 --- a/crates/smtp/src/reporting/send.rs +++ b/crates/smtp/src/reporting/send.rs @@ -98,7 +98,7 @@ impl MtaReportSend for Server { let dkim_signers = self .eval_signers(sign_config, &message.message, parent_session_id) .await; - message + let _ = message .queue( QueueParams::new(&report, parent_session_id, self).with_dkim_signers(dkim_signers), ) @@ -130,7 +130,7 @@ impl MtaReportSend for Server { } else { None }; - message + let _ = message .queue( QueueParams::new(&raw_message, parent_session_id, self) .with_dkim_signers(dkim_signers), diff --git a/crates/smtp/src/scripts/envelope.rs b/crates/smtp/src/scripts/envelope.rs index d6e9bbd..e8e860b 100644 --- a/crates/smtp/src/scripts/envelope.rs +++ b/crates/smtp/src/scripts/envelope.rs @@ -12,6 +12,7 @@ use smtp_proto::{ use utils::DomainPart; use crate::core::{SessionAddress, SessionData}; +use email::message::delivery::ORCPT_ADDR_TYPE; impl SessionData { pub fn apply_envelope_modification(&mut self, envelope: Envelope, value: String) { @@ -111,7 +112,11 @@ impl SessionData { } Envelope::Orcpt => { if let Some(rcpt_to) = self.rcpt_to.last_mut() { - rcpt_to.dsn_info = value.into(); + rcpt_to.dsn_info = value + .strip_prefix(ORCPT_ADDR_TYPE) + .map(str::to_string) + .unwrap_or(value) + .into(); } } Envelope::Envid => { diff --git a/crates/smtp/src/scripts/event_loop.rs b/crates/smtp/src/scripts/event_loop.rs index 25295b4..350dde2 100644 --- a/crates/smtp/src/scripts/event_loop.rs +++ b/crates/smtp/src/scripts/event_loop.rs @@ -297,7 +297,7 @@ impl RunScript for Server { None }; - message + let _ = message .queue( QueueParams::new(raw_message, session_id, self) .with_dkim_signers(dkim_signers) diff --git a/crates/smtp/src/scripts/exec.rs b/crates/smtp/src/scripts/exec.rs index 2900eef..2eeeaf1 100644 --- a/crates/smtp/src/scripts/exec.rs +++ b/crates/smtp/src/scripts/exec.rs @@ -95,10 +95,8 @@ impl Session { params .envelope .push((Envelope::To, rcpt.address_lcase.to_string().into())); - if let Some(orcpt) = &rcpt.dsn_info { - params - .envelope - .push((Envelope::Orcpt, orcpt.as_str().to_lowercase().into())); + if let Some(orcpt) = rcpt.orcpt_parameter() { + params.envelope.push((Envelope::Orcpt, orcpt.into())); } } } else { @@ -109,10 +107,10 @@ impl Session { for rcpt in &self.data.rcpt_to { recipients.push(Variable::from(rcpt.address_lcase.to_string())); - orcpts.push(match &rcpt.dsn_info { + orcpts.push(match rcpt.orcpt_parameter() { Some(orcpt) => { has_orcpts = true; - Variable::from(orcpt.as_str().to_lowercase()) + Variable::from(orcpt) } None => Variable::default(), }); diff --git a/crates/spam-filter/Cargo.toml b/crates/spam-filter/Cargo.toml index 221adc3..3006789 100644 --- a/crates/spam-filter/Cargo.toml +++ b/crates/spam-filter/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "spam-filter" -version = "0.16.22" +version = "0.16.23" edition = "2024" [dependencies] diff --git a/crates/store/Cargo.toml b/crates/store/Cargo.toml index 06e399a..347e426 100644 --- a/crates/store/Cargo.toml +++ b/crates/store/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "store" -version = "0.16.22" +version = "0.16.23" edition = "2024" [dependencies] diff --git a/crates/store/src/backend/foundationdb/mod.rs b/crates/store/src/backend/foundationdb/mod.rs index 09035df..f189a7e 100644 --- a/crates/store/src/backend/foundationdb/mod.rs +++ b/crates/store/src/backend/foundationdb/mod.rs @@ -71,7 +71,7 @@ impl ReadVersion { } fn expire(&self) { - self.obtained.store(0, Ordering::Release); + self.version.store(0, Ordering::Release); } fn try_begin_refresh(&self) -> Option> { diff --git a/crates/store/src/backend/meili/main.rs b/crates/store/src/backend/meili/main.rs index e6cb29f..6db9f62 100644 --- a/crates/store/src/backend/meili/main.rs +++ b/crates/store/src/backend/meili/main.rs @@ -11,14 +11,12 @@ use crate::{ CalendarSearchField, ContactSearchField, EmailSearchField, SearchField, SearchableField, TracingSearchField, }, - write::now, }; use registry::schema::structs; use reqwest::{Error, Response, Url}; use serde_json::{Value, json}; use std::{sync::Arc, time::Duration}; -const UNCONFIRMED_TASK_RECHECK_DELAY: u64 = 600; pub(crate) const MAX_TOTAL_HITS: u64 = 100_000; impl MeiliSearchStore { @@ -300,18 +298,13 @@ impl MeiliSearchStore { } } - let err = trc::StoreEvent::MeilisearchError - .reason("Timed out waiting for Meilisearch task") - .id(task_uid); - - Err(if self.task_fail_on_timeout { - err + if self.task_fail_on_timeout { + Err(trc::StoreEvent::MeilisearchError + .reason("Timed out waiting for Meilisearch task") + .id(task_uid)) } else { - err.ctx( - trc::Key::NextRetry, - now().saturating_add(UNCONFIRMED_TASK_RECHECK_DELAY), - ) - }) + Ok(true) + } } } diff --git a/crates/store/src/backend/redis/lookup.rs b/crates/store/src/backend/redis/lookup.rs index 9f681d9..300624b 100644 --- a/crates/store/src/backend/redis/lookup.rs +++ b/crates/store/src/backend/redis/lookup.rs @@ -6,36 +6,37 @@ use super::{RedisPool, RedisStore, into_error}; use crate::{Deserialize, write::now}; -use redis::AsyncCommands; +use deadpool::managed::{Manager, Object, Pool}; +use redis::{AsyncCommands, RedisError, RedisResult, RetryMethod, Script}; +use std::sync::LazyLock; + +static INCR_EXPIRE: LazyLock