From 3a272096c006273cd030fbf86a0354d1ba707d1d Mon Sep 17 00:00:00 2001 From: John Coffey Date: Tue, 22 Sep 2026 16:31:25 -0700 Subject: [PATCH] Import upstream v0.16.23, stripped Upstream commit: 9d1c75ab68435e4417337f768291e5f947686203 Enterprise-only files removed or emptied: 63 Enterprise-only snippets removed: 118 in 50 files Dangling module declarations removed: 5 Edits turning enterprise off: 25 Third-party code: 14 files, 0 not in THIRD-PARTY.md Verification: clean One snippet more than v0.16.22, in crates/common/src/auth/authentication.rs (3, was 2). --- CHANGELOG.md | 33 ++ Cargo.lock | 325 ++++++++--------- crates/common/Cargo.toml | 2 +- crates/common/src/auth/authentication.rs | 178 +++++++-- crates/common/src/auth/oauth/token.rs | 2 +- crates/common/src/config/smtp/resolver.rs | 2 + crates/common/src/expr/functions/misc.rs | 7 + crates/common/src/expr/functions/mod.rs | 1 + crates/common/src/manager/application.rs | 342 ++++++++++++------ crates/common/src/network/acme/order.rs | 67 +++- crates/common/src/network/dns/update.rs | 30 ++ crates/common/src/network/mta.rs | 57 ++- crates/coordinator/Cargo.toml | 2 +- crates/dav-proto/Cargo.toml | 2 +- crates/dav/Cargo.toml | 2 +- crates/directory/Cargo.toml | 2 +- crates/email/Cargo.toml | 2 +- crates/email/src/message/delivery.rs | 8 + crates/email/src/sieve/ingest.rs | 3 +- crates/groupware/Cargo.toml | 2 +- crates/http-proto/Cargo.toml | 2 +- crates/http/Cargo.toml | 2 +- crates/http/src/api/diagnose.rs | 16 +- crates/imap-proto/Cargo.toml | 2 +- crates/imap/Cargo.toml | 2 +- crates/imap/src/op/copy_move.rs | 32 ++ crates/jmap-proto/Cargo.toml | 2 +- crates/jmap/Cargo.toml | 2 +- crates/jmap/src/email/copy.rs | 31 +- crates/main/Cargo.toml | 2 +- crates/main/src/test_data.rs | 4 +- crates/managesieve/Cargo.toml | 2 +- crates/migration/Cargo.toml | 2 +- crates/nlp/Cargo.toml | 2 +- crates/pop3/Cargo.toml | 2 +- crates/pop3/src/op/fetch.rs | 10 +- crates/pop3/src/protocol/response.rs | 105 +++++- crates/registry/Cargo.toml | 2 +- crates/scim-proto/Cargo.toml | 2 +- crates/scim/Cargo.toml | 2 +- crates/services/Cargo.toml | 2 +- crates/services/src/broadcast/subscriber.rs | 15 +- .../src/task_manager/destroy_account.rs | 4 +- crates/services/src/task_manager/index.rs | 22 +- crates/services/src/task_manager/manager.rs | 1 + crates/services/src/task_manager/mod.rs | 16 + crates/smtp/Cargo.toml | 2 +- crates/smtp/src/core/mod.rs | 12 +- crates/smtp/src/inbound/data.rs | 2 +- crates/smtp/src/inbound/rcpt.rs | 5 +- crates/smtp/src/inbound/spam.rs | 12 +- crates/smtp/src/outbound/delivery.rs | 30 +- crates/smtp/src/outbound/local.rs | 2 +- crates/smtp/src/queue/dsn.rs | 58 ++- crates/smtp/src/queue/spool.rs | 16 +- crates/smtp/src/reporting/dmarc.rs | 3 - crates/smtp/src/reporting/inbound.rs | 2 +- crates/smtp/src/reporting/send.rs | 4 +- crates/smtp/src/scripts/envelope.rs | 7 +- crates/smtp/src/scripts/event_loop.rs | 2 +- crates/smtp/src/scripts/exec.rs | 10 +- crates/spam-filter/Cargo.toml | 2 +- crates/store/Cargo.toml | 2 +- crates/store/src/backend/foundationdb/mod.rs | 2 +- crates/store/src/backend/meili/main.rs | 19 +- crates/store/src/backend/redis/lookup.rs | 236 ++++++------ crates/trc/Cargo.toml | 2 +- crates/trc/event-macro/Cargo.toml | 2 +- crates/types/Cargo.toml | 2 +- crates/utils/Cargo.toml | 2 +- crates/utils/proc-macros/Cargo.toml | 2 +- tests/Cargo.toml | 2 +- tests/src/directory/issuer.rs | 115 ++++++ tests/src/directory/mod.rs | 2 + tests/src/imap/antispam.rs | 35 ++ tests/src/imap/pop.rs | 22 +- tests/src/smtp/lookup/expressions.rs | 4 + tests/src/smtp/outbound/fallback_relay.rs | 4 +- tests/src/smtp/outbound/tls.rs | 4 +- tests/src/smtp/queue/dsn.rs | 65 +++- tests/src/smtp/queue/manager.rs | 6 +- tests/src/system/delivery.rs | 138 ++++++- 82 files changed, 1544 insertions(+), 646 deletions(-) create mode 100644 tests/src/directory/issuer.rs 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 55cb635..4acca72 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]] @@ -597,12 +597,6 @@ version = "0.13.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "9e1b586273c5702936fe7b7d6896644d8be71e6314cfe09d3167c95f712589e8" -[[package]] -name = "base64" -version = "0.21.7" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "9d297deb1925b89f2ccc13d7635fa0714f12c87adce1c75356b39ca9b7178567" - [[package]] name = "base64" version = "0.22.1" @@ -880,7 +874,7 @@ dependencies = [ "log", "num", "pin-project-lite", - "rand 0.10.2", + "rand 0.10.3", "rustls", "rustls-native-certs", "rustls-pki-types", @@ -990,7 +984,7 @@ checksum = "46d07918caa9eeaaf06b7873925c53a61daac173539b4f7715090745e44e4e69" dependencies = [ "proc-macro2", "quote", - "syn 3.0.5", + "syn 3.0.6", ] [[package]] @@ -1041,16 +1035,16 @@ dependencies = [ [[package]] name = "calcard" -version = "0.3.13" +version = "0.3.14" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "75b779382e675380a1ff8a4873acee5ba60158be7a00458fb98e6b85aad1f2ee" +checksum = "c601473ec15a875626bce73db1a1f0fc81e9c949ca25f1463b714947c0a70a1f" dependencies = [ "ahash", "chrono", "chrono-tz", "hashify", "jmap-tools", - "mail-builder 0.5.0", + "mail-builder 1.0.0", "mail-parser", "rkyv", "serde", @@ -1116,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", @@ -1166,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" @@ -1308,7 +1302,7 @@ dependencies = [ [[package]] name = "common" -version = "0.16.22" +version = "0.16.23" dependencies = [ "aes-gcm-siv", "ahash", @@ -1407,9 +1401,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", @@ -1492,7 +1486,7 @@ checksum = "3d52eff69cd5e647efe296129160853a42795992097e8af39800e1060caeea9b" [[package]] name = "coordinator" -version = "0.16.22" +version = "0.16.23" dependencies = [ "async-nats", "futures", @@ -1854,7 +1848,7 @@ dependencies = [ "proc-macro2", "quote", "strsim", - "syn 3.0.5", + "syn 3.0.6", ] [[package]] @@ -1887,7 +1881,7 @@ checksum = "2ac7135c3ef02b2f7833bbeb1be5ba7f966dcde8a87c6b87f65a778d71a02785" dependencies = [ "darling_core 0.24.1", "quote", - "syn 3.0.5", + "syn 3.0.6", ] [[package]] @@ -1904,7 +1898,7 @@ checksum = "4583a4551df46e2792f82ceeac45e850d2e2d5debba0b91f102385cda5b11f06" [[package]] name = "dav" -version = "0.16.22" +version = "0.16.23" dependencies = [ "calcard", "chrono", @@ -1927,7 +1921,7 @@ dependencies = [ [[package]] name = "dav-proto" -version = "0.16.22" +version = "0.16.23" dependencies = [ "calcard", "chrono", @@ -2141,7 +2135,7 @@ dependencies = [ [[package]] name = "directory" -version = "0.16.22" +version = "0.16.23" dependencies = [ "ahash", "argon2 0.6.0", @@ -2198,7 +2192,7 @@ checksum = "c6232dd377dcc64799954cbd3a9bb882e9cdc1308ccd87b1c098f1fb2eaf82a8" dependencies = [ "proc-macro2", "quote", - "syn 3.0.5", + "syn 3.0.6", ] [[package]] @@ -2308,11 +2302,11 @@ dependencies = [ [[package]] name = "ece" -version = "2.3.1" +version = "2.4.2" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "c2ea1d2f2cc974957a4e2575d8e5bb494549bab66338d6320c2789abcfff5746" +checksum = "c2467bac73e5a36d75e16cab0fa8d40676f075db6afde7d78b35f033e1f66e37" dependencies = [ - "base64 0.21.7", + "base64 0.22.1", "byteorder", "hex", "hkdf 0.12.4", @@ -2321,7 +2315,7 @@ dependencies = [ "openssl", "serde", "sha2 0.10.9", - "thiserror 1.0.69", + "thiserror 2.0.20", ] [[package]] @@ -2382,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", @@ -2490,10 +2484,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]] @@ -2576,7 +2570,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", ] @@ -2599,9 +2593,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" @@ -2715,7 +2709,7 @@ dependencies = [ "foundationdb-sys", "foundationdb-tuple", "futures", - "rand 0.10.2", + "rand 0.10.3", "serde", "serde_bytes", "serde_json", @@ -2847,7 +2841,7 @@ checksum = "9fb9654ba8355388abeb8dcb4fc62f511300867002afc858860463bdd9fe0c44" dependencies = [ "proc-macro2", "quote", - "syn 3.0.5", + "syn 3.0.6", ] [[package]] @@ -3018,7 +3012,7 @@ dependencies = [ [[package]] name = "groupware" -version = "0.16.22" +version = "0.16.23" dependencies = [ "ahash", "calcard", @@ -3173,7 +3167,7 @@ dependencies = [ "jni", "lru-cache", "parking_lot", - "rand 0.10.2", + "rand 0.10.3", "rustls", "rustls-pki-types", "rustls-platform-verifier", @@ -3200,7 +3194,7 @@ dependencies = [ "jni", "once_cell", "prefix-trie", - "rand 0.10.2", + "rand 0.10.3", "ring", "rustls-pki-types", "thiserror 2.0.20", @@ -3227,7 +3221,7 @@ dependencies = [ "ndk-context", "once_cell", "parking_lot", - "rand 0.10.2", + "rand 0.10.3", "resolv-conf", "rustls", "smallvec", @@ -3306,7 +3300,7 @@ dependencies = [ [[package]] name = "http" -version = "0.16.22" +version = "0.16.23" dependencies = [ "async-stream", "base64 0.23.1", @@ -3400,7 +3394,7 @@ dependencies = [ [[package]] name = "http_proto" -version = "0.16.22" +version = "0.16.23" dependencies = [ "common", "compact_str", @@ -3490,9 +3484,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", @@ -3886,7 +3880,7 @@ checksum = "65b27460c2c92b037f3f94c538ed9a3342f3fdf923606781629ccb35f82d042a" [[package]] name = "imap" -version = "0.16.22" +version = "0.16.23" dependencies = [ "ahash", "common", @@ -3898,7 +3892,7 @@ dependencies = [ "md5", "nlp", "parking_lot", - "rand 0.10.2", + "rand 0.10.3", "registry", "store", "tokio", @@ -3910,7 +3904,7 @@ dependencies = [ [[package]] name = "imap_proto" -version = "0.16.22" +version = "0.16.23" dependencies = [ "ahash", "base64 0.23.1", @@ -4086,25 +4080,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", ] @@ -4164,7 +4157,7 @@ dependencies = [ [[package]] name = "jmap" -version = "0.16.22" +version = "0.16.23" dependencies = [ "async-stream", "base64 0.23.1", @@ -4187,7 +4180,7 @@ dependencies = [ "mail-parser", "nlp", "p256", - "rand 0.10.2", + "rand 0.10.3", "registry", "reqwest 0.13.5", "rkyv", @@ -4245,7 +4238,7 @@ dependencies = [ [[package]] name = "jmap_proto" -version = "0.16.22" +version = "0.16.23" dependencies = [ "ahash", "calcard", @@ -4345,9 +4338,9 @@ dependencies = [ [[package]] name = "jsonwebtoken" -version = "11.0.0" +version = "11.1.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "881733cbc631fc9e472e24447ce32a64bedf2da498d6d8570b08edc87de71f65" +checksum = "e75fe14a82d81e5f5af639997db37d8b96045938a7ac6ab18cdbe1c7467e05e1" dependencies = [ "aws-lc-rs", "base64 0.22.1", @@ -4649,9 +4642,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" @@ -4692,9 +4685,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", @@ -4707,7 +4700,7 @@ dependencies = [ "mail-parser", "memchr", "quick-xml 0.42.0", - "rand 0.10.2", + "rand 0.10.3", "rkyv", "rsa", "rustls-pki-types", @@ -4750,7 +4743,7 @@ dependencies = [ [[package]] name = "managesieve" -version = "0.16.22" +version = "0.16.23" dependencies = [ "common", "compact_str", @@ -4885,7 +4878,7 @@ checksum = "c797b9d6bb23aab2fc369c65f871be49214f5c759af65bde26ffaaa2b646b492" [[package]] name = "migration" -version = "0.16.22" +version = "0.16.23" dependencies = [ "common", "email", @@ -5016,7 +5009,7 @@ dependencies = [ "proc-macro2", "quote", "rustversion", - "syn 3.0.5", + "syn 3.0.6", ] [[package]] @@ -5081,7 +5074,7 @@ dependencies = [ "lru", "mysql_common", "percent-encoding", - "rand 0.10.2", + "rand 0.10.3", "rustls", "serde", "socket2 0.6.5", @@ -5156,14 +5149,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", @@ -5476,8 +5469,8 @@ dependencies = [ [[package]] name = "opentelemetry" -version = "0.31.0" -source = "git+https://github.com/stalwartlabs/opentelemetry-rust#274b4d324794280ce6f4def095a3428197a9e6e3" +version = "0.32.0" +source = "git+https://github.com/stalwartlabs/opentelemetry-rust#80a14a3b6846f62f85506d68d2600c948fccc9d2" dependencies = [ "futures-core", "futures-sink", @@ -5489,8 +5482,8 @@ dependencies = [ [[package]] name = "opentelemetry-http" -version = "0.31.0" -source = "git+https://github.com/stalwartlabs/opentelemetry-rust#274b4d324794280ce6f4def095a3428197a9e6e3" +version = "0.32.0" +source = "git+https://github.com/stalwartlabs/opentelemetry-rust#80a14a3b6846f62f85506d68d2600c948fccc9d2" dependencies = [ "async-trait", "bytes", @@ -5501,10 +5494,11 @@ dependencies = [ [[package]] name = "opentelemetry-otlp" -version = "0.31.0" -source = "git+https://github.com/stalwartlabs/opentelemetry-rust#274b4d324794280ce6f4def095a3428197a9e6e3" +version = "0.32.0" +source = "git+https://github.com/stalwartlabs/opentelemetry-rust#80a14a3b6846f62f85506d68d2600c948fccc9d2" dependencies = [ "http 1.5.0", + "httpdate", "opentelemetry", "opentelemetry-http", "opentelemetry-proto", @@ -5519,8 +5513,8 @@ dependencies = [ [[package]] name = "opentelemetry-proto" -version = "0.31.0" -source = "git+https://github.com/stalwartlabs/opentelemetry-rust#274b4d324794280ce6f4def095a3428197a9e6e3" +version = "0.32.0" +source = "git+https://github.com/stalwartlabs/opentelemetry-rust#80a14a3b6846f62f85506d68d2600c948fccc9d2" dependencies = [ "opentelemetry", "opentelemetry_sdk", @@ -5531,13 +5525,13 @@ dependencies = [ [[package]] name = "opentelemetry-semantic-conventions" -version = "0.31.0" -source = "git+https://github.com/stalwartlabs/opentelemetry-rust#274b4d324794280ce6f4def095a3428197a9e6e3" +version = "0.32.1" +source = "git+https://github.com/stalwartlabs/opentelemetry-rust#80a14a3b6846f62f85506d68d2600c948fccc9d2" [[package]] name = "opentelemetry_sdk" -version = "0.31.0" -source = "git+https://github.com/stalwartlabs/opentelemetry-rust#274b4d324794280ce6f4def095a3428197a9e6e3" +version = "0.32.1" +source = "git+https://github.com/stalwartlabs/opentelemetry-rust#80a14a3b6846f62f85506d68d2600c948fccc9d2" dependencies = [ "futures-channel", "futures-executor", @@ -5987,7 +5981,7 @@ dependencies = [ [[package]] name = "pop3" -version = "0.16.22" +version = "0.16.23" dependencies = [ "common", "directory", @@ -6031,7 +6025,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", ] @@ -6068,9 +6062,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" @@ -6155,7 +6149,7 @@ dependencies = [ "proc-macro-error-attr3", "proc-macro2", "quote", - "syn 3.0.5", + "syn 3.0.6", ] [[package]] @@ -6236,9 +6230,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", ] @@ -6266,7 +6260,7 @@ checksum = "1c8d9ca532f185d5d4db7a7c9d51420b452168ea1c2b913953281bd6fe1fcbd0" dependencies = [ "proc-macro2", "quote", - "syn 3.0.5", + "syn 3.0.6", ] [[package]] @@ -6337,9 +6331,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", @@ -6357,16 +6351,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", @@ -6483,9 +6477,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", @@ -6725,7 +6719,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", @@ -6749,11 +6743,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", ] @@ -6775,7 +6768,7 @@ checksum = "92ecd8964f8453721699a1ed72037b0db49ce2f5a5138486ee89bed6f67cdf3a" dependencies = [ "proc-macro2", "quote", - "syn 3.0.5", + "syn 3.0.6", ] [[package]] @@ -6809,7 +6802,7 @@ checksum = "d6f6ff9a378485b298a5286656da665ba74413d36db0979633275d2e708145d4" [[package]] name = "registry" -version = "0.16.22" +version = "0.16.23" dependencies = [ "ahash", "hashify", @@ -6993,7 +6986,7 @@ checksum = "1c25ef604ac7dd839d44d64648952ea23c97866f124ff671b0ed2cf3ad9bb06e" dependencies = [ "proc-macro2", "quote", - "syn 3.0.5", + "syn 3.0.6", ] [[package]] @@ -7168,9 +7161,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", @@ -7181,9 +7174,9 @@ dependencies = [ [[package]] name = "rustls" -version = "0.23.44" +version = "0.23.45" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "6725596c3f2c3a0aef021139e145d4eafe314a6623e4680ca83852b2c67ab2ba" +checksum = "0d41d731c7d2f962d1ccc364cec258de3c0e93b38c2fb3ba97ac74513048d634" dependencies = [ "aws-lc-rs", "log", @@ -7355,12 +7348,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 = [ "common", "directory", @@ -7382,7 +7375,7 @@ dependencies = [ [[package]] name = "scim-proto" -version = "0.16.22" +version = "0.16.23" dependencies = [ "hashify", "serde", @@ -7563,7 +7556,7 @@ checksum = "e7a5d71263a5a7d47b41f6b3f06ba276f10cc18b0931f1799f710578e2309348" dependencies = [ "proc-macro2", "quote", - "syn 3.0.5", + "syn 3.0.6", ] [[package]] @@ -7574,7 +7567,7 @@ checksum = "f852137cce035d6a4df67ccce505ff6b3e9fd3a10e3e52b24dc71e650bb1a9bd" dependencies = [ "proc-macro2", "quote", - "syn 3.0.5", + "syn 3.0.6", ] [[package]] @@ -7610,7 +7603,7 @@ checksum = "8d3b1629de253c70a0508c3899572da79ca359fdab27c7920ff00406df418906" dependencies = [ "proc-macro2", "quote", - "syn 3.0.5", + "syn 3.0.6", ] [[package]] @@ -7655,7 +7648,7 @@ dependencies = [ "darling 0.24.1", "proc-macro2", "quote", - "syn 3.0.5", + "syn 3.0.6", ] [[package]] @@ -7703,12 +7696,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", @@ -8021,7 +8014,7 @@ checksum = "ba467056f1b547ed52077911161fc86985becbc60e8e1857c8a144dab0def891" [[package]] name = "smtp" -version = "0.16.22" +version = "0.16.23" dependencies = [ "ahash", "base64 0.23.1", @@ -8035,7 +8028,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", @@ -8111,7 +8104,7 @@ dependencies = [ [[package]] name = "spam-filter" -version = "0.16.22" +version = "0.16.23" dependencies = [ "common", "compact_str", @@ -8224,7 +8217,7 @@ checksum = "6ce2be8dc25455e1f91df71bfa12ad37d7af1092ae736f3a6cd0e37bc7810596" [[package]] name = "stalwart" -version = "0.16.22" +version = "0.16.23" dependencies = [ "common", "coordinator", @@ -8232,7 +8225,7 @@ dependencies = [ "directory", "email", "groupware", - "http 0.16.22", + "http 0.16.23", "http_proto", "imap", "jmap", @@ -8262,7 +8255,7 @@ checksum = "a2eb9349b6444b326872e140eb1cf5e7c522154d69e7a0ffb0fb81c06b37543f" [[package]] name = "store" -version = "0.16.22" +version = "0.16.23" dependencies = [ "ahash", "arc-swap", @@ -8288,7 +8281,7 @@ dependencies = [ "parking_lot", "r2d2", "radsort", - "rand 0.10.2", + "rand 0.10.3", "rayon", "redis", "registry", @@ -8384,9 +8377,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", @@ -8413,6 +8406,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" @@ -8511,7 +8515,7 @@ dependencies = [ [[package]] name = "tests" -version = "0.16.22" +version = "0.16.23" dependencies = [ "ahash", "aws-lc-rs", @@ -8533,7 +8537,7 @@ dependencies = [ "form_urlencoded", "futures", "groupware", - "http 0.16.22", + "http 0.16.23", "http_proto", "hyper", "hyper-util", @@ -8618,7 +8622,7 @@ checksum = "bc04cd3e1236dd4a98afca4569f2deb3f120e5422a4023be2cb683f8486292af" dependencies = [ "proc-macro2", "quote", - "syn 3.0.5", + "syn 3.0.6", ] [[package]] @@ -8703,18 +8707,9 @@ dependencies = [ [[package]] name = "tinyvec" -version = "1.13.2" +version = "1.13.3" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "4cf0ded5c4e56918d8f8a339e1bb67d038d3bc6d144ac407904015ba2e4cde9b" -dependencies = [ - "tinyvec_macros", -] - -[[package]] -name = "tinyvec_macros" -version = "0.1.1" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "1f3ccbac311fea05f86f61904b462b55fb3df8837a366dfc601a0161d0532f20" +checksum = "fd3ca314f692efd6c868f8408f53fe444634a845f96c028b97d35f6a1f79f0ee" [[package]] name = "tls-listener" @@ -8765,7 +8760,7 @@ checksum = "78773a2a397f451582ce068015985c33193cf6dea8b74d2a639fe457b2f07b0e" dependencies = [ "proc-macro2", "quote", - "syn 3.0.5", + "syn 3.0.6", ] [[package]] @@ -8787,7 +8782,7 @@ dependencies = [ "pin-project-lite", "postgres-protocol", "postgres-types", - "rand 0.10.2", + "rand 0.10.3", "socket2 0.6.5", "tokio", "tokio-util", @@ -8972,7 +8967,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", @@ -9111,7 +9106,7 @@ dependencies = [ [[package]] name = "trc" -version = "0.16.22" +version = "0.16.23" dependencies = [ "ahash", "base64 0.23.1", @@ -9170,7 +9165,7 @@ dependencies = [ "http 1.5.0", "httparse", "log", - "rand 0.10.2", + "rand 0.10.3", "sha1 0.11.0", "thiserror 2.0.20", ] @@ -9220,7 +9215,7 @@ checksum = "b6f5e870be6c3b371b77fe0ee0bafb859fa4964b4404c27de1d380043c4dda20" [[package]] name = "types" -version = "0.16.22" +version = "0.16.23" dependencies = [ "blake3", "compact_str", @@ -9272,9 +9267,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" @@ -9389,7 +9384,7 @@ checksum = "b6c140620e7ffbb22c2dee59cafe6084a59b5ffc27a8859a5f0d494b5d52b6be" [[package]] name = "utils" -version = "0.16.22" +version = "0.16.23" dependencies = [ "ahash", "arcstr", @@ -9591,7 +9586,7 @@ dependencies = [ "bumpalo", "proc-macro2", "quote", - "syn 3.0.5", + "syn 3.0.6", "wasm-bindgen-shared", ] @@ -10100,14 +10095,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]] @@ -10630,14 +10625,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]] @@ -10692,7 +10687,7 @@ checksum = "34df6fc39dbd26ddc9c10e6a2984476e13acce22e64e4487636ef494369225da" dependencies = [ "proc-macro2", "quote", - "syn 3.0.5", + "syn 3.0.6", ] [[package]] @@ -10724,9 +10719,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 c584f85..7749b5e 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 d24568c..9d81f98 100644 --- a/crates/common/src/auth/authentication.rs +++ b/crates/common/src/auth/authentication.rs @@ -9,7 +9,7 @@ use crate::{ auth::{ AccessToken, AuthRequest, DomainCache, credential::{ApiKey, AppPassword}, - oauth::GrantType, + oauth::{GrantType, token::TOKEN_HEADER}, }, }; use base64::{Engine, engine::general_purpose}; @@ -21,7 +21,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; @@ -319,19 +320,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. @@ -520,31 +514,78 @@ impl Server { Ok(self.get_default_directory()) } + 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()), + } + } + + fn get_directory_for_issuer(&self, issuer: &str) -> Option<&Arc> { + + None + } + pub fn get_directory_for_cached_domain(&self, domain: &DomainCache) -> Option<&Arc> { self.get_default_directory() } } -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 { @@ -642,3 +683,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 a82bef8..2400113 100644 --- a/crates/common/src/manager/application.rs +++ b/crates/common/src/manager/application.rs @@ -11,8 +11,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 +39,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 +84,7 @@ impl WebApplications { Self { applications: ArcSwap::new(Arc::new(Vec::new())), routes: ArcSwap::new(Arc::new(AHashMap::new())), + generation: AtomicU64::new(0), } } @@ -128,48 +134,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 +195,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 +215,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 +289,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 +301,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 +395,6 @@ impl Resource> { } } -#[derive(Clone)] pub struct TempDir { pub path: PathBuf, } @@ -371,11 +404,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 +443,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 +573,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!("stalwart-app-{name}"))); - dir.clean().await.unwrap(); + dir.create().await.unwrap(); tokio::fs::write(dir.path.join("index.html"), INDEX) .await .unwrap(); @@ -544,6 +596,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 +606,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 +618,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 +637,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 +649,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 +661,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 +677,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("stalwart-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("stalwart-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 41f47c8..9073c21 100644 --- a/crates/common/src/network/acme/order.rs +++ b/crates/common/src/network/acme/order.rs @@ -96,7 +96,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(); @@ -108,7 +139,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; @@ -117,7 +148,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(), ); @@ -126,19 +157,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; @@ -149,7 +181,7 @@ impl AcmeRequestBuilder { trc::event!( Acme(AcmeEvent::OrderProcessing), Url = self.directory.new_order.to_string(), - Hostname = domains.as_slice(), + Hostname = domains, Total = i, ); @@ -177,7 +209,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| { @@ -190,10 +222,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, @@ -211,7 +243,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(), ); @@ -226,6 +258,7 @@ impl AcmeRequestBuilder { server: &Server, url: &String, dns_parameters: Option<&AcmeDnsParameters>, + published: Option<&mut BTreeSet<(String, String)>>, ) -> AcmeResult<()> { let response = self .auth(url) @@ -287,7 +320,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 @@ -308,6 +346,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 ee36fe9..5fcc998 100644 --- a/crates/common/src/network/mta.rs +++ b/crates/common/src/network/mta.rs @@ -21,6 +21,7 @@ use crate::{ manager::SPAM_CLASSIFIER_KEY, network::RcptResolution, }; +use ahash::AHashSet; use directory::Recipient; use mail_auth::IpLookupStrategy; use registry::schema::{enums::ExpressionVariable, structs::MaskedEmail}; @@ -36,6 +37,7 @@ use store::{ }; use trc::{AddContext, SpamEvent}; use types::id::Id; +use utils::DomainPart; impl Server { pub async fn rcpt_resolve( @@ -130,7 +132,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 @@ -162,6 +167,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 bc71398..b280094 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 1ae7ade..f65faec 100644 --- a/crates/email/src/message/delivery.rs +++ b/crates/email/src/message/delivery.rs @@ -17,6 +17,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, @@ -35,6 +37,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 e3e966a..12c17c0 100644 --- a/crates/email/src/sieve/ingest.rs +++ b/crates/email/src/sieve/ingest.rs @@ -124,6 +124,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 @@ -139,7 +140,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 1d557f7..46ac09b 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 b07ea9e..e902b01 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 8258a73..c445203 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 8f40ccb..3611e3a 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 66e1175..4e50be4 100644 --- a/crates/main/Cargo.toml +++ b/crates/main/Cargo.toml @@ -7,7 +7,7 @@ homepage = "https://stalw.art" keywords = ["imap", "jmap", "smtp", "email", "mail", "webdav", "server"] categories = ["email"] license = "AGPL-3.0-only OR LicenseRef-SEL" -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 007594b..b873188 100644 --- a/crates/pop3/src/protocol/response.rs +++ b/crates/pop3/src/protocol/response.rs @@ -14,7 +14,7 @@ pub enum Response<'x, T> { List(Vec), Message { bytes: SliceRange<'x>, - lines: u32, + lines: Option, }, Capability { mechanisms: Vec, @@ -52,40 +52,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 } => { @@ -206,9 +231,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 91cc716..e17b311 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 868380f..408367a 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 d12f24b..fbcd4bc 100644 --- a/crates/services/src/broadcast/subscriber.rs +++ b/crates/services/src/broadcast/subscriber.rs @@ -94,6 +94,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)) => { @@ -155,9 +162,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 @@ -166,9 +171,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 4a33f32..565ea22 100644 --- a/crates/services/src/task_manager/destroy_account.rs +++ b/crates/services/src/task_manager/destroy_account.rs @@ -4,7 +4,7 @@ * SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-SEL */ -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; @@ -39,7 +39,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 f45cba1..49cdb4d 100644 --- a/crates/services/src/task_manager/index.rs +++ b/crates/services/src/task_manager/index.rs @@ -4,7 +4,7 @@ * SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-SEL */ -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, message::metadata::MessageMetadata}; use groupware::{cache::GroupwareCache, calendar::CalendarEvent, contact::ContactCard}; @@ -255,7 +255,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; @@ -312,7 +312,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; @@ -426,22 +426,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 b7d2a00..c598ada 100644 --- a/crates/services/src/task_manager/mod.rs +++ b/crates/services/src/task_manager/mod.rs @@ -134,4 +134,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 2bdd6b7..dba9212 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 d9198c3..2bcc3ec 100644 --- a/crates/smtp/src/inbound/data.rs +++ b/crates/smtp/src/inbound/data.rs @@ -441,7 +441,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 561c73e..50ab9a5 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 e22c3a8..bfd33e8 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