Compare commits

..
100 changed files with 2482 additions and 2522 deletions
-109
View File
@@ -1,112 +1,5 @@
# Changelog
## [2.62.0] - 2026-09-22
### API-Changes
- Remove unused ImapMessageMoved event.
- [**breaking**] get rid of `Qr::FprWithoutAddr` and `Qr::FprMismatch` variants.
- removed `DC_QR_FPR_WITHOUT_ADDR` and `DC_QR_FPR_MISMATCH` constants from CFFI
- removed FprWithoutAddr and FprMismatch variants from JSON-RPC QrObject returned by `check_qr`
### Features / Changes
- Don't show "member added" messages in channels ([#8730](https://github.com/chatmail/core/pull/8730)).
### Fixes
- Don't notify about a reaction sent by a blocked contact ([#8729](https://github.com/chatmail/core/pull/8729)).
- keep a message read status if a failure arrives later.
- prefer login errors over connection errors in the log message ([#8727](https://github.com/chatmail/core/pull/8727)).
- blocked or email contacts contacts are never recently seen or old ([#8733](https://github.com/chatmail/core/pull/8733)).
### Refactor
- Remove code that moves messages between IMAP folders.
## [2.61.0] - 2026-09-21
### API-Changes
- add `init_transports()` for multi-relay onboarding.
- add JSON-RPC API `is_sending_finished()`.
- [**breaking**] remove verification methods from the FFI and JSON-RPC APIs.
- `dc_contact_is_verified()` and `dc_contact_get_verifier_id()` are removed.
- the JSON-RPC Contact object loses the `isVerified` and `verifierId` fields. A bot reading `snapshot.is_verified` now gets an `AttributeError` at runtime.
- the Python bindings lose `Contact.is_verified()` and `Contact.get_verifier()`.
- `DC_STR_CONTACT_VERIFIED` (35) is removed, so UIs should stop registering a translation for it. A stock id core does not know is logged and otherwise ignored, so an un-updated client keeps working.
- [**breaking**] remove default value for "addr" config.
- [**breaking**] remove `addr` field from Account objects in JSON-RPC APIs.
- `list_transports()` should be used instead.
- [**breaking**] remove `is_chatmail` and the XCHATMAIL capability.
- `is_chatmail` is no longer a known config key.
- [**breaking**] remove `Contact.get_name_n_addr()` and related APIs.
- `dc_contact_get_name_n_addr()` CFFI is removed
- JSON-RPC contact objects don't have nameAndAddr field anymore
- [**breaking**] replace `was_seen_recently` by `freshness` in contact object.
- use contact's `freshness` instead of `seen_recently`
### Fixes
- always emit `AccountsBackgroundFetchDone`.
- make `background_fetch` not wait on or trigger SMTP connections.
- do not send a sync message when changing `configured_addr`.
- emit `SmtpMessageSent` event after deleting the message from SMTP queue.
- don't use extra STUN nine server for fallback.
- use `max_smtp_rcpt_to` chunking for the actual transport we are sending from.
- use correct `From` address when sending MDNs.
- use correct address for Bcc-self in unencrypted mails.
- don't emit configure progress events during background relay additions.
### Features / Changes
- better quality of image recoding ([#8682](https://github.com/chatmail/core/pull/8682)).
- perform background fetch from all transports.
- queue messages for SMTP before encryption.
- [**breaking**] stop tracking contact verification.
- the statistics JSON sent to the self-reporting-bot on Android changes: Contacts have `encrypted` instead of `verified` and lose `transitive_chain` properties and message stats have `encrypted` instead of `verified` and `unverified_encrypted`, and securejoin invites lose `already_verified`. The collecting bot stores incoming reports verbatim but analysis will have to make sense of older and newer reports.
- remove last usage of XDELTAPUSH capability.
- base server-side message deletion on `force_encryption`.
- do not restart I/O when setting `configured_addr`.
- do not use ConfiguredAddr when connecting to SMTP.
- mark autorelays for relay operators ([#8701](https://github.com/chatmail/core/pull/8701)).
- try fasted relays to attempt first configure on.
- do not mark message as failed for which we got a read receipt before.
- a single NDN does not mark a group message as failed.
### Documentation
- update `sys.msgsize_max_recommended` documentation.
- add hint about how to reset an invitation ([#8160](https://github.com/chatmail/core/pull/8160)).
- mention `get_app_version()` API in the changelog for 2.59.0.
- remove `protect_autocrypt` setting.
### Miscellaneous Tasks
- cleanup "primary" wording in comments.
- update rustls to 0.23.45.
- fix nightly "cargo" warnings.
- fix some types in `deltachat_rpc_client`.
### Refactor
- add `Encryption.is_encrypted()`.
- separate QueuedEncryption.
- sql: disable double-quoted string literals.
- add smtp::queue module.
- substitute configure progress macro with simple function call.
### Tests
- cleanup `get_smtp_rows_for_msg()`.
- cross-core securejoin invites for every chat type.
- explicitly empty url in appversion updates ([#8702](https://github.com/chatmail/core/pull/8702)).
- do not talk about "inconsistent key state" in `test_securejoin_after_contact_resetup`.
- rename `test_aeap_transition_{0,1}` to `test_aeap_transition_{single,group}`.
- do not fetch all messages when `direct_imap` is created.
- `direct_imap`: always pass `mark_seen=False` to fetch().
- add pseudo transport explicitly rather than by setting ConfiguredAddr.
## [2.60.0] - 2026-09-11
### API-Changes
@@ -8901,5 +8794,3 @@ https://github.com/chatmail/core/pulls?q=is%3Apr+is%3Aclosed
[2.58.0]: https://github.com/chatmail/core/compare/v2.57.0..v2.58.0
[2.59.0]: https://github.com/chatmail/core/compare/v2.58.0..v2.59.0
[2.60.0]: https://github.com/chatmail/core/compare/v2.59.0..v2.60.0
[2.61.0]: https://github.com/chatmail/core/compare/v2.60.0..v2.61.0
[2.62.0]: https://github.com/chatmail/core/compare/v2.61.0..v2.62.0
Generated
+50 -81
View File
@@ -1327,7 +1327,7 @@ dependencies = [
[[package]]
name = "deltachat"
version = "2.63.0-dev"
version = "2.61.0-dev"
dependencies = [
"anyhow",
"astral-tokio-tar",
@@ -1435,7 +1435,7 @@ dependencies = [
[[package]]
name = "deltachat-jsonrpc"
version = "2.63.0-dev"
version = "2.61.0-dev"
dependencies = [
"anyhow",
"async-channel 2.5.0",
@@ -1456,14 +1456,14 @@ dependencies = [
[[package]]
name = "deltachat-jsonrpc-bindings"
version = "2.63.0-dev"
version = "2.61.0-dev"
dependencies = [
"deltachat-jsonrpc",
]
[[package]]
name = "deltachat-repl"
version = "2.63.0-dev"
version = "2.61.0-dev"
dependencies = [
"anyhow",
"deltachat",
@@ -1479,7 +1479,7 @@ dependencies = [
[[package]]
name = "deltachat-rpc-server"
version = "2.63.0-dev"
version = "2.61.0-dev"
dependencies = [
"anyhow",
"deltachat",
@@ -1508,7 +1508,7 @@ dependencies = [
[[package]]
name = "deltachat_ffi"
version = "2.63.0-dev"
version = "2.61.0-dev"
dependencies = [
"anyhow",
"deltachat",
@@ -2091,12 +2091,6 @@ version = "0.1.5"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "d9c4f5dac5e15c24eb999c26181a6ca40b39fe946cbe4c263c7209467bc83af2"
[[package]]
name = "foldhash"
version = "0.2.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "77ce24cb58228fbb8aa041425bb1050850ac19177686ea6e0f41a70416f56fdb"
[[package]]
name = "foreign-types"
version = "0.3.2"
@@ -2417,34 +2411,16 @@ checksum = "5971ac85611da7067dbfcabef3c70ebb5606018acd9e2a3903a0da507521e0d5"
dependencies = [
"allocator-api2",
"equivalent",
"foldhash 0.1.5",
]
[[package]]
name = "hashbrown"
version = "0.16.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "841d1cc9bed7f9236f321df977030373f4a4163ae1a7dbfe1a51a2c1a51d9100"
dependencies = [
"foldhash 0.2.0",
]
[[package]]
name = "hashbrown"
version = "0.17.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "ed5909b6e89a2db4456e54cd5f673791d7eca6732202bbf2a9cc504fe2f9b84a"
dependencies = [
"foldhash 0.2.0",
"foldhash",
]
[[package]]
name = "hashlink"
version = "0.12.2"
version = "0.10.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "a596f1b20ed2cc5ecac41a164aaebc7258057060f06c0cf7a2ba3991ee7990fb"
checksum = "7382cf6263419f2d8df38c55d7da83da5c18aef87fc7a7fc1fb1e344edfe14c1"
dependencies = [
"hashbrown 0.17.1",
"hashbrown",
]
[[package]]
@@ -2978,7 +2954,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "4b0f83760fb341a774ed326568e19f5a863af4a952def8c39f9ab92fd95b88e5"
dependencies = [
"equivalent",
"hashbrown 0.15.4",
"hashbrown",
]
[[package]]
@@ -3285,12 +3261,11 @@ checksum = "b1a46d1a171d865aa5f83f92695765caa047a9b4cbae2cbf37dbd613a793fd4c"
[[package]]
name = "js-sys"
version = "0.3.105"
version = "0.3.77"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "ce57d20d1ea864ce2ac172ab472d409214f4fd359f0b2a2775abdf522e2af99e"
checksum = "1cfaf33c695fc6e08064efbc1f72ec937429614f25eef83af942d0e227c3a28f"
dependencies = [
"cfg-if",
"futures-util",
"once_cell",
"wasm-bindgen",
]
@@ -3369,11 +3344,12 @@ dependencies = [
[[package]]
name = "libsqlite3-sys"
version = "0.38.2"
version = "0.35.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "f1d20bef17f513b9b3004532233187769cd072d790971f4e4da0e346eb6401e8"
checksum = "133c182a6a2c87864fe97778797e46c7e999672690dc9fa3ee8e241aa4a9c13f"
dependencies = [
"cc",
"openssl-sys",
"pkg-config",
"vcpkg",
]
@@ -3445,7 +3421,7 @@ version = "0.12.5"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "234cf4f4a04dc1f57e24b96cc0cd600cf2af460d4161ac5ecdd0af8e1f3b2a38"
dependencies = [
"hashbrown 0.15.4",
"hashbrown",
]
[[package]]
@@ -5234,21 +5210,11 @@ dependencies = [
"zeroize",
]
[[package]]
name = "rsqlite-vfs"
version = "0.1.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "c51c9ae4df8a7fba42103df5c621fa3c37eccf3a3c650879e90fc48b11cc192c"
dependencies = [
"hashbrown 0.16.1",
"thiserror 2.0.20",
]
[[package]]
name = "rusqlite"
version = "0.40.2"
version = "0.37.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "23f2a97da3e3873c73cb2a2e71b35c40ff95e0b1eefa8d72d8499a6928c3b5b3"
checksum = "165ca6e57b20e1351573e3729b958bc62f0e48025386970b6e4d29e7a7e71f3f"
dependencies = [
"bitflags 2.11.0",
"fallible-iterator",
@@ -5256,7 +5222,6 @@ dependencies = [
"hashlink",
"libsqlite3-sys",
"smallvec",
"sqlite-wasm-rs",
]
[[package]]
@@ -5961,18 +5926,6 @@ dependencies = [
"der",
]
[[package]]
name = "sqlite-wasm-rs"
version = "0.5.5"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "dc3efc0da82635d7e1ced0053bbbfa8c7ab9645d0bf36ceb4f7127bb85315d75"
dependencies = [
"cc",
"js-sys",
"rsqlite-vfs",
"wasm-bindgen",
]
[[package]]
name = "stable_deref_trait"
version = "1.2.0"
@@ -6880,32 +6833,48 @@ checksum = "b8dad83b4f25e74f184f64c43b150b91efe7647395b42289f38e50566d82855b"
[[package]]
name = "wasm-bindgen"
version = "0.2.128"
version = "0.2.100"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "aecb87a33d3b0c5e3b7aa46336eaf486cffafbd281b195e4c8b80d50df2351bf"
checksum = "1edc8929d7499fc4e8f0be2262a241556cfc54a0bea223790e71446f2aab1ef5"
dependencies = [
"cfg-if",
"once_cell",
"rustversion",
"wasm-bindgen-macro",
]
[[package]]
name = "wasm-bindgen-backend"
version = "0.2.100"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "2f0a0651a5c2bc21487bde11ee802ccaf4c51935d0d3d42a6101f98161700bc6"
dependencies = [
"bumpalo",
"log",
"proc-macro2",
"quote",
"syn 2.0.118",
"wasm-bindgen-shared",
]
[[package]]
name = "wasm-bindgen-futures"
version = "0.4.78"
version = "0.4.50"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "6ef4c5d3d2cdf5c54f4231181768f5510842e350db025faf1f7163b1030ed928"
checksum = "555d470ec0bc3bb57890405e5d4322cc9ea83cebb085523ced7be4144dac1e61"
dependencies = [
"cfg-if",
"js-sys",
"once_cell",
"wasm-bindgen",
"web-sys",
]
[[package]]
name = "wasm-bindgen-macro"
version = "0.2.128"
version = "0.2.100"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "a690d511e3c1a8b3a55e33511e3c2c00c78415cd23650f32b808627f5696b9ed"
checksum = "7fe63fc6d09ed3792bd0897b314f53de8e16568c2b3f7982f468c0bf9bd0b407"
dependencies = [
"quote",
"wasm-bindgen-macro-support",
@@ -6913,22 +6882,22 @@ dependencies = [
[[package]]
name = "wasm-bindgen-macro-support"
version = "0.2.128"
version = "0.2.100"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "411e4887f0071ef2d2164a9d5fdf2d20efbef78fccd3a78b0c10a1dc5295e48a"
checksum = "8ae87ea40c9f689fc23f209965b6fb8a99ad69aeeb0231408be24920604395de"
dependencies = [
"bumpalo",
"proc-macro2",
"quote",
"syn 3.0.4",
"syn 2.0.118",
"wasm-bindgen-backend",
"wasm-bindgen-shared",
]
[[package]]
name = "wasm-bindgen-shared"
version = "0.2.128"
version = "0.2.100"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "81941cd78d0c92026c33e5e01312845a4cb1e9af3407f9134b100dd03144103e"
checksum = "1a05d73b933a847d6cccdda8f838a22ff101ad9bf93e33684f39c1f5f0eece3d"
dependencies = [
"unicode-ident",
]
@@ -6948,9 +6917,9 @@ dependencies = [
[[package]]
name = "web-sys"
version = "0.3.105"
version = "0.3.77"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "9fbddc4a036f00ec4f18c83445bd3115cb306a91da554919a099d9222fe4a7f8"
checksum = "33b6dd2ef9186f1f2072e409e99cd22a975331a6b3591b12c764e0e55c60d5d2"
dependencies = [
"js-sys",
"wasm-bindgen",
+4 -4
View File
@@ -1,6 +1,6 @@
[package]
name = "deltachat"
version = "2.63.0-dev"
version = "2.61.0-dev"
edition = "2024"
license = "MPL-2.0"
rust-version = "1.89"
@@ -85,7 +85,7 @@ quick-xml = { version = "0.41", features = ["escape-html"] }
rand-old = { package = "rand", version = "0.8" }
rand = { workspace = true }
regex = { workspace = true }
rusqlite = { workspace = true, features = ["backup"] }
rusqlite = { workspace = true, features = ["sqlcipher"] }
sanitize-filename = { workspace = true }
sdp = "0.17.1"
serde_json = { workspace = true }
@@ -194,7 +194,7 @@ nu-ansi-term = "0.50"
num-traits = "0.2"
rand = "0.9"
regex = "1.12"
rusqlite = "0.40.2"
rusqlite = "0.37"
sanitize-filename = "0.6"
serde = "1.0"
serde_json = "1"
@@ -209,7 +209,7 @@ yerpc = "0.7"
default = ["vendored"]
internals = []
vendored = [
"rusqlite/bundled",
"rusqlite/bundled-sqlcipher-vendored-openssl",
"async-native-tls/vendored"
]
+1 -1
View File
@@ -104,7 +104,7 @@ pub fn sanitize_name_and_addr(name: &str, addr: &str) -> (String, String) {
let mut name = sanitize_name(name);
// If the 'display name' is just the address, remove it:
// Otherwise, the contact would sometimes be shown as "alice@example.com (alice@example.com)".
// Otherwise, the contact would sometimes be shown as "alice@example.com (alice@example.com)" (see `get_name_n_addr()`).
// If the display name is empty, DC will just show the address when it needs a display name.
if name == addr {
name = "".to_string();
+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "deltachat_ffi"
version = "2.63.0-dev"
version = "2.61.0-dev"
description = "Deltachat FFI"
edition = "2024"
license = "MPL-2.0"
+70 -84
View File
@@ -304,6 +304,21 @@ dc_context_t* dc_context_new_closed (const char* dbfile);
int dc_context_open (dc_context_t *context, const char* passphrase);
/**
* Changes the passphrase on the open database.
* Deprecated 2025-11, see `dc_context_open()` for reasoning.
*
* Existing database must already be encrypted and the passphrase cannot be NULL or empty.
* It is impossible to encrypt unencrypted database with this method and vice versa.
*
* @memberof dc_context_t
* @param context The context object.
* @param passphrase The new passphrase.
* @return 1 on success, 0 on error.
*/
int dc_context_change_passphrase (dc_context_t* context, const char* passphrase);
/**
* Returns 1 if database is open.
*
@@ -438,9 +453,18 @@ char* dc_get_blobdir (const dc_context_t* context);
* always auto-downloaded.
* 0 = no limit (default).
* Changes affect future messages only.
* - `protect_autocrypt` = Enable Header Protection for Autocrypt header.
* This is an experimental option not compatible to other MUAs
* and older Delta Chat versions.
* 1 = enable.
* 0 = disable (default).
* - `gossip_period` = How often to gossip Autocrypt keys in chats with multiple recipients, in
* seconds. 2 days by default.
* This is not supposed to be changed by UIs and only used for testing.
* - `is_chatmail` = (deprecated) 1 if the the server is a chatmail server, 0 otherwise.
* This is deprecated, UIs should not behave differently
* for chatmail relays and classical email servers.
* Most usages in UIs can be replaced by `force_encryption`.
* - `is_muted` = Whether a context is muted by the user.
* Muted contexts should not sound, vibrate or show notifications.
* In contrast to `dc_set_chat_mute_duration()`,
@@ -673,21 +697,6 @@ char* dc_get_connectivity_html (dc_context_t* context);
void dc_configure (dc_context_t* context);
/**
* Add fake transport that cannot be used to connect.
*
* Used for offline tests only.
*
* To add a transport, use JSON-RPC calls `add_or_update_transport`
* and `add_transport_from_qr` instead.
*
* @memberof dc_context_t
* @param context The context object.
* @param addr The email address of the new transport.
*/
void dc_add_pseudo_transport (dc_context_t* context, const char *addr);
/**
* Check if the context is already configured.
*
@@ -2440,6 +2449,8 @@ void dc_stop_ongoing_process (dc_context_t* context);
#define DC_QR_ASK_VERIFYGROUP 202 // text1=groupname
#define DC_QR_ASK_VERIFYBROADCAST 204 // text1=broadcast name
#define DC_QR_FPR_OK 210 // id=contact
#define DC_QR_FPR_MISMATCH 220 // id=contact
#define DC_QR_FPR_WITHOUT_ADDR 230 // test1=formatted fingerprint
#define DC_QR_ACCOUNT 250 // text1=domain
#define DC_QR_BACKUP2 252
#define DC_QR_BACKUP_TOO_NEW 255
@@ -2480,6 +2491,13 @@ void dc_stop_ongoing_process (dc_context_t* context);
* ask the user if they want to start chatting;
* if so, call dc_create_chat_by_contact_id().
*
* - DC_QR_FPR_MISMATCH with dc_lot_t::id=Contact ID:
* scanned fingerprint does not match last seen fingerprint.
*
* - DC_QR_FPR_WITHOUT_ADDR with dc_lot_t::text1=Formatted fingerprint
* the scanned QR code contains a fingerprint but no e-mail address;
* suggest the user to establish an encrypted connection first.
*
* - DC_QR_ACCOUNT dc_lot_t::text1=domain:
* ask the user if they want to create an account on the given domain,
* if so, call dc_set_config_from_qr() and then dc_configure().
@@ -3169,28 +3187,17 @@ void dc_accounts_maybe_network_lost (dc_accounts_t* accounts);
/**
* Perform a background fetch for all accounts in parallel with a timeout.
* Pauses the scheduler, fetches from all transports at once and then resumes the scheduler.
* The fetch for an account ends as soon as one of its transports received messages.
*
* For an account with IO stopped, the scheduler is paused
* and every transport is fetched concurrently on a dedicated connection.
* The account is done as soon as one transport received messages, the others stop.
* Only one batch of messages is fetched per transport this way,
* so a larger backlog is left to the next call or to started IO.
*
* For an account with IO running, IMAP IDLE is interrupted on every transport
* and the account is done once every transport is.
*
* The call never waits for outgoing messages and never triggers sending them itself.
* Received messages may still queue replies, securejoin handshakes for example,
* which go out only while IO is running.
* dc_accounts_background_fetch() was created for the iOS Background fetch.
*
* The `DC_EVENT_ACCOUNTS_BACKGROUND_FETCH_DONE` event is emitted at the end,
* also on timeout, when another background fetch is already running
* and when the call is ignored because the timeout is too small,
* so it is safe to wait for the event whenever `accounts` is not NULL.
* Process all events until you get this one and you can safely return to the background
* without forgetting to create a generic notification if no message was fetched.
* The event carries no data identifying the call it belongs to,
* so it marks your own call only if no concurrent background fetch is happening.
* without forgetting to create notifications caused by timing race conditions.
*
* @memberof dc_accounts_t
* @param accounts The account manager as created by dc_accounts_new().
@@ -4888,38 +4895,6 @@ uint32_t dc_msg_get_saved_msg_id (const dc_msg_t* msg);
int dc_msg_is_pinned (const dc_msg_t* msg);
/**
* @defgroup DC_FRESHNESS DC_FRESHNESS
*
* These constants describe the freshness of a contact,
* as returned by dc_contact_get_freshness().
*
* @addtogroup DC_FRESHNESS
* @{
*/
/**
* Contact shall not be highlighted, e.g. neither shown with a "seen recently" dot
* nor with a "not seen for a long time" hint.
*/
#define DC_FRESHNESS_NORMAL 0
/**
* Contact was seen recently, the UI shall highlight it e.g. with a little green dot on the avatar.
*/
#define DC_FRESHNESS_RECENTLY_SEEN 1
/**
* Contact was not seen for a long time, the UI shall highlight it e.g. with a string
* below the contact name (e.g. "Seen 2 months ago").
*/
#define DC_FRESHNESS_OLD 2
/**
* @}
*/
/**
* @class dc_contact_t
*
@@ -4935,6 +4910,7 @@ uint32_t dc_msg_get_saved_msg_id (const dc_msg_t* msg);
* By default, these names are equal,
* but functions working with contact names
* (e.g. dc_contact_get_name(), dc_contact_get_display_name(),
* dc_contact_get_name_n_addr(),
* dc_create_contact() or dc_add_address_book())
* only affect the given-name.
*/
@@ -4984,7 +4960,7 @@ char* dc_contact_get_addr (const dc_contact_t* contact);
* The function does not return the contact name as received from the network.
*
* This name is typically used in a form where the user can edit the name of a contact.
* To get a fine name to display in lists etc., use dc_contact_get_display_name().
* To get a fine name to display in lists etc., use dc_contact_get_display_name() or dc_contact_get_name_n_addr().
*
* @memberof dc_contact_t
* @param contact The contact object.
@@ -5039,6 +5015,23 @@ char* dc_contact_get_display_name (const dc_contact_t* contact);
#define dc_contact_get_first_name dc_contact_get_display_name
/**
* Get a summary of name and address.
*
* The returned string is either "Name (email@domain.com)" or just
* "email@domain.com" if the name is unset.
*
* The summary is typically used when asking the user something about the contact.
* The attached e-mail address makes the question unique, e.g. "Chat with Alan Miller (am@uniquedomain.com)?"
*
* @memberof dc_contact_t
* @param contact The contact object.
* @return A summary string, must be released using dc_str_unref().
* Never returns NULL.
*/
char* dc_contact_get_name_n_addr (const dc_contact_t* contact);
/**
* Get the contact's profile image.
* This is the image set by each remote user on their own
@@ -5092,18 +5085,18 @@ int64_t dc_contact_get_last_seen (const dc_contact_t* contact);
/**
* Get the contact's freshness.
*
* The UI shall hightlight contacts that are recently seen by a little green dot on the avatar
* and contacts that were not seen for a long time by a string below the contact name (e.g. "Seen 2 months ago")
* Check if the contact was seen recently.
*
* The UI may highlight these contacts,
* eg. draw a little green dot on the avatars of the users recently seen.
* DC_CONTACT_ID_SELF and other special contact IDs are defined as never seen recently (they should not get a dot).
* To get the time a contact was seen, use dc_contact_get_last_seen().
*
* @memberof dc_contact_t
* @param contact The contact object.
* @return One of the @ref DC_FRESHNESS constants.
* @return 1=contact seen recently, 0=contact not seen recently.
*/
int dc_contact_get_freshness (const dc_contact_t* contact);
int dc_contact_was_seen_recently (const dc_contact_t* contact);
/**
@@ -5915,6 +5908,14 @@ void dc_event_unref(dc_event_t* event);
*/
#define DC_EVENT_IMAP_MESSAGE_DELETED 104
/**
* Emitted when a message was successfully moved on IMAP.
*
* @param data1 0
* @param data2 (char*) Info string in English language.
*/
#define DC_EVENT_IMAP_MESSAGE_MOVED 105
/**
* Emitted before going into IDLE on the Inbox folder.
*
@@ -6163,17 +6164,6 @@ void dc_event_unref(dc_event_t* event);
#define DC_EVENT_CHAT_DELETED 2023
/**
* The list of pinned messages for the chat has changed.
*
* Some message got pinned, or pinned message is unpinned or deleted.
*
* @param data1 (int) chat_id
* @param data2 (int) 0
*/
#define DC_EVENT_PINNED_MESSAGES_CHANGED 2024
/**
* Contact(s) created, renamed, blocked or deleted.
*
@@ -6338,10 +6328,6 @@ void dc_event_unref(dc_event_t* event);
* A call made while another background fetch is running gets the event immediately,
* and the running fetch keeps emitting events until its own marker.
*
* The event carries no data identifying the call it belongs to,
* so it marks your own call only if no concurrent background fetch is happening.
* Your own call has finished when dc_accounts_background_fetch() returns.
*
* This event is only emitted by the account manager
*/
+40 -26
View File
@@ -31,7 +31,6 @@ use deltachat::key::preconfigure_keypair;
use deltachat::message::MsgId;
use deltachat::qr_code_generator::{create_qr_svg, generate_backup_qr, get_securejoin_qr_svg};
use deltachat::stock_str::StockMessage;
use deltachat::transport::add_pseudo_transport;
use deltachat::webxdc::StatusUpdateSerial;
use deltachat::*;
use deltachat::{accounts::Accounts, log::LogExt};
@@ -160,6 +159,24 @@ pub unsafe extern "C" fn dc_context_open(
.unwrap_or(0)
}
#[unsafe(no_mangle)]
pub unsafe extern "C" fn dc_context_change_passphrase(
context: *mut dc_context_t,
passphrase: *const libc::c_char,
) -> libc::c_int {
if context.is_null() {
eprintln!("ignoring careless call to dc_context_change_passphrase()");
return 0;
}
let ctx = unsafe { &*context };
let passphrase = to_string_lossy(passphrase);
block_on(ctx.change_passphrase(passphrase))
.context("dc_context_change_passphrase() failed")
.log_err(ctx)
.is_ok() as libc::c_int
}
#[unsafe(no_mangle)]
pub unsafe extern "C" fn dc_context_is_open(context: *mut dc_context_t) -> libc::c_int {
if context.is_null() {
@@ -397,21 +414,6 @@ pub unsafe extern "C" fn dc_configure(context: *mut dc_context_t) {
spawn_configure(ctx.clone());
}
#[unsafe(no_mangle)]
pub unsafe extern "C" fn dc_add_pseudo_transport(
context: *mut dc_context_t,
addr: *const libc::c_char,
) {
if context.is_null() {
eprintln!("ignoring careless call to dc_add_pseudo_transport()");
return;
}
let ctx = unsafe { &*context };
let addr = to_string_lossy(addr);
block_on(add_pseudo_transport(ctx, &addr)).log_err(ctx).ok();
}
#[unsafe(no_mangle)]
pub unsafe extern "C" fn dc_is_configured(context: *mut dc_context_t) -> libc::c_int {
if context.is_null() {
@@ -475,6 +477,7 @@ pub unsafe extern "C" fn dc_event_get_id(event: *mut dc_event_t) -> libc::c_int
EventType::ImapConnected(_) => 102,
EventType::SmtpMessageSent(_) => 103,
EventType::ImapMessageDeleted(_) => 104,
EventType::ImapMessageMoved(_) => 105,
EventType::ImapInboxIdle => 106,
EventType::NewBlobFile(_) => 150,
EventType::DeletedBlobFile(_) => 151,
@@ -496,7 +499,6 @@ pub unsafe extern "C" fn dc_event_get_id(event: *mut dc_event_t) -> libc::c_int
EventType::ChatModified(_) => 2020,
EventType::ChatEphemeralTimerModified { .. } => 2021,
EventType::ChatDeleted { .. } => 2023,
EventType::PinnedMessagesChanged { .. } => 2024,
EventType::ContactsChanged(_) => 2030,
EventType::LocationChanged(_) => 2035,
EventType::ConfigureProgress { .. } => 2041,
@@ -542,6 +544,7 @@ pub unsafe extern "C" fn dc_event_get_data1_int(event: *mut dc_event_t) -> libc:
| EventType::ImapConnected(_)
| EventType::SmtpMessageSent(_)
| EventType::ImapMessageDeleted(_)
| EventType::ImapMessageMoved(_)
| EventType::ImapInboxIdle
| EventType::NewBlobFile(_)
| EventType::DeletedBlobFile(_)
@@ -570,8 +573,7 @@ pub unsafe extern "C" fn dc_event_get_data1_int(event: *mut dc_event_t) -> libc:
| EventType::MsgReadCountChanged { chat_id, .. }
| EventType::ChatModified(chat_id)
| EventType::ChatEphemeralTimerModified { chat_id, .. }
| EventType::ChatDeleted { chat_id }
| EventType::PinnedMessagesChanged { chat_id } => chat_id.to_u32() as libc::c_int,
| EventType::ChatDeleted { chat_id } => chat_id.to_u32() as libc::c_int,
EventType::ContactsChanged(id) | EventType::LocationChanged(id) => {
let id = id.unwrap_or_default();
id.to_u32() as libc::c_int
@@ -617,6 +619,7 @@ pub unsafe extern "C" fn dc_event_get_data2_int(event: *mut dc_event_t) -> libc:
| EventType::ImapConnected(_)
| EventType::SmtpMessageSent(_)
| EventType::ImapMessageDeleted(_)
| EventType::ImapMessageMoved(_)
| EventType::ImapInboxIdle
| EventType::NewBlobFile(_)
| EventType::DeletedBlobFile(_)
@@ -645,8 +648,7 @@ pub unsafe extern "C" fn dc_event_get_data2_int(event: *mut dc_event_t) -> libc:
| EventType::OutgoingCallAccepted { .. }
| EventType::CallEnded { .. }
| EventType::EventChannelOverflow { .. }
| EventType::TransportsModified
| EventType::PinnedMessagesChanged { .. } => 0,
| EventType::TransportsModified => 0,
EventType::MsgsChanged { msg_id, .. }
| EventType::ReactionsChanged { msg_id, .. }
| EventType::IncomingReaction { msg_id, .. }
@@ -712,6 +714,7 @@ pub unsafe extern "C" fn dc_event_get_data2_str(event: *mut dc_event_t) -> *mut
| EventType::ImapConnected(msg)
| EventType::SmtpMessageSent(msg)
| EventType::ImapMessageDeleted(msg)
| EventType::ImapMessageMoved(msg)
| EventType::NewBlobFile(msg)
| EventType::DeletedBlobFile(msg)
| EventType::Warning(msg)
@@ -747,8 +750,7 @@ pub unsafe extern "C" fn dc_event_get_data2_str(event: *mut dc_event_t) -> *mut
| EventType::AccountsItemChanged
| EventType::IncomingCallAccepted { .. }
| EventType::WebxdcRealtimeAdvertisementReceived { .. }
| EventType::TransportsModified
| EventType::PinnedMessagesChanged { .. } => ptr::null_mut(),
| EventType::TransportsModified => ptr::null_mut(),
EventType::IncomingCall {
place_call_info, ..
} => place_call_info.strdup(),
@@ -3961,6 +3963,18 @@ pub unsafe extern "C" fn dc_contact_get_display_name(
ffi_contact.contact.get_display_name().strdup()
}
#[unsafe(no_mangle)]
pub unsafe extern "C" fn dc_contact_get_name_n_addr(
contact: *mut dc_contact_t,
) -> *mut libc::c_char {
if contact.is_null() {
eprintln!("ignoring careless call to dc_contact_get_name_n_addr()");
return "".strdup();
}
let ffi_contact = unsafe { &*contact };
ffi_contact.contact.get_name_n_addr().strdup()
}
#[unsafe(no_mangle)]
pub unsafe extern "C" fn dc_contact_get_profile_image(
contact: *mut dc_contact_t,
@@ -4016,13 +4030,13 @@ pub unsafe extern "C" fn dc_contact_get_last_seen(contact: *mut dc_contact_t) ->
}
#[unsafe(no_mangle)]
pub unsafe extern "C" fn dc_contact_get_freshness(contact: *mut dc_contact_t) -> libc::c_int {
pub unsafe extern "C" fn dc_contact_was_seen_recently(contact: *mut dc_contact_t) -> libc::c_int {
if contact.is_null() {
eprintln!("ignoring careless call to dc_contact_get_freshness()");
eprintln!("ignoring careless call to dc_contact_was_seen_recently()");
return 0;
}
let ffi_contact = unsafe { &*contact };
u32::from(ffi_contact.contact.get_freshness()) as libc::c_int
ffi_contact.contact.was_seen_recently() as libc::c_int
}
#[unsafe(no_mangle)]
+12
View File
@@ -47,6 +47,8 @@ impl Lot {
Qr::AskVerifyGroup { grpname, .. } => Some(Cow::Borrowed(grpname)),
Qr::AskJoinBroadcast { name, .. } => Some(Cow::Borrowed(name)),
Qr::FprOk { .. } => None,
Qr::FprMismatch { .. } => None,
Qr::FprWithoutAddr { fingerprint, .. } => Some(Cow::Borrowed(fingerprint)),
Qr::Account { domain } => Some(Cow::Borrowed(domain)),
Qr::Backup2 { .. } => None,
Qr::BackupTooNew { .. } => None,
@@ -101,6 +103,8 @@ impl Lot {
Qr::AskVerifyGroup { .. } => LotState::QrAskVerifyGroup,
Qr::AskJoinBroadcast { .. } => LotState::QrAskJoinBroadcast,
Qr::FprOk { .. } => LotState::QrFprOk,
Qr::FprMismatch { .. } => LotState::QrFprMismatch,
Qr::FprWithoutAddr { .. } => LotState::QrFprWithoutAddr,
Qr::Account { .. } => LotState::QrAccount,
Qr::Backup2 { .. } => LotState::QrBackup2,
Qr::BackupTooNew { .. } => LotState::QrBackupTooNew,
@@ -128,6 +132,8 @@ impl Lot {
Qr::AskVerifyGroup { .. } => Default::default(),
Qr::AskJoinBroadcast { .. } => Default::default(),
Qr::FprOk { contact_id } => contact_id.to_u32(),
Qr::FprMismatch { contact_id } => contact_id.unwrap_or_default().to_u32(),
Qr::FprWithoutAddr { .. } => Default::default(),
Qr::Account { .. } => Default::default(),
Qr::Backup2 { .. } => Default::default(),
Qr::BackupTooNew { .. } => Default::default(),
@@ -175,6 +181,12 @@ pub enum LotState {
/// id=contact
QrFprOk = 210,
/// id=contact
QrFprMismatch = 220,
/// text1=formatted fingerprint
QrFprWithoutAddr = 230,
/// text1=domain
QrAccount = 250,
+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "deltachat-jsonrpc-bindings"
version = "2.63.0-dev"
version = "2.61.0-dev"
description = "Autogenerate DeltaChat JSON-RPC API bindings at build time"
edition = "2024"
license = "MPL-2.0"
@@ -54,5 +54,5 @@
},
"type": "module",
"types": "dist/deltachat.d.ts",
"version": "2.63.0-dev"
"version": "2.61.0-dev"
}
+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "deltachat-jsonrpc"
version = "2.63.0-dev"
version = "2.61.0-dev"
description = "DeltaChat JSON-RPC API"
edition = "2024"
license = "MPL-2.0"
+7 -28
View File
@@ -278,26 +278,10 @@ impl CommandApi {
/// Performs a background fetch for all accounts in parallel with a timeout.
///
/// For an account with IO stopped, the scheduler is paused
/// and every transport is fetched concurrently on a dedicated connection.
/// The account is done as soon as one transport received messages, the others stop.
/// Only one batch of messages is fetched per transport this way,
/// so a larger backlog is left to the next call or to started IO.
///
/// For an account with IO running, IMAP IDLE is interrupted on every transport
/// and the account is done once every transport is.
///
/// The call never waits for outgoing messages and never triggers sending them itself.
/// Received messages may still queue replies, securejoin handshakes for example,
/// which go out only while IO is running.
/// Use `is_sending_finished()` to tell whether the outgoing queue is empty.
///
/// The `AccountsBackgroundFetchDone` event is emitted at the end even in case of timeout,
/// and immediately if another background fetch is already running.
/// Process all events until you get this one and you can safely return to the background
/// without forgetting to create a generic notification if no message was fetched.
/// The event carries no data identifying the call it belongs to,
/// so it marks your own call only if no concurrent background fetch is happening.
/// without forgetting to create notifications caused by timing race conditions.
async fn background_fetch(&self, timeout_in_seconds: f64) -> Result<()> {
let future = {
let lock = self.accounts.read().await;
@@ -308,11 +292,6 @@ impl CommandApi {
Ok(())
}
/// Stops an ongoing `background_fetch()` call, making it return early
/// without waiting for the remaining transports or for the timeout.
///
/// The `AccountsBackgroundFetchDone` event is emitted as usual.
/// Does nothing if no background fetch is running.
async fn stop_background_fetch(&self) -> Result<()> {
self.accounts.read().await.stop_background_fetch();
Ok(())
@@ -546,13 +525,13 @@ impl CommandApi {
ctx.add_transport_from_qr(&qr).await
}
/// Adds an initial transport on the chatmail relay that answers fastest
/// and lets the profile add further ones in the background.
/// Automatically adds up to three transports.
///
/// A `DCACCOUNT:` or `DCLOGIN:` `qr` code adds a single transport
/// while securejoin codes add the inviter's relays to the candidates.
///
/// Does nothing if the profile already has a transport.
/// If the user just scanned a QR code of type `Account`, `Login`,
/// `AskVerifyContact`, `AskVerifyGroup`, or `AskJoinBroadcast`,
/// then UI implementations should pass it as the `qr` parameter.
/// The host(s) from the QR code will then also be considered
/// for creating an account there.
async fn init_transports(&self, account_id: u32, qr: Option<String>) -> Result<()> {
let ctx = self.get_context(account_id).await?;
ctx.init_transports(qr.as_deref()).await
@@ -13,6 +13,7 @@ pub enum Account {
Configured {
id: u32,
display_name: Option<String>,
addr: Option<String>,
// size: u32,
profile_image: Option<String>,
color: String,
@@ -28,6 +29,7 @@ impl Account {
pub async fn from_context(ctx: &deltachat::context::Context, id: u32) -> Result<Self> {
if ctx.is_configured().await? {
let display_name = ctx.get_config(Config::Displayname).await?;
let addr = ctx.get_config(Config::Addr).await?;
let profile_image = ctx.get_config(Config::Selfavatar).await?;
let color = color_int_to_hex_string(
Contact::get_by_id(ctx, ContactId::SELF)
@@ -39,6 +41,7 @@ impl Account {
Ok(Account::Configured {
id,
display_name,
addr,
profile_image,
color,
private_tag,
+8 -11
View File
@@ -9,8 +9,6 @@ use deltachat::context::Context;
use serde::{Deserialize, Serialize};
use typescript_type_def::TypeDef;
use crate::api::types::contact::ContactFreshness;
use super::color_int_to_hex_string;
#[derive(Serialize, TypeDef, schemars::JsonSchema)]
@@ -71,7 +69,7 @@ pub struct FullChat {
is_muted: bool,
ephemeral_timer: u32,
can_send: bool,
freshness: ContactFreshness,
was_seen_recently: bool,
mailing_list_address: Option<String>,
}
@@ -94,17 +92,16 @@ impl FullChat {
let can_send = chat.can_send(context).await?;
let freshness = if chat.get_type() == Chattype::Single {
let was_seen_recently = if chat.get_type() == Chattype::Single {
match contact_ids.first() {
Some(contact) => Contact::get_by_id(context, *contact)
.await
.context("failed to load contact for get_freshness")?
.get_freshness()
.into(),
None => ContactFreshness::Normal,
.context("failed to load contact for was_seen_recently")?
.was_seen_recently(),
None => false,
}
} else {
ContactFreshness::Normal
false
};
let mailing_list_address = chat.get_mailinglist_addr().map(|s| s.to_string());
@@ -129,7 +126,7 @@ impl FullChat {
is_muted: chat.is_muted(),
ephemeral_timer,
can_send,
freshness,
was_seen_recently,
mailing_list_address,
})
}
@@ -140,7 +137,7 @@ impl FullChat {
/// - fresh_message_counter
/// - ephemeral_timer
/// - self_in_group
/// - freshness
/// - was_seen_recently
/// - can_send
///
/// used when you only need the basic metadata of a chat like type, name, profile picture
+11 -11
View File
@@ -11,8 +11,6 @@ use num_traits::cast::ToPrimitive;
use serde::Serialize;
use typescript_type_def::TypeDef;
use crate::api::types::contact::ContactFreshness;
use super::chat::JsonrpcChatType;
use super::color_int_to_hex_string;
use super::message::MessageViewtype;
@@ -70,7 +68,7 @@ pub enum ChatListItemFetchResult {
is_contact_request: bool,
/// contact id if this is a dm chat (for view profile entry in context menu)
dm_chat_contact: Option<u32>,
freshness: ContactFreshness,
was_seen_recently: bool,
last_message_type: Option<MessageViewtype>,
last_message_id: Option<u32>,
},
@@ -129,20 +127,22 @@ pub(crate) async fn get_chat_list_item_by_id(
None => (None, None),
};
let (dm_chat_contact, freshness) = if chat.get_type() == Chattype::Single {
let (dm_chat_contact, was_seen_recently) = if chat.get_type() == Chattype::Single {
let chat_contacts = get_chat_contacts(ctx, chat_id).await?;
let contact = chat_contacts.first();
let freshness = match contact {
let was_seen_recently = match contact {
Some(contact) => Contact::get_by_id(ctx, *contact)
.await
.context("contact")?
.get_freshness()
.into(),
None => ContactFreshness::Normal,
.was_seen_recently(),
None => false,
};
(contact.map(|contact_id| contact_id.to_u32()), freshness)
(
contact.map(|contact_id| contact_id.to_u32()),
was_seen_recently,
)
} else {
(None, ContactFreshness::Normal)
(None, false)
};
let color = color_int_to_hex_string(chat.get_color(ctx).await?);
@@ -170,7 +170,7 @@ pub(crate) async fn get_chat_list_item_by_id(
is_muted: chat.is_muted(),
is_contact_request: chat.is_contact_request(),
dm_chat_contact,
freshness,
was_seen_recently,
last_message_type: message_type,
last_message_id: last_msgid.map(|id| id.to_u32()),
})
+4 -24
View File
@@ -1,5 +1,4 @@
use anyhow::Result;
use deltachat::contact;
use deltachat::context::Context;
use deltachat::key::{DcKey, SignedPublicKey};
use serde::Serialize;
@@ -7,27 +6,6 @@ use typescript_type_def::TypeDef;
use super::color_int_to_hex_string;
/// Freshness of a contact, based on when it was last seen.
#[derive(Serialize, TypeDef, schemars::JsonSchema)]
pub enum ContactFreshness {
/// Contact shall not be highlighted.
Normal,
/// Contact was seen recently.
RecentlySeen,
/// Contact was not seen for a long time.
Old,
}
impl From<contact::Freshness> for ContactFreshness {
fn from(freshness: contact::Freshness) -> Self {
match freshness {
contact::Freshness::Normal => ContactFreshness::Normal,
contact::Freshness::RecentlySeen => ContactFreshness::RecentlySeen,
contact::Freshness::Old => ContactFreshness::Old,
}
}
}
#[derive(Serialize, TypeDef, schemars::JsonSchema)]
#[serde(rename = "Contact", rename_all = "camelCase")]
pub struct ContactObject {
@@ -39,6 +17,7 @@ pub struct ContactObject {
id: u32,
name: String,
profile_image: Option<String>, // BLOBS
name_and_addr: String,
is_blocked: bool,
/// Is the contact a key contact.
@@ -54,7 +33,7 @@ pub struct ContactObject {
/// the contact's last seen timestamp
last_seen: i64,
freshness: ContactFreshness,
was_seen_recently: bool,
/// If the contact is a bot.
is_bot: bool,
@@ -78,11 +57,12 @@ impl ContactObject {
id: contact.id.to_u32(),
name: contact.get_name().to_owned(),
profile_image, //BLOBS
name_and_addr: contact.get_name_n_addr(),
is_blocked: contact.is_blocked(),
is_key_contact: contact.is_key_contact(),
e2ee_avail: contact.e2ee_avail(context).await?,
last_seen: contact.last_seen(),
freshness: contact.get_freshness().into(),
was_seen_recently: contact.was_seen_recently(),
is_bot: contact.is_bot(),
})
}
+4 -12
View File
@@ -44,6 +44,9 @@ pub enum EventType {
/// Emitted when an IMAP message has been marked as deleted
ImapMessageDeleted { msg: String },
/// Emitted when an IMAP message has been moved
ImapMessageMoved { msg: String },
/// Emitted before going into IDLE on the Inbox folder.
ImapInboxIdle,
@@ -258,15 +261,6 @@ pub enum EventType {
chat_id: u32,
},
/// The list of pinned messages for the chat has changed.
///
/// Some message got pinned, or pinned message is unpinned or deleted.
#[serde(rename_all = "camelCase")]
PinnedMessagesChanged {
/// ID of the chat where the list of pinned messages changed.
chat_id: u32,
},
/// Contact(s) created, renamed, blocked or deleted.
#[serde(rename_all = "camelCase")]
ContactsChanged {
@@ -502,6 +496,7 @@ impl From<CoreEventType> for EventType {
CoreEventType::ImapConnected(msg) => ImapConnected { msg },
CoreEventType::SmtpMessageSent(msg) => SmtpMessageSent { msg },
CoreEventType::ImapMessageDeleted(msg) => ImapMessageDeleted { msg },
CoreEventType::ImapMessageMoved(msg) => ImapMessageMoved { msg },
CoreEventType::ImapInboxIdle => ImapInboxIdle,
CoreEventType::NewBlobFile(file) => NewBlobFile { file },
CoreEventType::DeletedBlobFile(file) => DeletedBlobFile { file },
@@ -521,9 +516,6 @@ impl From<CoreEventType> for EventType {
msg_id: msg_id.to_u32(),
contact_id: contact_id.to_u32(),
},
CoreEventType::PinnedMessagesChanged { chat_id } => PinnedMessagesChanged {
chat_id: chat_id.to_u32(),
},
CoreEventType::IncomingReaction {
chat_id,
contact_id,
+15
View File
@@ -68,6 +68,16 @@ pub enum QrObject {
/// Contact ID.
contact_id: u32,
},
/// Scanned fingerprint does not match the last seen fingerprint.
FprMismatch {
/// Contact ID.
contact_id: Option<u32>,
},
/// The scanned QR code contains a fingerprint but no e-mail address.
FprWithoutAddr {
/// Key fingerprint.
fingerprint: String,
},
/// Ask the user if they want to create an account on the given domain.
Account {
/// Server domain name.
@@ -286,6 +296,11 @@ impl From<Qr> for QrObject {
let contact_id = contact_id.to_u32();
QrObject::FprOk { contact_id }
}
Qr::FprMismatch { contact_id } => {
let contact_id = contact_id.map(|contact_id| contact_id.to_u32());
QrObject::FprMismatch { contact_id }
}
Qr::FprWithoutAddr { fingerprint } => QrObject::FprWithoutAddr { fingerprint },
Qr::Account { domain } => QrObject::Account { domain },
Qr::Backup2 {
ref node_addr,
+1 -1
View File
@@ -12,7 +12,7 @@ pub struct JsonrpcReaction {
emoji: String,
/// Emoji frequency.
count: u32,
count: usize,
/// True if we reacted with this emoji.
is_from_self: bool,
+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "deltachat-repl"
version = "2.63.0-dev"
version = "2.61.0-dev"
license = "MPL-2.0"
edition = "2024"
repository = "https://github.com/chatmail/core"
+3 -3
View File
@@ -1113,11 +1113,11 @@ pub async fn cmdline(context: Context, line: &str, chat_id: &mut ChatId) -> Resu
let contact_id = ContactId::new(arg1.parse()?);
let contact = Contact::get_by_id(&context, contact_id).await?;
let name = contact.get_display_name();
let addr = contact.get_addr();
let name_n_addr = contact.get_name_n_addr();
let mut res = format!(
"Contact info for: {name} ({addr}):\nIcon: {}\n",
"Contact info for: {}:\nIcon: {}\n",
name_n_addr,
match contact.get_profile_image(&context).await? {
Some(image) => image.to_str().unwrap().to_string(),
None => "NoIcon".to_string(),
+1 -1
View File
@@ -4,7 +4,7 @@ build-backend = "setuptools.build_meta"
[project]
name = "deltachat-rpc-client"
version = "2.63.0-dev"
version = "2.61.0-dev"
license = "MPL-2.0"
description = "Python client for Delta Chat core JSON-RPC interface"
classifiers = [
@@ -141,13 +141,13 @@ class Account:
@futuremethod
def init_transports(self, qr: Optional[str] = None):
"""Add an initial transport on the chatmail relay that answers fastest.
"""Automatically adds up to three transports.
The profile then adds further ones in the background.
A ``DCACCOUNT:`` or ``DCLOGIN:`` ``qr`` code adds a single transport
while securejoin codes add the inviter's relays to the candidates.
Does nothing if the profile already has a transport.
If the user just scanned a QR code of type `Account`, `Login`,
`AskVerifyContact`, `AskVerifyGroup`, or `AskJoinBroadcast`,
then UI implementations should pass it as the `qr` parameter.
The host(s) from the QR code will then also be considered
for creating an account there.
"""
yield self._rpc.init_transports.future(self.id, qr)
@@ -38,6 +38,7 @@ class EventType(str, Enum):
IMAP_CONNECTED = "ImapConnected"
SMTP_MESSAGE_SENT = "SmtpMessageSent"
IMAP_MESSAGE_DELETED = "ImapMessageDeleted"
IMAP_MESSAGE_MOVED = "ImapMessageMoved"
IMAP_INBOX_IDLE = "ImapInboxIdle"
NEW_BLOB_FILE = "NewBlobFile"
DELETED_BLOB_FILE = "DeletedBlobFile"
@@ -69,7 +70,6 @@ class EventType(str, Enum):
SELFAVATAR_CHANGED = "SelfavatarChanged"
WEBXDC_STATUS_UPDATE = "WebxdcStatusUpdate"
WEBXDC_INSTANCE_DELETED = "WebxdcInstanceDeleted"
ACCOUNTS_BACKGROUND_FETCH_DONE = "AccountsBackgroundFetchDone"
CHATLIST_CHANGED = "ChatlistChanged"
CHATLIST_ITEM_CHANGED = "ChatlistItemChanged"
ACCOUNTS_CHANGED = "AccountsChanged"
@@ -48,13 +48,6 @@ class DeltaChat:
"""Stop ongoing background fetch."""
self.rpc.stop_background_fetch()
def wait_for_event(self, event_type=None) -> AttrDict:
"""Wait until the next account manager event and return it."""
while True:
next_event = AttrDict(self.rpc.wait_for_event(0))
if event_type is None or next_event.kind == event_type:
return next_event
def maybe_network(self) -> None:
"""Indicate that the network conditions might have changed."""
self.rpc.maybe_network()
@@ -92,9 +92,10 @@ class RPCAccountFactory:
"""Create a new configured account."""
account = self.get_unconfigured_account()
qr = self.get_account_qr()
yield account.add_transport_from_qr.future(qr)
yield account.init_transports.future(qr)
assert account.is_configured()
assert len(account.list_transports()) == 1
return account
def new_configured_bot(self) -> Bot:
+2 -10
View File
@@ -24,16 +24,8 @@ class DirectImap:
def __init__(self, account: Account, addr=None, password=None) -> None:
self.account = account
if addr is None or password is None:
transport = account.list_transports()[-1]
if addr is None:
self.addr = transport["addr"]
else:
self.addr = addr
if password is None:
self.password = transport["password"]
else:
self.password = password
self.addr = addr or account.get_config("addr")
self.password = password or account.get_config("mail_pw")
self.logid = account.get_config("displayname") or id(account)
self._idling = False
self.connect()
@@ -24,12 +24,6 @@ def wait_for_imap_message(imap):
time.sleep(1)
def test_init_transports(acf):
account = acf.get_unconfigured_account()
account.init_transports(acf.get_account_qr())
assert len(account.list_transports()) == 1
def test_add_second_address(acf) -> None:
account = acf.new_configured_account()
assert len(account.list_transports()) == 1
@@ -38,6 +32,12 @@ def test_add_second_address(acf) -> None:
account.add_transport_from_qr(qr)
assert len(account.list_transports()) == 2
# init_transports() only works on an unconfigured profile:
with pytest.raises(JsonRpcError):
account.init_transports(qr)
with pytest.raises(JsonRpcError):
account.init_transports()
account.add_transport_from_qr(qr)
assert len(account.list_transports()) == 3
@@ -58,9 +58,10 @@ def test_add_second_address(acf) -> None:
def test_change_address(acf) -> None:
"""Test Alice configuring a second transport and removing the first one."""
"""Test Alice configuring a second transport and setting it as a primary one."""
alice, bob = acf.get_online_accounts(2)
bob_addr = bob.get_config("configured_addr")
bob.create_chat(alice)
alice_chat_bob = alice.create_chat(bob)
@@ -70,14 +71,22 @@ def test_change_address(acf) -> None:
sender_addr1 = msg1.sender.get_snapshot().address
alice.stop_io()
old_alice_addr = alice.list_transports()[0]["addr"]
old_alice_addr = alice.get_config("configured_addr")
alice_vcard = alice.self_contact.make_vcard()
assert old_alice_addr in alice_vcard
qr = acf.get_account_qr()
alice.add_transport_from_qr(qr)
new_alice_addr = alice.list_transports()[1]["addr"]
with pytest.raises(JsonRpcError):
# Cannot use the address that is not
# configured for any transport.
alice.set_config("configured_addr", bob_addr)
alice.delete_transport(old_alice_addr)
# Load old address so it is cached.
assert alice.get_config("configured_addr") == old_alice_addr
alice.set_config("configured_addr", new_alice_addr)
# Make sure that setting `configured_addr` invalidated the cache.
assert alice.get_config("configured_addr") == new_alice_addr
alice_vcard = alice.self_contact.make_vcard()
assert old_alice_addr not in alice_vcard
@@ -225,6 +234,10 @@ def test_transport_sync_new_as_primary(acf, log) -> None:
log.section("ac1 changes the primary transport")
ac1.set_config("configured_addr", transport2["addr"])
ac1.wait_for_event(EventType.TRANSPORTS_MODIFIED)
ac1_clone.wait_for_event(EventType.TRANSPORTS_MODIFIED)
assert ac1_clone.get_config("configured_addr") == transport1["addr"]
log.section("ac1_clone receives a message via the new transport")
ac1_chat = ac1.create_chat(bob)
@@ -395,14 +408,3 @@ def test_background_fetch_no_duplicates(acf, direct_imap, dc):
dc.background_fetch(300)
assert len(messages_with_text(alice_chat, "hello")) == 1
def test_multitransport_mdn(acf):
"""Test sending an MDN right after configuring two transports."""
alice, bob = acf.get_online_accounts(2)
alice.add_transport_from_qr(acf.get_account_qr())
alice.bring_online()
alice.create_chat(bob)
bob_msg = bob.create_chat(alice).send_text("Hello!")
alice.wait_for_incoming_msg().mark_seen()
assert bob.wait_for_event(EventType.MSG_READ).msg_id == bob_msg.id
+17 -15
View File
@@ -27,7 +27,7 @@ def test_qr_setup_contact(acf) -> None:
def test_qr_setup_contact_svg(acf) -> None:
alice = acf.new_configured_account()
_, _, domain = alice.list_transports()[0]["addr"].rpartition("@")
_, _, domain = alice.get_config("addr").rpartition("@")
_qr_code, svg = alice.get_qr_code_svg()
@@ -43,7 +43,6 @@ def test_qr_setup_contact_svg(acf) -> None:
def test_qr_securejoin(acf):
alice, bob, fiona = acf.get_online_accounts(3)
alice.set_config("displayname", "Alice")
# Setup second device for Alice
# to test observing securejoin protocol.
alice2 = alice.clone()
@@ -68,7 +67,7 @@ def test_qr_securejoin(acf):
assert alice_contact_bob_snapshot.e2ee_avail
snapshot = bob.wait_for_incoming_msg().get_snapshot()
assert snapshot.text == "You were added by Alice."
assert snapshot.text == "You were added by {}.".format(alice.get_config("addr"))
bob_contact_alice = bob.create_contact(alice)
bob_contact_alice_snapshot = bob_contact_alice.get_snapshot()
@@ -147,13 +146,17 @@ def test_qr_securejoin_broadcast(acf, all_devices_online):
assert "invited you to join this channel" in first_msg.text
assert first_msg.is_info
if not inviter_side:
if inviter_side:
member_added_msg = chat_msgs.pop(0).get_snapshot()
assert member_added_msg.text == f"Member {contact_snapshot.display_name} added."
assert member_added_msg.info_contact_id == contact_snapshot.id
else:
if chat_msgs[0].get_snapshot().text == "You joined the channel.":
member_added_msg = chat_msgs.pop(0).get_snapshot()
else:
member_added_msg = chat_msgs.pop(1).get_snapshot()
assert member_added_msg.text == "You joined the channel."
assert member_added_msg.is_info
assert member_added_msg.is_info
hello_msg = chat_msgs.pop(0).get_snapshot()
assert hello_msg.text == "Hello everyone!"
@@ -218,7 +221,7 @@ def test_qr_securejoin_broadcast(acf, all_devices_online):
snapshot = fiona.wait_for_incoming_msg().get_snapshot()
assert snapshot.text == "You joined the channel."
get_broadcast(alice2).get_messages()[1].resend()
get_broadcast(alice2).get_messages()[2].resend()
snapshot = fiona.wait_for_incoming_msg().get_snapshot()
assert snapshot.text == "Hello everyone!"
@@ -500,9 +503,11 @@ def test_aeap_flow(acf):
assert msg_in_1.text == msg_out.text
logging.info("changing email account")
old_addr = ac1.list_transports()[0]["addr"]
ac1.add_transport_from_qr(acf.get_account_qr())
ac1.delete_transport(old_addr)
ac1.set_config("addr", addr)
ac1.set_config("mail_pw", password)
ac1.stop_io()
ac1.configure()
ac1.start_io()
logging.info("sending second message")
msg_out = chat.send_text("changed address").get_snapshot()
@@ -524,7 +529,6 @@ def test_securejoin_after_contact_resetup(acf) -> None:
but different key fingerprint while a securejoin with that contact is still pending.
"""
ac1, ac2, ac3 = acf.get_online_accounts(3)
ac3.set_config("displayname", "ac3")
# ac3 creates a group with ac1.
ac3_chat = ac3.create_group("Group")
@@ -536,7 +540,7 @@ def test_securejoin_after_contact_resetup(acf) -> None:
# ac1 waits for member added message and creates a QR code.
snapshot = ac1.wait_for_incoming_msg().get_snapshot()
assert snapshot.text == "You were added by ac3."
assert snapshot.text == "You were added by {}.".format(ac3.get_config("addr"))
ac1_qr_code = snapshot.chat.get_qr_code()
# ac2 sets up contact with ac1
@@ -576,8 +580,6 @@ def test_securejoin_after_contact_resetup(acf) -> None:
def test_withdraw_securejoin_qr(acf):
alice, bob = acf.get_online_accounts(2)
alice.set_config("displayname", "Alice")
bob.set_config("displayname", "Bob")
logging.info("Alice creates a group")
alice_chat = alice.create_group("Group")
@@ -590,11 +592,11 @@ def test_withdraw_securejoin_qr(acf):
alice.clear_all_events()
snapshot = bob.wait_for_incoming_msg().get_snapshot()
assert snapshot.text == "You were added by Alice."
assert snapshot.text == "You were added by {}.".format(alice.get_config("addr"))
bob_chat.leave()
snapshot = alice.get_message_by_id(alice.wait_for_msgs_changed_event().msg_id).get_snapshot()
assert snapshot.text == "Group left by Bob."
assert snapshot.text == "Group left by {}.".format(bob.get_config("addr"))
logging.info("Alice withdraws QR code.")
qr = alice.check_qr(qr_code)
+12 -21
View File
@@ -155,7 +155,7 @@ def test_list_transports(acf) -> None:
def test_account(acf) -> None:
alice, bob = acf.get_online_accounts(2)
bob_addr = bob.get_config("configured_addr")
bob_addr = bob.get_config("addr")
alice_contact_bob = alice.create_contact(bob, "Bob")
alice_chat_bob = alice_contact_bob.create_chat()
alice_chat_bob.send_text("Hello!")
@@ -318,7 +318,7 @@ def test_chat(acf) -> None:
def test_contact(acf) -> None:
alice, bob = acf.get_online_accounts(2)
bob_addr = bob.get_config("configured_addr")
bob_addr = bob.get_config("addr")
alice_contact_bob = alice.create_contact(bob, "Bob")
assert alice_contact_bob == alice.get_contact_by_id(alice_contact_bob.id)
@@ -421,6 +421,7 @@ def test_dont_move_sync_msgs(acf, direct_imap):
addr, password = acf.get_credentials()
ac1 = acf.get_unconfigured_account()
ac1.set_config("bcc_self", "1")
ac1.set_config("fix_is_chatmail", "1")
ac1.add_or_update_transport({"addr": addr, "password": password})
ac1.start_io()
ac1_direct_imap = direct_imap(ac1)
@@ -604,6 +605,7 @@ def test_import_export_online_all(acf, tmp_path, rpcdata, log) -> None:
(ac1, some1) = acf.get_online_accounts(2)
log.section("create some chat content")
some1_addr = some1.get_config("addr")
chat1 = ac1.create_contact(some1).create_chat()
chat1.send_text("msg1")
assert len(ac1.get_contacts()) == 1
@@ -621,6 +623,7 @@ def test_import_export_online_all(acf, tmp_path, rpcdata, log) -> None:
contacts = ac.get_contacts()
assert len(contacts) == 1
contact2 = contacts[0]
assert contact2.get_snapshot().address == some1_addr
chat2 = contact2.create_chat()
messages = chat2.get_messages()
assert len(messages) == 3 + E2EE_INFO_MSGS
@@ -1185,6 +1188,7 @@ def test_leave_broadcast(acf, all_devices_online):
def check_account(ac, contact, inviter_side, please_wait_info_msg=False):
chat = get_broadcast(ac)
contact_snapshot = contact.get_snapshot()
chat_msgs = chat.get_messages()
encrypted_msg = chat_msgs.pop(0).get_snapshot()
@@ -1196,11 +1200,14 @@ def test_leave_broadcast(acf, all_devices_online):
assert "invited you to join this channel" in first_msg.text
assert first_msg.is_info
if not inviter_side:
member_added_msg = chat_msgs.pop(0).get_snapshot()
member_added_msg = chat_msgs.pop(0).get_snapshot()
if inviter_side:
assert member_added_msg.text == f"Member {contact_snapshot.display_name} added."
else:
assert member_added_msg.text == "You joined the channel."
assert member_added_msg.is_info
assert member_added_msg.is_info
if not inviter_side:
leave_msg = chat_msgs.pop(0).get_snapshot()
assert leave_msg.text == "You left the channel."
@@ -1347,22 +1354,6 @@ def test_background_fetch(acf, dc):
break
def test_background_fetch_does_not_wait_for_sending(dc, acf):
alice, bob = acf.get_online_accounts(2)
alice_chat_bob = alice.create_chat(bob)
alice.stop_io()
text = "x" * 200_000
for _ in range(50):
alice_chat_bob.send_text(text)
assert not dc.is_sending_finished()
alice.start_io()
dc.background_fetch(50)
dc.wait_for_event(EventType.ACCOUNTS_BACKGROUND_FETCH_DONE)
assert not dc.is_sending_finished()
def test_message_exists(acf):
ac1, ac2 = acf.get_online_accounts(2)
chat = ac1.create_chat(ac2)
+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "deltachat-rpc-server"
version = "2.63.0-dev"
version = "2.61.0-dev"
description = "DeltaChat JSON-RPC server"
edition = "2024"
license = "MPL-2.0"
@@ -15,5 +15,5 @@
},
"type": "module",
"types": "index.d.ts",
"version": "2.63.0-dev"
"version": "2.61.0-dev"
}
-3
View File
@@ -69,11 +69,8 @@ skip = [
{ name = "derive_more-impl", version = "1.0.0" },
{ name = "derive_more", version = "1.0.0" },
{ name = "event-listener", version = "2.5.3" },
{ name = "foldhash", version = "0.1.5" },
{ name = "getrandom", version = "0.2.12" },
{ name = "getrandom", version = "0.3.3" },
{ name = "hashbrown", version = "0.15.4" },
{ name = "hashbrown", version = "0.16.1" },
{ name = "heck", version = "0.4.1" },
{ name = "http", version = "0.2.12" },
{ name = "hybrid-array", version = "0.2.3" },
+11 -9
View File
@@ -630,12 +630,18 @@ CREATE TABLE broadcast_secrets(
-- Candidate chatmail relays for automatic relay management.
-- Holds the hosts a QR code contributed and the default relays already tried;
-- the default list itself lives in `autorelay.rs` and is not stored here.
CREATE TABLE relay_candidates(
host TEXT PRIMARY KEY NOT NULL,
-- Timestamp of the last connection attempt, 0 if the host was never tried.
last_tried INTEGER NOT NULL DEFAULT 0
last_tried INTEGER NOT NULL DEFAULT 0 -- Deprecated 2026-09, replaced with separate relay_candidates_last_tried table.
) STRICT;
-- This table is used for storing the timestamp of the last
-- connection attempt per chatmail relay candidate.
-- This table can contain relays that were removed from the list of candidates,
-- and it does not contain the default relays.
CREATE TABLE relay_candidates_last_tried(
host TEXT PRIMARY KEY NOT NULL,
last_tried INTEGER NOT NULL DEFAULT 0 -- Timestamp of the last connection attempt.
) STRICT;
CREATE TABLE transports (
@@ -683,11 +689,7 @@ CREATE TABLE imap (
transport_id INTEGER NOT NULL, -- ID of the transport in the `transports` table.
rfc724_mid TEXT NOT NULL, -- Message-ID header
folder TEXT NOT NULL, -- IMAP folder
-- Destination folder. Empty string means that the message shall be deleted.
-- Since we don't move messages between IMAP folders anymore,
-- this is always either empty or equal to `folder`.
target TEXT NOT NULL,
target TEXT NOT NULL, -- Destination folder. Empty string means that the message shall be deleted.
uid INTEGER NOT NULL, -- UID
uidvalidity INTEGER NOT NULL,
UNIQUE (transport_id, folder, uid, uidvalidity)
+1 -1
View File
@@ -4,7 +4,7 @@ build-backend = "setuptools.build_meta"
[project]
name = "deltachat"
version = "2.63.0-dev"
version = "2.61.0-dev"
license = "MPL-2.0"
description = "Python bindings for the Delta Chat Core library using CFFI against the Rust-implemented libdeltachat"
readme = "README.rst"
-5
View File
@@ -11,9 +11,6 @@ import random
from queue import Queue
from typing import Callable, Dict, List, Optional
from .capi import lib
from .cutil import as_dc_charpointer
import pytest
from _pytest._code import Source
@@ -369,7 +366,6 @@ class ACFactory:
ac.open(passphrase)
acname = ac._logid
addr = f"{acname}@offline.org"
lib.dc_add_pseudo_transport(ac._dc_context, as_dc_charpointer(addr))
ac.update_config(
{
"configured_addr": addr,
@@ -378,7 +374,6 @@ class ACFactory:
)
self._preconfigure_key(ac)
self._acsetup.init_logging(ac)
assert ac.is_configured(), "Pseudo configured account should look like if it is configured"
return ac
def new_online_configuring_account(self, cloned_from=None, **kwargs) -> Account:
+1 -1
View File
@@ -1 +1 @@
2026-09-22
2026-09-11
+53 -7
View File
@@ -169,7 +169,9 @@ impl Accounts {
.with_push_subscriber(self.push_subscriber.clone())
.build()
.await?;
ctx.open().await?;
// Try to open without a passphrase,
// but do not return an error if account is passphare-protected.
ctx.open("".to_string()).await?;
self.accounts.insert(account_config.id, ctx);
self.emit_event(EventType::AccountsChanged);
@@ -480,15 +482,11 @@ impl Accounts {
/// return immediately even before the timeout expiration
/// or finishing fetching.
///
/// Pending outgoing messages are not waited for and not triggered.
///
/// The `AccountsBackgroundFetchDone` event is emitted at the end,
/// process all events until you get this one and you can safely return to the background
/// without forgetting to create notifications caused by timing race conditions.
/// If another background fetch is already running,
/// nothing is fetched and the event is emitted immediately.
/// The event carries no data identifying the call it belongs to,
/// so it only safely refers to your call if no concurrent background fetch is happening.
///
/// Returns a future that resolves when background fetch is done,
/// but does not capture `&self`.
@@ -526,7 +524,7 @@ impl Accounts {
pub async fn is_sending_finished(&self) -> Result<bool> {
let accounts: Vec<Context> = self.accounts.values().cloned().collect();
for account in accounts {
if !smtp::queue::is_empty(&account).await? {
if !smtp::is_queue_empty(&account).await? {
return Ok(false);
}
}
@@ -819,7 +817,9 @@ impl Config {
.build()
.await
.with_context(|| format!("failed to create context from file {dbfile:?}"))?;
ctx.open().await?;
// Try to open without a passphrase,
// but do not return an error if account is passphare-protected.
ctx.open("".to_string()).await?;
accounts.insert(account_config.id, ctx);
}
@@ -1268,6 +1268,52 @@ mod tests {
Ok(())
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn test_encrypted_account() -> Result<()> {
let dir = tempfile::tempdir().context("failed to create tempdir")?;
let p: PathBuf = dir.path().join("accounts");
let writable = true;
let mut accounts = Accounts::new(p.clone(), writable)
.await
.context("failed to create accounts manager")?;
assert_eq!(accounts.accounts.len(), 0);
let account_id = accounts
.add_closed_account()
.await
.context("failed to add closed account")?;
let account = accounts
.get_selected_account()
.context("failed to get account")?;
assert_eq!(account.id, account_id);
let passphrase_set_success = account
.open("foobar".to_string())
.await
.context("failed to set passphrase")?;
assert!(passphrase_set_success);
drop(accounts);
let writable = false;
let accounts = Accounts::new(p.clone(), writable)
.await
.context("failed to create second accounts manager")?;
let account = accounts
.get_selected_account()
.context("failed to get account")?;
assert_eq!(account.is_open().await, false);
// Try wrong passphrase.
assert_eq!(account.open("barfoo".to_string()).await?, false);
assert_eq!(account.open("".to_string()).await?, false);
assert_eq!(account.open("foobar".to_string()).await?, true);
assert_eq!(account.is_open().await, true);
Ok(())
}
/// Tests that accounts share stock string translations.
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn test_accounts_share_translations() -> Result<()> {
+107 -127
View File
@@ -1,24 +1,31 @@
//! # Automatic multi-relay onboarding
//! # Automatic relay handling (experimental, still in development)
//!
//! Support for automatically onboarding a profile on transport
//! candidates without the user choosing a relay.
//! Chatmail relays create an account on first login,
//! so a profile can add further transports on its own without user interaction.
//! Candidate hosts come from the `relay_candidates` table
//! as well as the [`DEFAULT_RELAY_CANDIDATES`] list.
//!
//! Status of implementation:
//! Additions are attempted right before going into IMAP IDLE,
//! i.e. only while connected and with nothing more important to do,
//! and only if a UI opted in via [`Config::Autorelay`].
//! Once a profile has reached `NUM_TRANSPORTS_TARGET` transports,
//! [`Config::AutorelayFinished`] is set and nothing is ever added again,
//! so deleting a transport later does not pull in a replacement.
use std::collections::BTreeMap;
use std::collections::BTreeSet;
use std::pin::Pin;
use anyhow::{Result, format_err};
use anyhow::Result;
use deltachat_contact_tools::addr_normalize;
use rand::distr::{Alphanumeric, SampleString};
use rand::seq::{IndexedRandom, SliceRandom};
use rusqlite::Transaction;
use tokio::task::JoinSet;
use rand::seq::IndexedRandom;
use crate::config::{self, Config};
use crate::configure::{EnteredLoginParam, SILENT_PROGRESS, configure};
use crate::log::{LogExt, warn};
use crate::login_param::{EnteredCertificateChecks, EnteredImapLoginParam};
use crate::net::{connect_tcp, proxy::ProxyConfig};
use crate::{context::Context, tools::time};
use crate::sql::TransactionExt as _;
use crate::{configure::EnteredLoginParam, context::Context, tools::time};
/// The target number of transports.
const NUM_TRANSPORTS_TARGET: usize = 3;
@@ -27,9 +34,12 @@ const AUTOMATIC_ADDITION_DEBOUNCE_SECONDS: i64 = 60 * 60; // one hour
/// How long we ignore a relay candidate after failing to connect to it:
const BACKOFF_PERIOD_FOR_NOT_WORKING_RELAY: i64 = 60 * 60 * 24 * 7; // one week
/// Sorted relay list a profile can attempt to onboard on without the user choosing one.
/// The list of relays to which we onboard
/// if no other relays are provided via QR codes.
/// Please keep this list alphabetically sorted.
const DEFAULT_RELAY_CANDIDATES: &[&str] = &[
"chat.adminforge.de", // iroh relay 404s
"chat.adminforge.de",
"chat.feld.me",
"chat.me.ke",
"chat.nuvon.app",
"chat.tinydispatch.org",
@@ -38,7 +48,7 @@ const DEFAULT_RELAY_CANDIDATES: &[&str] = &[
"chtml.ca",
"deltachat.me",
"e2e.sus.fr",
"e2ee.wang",
"jp.deltachat.me",
"mailchat.pl",
"nchrcht.la10cy.net",
"nine.testrun.org",
@@ -46,66 +56,35 @@ const DEFAULT_RELAY_CANDIDATES: &[&str] = &[
"tarpit.fun",
];
/// Records the hosts of `addrs` as relay candidates.
pub(crate) async fn add_relay_candidates(context: &Context, addrs: &[String]) -> Result<()> {
context
.sql
.transaction(|tx| hosts_of(addrs).try_for_each(|host| save_relay_candidate(tx, host, 0)))
.await
}
/// Adds a first transport on the relay candidate that answers fastest.
///
/// All candidates are probed at once with a TCP connection to their HTTPS port
/// and configured in the order in which the connections complete,
/// stopping at the first success. Candidates that fail the probe are skipped.
/// Answering TCP fastest is used as a network proximity measure,
/// which keeps latency low for initial onboarding,
/// and it avoids relays that are down or black-holing traffic.
pub(crate) async fn add_transport_from_candidates(
pub(crate) async fn init_transports_inner(
context: &Context,
addrs_from_qr: Vec<String>,
skip_network: bool,
) -> Result<()> {
let mut candidates = triable_relay_candidates(context, time()).await?;
candidates.shuffle(&mut rand::rng());
let mut probes = JoinSet::new();
let proxy_config = ProxyConfig::load(context).await?;
let load_cache = false;
for host in candidates {
let ctx = context.clone();
let proxy_config = proxy_config.clone();
probes.spawn(async move {
let res = match proxy_config {
_ if skip_network => Ok(()),
Some(proxy) => proxy.connect(&ctx, &host, 443, load_cache).await.map(drop),
None => connect_tcp(&ctx, &host, 443, load_cache).await.map(drop),
};
(host, res)
});
}
context
.sql
.transaction(|transaction| {
let mut stmt = transaction.prepare("INSERT INTO relay_candidates(host) VALUES(?)")?;
for addr in addrs_from_qr {
stmt.execute((addr,))?;
}
Ok(())
})
.await?;
let mut last_err = format_err!("No relay candidates");
let mark_as_autorelay = true;
while let Some(res) = probes.join_next().await {
let (host, res) = res?;
if let Err(err) = res {
warn!(context, "Failed to connect to relay {host}: {err:#}.");
last_err = err;
continue;
}
let param = login_param_from_host(&host, mark_as_autorelay);
match configure(context, &param, skip_network).await {
Ok(()) => {
info!(context, "Added a transport on relay {host}.");
return Ok(());
}
Err(err) => {
warn!(context, "Failed to add relay {host}: {err:#}.");
last_err = err;
}
}
let host = "nine.testrun.org";
let param = login_param_from_host(host);
let res = crate::configure::configure(context, &param, skip_network).await;
if let Err(err) = &res {
warn!(context, "Failed to init transports: {err:#}.");
} else {
info!(context, "Initialized with transport {host}");
}
Err(last_err)
res?;
context.set_config_bool(Config::Autorelay, true).await?;
Ok(())
}
pub(crate) fn maybe_add_additional_relays(
@@ -161,15 +140,22 @@ async fn maybe_add_additional_relays_inner(context: &Context, skip_network: bool
let mut relay_added = false;
// Using `for` instead of `while` to prevent infinite loop
for _ in 0..NUM_TRANSPORTS_TARGET {
if context.count_transports().await? >= NUM_TRANSPORTS_TARGET {
let num_transports = context.count_transports().await?;
if num_transports >= NUM_TRANSPORTS_TARGET {
info!(context, "Transports target reached at {num_transports}");
context
.set_config_internal(Config::AutorelayFinished, config::from_bool(true))
.await?;
return Ok(relay_added);
} else {
info!(context, "There are {num_transports} relays, will add more");
}
let candidates = triable_relay_candidates(context, now).await?;
// First, query all candidates that were not tried since `BACKOFF_PERIOD_FOR_NOT_WORKING_RELAY` seconds.
// Hosts that are already used are excluded.
let candidates = load_relay_candidates(context, now).await?;
let Some(host) = candidates.choose(&mut rand::rng()) else {
info!(
context,
@@ -184,15 +170,9 @@ async fn maybe_add_additional_relays_inner(context: &Context, skip_network: bool
candidates.len(),
);
context
.sql
.transaction(|tx| save_relay_candidate(tx, host, now))
.await?;
let mark_as_autorelay = true;
let param = login_param_from_host(host, mark_as_autorelay);
let res = SILENT_PROGRESS
.scope((), configure(context, &param, skip_network))
.await;
set_relay_candidate_last_tried(context, host, now).await?;
let param = login_param_from_host(host);
let res = crate::configure::configure(context, &param, skip_network).await;
if let Err(e) = res {
warn!(
context,
@@ -207,63 +187,63 @@ async fn maybe_add_additional_relays_inner(context: &Context, skip_network: bool
Ok(relay_added)
}
async fn triable_relay_candidates(context: &Context, now: i64) -> Result<Vec<String>> {
let cutoff_timestamp = now.saturating_sub(BACKOFF_PERIOD_FOR_NOT_WORKING_RELAY);
let mut last_tried: BTreeMap<String, i64> = context
async fn set_relay_candidate_last_tried(
context: &Context,
host: &str,
now: i64,
) -> Result<(), anyhow::Error> {
context
.sql
.query_map_collect("SELECT host, last_tried FROM relay_candidates", (), |row| {
Ok((row.get(0)?, row.get(1)?))
})
.execute(
"INSERT OR REPLACE INTO relay_candidates_last_tried(host, last_tried) VALUES(?, ?)",
(host, now),
)
.await?;
for host in DEFAULT_RELAY_CANDIDATES {
last_tried.entry(host.to_string()).or_insert(0);
}
let self_addrs = context.get_self_addrs().await?;
let used_hosts: Vec<&str> = hosts_of(&self_addrs).collect();
// We also try candidates which have `last_tried` in the future,
// which on next failure get `last_tried` reset to the current time.
let candidates = last_tried
.into_iter()
.filter(|(host, last_tried)| {
(*last_tried < cutoff_timestamp || *last_tried > now)
&& !used_hosts.contains(&host.as_str())
})
.map(|(host, _)| host)
.collect();
Ok(candidates)
}
/// Returns the host of each address in `addrs`.
fn hosts_of(addrs: &[String]) -> impl Iterator<Item = &str> {
addrs.iter().filter_map(|a| Some(a.rsplit_once('@')?.1))
}
/// Records `host` as a relay candidate, overwriting a stored `last_tried`.
fn save_relay_candidate(tx: &Transaction, host: &str, last_tried: i64) -> Result<()> {
tx.execute(
"INSERT INTO relay_candidates (host, last_tried) VALUES (?, ?)
ON CONFLICT(host) DO UPDATE SET last_tried=excluded.last_tried",
(host, last_tried),
)?;
Ok(())
}
pub(crate) fn login_param_from_host(host: &str, mark_as_autorelay: bool) -> EnteredLoginParam {
async fn load_relay_candidates(context: &Context, now: i64) -> Result<Vec<String>> {
let res = context
.sql
.transaction(|transaction| {
let mut candidates: BTreeSet<String> =
transaction.query_map_collect("SELECT host FROM relay_candidates", (), |row| {
Ok(row.get(0)?)
})?;
candidates.extend(DEFAULT_RELAY_CANDIDATES.iter().map(|s| s.to_string()));
let cutoff_timestamp = now.saturating_sub(BACKOFF_PERIOD_FOR_NOT_WORKING_RELAY);
// This does not select candidates which have last_tried in the future,
// essentially treating them as never tried,
// so if some timestamp far in the future is accidentally stored,
// we are not stuck never trying the candidate.
// After trying the candidate, last_tried will be corrected to the current time.
let exclude: BTreeSet<String> = transaction.query_map_collect(
"SELECT host FROM relay_candidates_last_tried WHERE (last_tried>=? AND last_tried<=?)
UNION
SELECT substr(addr, instr(addr, '@') + 1) FROM transports",
(cutoff_timestamp, now),
|row| Ok(row.get(0)?),
)?;
Ok(candidates
.difference(&exclude)
.map(|s| s.to_string())
.collect::<Vec<String>>())
})
.await?;
Ok(res)
}
pub(crate) fn login_param_from_host(host: &str) -> EnteredLoginParam {
let rng = &mut rand::rng();
let username = Alphanumeric.sample_string(rng, 9);
let addr = username + "@" + host;
let addr = addr_normalize(&addr);
// `mark_as_autorelay` is a temporary precaution hack
// while introducing onboarding on multiple community relays from a list:
// though relay operators were asked to get on that list, unexpected things can happen,
// and they want to return to allow only manual onboarding.
// this is possible by failing on `password_len == 23`.
// 22 * log2(26 * 2 + 10) = 130 bits of entropy
let password = Alphanumeric.sample_string(rng, if mark_as_autorelay { 23 } else { 22 });
let password = Alphanumeric.sample_string(rng, 22);
EnteredLoginParam {
addr,
+106 -149
View File
@@ -1,176 +1,107 @@
use std::time::Duration;
use super::*;
use crate::EventType;
use crate::test_utils::{TestContext, TestContextManager};
use crate::test_utils::TestContext;
use crate::tools::SystemTime;
/// Tests that the default relays are candidates without a row in the table.
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn test_triable_relay_candidates_defaults() -> Result<()> {
let mut tcm = TestContextManager::new();
let t = &tcm.unconfigured().await;
let now = time();
assert!(DEFAULT_RELAY_CANDIDATES.is_sorted());
let mut candidates = triable_relay_candidates(t, now).await?;
candidates.sort();
assert_eq!(candidates, DEFAULT_RELAY_CANDIDATES);
let tried = DEFAULT_RELAY_CANDIDATES[0];
save_relay_candidates(t, &[tried], now).await?;
let candidates = triable_relay_candidates(t, now).await?;
assert_eq!(candidates.len(), DEFAULT_RELAY_CANDIDATES.len() - 1);
assert!(!candidates.contains(&tried.to_string()));
Ok(())
}
/// Tests that a transport is added on a candidate from the given addresses.
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn test_add_transport_from_candidates() -> Result<()> {
let mut tcm = TestContextManager::new();
let t = &tcm.unconfigured().await;
mark_defaults_tried(t, time()).await?;
let addrs_from_qr = [
"alice@example.org".to_string(),
"bob@example.org".to_string(),
];
let skip_network = true;
add_relay_candidates(t, &addrs_from_qr).await?;
add_transport_from_candidates(t, skip_network).await?;
let transports = t.list_transports().await?;
assert_eq!(transports.len(), 1);
assert!(transports[0].addr.ends_with("@example.org"));
let untried = untried_relay_candidates(t).await?;
assert_eq!(untried, ["example.org"]);
assert!(configure_progress_emitted(t).await);
Ok(())
}
/// Tests correct add_transport_from_candidates error handling.
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn test_add_transport_from_candidates_failure() -> Result<()> {
let mut tcm = TestContextManager::new();
let t = &tcm.unconfigured().await;
mark_defaults_tried(t, time()).await?;
save_relay_candidates(t, &["bad host", "worse host"], 0).await?;
async fn test_init_transports_basic() -> Result<()> {
let t = &TestContext::new().await;
assert!(t.list_transports().await?.is_empty());
let skip_network = true;
let err = add_transport_from_candidates(t, skip_network)
.await
.unwrap_err();
assert!(format!("{err:#}").contains("Bad email-address"));
assert!(!t.is_configured().await?);
let untried = untried_relay_candidates(t).await?;
assert_eq!(untried, ["bad host", "worse host"]);
t.assert_warns_or_errors(&[
"Failed to add relay bad host",
"Failed to add relay worse host",
])
.await;
init_transports_inner(t, vec![], skip_network).await?;
Ok(())
}
async fn untried_relay_candidates(t: &TestContext) -> Result<Vec<String>> {
t.sql
.query_map_vec(
"SELECT host FROM relay_candidates WHERE last_tried=0 ORDER BY host",
(),
|row| Ok(row.get(0)?),
)
.await
}
async fn save_relay_candidates(t: &TestContext, hosts: &[&str], last_tried: i64) -> Result<()> {
t.sql
.transaction(|tx| {
for host in hosts {
save_relay_candidate(tx, host, last_tried)?;
}
Ok(())
})
.await
}
/// Keeps the default relays out of `triable_relay_candidates()`.
async fn mark_defaults_tried(t: &TestContext, now: i64) -> Result<()> {
save_relay_candidates(t, DEFAULT_RELAY_CANDIDATES, now).await
}
/// Consumes emitted events, telling whether a configure progress is among them.
async fn configure_progress_emitted(t: &TestContext) -> bool {
t.evtracker
.get_matching_opt(t, |evt| matches!(evt, EventType::ConfigureProgress { .. }))
.await
.is_some()
}
/// Tests that saving a candidate overwrites its stored timestamp.
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn test_save_relay_candidate() -> Result<()> {
let mut tcm = TestContextManager::new();
let t = &tcm.unconfigured().await;
let now = time();
for last_tried in [0, now, 0] {
t.sql
.transaction(|tx| save_relay_candidate(tx, "relay.example", last_tried))
.await?;
let stored: Option<i64> = t
.sql
.query_get_value(
"SELECT last_tried FROM relay_candidates WHERE host=?",
("relay.example",),
)
.await?;
assert_eq!(stored, Some(last_tried));
let relays = get_configured_relays(t).await;
assert_eq!(relays.len(), 1);
for relay in &relays {
assert!(DEFAULT_RELAY_CANDIDATES.contains(&relay.as_ref()));
}
assert_eq!(t.get_config_bool(Config::Autorelay).await?, true);
Ok(())
}
async fn get_configured_relays(t: &TestContext) -> Vec<String> {
let transports = t.list_transports().await.unwrap();
let mut relays: Vec<_> = transports
.iter()
.map(|t| t.addr.split_once('@').unwrap().1)
.collect();
// Check that every relay is used only once:
relays.sort();
relays.dedup();
assert_eq!(relays.len(), transports.len());
relays.into_iter().map(|s| s.to_string()).collect()
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn test_triable_relay_candidates_single() -> Result<()> {
async fn test_load_relay_candidates_single() -> Result<()> {
let t = &TestContext::new_alice().await;
enable_config(t).await;
let now = time();
mark_defaults_tried(t, now).await?;
// This host should be returned by load_relay_candidates():
t.sql
.execute(
"INSERT INTO relay_candidates (host) VALUES (?)",
("never_tried.example",),
)
.await?;
save_relay_candidates(t, &["never_tried.example", "example.org"], 0).await?;
save_relay_candidates(t, &["recent.example"], now).await?;
// This host was recently tried and should not be returned:
t.sql
.execute(
"INSERT INTO relay_candidates (host) VALUES (?)",
("recent.example",),
)
.await?;
set_relay_candidate_last_tried(t, "recent.example", now).await?;
let candidates = triable_relay_candidates(t, now).await?;
// This host is already in use (alice@example.org) and should not be returned:
t.sql
.execute(
"INSERT INTO relay_candidates (host) VALUES (?)",
("example.org",),
)
.await?;
assert_eq!(candidates, vec!["never_tried.example".to_string()]);
let candidates = load_relay_candidates(t, now).await?;
assert!(candidates.contains(&"never_tried.example".to_string()));
assert_eq!(candidates.contains(&"recent.example".to_string()), false);
assert_eq!(candidates.contains(&"example.org".to_string()), false);
assert_eq!(candidates.len(), DEFAULT_RELAY_CANDIDATES.len() + 1);
Ok(())
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn test_triable_relay_candidates_multiple() -> Result<()> {
async fn test_load_relay_candidates_multiple() -> Result<()> {
let t = &TestContext::new().await;
enable_config(t).await;
let now = time();
mark_defaults_tried(t, now).await?;
save_relay_candidates(t, &["a.example", "b.example", "c.example"], 0).await?;
const EXAMPLE_CANDIDATES: &[&str] = &["a.example", "b.example", "c.example"];
let mut candidates = triable_relay_candidates(t, now).await?;
for host in EXAMPLE_CANDIDATES {
t.sql
.execute("INSERT INTO relay_candidates (host) VALUES (?)", (host,))
.await?;
}
let mut candidates = load_relay_candidates(t, now).await?;
candidates.sort();
assert_eq!(
candidates,
vec![
"a.example".to_string(),
"b.example".to_string(),
"c.example".to_string()
]
);
let mut expected = EXAMPLE_CANDIDATES.to_vec();
expected.extend(DEFAULT_RELAY_CANDIDATES);
expected.sort();
assert_eq!(candidates, expected);
Ok(())
}
@@ -261,8 +192,23 @@ async fn test_maybe_add_additional_relays_add_one() -> Result<()> {
enable_config(t).await;
let now = time();
mark_defaults_tried(t, now).await?;
save_relay_candidates(t, &["relay.example"], 0).await?;
// Make sure that default relay candidates
// are not used by setting last_used to now:
for candidate in DEFAULT_RELAY_CANDIDATES {
t.sql
.execute(
"INSERT INTO relay_candidates_last_tried(host, last_tried) VALUES(?,?)",
(candidate, now),
)
.await?;
}
t.sql
.execute(
"INSERT INTO relay_candidates (host) VALUES (?)",
("relay.example",),
)
.await?;
let transports_before = t.count_transports().await?;
@@ -275,7 +221,6 @@ async fn test_maybe_add_additional_relays_add_one() -> Result<()> {
let transports_after = t.count_transports().await?;
assert_eq!(transports_after, transports_before + 1);
assert!(!configure_progress_emitted(t).await);
Ok(())
}
@@ -286,9 +231,6 @@ async fn test_maybe_add_additional_relays_add_multiple() -> Result<()> {
enable_config(t).await;
let now = time();
mark_defaults_tried(t, now).await?;
save_relay_candidates(t, &["a.example", "b.example", "c.example", "d.example"], 0).await?;
let skip_network = true;
let relay_added = maybe_add_additional_relays_inner(t, skip_network).await?;
assert!(relay_added);
@@ -308,9 +250,24 @@ async fn test_maybe_add_additional_relays_failure() -> Result<()> {
enable_config(t).await;
let now = time();
mark_defaults_tried(t, now).await?;
// Make sure that default relay candidates
// are not used by setting last_used to now:
for candidate in DEFAULT_RELAY_CANDIDATES {
t.sql
.execute(
"INSERT INTO relay_candidates_last_tried(host, last_tried) VALUES(?,?)",
(candidate, now - 2),
)
.await?;
}
for i in 1..10 {
save_relay_candidates(t, &[format!("{i}.invalid.example").as_str()], 0).await?;
t.sql
.execute(
"INSERT INTO relay_candidates (host) VALUES (?)",
(format!("{i}.invalid.example"),),
)
.await?;
}
let transports_before = t.count_transports().await?;
@@ -331,7 +288,7 @@ async fn test_maybe_add_additional_relays_failure() -> Result<()> {
assert!(
t.sql
.exists(
"SELECT COUNT(*) FROM relay_candidates WHERE last_tried>=?",
"SELECT COUNT(*) FROM relay_candidates_last_tried WHERE last_tried>=?",
(now,)
)
.await?
@@ -339,7 +296,7 @@ async fn test_maybe_add_additional_relays_failure() -> Result<()> {
// ...but not all, because there might be many relay candidates
// and we don't want to try all of them in a single call:
assert_eq!(triable_relay_candidates(t, now).await?.is_empty(), false);
assert_eq!(load_relay_candidates(t, now).await?.is_empty(), false);
t.assert_warns_or_errors(&[
"DNS lookup with memory cache failure",
+46 -13
View File
@@ -42,7 +42,7 @@ const RINGING_SECONDS: i64 = 120;
const CALL_ACCEPTED_TIMESTAMP: Param = Param::Arg;
const CALL_ENDED_TIMESTAMP: Param = Param::Arg4;
const TURN_PORT: u16 = 3478;
const STUN_PORT: u16 = 3478;
/// Set if incoming call was ended explicitly
/// by the other side before we accepted it.
@@ -658,9 +658,12 @@ pub(crate) async fn create_ice_servers_from_metadata(
Ok((expiration_timestamp, ice_servers))
}
/// TURN server with unresolved DNS name.
/// STUN or TURN server with unresolved DNS name.
#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord)]
pub(crate) enum UnresolvedIceServer {
/// STUN server.
Stun { hostname: String, port: u16 },
/// TURN server with the username and password.
Turn {
hostname: String,
@@ -685,6 +688,28 @@ pub(crate) async fn resolve_ice_servers(
for unresolved_ice_server in unresolved_ice_servers {
match unresolved_ice_server {
UnresolvedIceServer::Stun { hostname, port } => {
match lookup_host_with_cache(context, &hostname, port, "", load_cache).await {
Ok(addrs) => {
let urls: Vec<String> = addrs
.into_iter()
.map(|addr| format!("stun:{addr}"))
.collect();
let stun_server = IceServer {
urls,
username: None,
credential: None,
};
result.push(stun_server);
}
Err(err) => {
warn!(
context,
"Failed to resolve STUN {hostname}:{port}: {err:#}."
);
}
}
}
UnresolvedIceServer::Turn {
hostname,
port,
@@ -717,17 +742,25 @@ pub(crate) async fn resolve_ice_servers(
/// Creates JSON with ICE servers when no TURN servers are known.
pub(crate) fn create_fallback_ice_servers() -> Vec<UnresolvedIceServer> {
// As long as we can't rely on most chat profiles
// having a multi-relay setup and offering TURN servers,
// fall back to a service run by Delta Chat developers
// which ensures that no IP addresses or other metadata is
// persistently logged.
vec![UnresolvedIceServer::Turn {
hostname: "turn.delta.chat".to_string(),
port: TURN_PORT,
username: "public".to_string(),
credential: "o4tR7yG4rG2slhXqRUf9zgmHz".to_string(),
}]
// Do not use public STUN server from https://stunprotocol.org/.
// It changes the hostname every year
// (e.g. stunserver2025.stunprotocol.org
// which was previously stunserver2024.stunprotocol.org)
// because of bandwidth costs:
// <https://github.com/jselbie/stunserver/issues/50>
vec![
UnresolvedIceServer::Stun {
hostname: "nine.testrun.org".to_string(),
port: STUN_PORT,
},
UnresolvedIceServer::Turn {
hostname: "turn.delta.chat".to_string(),
port: STUN_PORT,
username: "public".to_string(),
credential: "o4tR7yG4rG2slhXqRUf9zgmHz".to_string(),
},
]
}
/// Returns JSON with ICE servers.
-10
View File
@@ -779,13 +779,3 @@ async fn test_end_text_call() -> Result<()> {
Ok(())
}
/// Tests that fallback ice servers just carry the Turn one
#[test]
fn test_fallback_ice_servers() {
let hostnames: Vec<String> = create_fallback_ice_servers()
.into_iter()
.map(|UnresolvedIceServer::Turn { hostname, .. }| hostname)
.collect();
assert_eq!(hostnames, ["turn.delta.chat"]);
}
+129 -14
View File
@@ -30,17 +30,17 @@ use crate::download::{DownloadState, PRE_MSG_ATTACHMENT_SIZE_THRESHOLD};
use crate::ensure_and_debug_assert_eq;
use crate::ephemeral::{Timer as EphemeralTimer, start_chat_ephemeral_timers};
use crate::events::EventType;
use crate::key::{Fingerprint, self_fingerprint};
use crate::key::{DcKey as _, Fingerprint, self_fingerprint};
use crate::log::{LogExt, warn};
use crate::logged_debug_assert;
use crate::message::{self, Message, MessageState, MsgId, Viewtype};
use crate::mimefactory::MimeFactory;
use crate::mimefactory;
use crate::mimefactory::{MimeFactory, QueueSideEffects, QueuedMail, ToBeQueuedMail};
use crate::mimeparser::SystemMessage;
use crate::param::{Param, Params};
use crate::pgp::addresses_from_public_key;
use crate::reaction::broadcast_reactions;
use crate::receive_imf::ReceivedMsg;
use crate::smtp::queue::{ToBeQueuedMail, enqueue_mail};
use crate::smtp::send_msg_to_smtp;
use crate::stock_str;
use crate::sync::{self, Sync::*, SyncData};
@@ -331,7 +331,7 @@ impl ChatId {
Ok(chat_id)
}
pub(crate) fn set_selfavatar_timestamp(
fn set_selfavatar_timestamp(
self,
transaction: &mut rusqlite::Transaction<'_>,
timestamp: i64,
@@ -2815,6 +2815,99 @@ async fn render_mime_message_and_pre_message(
}
}
/// Process side effects and store queued mail.
pub(crate) fn enqueue_mail(
transaction: &mut rusqlite::Transaction<'_>,
now: i64,
msg_id: MsgId,
queued_mail: &QueuedMail,
side_effects: Option<&QueueSideEffects>,
) -> Result<i64> {
if let Some(side_effects) = side_effects {
if let Some(last_added_location_timestamp) = side_effects.last_added_location_timestamp {
transaction.execute(
"UPDATE chats SET locations_last_sent=? WHERE id=?;",
(last_added_location_timestamp, side_effects.chat_id),
)?;
}
if side_effects.avatar_is_attached {
side_effects
.chat_id
.set_selfavatar_timestamp(transaction, now)
.context("Failed to set selfavatar timestamp")?;
}
if let Some(ref sync_ids) = side_effects.sync_ids_to_delete {
transaction.execute(
&format!("DELETE FROM multi_device_sync WHERE id IN ({sync_ids})"),
(),
)?;
}
}
// Store mail into queue.
let all_recipients = queued_mail.recipients.join(" ");
let is_encrypted = queued_mail.encryption.is_encrypted();
transaction
.execute(
"
INSERT INTO smtp2 (
display_name,
rfc724_mid,
mime,
should_attach_pubkey,
should_compress,
should_sign,
msg_id,
recipients,
bcc_self,
is_encrypted,
shared_secret,
encryption_fingerprints
)
VALUES (
?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?
)
",
(
&queued_mail.display_name,
&queued_mail.rfc724_mid,
&queued_mail.raw_message,
queued_mail.should_attach_pubkey,
queued_mail.should_compress,
queued_mail.should_sign,
msg_id,
&all_recipients,
queued_mail.bcc_self,
is_encrypted,
if let mimefactory::QueuedEncryption::Symmetric { ref shared_secret } =
queued_mail.encryption
{
shared_secret
} else {
""
},
if let mimefactory::QueuedEncryption::Asymmetric {
ref encryption_pubkeys,
} = queued_mail.encryption
{
let res: Vec<String> = encryption_pubkeys
.iter()
.map(|pubkey| pubkey.dc_fingerprint().hex())
.collect();
res.join(" ")
} else {
"".to_string()
},
),
)
.context("Failed to insert a row into smtp2 table")?;
let row_id = transaction.last_insert_rowid();
Ok(row_id)
}
/// Constructs jobs for sending a message and inserts them into the `smtp` table.
///
/// Updates the message `GuaranteeE2ee` parameter and persists it
@@ -3918,7 +4011,15 @@ pub(crate) async fn add_contact_to_chat_ext(
msg.viewtype = Viewtype::Text;
let contact_addr = contact.get_addr().to_lowercase();
let added_by = ContactId::SELF;
let added_by = if from_handshake && chat.typ == Chattype::OutBroadcast {
// The contact was added via a QR code rather than explicit user action,
// so it could be confusing to say 'You added member Alice'.
// And in a broadcast, SELF is the only one who can add members,
// so, no information is lost by just writing 'Member Alice added' instead.
ContactId::UNDEFINED
} else {
ContactId::SELF
};
msg.text = stock_str::msg_add_member_local(context, contact.id, added_by).await;
msg.param.set_cmd(SystemMessage::MemberAddedToGroup);
msg.param.set(Param::Arg, contact_addr);
@@ -3932,11 +4033,6 @@ pub(crate) async fn add_contact_to_chat_ext(
.await?
.context("Failed to find broadcast shared secret")?;
msg.param.set(PARAM_BROADCAST_SECRET, secret);
// We don't show "member added" info-messages in channels,
// because there can be a lot members added,
// and these messages would clutter the timeline.
msg.hidden = true;
}
send_msg(context, chat_id, &mut msg).await?;
@@ -3985,7 +4081,7 @@ ORDER BY timestamp DESC, id DESC -- final ORDER BY is needed as UNION does not g
(
chat_id,
Viewtype::Webxdc,
constants::N_MSGS_TO_NEW_BROADCAST_MEMBER as u32,
constants::N_MSGS_TO_NEW_BROADCAST_MEMBER,
ContactId::INFO,
),
|row: &rusqlite::Row| Ok(row.get::<_, MsgId>(0)?),
@@ -5104,7 +5200,7 @@ async fn set_contacts_by_fingerprints(
if contacts == contacts_old {
return Ok(());
}
context
let broadcast_contacts_added = context
.sql
.transaction(move |transaction| {
// For broadcast channels, we only add members,
@@ -5121,12 +5217,31 @@ async fn set_contacts_by_fingerprints(
let mut statement = transaction.prepare(
"INSERT OR IGNORE INTO chats_contacts (chat_id, contact_id) VALUES (?, ?)",
)?;
let mut broadcast_contacts_added = Vec::new();
for contact_id in &contacts {
statement.execute((id, contact_id))?;
if statement.execute((id, contact_id))? > 0 && chat.typ == Chattype::OutBroadcast {
broadcast_contacts_added.push(*contact_id);
}
}
Ok(())
Ok(broadcast_contacts_added)
})
.await?;
let timestamp = time();
for added_id in broadcast_contacts_added {
let msg = stock_str::msg_add_member_local(context, added_id, ContactId::UNDEFINED).await;
add_info_msg_with_cmd(
context,
id,
&msg,
SystemMessage::MemberAddedToGroup,
Some(timestamp),
timestamp,
None,
Some(ContactId::SELF),
Some(added_id),
)
.await?;
}
context.emit_event(EventType::ChatModified(id));
Ok(())
}
+29 -15
View File
@@ -3003,13 +3003,13 @@ async fn test_broadcast_change_name() -> Result<()> {
tcm.section("Bob receives the name-change system message");
let msg = bob.recv_msg(&sent).await;
assert_eq!(msg.subject, "My great broadcast");
assert_eq!(msg.subject, "Re: My great broadcast");
let bob_chat = Chat::load_from_db(bob, msg.chat_id).await?;
assert_eq!(bob_chat.name, "My great broadcast");
tcm.section("Fiona receives the name-change system message");
let msg = fiona.recv_msg(&sent).await;
assert_eq!(msg.subject, "My great broadcast");
assert_eq!(msg.subject, "Re: My great broadcast");
let fiona_chat = Chat::load_from_db(fiona, msg.chat_id).await?;
assert_eq!(fiona_chat.name, "My great broadcast");
}
@@ -3343,27 +3343,41 @@ async fn test_broadcast_recipients_sync1() -> Result<()> {
alice2.assert_warn("unknown grpid").await;
let member_added = alice1.pop_sent_msg().await;
alice2.recv_msg_trash(&member_added).await;
let a2_charlie_added = alice2.recv_msg(&member_added).await;
let _c_member_added = charlie.recv_msg(&member_added).await;
let a2_chatlist = Chatlist::try_load(alice2, 0, Some("Channel"), None).await?;
assert_eq!(a2_chatlist.get_msg_id(0)?.unwrap(), a2_charlie_added.id);
// Alice1 will now sync the full member list to Alice2:
sync(alice1, alice2).await;
let a2_bob_contact = alice2.add_or_lookup_contact_id(bob).await;
let a2_charlie_contact = alice2.add_or_lookup_contact_id(charlie).await;
let a2_chatlist = Chatlist::try_load(alice2, 0, Some("Channel"), None).await?;
let a2_chat_id = a2_chatlist.get_chat_id(0).unwrap();
let msg_id = a2_chatlist.get_msg_id(0)?.unwrap();
let a2_bob_added = Message::load_from_db(alice2, msg_id).await?;
assert_ne!(a2_bob_added.id, a2_charlie_added.id);
assert_eq!(
a2_bob_added.text,
stock_str::msg_add_member_local(alice2, a2_bob_contact, ContactId::UNDEFINED).await
);
assert_eq!(a2_bob_added.from_id, ContactId::SELF);
assert_eq!(
a2_bob_added.param.get_cmd(),
SystemMessage::MemberAddedToGroup
);
assert_eq!(
ContactId::new(
a2_bob_added
.param
.get_int(Param::ContactAddedRemoved)
.unwrap()
.try_into()
.unwrap()
),
a2_bob_contact
);
// Also for Alice2, no info message should be shown;
// she should see only the "Messages are end-to-end encrypted" message.
let a2_chat_msgs = get_chat_msgs(alice2, a2_chat_id).await?;
assert_eq!(a2_chat_msgs.len(), 1);
let ChatItem::Message { msg_id } = a2_chat_msgs[0] else {
unreachable!()
};
let a2_msg = Message::load_from_db(alice2, msg_id).await?;
assert_eq!(a2_msg.get_info_type(), SystemMessage::ChatE2ee);
let a2_chat_members = get_chat_contacts(alice2, a2_chat_id).await?;
let a2_chat_members = get_chat_contacts(alice2, a2_charlie_added.chat_id).await?;
assert!(a2_chat_members.contains(&a2_bob_contact));
assert!(a2_chat_members.contains(&a2_charlie_contact));
assert_eq!(a2_chat_members.len(), 2);
+78 -30
View File
@@ -18,8 +18,8 @@ use crate::events::EventType;
use crate::log::LogExt;
use crate::mimefactory::RECOMMENDED_FILE_SIZE;
use crate::sync::{self, Sync::*, SyncData};
use crate::tools::get_abs_path;
use crate::transport::transport_addrs;
use crate::tools::{get_abs_path, time};
use crate::transport::{add_pseudo_transport, send_sync_transports, transport_addrs};
use crate::{constants, stats};
/// The available configuration keys.
@@ -42,12 +42,10 @@ use crate::{constants, stats};
#[strum(serialize_all = "snake_case")]
pub enum Config {
/// Deprecated(2026-04).
/// Use ConfiguredAddr, [`crate::login_param::EnteredLoginParam`],
/// or add_transport{from_qr}()/list_transports() instead.
///
/// Email address used by the deprecated configure() procedure.
///
/// Use add_transport{from_qr}() to configure new transports,
/// Use list_transports() to learn about configured transports,
/// including their addresses.
/// Email address, used in the `From:` field.
Addr,
/// Deprecated(2026-04).
@@ -197,9 +195,9 @@ pub enum Config {
#[strum(props(default = "0"))]
DeleteDeviceAfter,
/// Deprecated(2026-09).
/// The address of the transport used for sending.
///
/// Use ConfiguredLoginParam and list_transports() instead.
/// Device-local, other devices choose their own sending transport.
ConfiguredAddr,
/// Deprecated(2026-04).
@@ -307,6 +305,18 @@ pub enum Config {
/// True if account is configured.
Configured,
/// Deprecated, we are trying to get rid of this global setting.
/// It is possible to configure a profile with both chatmail relays
/// and classical email servers.
///
/// Most usages in UIs can be replaced by `force_encryption`.
///
/// True if account is a chatmail account.
IsChatmail,
/// True if `IsChatmail` mustn't be autoconfigured. For tests.
FixIsChatmail,
/// True if account is muted.
IsMuted,
@@ -489,6 +499,11 @@ impl Config {
| Self::ForceEncryption,
)
}
/// Whether the config option needs an IO scheduler restart to take effect.
pub(crate) fn needs_io_restart(&self) -> bool {
matches!(self, Config::ConfiguredAddr)
}
}
impl Context {
@@ -667,6 +682,10 @@ impl Context {
pub async fn set_config(&self, key: Config, value: Option<&str>) -> Result<()> {
Self::check_config(key, value)?;
let _pause = match key.needs_io_restart() {
true => self.scheduler.pause(self).await?,
_ => Default::default(),
};
if key == Config::StatsSending {
let old_value = self.get_config(key).await?;
let old_value = bool_from_config(old_value.as_deref());
@@ -746,28 +765,57 @@ impl Context {
bail!("Cannot unset configured_addr");
};
self.sql
.transaction(|transaction| {
if transaction.query_row(
"SELECT COUNT(*) FROM transports WHERE addr=?",
(addr,),
|row| {
let res: i64 = row.get(0)?;
Ok(res)
},
)? == 0
{
bail!("Address does not belong to any transport.");
}
transaction.execute(
"INSERT OR REPLACE INTO config (keyname, value) VALUES ('configured_addr', ?)",
(addr,),
)?;
if !self.is_configured().await? {
info!(
self,
"Creating a pseudo configured account which will not be able to send or receive messages. Only meant for tests!"
);
add_pseudo_transport(self, addr).await?;
self.sql
.set_raw_config(Config::ConfiguredAddr.as_ref(), Some(addr))
.await?;
} else {
self.sql
.transaction(|transaction| {
if transaction.query_row(
"SELECT COUNT(*) FROM transports WHERE addr=?",
(addr,),
|row| {
let res: i64 = row.get(0)?;
Ok(res)
},
)? == 0
{
bail!("Address does not belong to any transport.");
}
transaction.execute(
"UPDATE config SET value=? WHERE keyname='configured_addr'",
(addr,),
)?;
Ok(())
})
.await?;
self.sql.uncache_raw_config("configured_addr").await;
// The timestamp must strictly increase because
// other devices ignore the row update otherwise,
// and contacts only adopt the re-signed key
// if its signature timestamp increases.
transaction
.execute(
"UPDATE transports
SET add_timestamp=MAX(?, add_timestamp+1)
WHERE addr=?",
(time(), addr),
)
.context(
"Failed to update add_timestamp for the new sending transport",
)?;
Ok(())
})
.await?;
// Invalidate the cache so the sync message
// cannot read a stale sending address.
self.sql.uncache_raw_config("configured_addr").await;
send_sync_transports(self).await?;
}
}
_ => {
self.sql.set_raw_config(key.as_ref(), value).await?;
+75 -98
View File
@@ -31,7 +31,7 @@ use crate::login_param::EnteredCertificateChecks;
pub use crate::login_param::EnteredLoginParam;
use crate::net::proxy::ProxyConfig;
use crate::provider::{self, Protocol, Socket};
use crate::qr::{Qr, check_qr, login_param_from_account_qr, login_param_from_login_qr};
use crate::qr::{login_param_from_account_qr, login_param_from_login_qr};
use crate::smtp::Smtp;
use crate::sync::Sync::Nosync;
use crate::tools::time;
@@ -47,23 +47,20 @@ use crate::{EventType, autorelay, stock_str};
/// See <https://github.com/chatmail/core/issues/7608>.
pub(crate) const MAX_RELAYS: usize = 5;
tokio::task_local! {
pub(crate) static SILENT_PROGRESS: ();
}
#[track_caller]
fn emit_progress(ctx: &Context, progress: u16) {
assert!(
progress <= 1000,
"value in range 0..1000 expected with: 0=error, 1..999=progress, 1000=success"
);
if SILENT_PROGRESS.try_with(|_| ()).is_ok() {
return;
}
ctx.emit_event(EventType::ConfigureProgress {
progress,
comment: None,
});
macro_rules! progress {
($context:tt, $progress:expr, $comment:expr) => {
assert!(
$progress <= 1000,
"value in range 0..1000 expected with: 0=error, 1..999=progress, 1000=success"
);
$context.emit_event($crate::events::EventType::ConfigureProgress {
progress: $progress,
comment: $comment,
});
};
($context:tt, $progress:expr) => {
progress!($context, $progress, None);
};
}
impl Context {
@@ -127,17 +124,14 @@ impl Context {
pub(crate) async fn add_transport_inner(&self, param: &mut EnteredLoginParam) -> Result<()> {
match self.add_transport_unreported(param).await {
Ok(()) => {
emit_progress(self, 1000);
progress!(self, 1000);
Ok(())
}
Err(err) => {
// We are using Anyhow's .context() and to show the
// inner error, too, we need the {:#}:
let error_msg = stock_str::configuration_failed(self, &format!("{err:#}"));
self.emit_event(EventType::ConfigureProgress {
progress: 0,
comment: Some(error_msg.clone()),
});
progress!(self, 0, Some(error_msg.clone()));
bail!(error_msg);
}
}
@@ -162,7 +156,9 @@ impl Context {
.await;
self.free_ongoing().await;
res
res?;
param.save_legacy(self).await
}
/// Adds a new email account as a transport
@@ -194,59 +190,52 @@ impl Context {
Ok(())
}
/// Adds an initial transport on the chatmail relay that answers fastest
/// and lets the profile add further ones in the background.
/// Automatically adds up to three transports.
///
/// A `DCACCOUNT:` or `DCLOGIN:` `qr` code adds a single transport
/// while securejoin codes add the inviter's relays to the candidates.
///
/// Does nothing if the profile already has a transport.
/// If the user just scanned a QR code of type `Account`, `Login`,
/// `AskVerifyContact`, `AskVerifyGroup`, or `AskJoinBroadcast`,
/// then UI implementations should pass it as the `qr` parameter.
/// The host(s) from the QR code will then also be considered
/// for creating an account there.
pub async fn init_transports(&self, qr: Option<&str>) -> Result<()> {
if self.is_configured().await? {
return Ok(());
bail!("Transports are already initialized");
}
let mut addrs_from_qr = vec![];
if let Some(qr) = qr {
match check_qr(self, qr).await? {
Qr::Account { .. } | Qr::Login { .. } => {
match crate::qr::check_qr(self, qr).await? {
crate::qr::Qr::Account { .. } | crate::qr::Qr::Login { .. } => {
return self.add_transport_from_qr(qr).await;
}
Qr::AskVerifyContact { addrs, .. }
| Qr::AskVerifyGroup { addrs, .. }
| Qr::AskJoinBroadcast { addrs, .. } => {
autorelay::add_relay_candidates(self, &addrs).await?
}
_ => bail!("QR code does not contain a relay"),
crate::qr::Qr::AskVerifyContact { addrs, .. }
| crate::qr::Qr::AskVerifyGroup { addrs, .. }
| crate::qr::Qr::AskJoinBroadcast { addrs, .. } => addrs_from_qr = addrs,
_ => bail!("This QR code can't be used to initialize transports"),
}
}
self.stop_io().await;
let cancel_channel = self.alloc_ongoing().await?;
let skip_network = false;
let res = autorelay::add_transport_from_candidates(self, skip_network)
let res = autorelay::init_transports_inner(self, addrs_from_qr, skip_network)
.race(cancel_channel.recv().map(|_| Err(format_err!("Canceled"))))
.await;
self.free_ongoing().await;
let configured = self.is_configured().await?;
match res {
Ok(()) => {}
Err(err) if configured => {
warn!(
self,
"Onboarding interrupted after adding a transport: {err:#}."
);
}
Err(err) => {
let error_msg = stock_str::configuration_failed(self, &format!("{err:#}"));
self.emit_event(EventType::ConfigureProgress {
progress: 0,
comment: Some(error_msg.clone()),
});
bail!(error_msg);
}
if let Err(err) = res {
let error_msg = stock_str::configuration_failed(self, &format!("{err:#}"));
progress!(self, 0, Some(error_msg.clone()));
bail!(error_msg);
}
self.set_config_bool(Config::Autorelay, true).await?;
emit_progress(self, 1000);
progress!(self, 1000);
self.start_io().await;
Ok(())
}
@@ -381,7 +370,7 @@ async fn get_configured_param(
let parsed = EmailAddress::new(&param.addr).context("Bad email-address")?;
let param_domain = parsed.domain;
emit_progress(ctx, 200);
progress!(ctx, 200);
let param_autoconfig = if param.imap.server.is_empty()
&& param.imap.port == 0
@@ -403,7 +392,7 @@ async fn get_configured_param(
None
};
emit_progress(ctx, 500);
progress!(ctx, 500);
let mut servers = param_autoconfig.unwrap_or_default();
if !servers
@@ -497,13 +486,13 @@ pub(crate) async fn configure(
param: &EnteredLoginParam,
skip_network: bool,
) -> Result<()> {
emit_progress(ctx, 1);
progress!(ctx, 1);
let configured_param = get_configured_param(ctx, param, skip_network).await?;
let proxy_config = ProxyConfig::load(ctx).await?;
let strict_tls = configured_param.strict_tls(proxy_config.is_some())?;
emit_progress(ctx, 550);
progress!(ctx, 550);
if !skip_network {
// Spawn SMTP configuration task
@@ -529,7 +518,7 @@ pub(crate) async fn configure(
Ok::<(), anyhow::Error>(())
});
emit_progress(ctx, 600);
progress!(ctx, 600);
// Configure IMAP
@@ -543,12 +532,23 @@ pub(crate) async fn configure(
}
};
emit_progress(ctx, 850);
progress!(ctx, 850);
// Wait for SMTP configuration
smtp_config_task.await??;
emit_progress(ctx, 900);
progress!(ctx, 900);
let is_configured = ctx.is_configured().await?;
if !ctx.get_config_bool(Config::FixIsChatmail).await? {
if imap_session.is_chatmail() {
ctx.sql.set_raw_config("is_chatmail", Some("1")).await?;
} else if !is_configured {
// Reset the setting that may have been set
// during failed configuration.
ctx.sql.set_raw_config("is_chatmail", Some("0")).await?;
}
}
// Drop the imap connection explicitly
// to make sure that it's not forgotten in a future refactoring
@@ -556,7 +556,7 @@ pub(crate) async fn configure(
drop(imap);
}
emit_progress(ctx, 910);
progress!(ctx, 910);
configured_param
.clone()
@@ -567,11 +567,11 @@ pub(crate) async fn configure(
ctx.set_config_internal(Config::ConfiguredTimestamp, Some(&time().to_string()))
.await?;
emit_progress(ctx, 920);
progress!(ctx, 920);
ctx.scheduler.interrupt_inbox().await;
emit_progress(ctx, 940);
progress!(ctx, 940);
ctx.update_device_chats()
.await
.context("Failed to update device chats")?;
@@ -615,7 +615,7 @@ async fn get_autoconfig(
{
return Some(res);
}
emit_progress(ctx, 300);
progress!(ctx, 300);
// `?emailaddress=` query string is excluded on purpose.
// It is not part of the URL according to <https://datatracker.ietf.org/doc/draft-ietf-mailmaint-autoconfig/06/>.
@@ -630,7 +630,7 @@ async fn get_autoconfig(
{
return Some(res);
}
emit_progress(ctx, 310);
progress!(ctx, 310);
// Outlook uses always SSL but different domains (this comment describes the next two steps)
if let Ok(res) = outlk_autodiscover(
@@ -642,7 +642,7 @@ async fn get_autoconfig(
{
return Some(res);
}
emit_progress(ctx, 320);
progress!(ctx, 320);
if let Ok(res) = outlk_autodiscover(
ctx,
@@ -653,7 +653,7 @@ async fn get_autoconfig(
{
return Some(res);
}
emit_progress(ctx, 330);
progress!(ctx, 330);
// always SSL for Thunderbird's database
if let Ok(res) = moz_autoconfigure(
@@ -730,8 +730,7 @@ mod tests {
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn test_early_configure_failure_is_reported() -> Result<()> {
let t = TestContext::new().await;
let mark_as_autorelay = false;
let mut param = login_param_from_host("example.org", mark_as_autorelay);
let mut param = login_param_from_host("example.org");
// An ongoing process, e.g. a backup import,
// makes configuration fail without ever contacting a relay.
@@ -750,28 +749,6 @@ mod tests {
Ok(())
}
/// Tests that init_transports() fails on a bad code
/// and does nothing on a profile that already has a transport.
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn test_init_transports() -> Result<()> {
let mut tcm = TestContextManager::new();
let t = &tcm.unconfigured().await;
assert!(t.init_transports(Some("not a qr code")).await.is_err());
assert!(!t.is_configured().await?);
let alice = &tcm.alice().await;
let invite = "openpgp4fpr:79252762C34C5096AF57958F4FC3D21A81B0F0A7#a=cli%40invite.example&i=TbnwJ6lSvD5&s=0ejvbdFSQxB";
alice.init_transports(Some(invite)).await?;
let candidates = alice
.sql
.count("SELECT COUNT(*) FROM relay_candidates", ())
.await?;
assert_eq!(candidates, 0);
assert_eq!(alice.count_transports().await?, 1);
Ok(())
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn test_get_configured_param() -> Result<()> {
let t = &TestContext::new().await;
@@ -798,7 +775,7 @@ mod tests {
let mut tcm = TestContextManager::new();
let t = &tcm.unconfigured().await;
add_pseudo_transport(t, "primary@example.org").await?;
// Setting ConfiguredAddr on an unconfigured account creates a pseudo transport
t.set_config(Config::ConfiguredAddr, Some("primary@example.org"))
.await?;
assert_eq!(t.count_transports().await?, 1);
+28 -45
View File
@@ -40,33 +40,9 @@ use crate::sync::{self, Sync::*};
use crate::tools::{SystemTime, duration_to_str, get_abs_path, normalize_text, time, to_lowercase};
use crate::{chat, chatlist_events, ensure_and_debug_assert, stock_str};
/// If the contact's "last seen" is newer, the contact freshness is set to "recently seen".
/// Time during which a contact is considered as seen recently.
const SEEN_RECENTLY_SECONDS: i64 = 600;
/// If the contact's "last seen" is older, the contact freshness is set to "old".
const CONTACT_OLD_SECONDS: i64 = 60 * 24 * 60 * 60;
/// Freshness of a contact, based on when it was last seen.
///
/// Used by the UI to highlight contacts:
/// recently seen contacts get a little green dot on the avatar,
/// contacts not seen for a long time get a string below the name (e.g. "Seen 2 months ago").
#[derive(Debug, PartialEq, Eq)]
pub enum Freshness {
/// Contact shall not be highlighted.
Normal = 0,
/// Contact was seen recently.
RecentlySeen = 1,
/// Contact was not seen for a long time.
Old = 2,
}
impl From<Freshness> for u32 {
fn from(freshness: Freshness) -> Self {
freshness as u32
}
}
/// Contact ID, including reserved IDs.
///
/// Some contact IDs are reserved to identify special contacts. This
@@ -510,8 +486,8 @@ pub struct Contact {
/// The contact ID.
pub id: ContactId,
/// Contact name. It is recommended to use `Contact::get_name`
/// or `Contact::get_display_name` to access this field.
/// Contact name. It is recommended to use `Contact::get_name`,
/// `Contact::get_display_name` or `Contact::get_name_n_addr` to access this field.
/// May be empty, initially set to `authname`.
name: String,
@@ -749,23 +725,10 @@ impl Contact {
self.last_seen
}
/// Returns freshness of the contact.
pub fn get_freshness(&self) -> Freshness {
if self.id.is_special() || !self.is_key_contact() || self.is_blocked() {
return Freshness::Normal;
}
let is_old = time().saturating_sub(self.last_seen) > CONTACT_OLD_SECONDS;
if is_old || self.last_seen <= 0 {
return Freshness::Old;
}
let seen_recently = time().saturating_sub(self.last_seen) <= SEEN_RECENTLY_SECONDS;
if seen_recently {
return Freshness::RecentlySeen;
}
Freshness::Normal
/// Returns `true` if this contact was seen recently.
#[expect(clippy::arithmetic_side_effects)]
pub fn was_seen_recently(&self) -> bool {
time() - self.last_seen <= SEEN_RECENTLY_SECONDS
}
/// Check if a contact is blocked.
@@ -1613,7 +1576,7 @@ WHERE addr=?
/// May be an empty string.
///
/// This name is typically used in a form where the user can edit the name of a contact.
/// To get a fine name to display in lists etc., use `Contact::get_display_name`.
/// To get a fine name to display in lists etc., use `Contact::get_display_name` or `Contact::get_name_n_addr`.
pub fn get_name(&self) -> &str {
&self.name
}
@@ -1633,6 +1596,26 @@ WHERE addr=?
&self.addr
}
/// Get a summary of name and address.
///
/// The returned string is either "Name (email@domain.com)" or just
/// "email@domain.com" if the name is unset.
///
/// The result should only be used locally and never sent over the network
/// as it leaks the local contact name.
///
/// The summary is typically used when asking the user something about the contact.
/// The attached email address makes the question unique, eg. "Chat with Alan Miller (am@uniquedomain.com)?"
pub fn get_name_n_addr(&self) -> String {
if !self.name.is_empty() {
format!("{} ({})", self.name, self.addr)
} else if !self.authname.is_empty() {
format!("{} ({})", self.authname, self.addr)
} else {
(&self.addr).into()
}
}
/// Get the contact's profile image.
/// This is the image set by each remote user on their own
/// using set_config(context, "selfavatar", image).
+17 -50
View File
@@ -253,6 +253,7 @@ async fn test_add_or_lookup() {
assert_eq!(contact.get_authname(), "bla foo");
assert_eq!(contact.get_display_name(), "Name one");
assert_eq!(contact.get_addr(), "one@eins.org");
assert_eq!(contact.get_name_n_addr(), "Name one (one@eins.org)");
// modify first added contact
let (contact_id_test, sth_modified) = Contact::add_or_lookup(
@@ -285,6 +286,7 @@ async fn test_add_or_lookup() {
assert_eq!(contact.get_name(), "");
assert_eq!(contact.get_display_name(), "three@drei.sam");
assert_eq!(contact.get_addr(), "three@drei.sam");
assert_eq!(contact.get_name_n_addr(), "three@drei.sam");
// add name to third contact from incoming message (this becomes authorized name)
let (contact_id_test, sth_modified) = Contact::add_or_lookup(
@@ -298,6 +300,7 @@ async fn test_add_or_lookup() {
assert_eq!(contact_id, contact_id_test);
assert_eq!(sth_modified, Modifier::Modified);
let contact = Contact::get_by_id(&t, contact_id).await.unwrap();
assert_eq!(contact.get_name_n_addr(), "m. serious (three@drei.sam)");
assert!(!contact.is_blocked());
// manually edit name of third contact (does not changed authorized name)
@@ -313,6 +316,7 @@ async fn test_add_or_lookup() {
assert_eq!(sth_modified, Modifier::Modified);
let contact = Contact::get_by_id(&t, contact_id).await.unwrap();
assert_eq!(contact.get_authname(), "m. serious");
assert_eq!(contact.get_name_n_addr(), "schnucki (three@drei.sam)");
assert!(!contact.is_blocked());
// Fourth contact:
@@ -330,6 +334,7 @@ async fn test_add_or_lookup() {
assert_eq!(contact.get_name(), "Wonderland, Alice");
assert_eq!(contact.get_display_name(), "Wonderland, Alice");
assert_eq!(contact.get_addr(), "alice@w.de");
assert_eq!(contact.get_name_n_addr(), "Wonderland, Alice (alice@w.de)");
// check SELF
let contact = Contact::get_by_id(&t, ContactId::SELF).await.unwrap();
@@ -368,6 +373,7 @@ async fn test_contact_name_changes() -> Result<()> {
assert_eq!(contact.get_authname(), "");
assert_eq!(contact.get_name(), "");
assert_eq!(contact.get_display_name(), "f@example.org");
assert_eq!(contact.get_name_n_addr(), "f@example.org");
let contacts = Contact::get_all(&t, 0, Some("f@example.org")).await?;
assert_eq!(contacts.len(), 0);
@@ -393,6 +399,7 @@ async fn test_contact_name_changes() -> Result<()> {
assert_eq!(contact.get_authname(), "Flobbyfoo");
assert_eq!(contact.get_name(), "");
assert_eq!(contact.get_display_name(), "Flobbyfoo");
assert_eq!(contact.get_name_n_addr(), "Flobbyfoo (f@example.org)");
let contacts = Contact::get_all(&t, 0, Some("f@example.org")).await?;
assert_eq!(contacts.len(), 0);
let contacts = Contact::get_all(&t, 0, Some("flobbyfoo")).await?;
@@ -422,6 +429,7 @@ async fn test_contact_name_changes() -> Result<()> {
assert_eq!(contact.get_authname(), "Foo Flobby");
assert_eq!(contact.get_name(), "");
assert_eq!(contact.get_display_name(), "Foo Flobby");
assert_eq!(contact.get_name_n_addr(), "Foo Flobby (f@example.org)");
let contacts = Contact::get_all(&t, 0, Some("f@example.org")).await?;
assert_eq!(contacts.len(), 0);
let contacts = Contact::get_all(&t, 0, Some("flobbyfoo")).await?;
@@ -439,6 +447,7 @@ async fn test_contact_name_changes() -> Result<()> {
assert_eq!(contact.get_authname(), "Foo Flobby");
assert_eq!(contact.get_name(), "Falk");
assert_eq!(contact.get_display_name(), "Falk");
assert_eq!(contact.get_name_n_addr(), "Falk (f@example.org)");
let contacts = Contact::get_all(&t, 0, Some("f@example.org")).await?;
assert_eq!(contacts.len(), 0);
let contacts = Contact::get_all(&t, 0, Some("falk")).await?;
@@ -1048,7 +1057,7 @@ async fn test_last_seen() -> Result<()> {
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn test_contact_freshness() -> Result<()> {
async fn test_was_seen_recently() -> Result<()> {
let _n = TimeShiftFalsePositiveNote;
let mut tcm = TestContextManager::new();
@@ -1061,57 +1070,15 @@ async fn test_contact_freshness() -> Result<()> {
let chat = bob.create_chat(&alice).await;
let contacts = chat::get_chat_contacts(&bob, chat.id).await?;
let contact = Contact::get_by_id(&bob, *contacts.first().unwrap()).await?;
assert_eq!(contact.get_freshness(), Freshness::Old);
assert!(!contact.was_seen_recently());
bob.recv_msg(&sent_msg).await;
let contact = Contact::get_by_id(&bob, *contacts.first().unwrap()).await?;
assert_eq!(contact.get_freshness(), Freshness::RecentlySeen);
assert!(contact.was_seen_recently());
let self_contact = Contact::get_by_id(&bob, ContactId::SELF).await?;
assert_eq!(self_contact.get_freshness(), Freshness::Normal);
Ok(())
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn test_contact_freshness_blocked() -> Result<()> {
let mut tcm = TestContextManager::new();
let alice = tcm.alice().await;
let bob = tcm.bob().await;
// Alice sends message to Bob. Bob receives message, contact's freshness is "recently seen"
let msg = tcm.send_recv(&alice, &bob, "moin").await;
let contact = Contact::get_by_id(&bob, msg.from_id).await?;
assert_eq!(contact.get_freshness(), Freshness::RecentlySeen);
// Bob blocks Alice, contact's freshness is "normal" now
Contact::block(&bob, msg.from_id).await?;
let contact = Contact::get_by_id(&bob, msg.from_id).await?;
assert!(contact.is_blocked());
assert_eq!(contact.get_freshness(), Freshness::Normal);
Ok(())
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn test_contact_freshness_address_contact() -> Result<()> {
let mut tcm = TestContextManager::new();
let alice = tcm.alice().await;
alice.allow_unencrypted().await?;
let bob = tcm.bob().await;
bob.allow_unencrypted().await?;
// Alice sends a message to Bob
let alice_chat = alice.create_email_chat(&bob).await;
let sent_msg = alice.send_text(alice_chat.id, "moin").await;
// Bob receives message, contact's freshness is "normal", even though the messages was just received
let msg = bob.recv_msg(&sent_msg).await;
let contact = Contact::get_by_id(&bob, msg.from_id).await?;
assert!(!contact.is_key_contact());
assert_eq!(contact.get_freshness(), Freshness::Normal);
assert!(!self_contact.was_seen_recently());
Ok(())
}
@@ -1129,11 +1096,11 @@ async fn test_was_seen_recently_event() -> Result<()> {
let chat = alice.create_chat(&bob).await;
let sent_msg = alice.send_text(chat.id, "moin").await;
let contact = Contact::get_by_id(&bob, *contacts.first().unwrap()).await?;
assert_ne!(contact.get_freshness(), Freshness::RecentlySeen);
assert!(!contact.was_seen_recently());
bob.evtracker.clear_events();
bob.recv_msg(&sent_msg).await;
let contact = Contact::get_by_id(&bob, *contacts.first().unwrap()).await?;
assert_eq!(contact.get_freshness(), Freshness::RecentlySeen);
assert!(contact.was_seen_recently());
bob.evtracker
.get_matching(|evt| matches!(evt, EventType::ContactsChanged { .. }))
.await;
@@ -1141,12 +1108,12 @@ async fn test_was_seen_recently_event() -> Result<()> {
.interrupt(contact.id, contact.last_seen)
.await;
// Wait for "seen recently" to turn off.
// Wait for `was_seen_recently()` to turn off.
bob.evtracker.clear_events();
SystemTime::shift(Duration::from_secs(SEEN_RECENTLY_SECONDS as u64 * 2));
recently_seen_loop.interrupt(ContactId::UNDEFINED, 0).await;
let contact = Contact::get_by_id(&bob, *contacts.first().unwrap()).await?;
assert_eq!(contact.get_freshness(), Freshness::Normal);
assert!(!contact.was_seen_recently());
bob.evtracker
.get_matching(|evt| matches!(evt, EventType::ContactsChanged { .. }))
.await;
+82 -21
View File
@@ -67,6 +67,7 @@ pub struct ContextBuilder {
id: u32,
events: Events,
stock_strings: StockStrings,
password: Option<String>,
push_subscriber: Option<PushSubscriber>,
}
@@ -83,6 +84,7 @@ impl ContextBuilder {
id: rand::random(),
events: Events::new(),
stock_strings: StockStrings::new(),
password: None,
push_subscriber: None,
}
}
@@ -129,6 +131,19 @@ impl ContextBuilder {
self
}
/// Sets the password to unlock the database.
/// Deprecated 2025-11:
/// - Db encryption does nothing with blobs, so fs/disk encryption is recommended.
/// - Isolation from other apps is needed anyway.
///
/// If an encrypted database is used it must be opened with a password. Setting a
/// password on a new database will enable encryption.
#[deprecated(since = "TBD")]
pub fn with_password(mut self, password: String) -> Self {
self.password = Some(password);
self
}
/// Sets push subscriber.
pub(crate) fn with_push_subscriber(mut self, push_subscriber: PushSubscriber) -> Self {
self.push_subscriber = Some(push_subscriber);
@@ -153,10 +168,11 @@ impl ContextBuilder {
///
/// Returns error if context cannot be opened.
pub async fn open(self) -> Result<Context> {
let password = self.password.clone().unwrap_or_default();
let context = self.build().await?;
match context.open().await? {
match context.open(password).await? {
true => Ok(context),
false => bail!("FIXME database could not be decrypted, incorrect or missing password"),
false => bail!("database could not be decrypted, incorrect or missing password"),
}
}
}
@@ -215,7 +231,7 @@ pub struct InnerContext {
/// This is a global mutex-like state for operations which should be modal in the
/// clients.
running_state: RwLock<RunningState>,
/// Mutex to prevent running housekeeping or relay management from multiple threads at once.
/// Lock to prevent running housekeeping or relay management from multiple threads at once.
pub(crate) background_task_mutex: Mutex<()>,
/// Mutex to prevent multiple IMAP loops from fetching the messages at once.
@@ -370,7 +386,10 @@ impl Context {
let context =
Self::new_closed(dbfile, id, events, stock_strings, Default::default()).await?;
context.sql.open(&context).await?;
// Open the database if is not encrypted.
if context.check_passphrase("".to_string()).await? {
context.sql.open(&context, "".to_string()).await?;
}
Ok(context)
}
@@ -414,9 +433,20 @@ impl Context {
/// Returns true if passphrase is correct, false is passphrase is not correct. Fails on other
/// errors.
#[deprecated(since = "TBD")]
pub async fn open(&self) -> Result<bool> {
self.sql.open(self).await?;
Ok(true)
pub async fn open(&self, passphrase: String) -> Result<bool> {
if self.sql.check_passphrase(passphrase.clone()).await? {
self.sql.open(self, passphrase).await?;
Ok(true)
} else {
Ok(false)
}
}
/// Changes encrypted database passphrase.
/// Deprecated 2025-11, see [`ContextBuilder::with_password()`] for reasoning.
pub async fn change_passphrase(&self, passphrase: String) -> Result<()> {
self.sql.change_passphrase(passphrase).await?;
Ok(())
}
/// Returns true if database is open.
@@ -424,6 +454,15 @@ impl Context {
self.sql.is_open().await
}
/// Tests the database passphrase.
///
/// Returns true if passphrase is correct.
///
/// Fails if database is already open.
pub(crate) async fn check_passphrase(&self, passphrase: String) -> Result<bool> {
self.sql.check_passphrase(passphrase).await
}
pub(crate) fn with_blobdir(
dbfile: PathBuf,
blobdir: PathBuf,
@@ -529,9 +568,20 @@ impl Context {
self.scheduler.maybe_network().await;
}
/// Returns maximum number of recipients a single email can be sent to
/// over the transport `transport_id`, which sends from `addr`.
pub(crate) async fn get_max_smtp_rcpt_to(&self, transport_id: u32, addr: &str) -> Result<u32> {
/// Deprecated, we are trying to get rid of this global setting.
/// It is possible to configure a profile with both chatmail relays
/// and classical email servers.
///
/// Returns true if an account is on a chatmail server.
pub async fn is_chatmail(&self) -> Result<bool> {
self.get_config_bool(Config::IsChatmail).await
}
/// Returns maximum number of recipients a single email can be sent to.
pub(crate) async fn get_max_smtp_rcpt_to(&self) -> Result<u32> {
let Some((transport_id, param)) = ConfiguredLoginParam::load(self).await? else {
bail!("Not configured");
};
let metadata_limit = self
.metadata
.read()
@@ -541,7 +591,9 @@ impl Context {
if let Some(limit) = metadata_limit {
return Ok(limit);
}
if let Some(limit) = crate::provider::legacy_settings_for_addr(addr)?.max_smtp_rcpt_to {
if let Some(limit) =
crate::provider::legacy_settings_for_addr(&param.addr)?.max_smtp_rcpt_to
{
return Ok(limit);
}
Ok(constants::DEFAULT_MAX_SMTP_RCPT_TO)
@@ -549,13 +601,9 @@ impl Context {
/// Does a single round of fetching messages from all transports and returns.
///
/// If IO is stopped, pauses the scheduler and fetches over a dedicated connection
/// per transport, returning as soon as one of them fetched messages.
/// If IO is running, interrupts IMAP IDLE on all transports
/// and waits until they are done fetching.
///
/// Does not wait for outgoing messages to be sent out,
/// use [`crate::accounts::Accounts::is_sending_finished`] for that.
/// Can be used even if I/O is currently stopped.
/// If I/O is stopped, fetches over a dedicated connection per transport
/// and returns as soon as one of them fetched messages.
pub async fn background_fetch(&self) -> Result<()> {
if !(self.is_configured().await?) {
return Ok(());
@@ -565,9 +613,8 @@ impl Context {
info!(self, "background_fetch started.");
if self.scheduler.is_running().await {
self.scheduler.interrupt_inbox_idle().await;
let include_smtp = false;
self.wait_for_work_done(include_smtp).await;
self.scheduler.maybe_network().await;
self.wait_for_all_work_done().await;
} else {
self.scheduler.background_fetch_any(self).await?;
}
@@ -811,6 +858,13 @@ impl Context {
res.insert("number_of_contacts", contacts.to_string());
res.insert("database_dir", self.get_dbfile().display().to_string());
res.insert("database_version", dbversion.to_string());
res.insert(
"database_encrypted",
self.sql
.is_encrypted()
.await
.map_or_else(|| "closed".to_string(), |b| b.to_string()),
);
res.insert("journal_mode", journal_mode);
res.insert("blobdir", self.get_blobdir().display().to_string());
res.insert(
@@ -826,6 +880,13 @@ impl Context {
res.insert("imap_server_id", format!("{server_id:?}"));
}
res.insert("is_chatmail", self.is_chatmail().await?.to_string());
res.insert(
"fix_is_chatmail",
self.get_config_bool(Config::FixIsChatmail)
.await?
.to_string(),
);
res.insert(
"is_muted",
self.get_config_bool(Config::IsMuted).await?.to_string(),
+60 -24
View File
@@ -1,37 +1,16 @@
use anyhow::Context as _;
use strum::IntoEnumIterator;
use tempfile::tempdir;
use super::*;
use crate::chat::{Chat, MuteDuration, get_chat_contacts, get_chat_msgs, send_msg, set_muted};
use crate::chatlist::Chatlist;
use crate::constants::{Chattype, DEFAULT_MAX_SMTP_RCPT_TO};
use crate::constants::Chattype;
use crate::message::Message;
use crate::receive_imf::receive_imf;
use crate::test_utils::{E2EE_INFO_MSGS, TestContext, TestContextManager};
use crate::tools::{SystemTime, create_outgoing_rfc724_mid};
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn test_get_max_smtp_rcpt_to() -> Result<()> {
let t = TestContext::new().await;
assert_eq!(
t.get_max_smtp_rcpt_to(2, "alice@example.org").await?,
DEFAULT_MAX_SMTP_RCPT_TO
);
for (transport_id, limit) in [(1, 3), (2, 7)] {
t.metadata.write().await.insert(
transport_id,
ServerMetadata {
max_smtp_rcpt_to: Some(limit),
..Default::default()
},
);
}
assert_eq!(t.get_max_smtp_rcpt_to(2, "alice@example.org").await?, 7);
assert_eq!(t.get_max_smtp_rcpt_to(1, "alice@example.org").await?, 3);
Ok(())
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn test_wrong_db() -> Result<()> {
let tmp = tempfile::tempdir()?;
@@ -484,6 +463,63 @@ async fn test_limit_search_msgs() -> Result<()> {
Ok(())
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn test_check_passphrase() -> Result<()> {
let dir = tempdir()?;
let dbfile = dir.path().join("db.sqlite");
let context = ContextBuilder::new(dbfile.clone())
.with_id(1)
.build()
.await
.context("failed to create context")?;
assert_eq!(context.open("foo".to_string()).await?, true);
assert_eq!(context.is_open().await, true);
drop(context);
let context = ContextBuilder::new(dbfile)
.with_id(2)
.build()
.await
.context("failed to create context")?;
assert_eq!(context.is_open().await, false);
assert_eq!(context.check_passphrase("bar".to_string()).await?, false);
assert_eq!(context.open("false".to_string()).await?, false);
assert_eq!(context.open("foo".to_string()).await?, true);
Ok(())
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn test_context_change_passphrase() -> Result<()> {
let dir = tempdir()?;
let dbfile = dir.path().join("db.sqlite");
let context = ContextBuilder::new(dbfile)
.with_id(1)
.build()
.await
.context("failed to create context")?;
assert_eq!(context.open("foo".to_string()).await?, true);
assert_eq!(context.is_open().await, true);
context
.set_config(Config::Addr, Some("alice@example.org"))
.await?;
context
.change_passphrase("bar".to_string())
.await
.context("Failed to change passphrase")?;
assert_eq!(
context.get_config(Config::Addr).await?.unwrap(),
"alice@example.org"
);
Ok(())
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn test_ongoing() -> Result<()> {
let context = TestContext::new().await;
+3 -1
View File
@@ -6,6 +6,7 @@ use anyhow::{Result, anyhow, bail, ensure};
use deltachat_derive::{FromSql, ToSql};
use serde::{Deserialize, Serialize};
use crate::config::Config;
use crate::context::Context;
use crate::imap::session::Session;
use crate::log::warn;
@@ -163,7 +164,8 @@ pub(crate) async fn download_msg(
.fetch_single_msg(context, &server_folder, server_uid, rfc724_mid)
.await?;
if ephemeral::should_delete_all_downloaded_messages(context).await? {
let bcc_self = context.get_config_bool(Config::BccSelf).await?;
if ephemeral::should_delete_all_downloaded_messages(bcc_self, session.is_chatmail()) {
// Now that the message was downloaded, it likely needs to be deleted;
// trigger a re-check by interrupting the inbox folder.
// This is mainly needed to make the tests pass;
+2
View File
@@ -6,6 +6,7 @@ mod tests {
use crate::chat;
use crate::chat::send_text_msg;
use crate::config::Config;
use crate::message::Message;
use crate::mimeparser::SystemMessage;
use crate::receive_imf::receive_imf;
@@ -60,6 +61,7 @@ Sent with my Delta Chat Messenger: https://delta.chat";
async fn test_chatmail_can_send_unencrypted() -> Result<()> {
let mut tcm = TestContextManager::new();
let bob = &tcm.bob().await;
bob.set_config_bool(Config::IsChatmail, true).await?;
bob.allow_unencrypted().await?;
let bob_chat_id = receive_imf(
bob,
+14 -28
View File
@@ -363,7 +363,7 @@ pub(crate) async fn start_chat_ephemeral_timers(context: &Context, chat_id: Chat
/// Selects messages which are expired according to
/// `delete_device_after` setting or `ephemeral_timestamp` column.
///
/// For each message a row ID, chat id, viewtype, whether the message is pinned and location ID is returned.
/// For each message a row ID, chat id, viewtype and location ID is returned.
///
/// Unknown viewtypes are returned as `Viewtype::Unknown`
/// and not as errors bubbled up, easily resulting in infinite loop or leaving messages undeleted.
@@ -371,12 +371,12 @@ pub(crate) async fn start_chat_ephemeral_timers(context: &Context, chat_id: Chat
async fn select_expired_messages(
context: &Context,
now: i64,
) -> Result<Vec<(MsgId, ChatId, Viewtype, bool, u32)>> {
) -> Result<Vec<(MsgId, ChatId, Viewtype, u32)>> {
let mut rows = context
.sql
.query_map_vec(
r#"
SELECT id, chat_id, type, pinned, location_id
SELECT id, chat_id, type, location_id
FROM msgs
WHERE
ephemeral_timestamp != 0
@@ -392,9 +392,8 @@ WHERE
.context("Using default viewtype for ephemeral handling.")
.log_err(context)
.unwrap_or_default();
let pinned: bool = row.get("pinned")?;
let location_id: u32 = row.get("location_id")?;
Ok((id, chat_id, viewtype, pinned, location_id))
Ok((id, chat_id, viewtype, location_id))
},
)
.await?;
@@ -415,7 +414,7 @@ WHERE
.sql
.query_map_vec(
r#"
SELECT id, chat_id, type, pinned, location_id
SELECT id, chat_id, type, location_id
FROM msgs
WHERE
timestamp < ?1
@@ -438,9 +437,8 @@ WHERE
.context("Using default viewtype for delete-old handling.")
.log_err(context)
.unwrap_or_default();
let pinned: bool = row.get("pinned")?;
let location_id: u32 = row.get("location_id")?;
Ok((id, chat_id, viewtype, pinned, location_id))
Ok((id, chat_id, viewtype, location_id))
},
)
.await?;
@@ -465,15 +463,11 @@ pub(crate) async fn delete_expired_messages(context: &Context, now: i64) -> Resu
if !rows.is_empty() {
info!(context, "Attempting to delete {} messages.", rows.len());
let (msgs_changed, webxdc_deleted, pinned_chat_ids) = context
let (msgs_changed, webxdc_deleted) = context
.sql
.transaction(|transaction| {
let mut msgs_changed = Vec::with_capacity(rows.len());
let mut webxdc_deleted = Vec::new();
// IDs of the chats in which pinned messages were deleted.
let mut pinned_chat_ids = BTreeSet::new();
// If you change which information is preserved here, also change `MsgId::trash()`
// and other places it references.
let mut del_msg_stmt = transaction.prepare(
@@ -484,7 +478,7 @@ SELECT ?1, rfc724_mid, pre_rfc724_mid, timestamp, ? FROM msgs WHERE id=?1
)?;
let mut del_location_stmt =
transaction.prepare("DELETE FROM locations WHERE independent=1 AND id=?")?;
for (msg_id, chat_id, viewtype, is_pinned, location_id) in rows {
for (msg_id, chat_id, viewtype, location_id) in rows {
del_msg_stmt.execute((msg_id, ChatId::TRASH))?;
if location_id > 0 {
del_location_stmt.execute((location_id,))?;
@@ -494,12 +488,8 @@ SELECT ?1, rfc724_mid, pre_rfc724_mid, timestamp, ? FROM msgs WHERE id=?1
if viewtype == Viewtype::Webxdc {
webxdc_deleted.push(msg_id)
}
if is_pinned {
pinned_chat_ids.insert(chat_id);
}
}
Ok((msgs_changed, webxdc_deleted, pinned_chat_ids))
Ok((msgs_changed, webxdc_deleted))
})
.await?;
@@ -514,10 +504,6 @@ SELECT ?1, rfc724_mid, pre_rfc724_mid, timestamp, ? FROM msgs WHERE id=?1
context.emit_msgs_changed_without_msg_id(modified_chat_id);
}
for chat_id in pinned_chat_ids {
context.emit_event(EventType::PinnedMessagesChanged { chat_id });
}
for msg_id in webxdc_deleted {
context.emit_event(EventType::WebxdcInstanceDeleted { msg_id });
}
@@ -668,11 +654,12 @@ pub(crate) async fn ephemeral_loop(context: &Context, interrupt_receiver: Receiv
pub(crate) async fn delete_expired_imap_messages(
context: &Context,
transport_id: u32,
is_chatmail: bool,
) -> Result<()> {
let now = time();
let bcc_self = context.get_config_bool(Config::BccSelf).await?;
if should_delete_all_downloaded_messages(context).await? {
if should_delete_all_downloaded_messages(bcc_self, is_chatmail) {
// This is the only device using this relay.
// Mark all downloaded messages for deletion, because they are not needed anymore.
//
@@ -721,7 +708,7 @@ pub(crate) async fn delete_expired_imap_messages(
)
.await?;
} else {
// Single device, but unencrypted messages are allowed.
// Single device.
// Delete all expired and encrypted messages.
context
.sql
@@ -749,9 +736,8 @@ pub(crate) async fn delete_expired_imap_messages(
Ok(())
}
pub(crate) async fn should_delete_all_downloaded_messages(context: &Context) -> Result<bool> {
Ok(!context.get_config_bool(Config::BccSelf).await?
&& context.get_config_bool(Config::ForceEncryption).await?)
pub(crate) fn should_delete_all_downloaded_messages(bcc_self: bool, is_chatmail: bool) -> bool {
!bcc_self && is_chatmail
}
/// Start ephemeral timers for seen messages if they are not started
+6 -12
View File
@@ -509,7 +509,7 @@ async fn test_delete_expired_imap_messages() -> Result<()> {
.await?;
}
for (force_encryption, other_transport, bcc_self) in [
for (is_chatmail, other_transport, bcc_self) in [
(false, false, false),
(false, false, true),
(false, true, false),
@@ -520,12 +520,10 @@ async fn test_delete_expired_imap_messages() -> Result<()> {
(true, true, true),
] {
println!(
"Testing combination force_encryption={force_encryption}, other_transport={other_transport}, bcc_self={bcc_self}"
"Testing combination is_chatmail={is_chatmail}, other_transport={other_transport}, bcc_self={bcc_self}"
);
t.set_config_bool(Config::BccSelf, bcc_self).await?;
t.set_config_bool(Config::ForceEncryption, force_encryption)
.await?;
delete_expired_imap_messages(
&t,
@@ -534,6 +532,7 @@ async fn test_delete_expired_imap_messages() -> Result<()> {
} else {
transport_id
},
is_chatmail,
)
.await?;
@@ -549,7 +548,7 @@ async fn test_delete_expired_imap_messages() -> Result<()> {
assert_eq!(is_deleted(&t, "no_expire@localhost").await?, !bcc_self);
assert_eq!(
is_deleted(&t, "no_expire_unencrypted@localhost").await?,
force_encryption && !bcc_self
is_chatmail && !bcc_self
);
assert_eq!(is_deleted(&t, "future@localhost").await?, !bcc_self);
assert_eq!(is_deleted(&t, "expired_post@localhost").await?, true);
@@ -564,10 +563,9 @@ async fn test_delete_expired_imap_messages() -> Result<()> {
reset_targets(&t).await;
}
// With BccSelf=true, non-expired messages are kept even if `force_encryption` is true
// With BccSelf=true, non-expired messages are kept even if `is_chatmail` is true
t.set_config_bool(Config::BccSelf, true).await?;
t.set_config_bool(Config::ForceEncryption, true).await?;
delete_expired_imap_messages(&t, transport_id).await?;
delete_expired_imap_messages(&t, transport_id, true).await?;
assert_eq!(is_deleted(&t, "expired@localhost").await?, true);
assert_eq!(is_deleted(&t, "no_expire@localhost").await?, false);
assert_eq!(is_deleted(&t, "done_pre@localhost").await?, false);
@@ -684,10 +682,6 @@ async fn test_ephemeral_msg_offline() -> Result<()> {
check_msg_will_be_deleted(alice, msg.id, &chat, now, now + i64::from(duration) + 1).await?;
assert!(alice.sql.exists(stmt, (msg.id,)).await?);
alice
.assert_warn("No SMTP connection candidates provided")
.await;
Ok(())
}
+3 -11
View File
@@ -33,6 +33,9 @@ pub enum EventType {
/// Emitted when an IMAP message has been marked as deleted
ImapMessageDeleted(String),
/// Emitted when an IMAP message has been moved
ImapMessageMoved(String),
/// Emitted before going into IDLE on the Inbox folder.
ImapInboxIdle,
@@ -93,14 +96,6 @@ pub enum EventType {
contact_id: ContactId,
},
/// The list of pinned messages for the chat has changed.
///
/// Some message got pinned, or pinned message is unpinned or deleted.
PinnedMessagesChanged {
/// ID of the chat where the list of pinned messages changed.
chat_id: ChatId,
},
/// A reaction to one's own sent message received.
/// Typically, the UI will show a notification for that.
///
@@ -368,9 +363,6 @@ pub enum EventType {
/// A call made while another background fetch is running gets the event immediately,
/// and the running fetch keeps emitting events until its own marker.
///
/// The event carries no data identifying the call it belongs to,
/// so it is unambiguous only if there are no concurrent background fetch calls.
///
/// This event is only emitted by the account manager.
AccountsBackgroundFetchDone,
/// Inform that set of chats or the order of the chats in the chatlist has changed.
+104 -28
View File
@@ -129,8 +129,6 @@ pub(crate) struct ServerMetadata {
/// True if we think the relay supports push notifications.
/// This gates wether we attempt to write an encrypted device token
/// to per-transport IMAP metadata key `/private/devicetoken`.
///
/// Any `/shared/vendor/deltachat/` entry identifies a chatmail relay.
pub supports_push: bool,
/// ICE servers for WebRTC calls.
@@ -301,8 +299,7 @@ impl Imap {
self.conn_backoff_ms = max(BACKOFF_MIN_MS, self.conn_backoff_ms);
let login_params = prioritize_server_login_params(&context.sql, &self.lp, "imap").await?;
let mut first_connection_error = None;
let mut first_login_error = None;
let mut first_error = None;
'candidate: for lp in login_params {
info!(context, "IMAP trying to connect to {}.", lp.connection);
let connection_candidate = lp.connection.clone();
@@ -318,7 +315,7 @@ impl Imap {
Ok(client) => client,
Err(err) => {
warn!(context, "{err:#}.");
first_connection_error.get_or_insert(err);
first_error.get_or_insert(err);
continue 'candidate;
}
};
@@ -397,14 +394,12 @@ impl Imap {
Err(err) => {
warn!(context, "{err:#}.");
first_login_error.get_or_insert(err);
first_error.get_or_insert(err);
}
}
}
Err(first_login_error
.or(first_connection_error)
.unwrap_or_else(|| format_err!("No IMAP connection candidates provided")))
Err(first_error.unwrap_or_else(|| format_err!("No IMAP connection candidates provided")))
}
/// Prepare a new IMAP session.
@@ -423,13 +418,13 @@ impl Imap {
Ok(session)
}
/// FETCH-and-DELETE iteration.
/// FETCH-MOVE-DELETE iteration.
///
/// Prefetches headers and downloads new message from the folder
/// and deletes messages in the folder.
/// Prefetches headers and downloads new message from the folder, moves messages away from the
/// folder and deletes messages in the folder.
///
/// Returns true if at least one message was fetched.
pub async fn fetch_delete(
pub async fn fetch_move_delete(
&mut self,
context: &Context,
session: &mut Session,
@@ -455,14 +450,14 @@ impl Imap {
// Mark expired messages for deletion. Note that `delete_expired_imap_messages` is
// not well optimized and should not be called before fetching.
delete_expired_imap_messages(context, session.transport_id())
delete_expired_imap_messages(context, session.transport_id(), session.is_chatmail())
.await
.context("delete_expired_imap_messages")?;
session
.delete_messages(context, watch_folder)
.move_delete_messages(context, watch_folder)
.await
.context("delete_messages")?;
.context("move_delete_messages")?;
Ok(msgs_fetched)
}
@@ -573,6 +568,16 @@ impl Imap {
.size
.context("imap fetch response does not contain size")?;
// Determine the target folder where the message should be moved to.
//
// We only move the messages from the INBOX and Spam folders.
// This is required to avoid infinite MOVE loop on IMAP servers
// that alias `DeltaChat` folder to other names.
// For example, some Dovecot servers alias `DeltaChat` folder to `INBOX.DeltaChat`.
// In this case moving from `INBOX.DeltaChat` to `DeltaChat`
// results in the messages getting a new UID,
// so the messages will be detected as new
// in the `INBOX.DeltaChat` folder again.
let delete = if let Some(message_id) = &message_id {
message::rfc724_mid_exists_ext(context, message_id, "deleted=1")
.await?
@@ -610,7 +615,13 @@ impl Imap {
)
.await?;
if !delete
// Download only the messages which have reached their target folder if there are
// multiple devices. This prevents race conditions in multidevice case, where one
// device tries to download the message while another device moves the message at the
// same time. Even in single device case it is possible to fail downloading the first
// message, move it to the movebox and then download the second message before
// downloading the first one, if downloading from inbox before moving is allowed.
if folder == target
&& prefetch_should_download(context, &headers, &message_id, fetch_response.flags())
.await
.context("prefetch_should_download")?
@@ -847,10 +858,75 @@ impl Session {
Ok(())
}
/// Deletes messages as planned in the `imap` table.
/// Moves batch of messages identified by their UID from the currently
/// selected folder to the target folder.
async fn move_message_batch(
&mut self,
context: &Context,
set: &str,
row_ids: Vec<i64>,
target: &str,
) -> Result<()> {
if self.can_move() {
match self.uid_mv(set, &target).await {
Ok(()) => {
// Messages are moved or don't exist, IMAP returns OK response in both cases.
context
.sql
.transaction(|transaction| {
let mut stmt = transaction.prepare("DELETE FROM imap WHERE id = ?")?;
for row_id in row_ids {
stmt.execute((row_id,))?;
}
Ok(())
})
.await
.context("Cannot delete moved messages from imap table")?;
context.emit_event(EventType::ImapMessageMoved(format!(
"IMAP messages {set} moved to {target}"
)));
return Ok(());
}
Err(err) => {
warn!(
context,
"Cannot move messages, fallback to COPY/DELETE {} to {}: {}",
set,
target,
err
);
}
}
}
// Server does not support MOVE or MOVE failed.
// Copy messages to the destination folder if needed and mark records for deletion.
info!(
context,
"Server does not support MOVE, fallback to COPY/DELETE {} to {}", set, target
);
self.uid_copy(&set, &target).await?;
context
.sql
.transaction(|transaction| {
let mut stmt = transaction.prepare("UPDATE imap SET target='' WHERE id = ?")?;
for row_id in row_ids {
stmt.execute((row_id,))?;
}
Ok(())
})
.await
.context("Cannot plan deletion of messages")?;
context.emit_event(EventType::ImapMessageMoved(format!(
"IMAP messages {set} copied to {target}"
)));
Ok(())
}
/// Moves and deletes messages as planned in the `imap` table.
///
/// This is the only place where messages are deleted on the IMAP server.
async fn delete_messages(&mut self, context: &Context, folder: &str) -> Result<()> {
/// This is the only place where messages are moved or deleted on the IMAP server.
async fn move_delete_messages(&mut self, context: &Context, folder: &str) -> Result<()> {
let transport_id = self.transport_id();
let rows = context
.sql
@@ -872,21 +948,23 @@ impl Session {
for (target, rowid_set, uid_set) in UidGrouper::from(rows) {
// Select folder inside the loop to avoid selecting it if there are no pending
// DELETE operations. This does not result in multiple SELECT commands
// MOVE/DELETE operations. This does not result in multiple SELECT commands
// being sent because `select_folder()` does nothing if the folder is already
// selected.
let folder_exists = self.select_with_uidvalidity(context, folder).await?;
ensure!(folder_exists, "No folder {folder}");
// Empty target folder name means messages should be deleted.
// Since we don't move messages between IMAP folders anymore,
// `target` is always either empty or equal to `folder`.
debug_assert!(target.is_empty() || folder == target);
if target.is_empty() {
self.delete_message_batch(context, &uid_set, rowid_set)
.await
.with_context(|| format!("cannot delete batch of messages {uid_set:?}"))?;
} else {
self.move_message_batch(context, &uid_set, rowid_set, &target)
.await
.with_context(|| {
format!("cannot move batch of messages {uid_set:?} to folder {target:?}",)
})?;
}
}
@@ -1306,8 +1384,6 @@ impl Session {
_ => {}
}
}
let supports_push =
max_smtp_rcpt_to.is_some() || iroh_relay.is_some() || ice_servers.is_some();
let ice_servers = if let Some(ice_servers) = ice_servers {
ice_servers
} else {
@@ -1323,7 +1399,7 @@ impl Session {
admin,
iroh_relay,
max_smtp_rcpt_to,
supports_push,
supports_push: max_smtp_rcpt_to.is_some() || self.capabilities.has_xdeltapush,
ice_servers,
ice_servers_expiration_timestamp,
app_versions,
+15
View File
@@ -9,6 +9,10 @@ pub(crate) struct Capabilities {
/// <https://tools.ietf.org/html/rfc2177>
pub can_idle: bool,
/// True if the server has MOVE capability as defined in
/// <https://tools.ietf.org/html/rfc6851>
pub can_move: bool,
/// True if the server has QUOTA capability as defined in
/// <https://tools.ietf.org/html/rfc2087>
pub can_check_quota: bool,
@@ -21,6 +25,17 @@ pub(crate) struct Capabilities {
/// <https://tools.ietf.org/html/rfc4978>
pub can_compress: bool,
/// True if the server advertises the legacy `XDELTAPUSH` capability.
pub has_xdeltapush: bool,
/// True if the server has an XCHATMAIL capability
/// indicating that it is a <https://github.com/deltachat/chatmail> server.
///
/// This can be used to hide some advanced settings in the UI
/// that are only interesting for normal email accounts,
/// e.g. the ability to move messages to Delta Chat folder.
pub is_chatmail: bool,
/// Server ID if the server supports ID capability.
pub server_id: Option<HashMap<String, String>>,
}
+3
View File
@@ -78,9 +78,12 @@ pub(crate) async fn identify_server(
};
let capabilities = Capabilities {
can_idle: caps.has_str("IDLE"),
can_move: caps.has_str("MOVE"),
can_check_quota: caps.has_str("QUOTA"),
can_metadata: caps.has_str("METADATA"),
can_compress: caps.has_str("COMPRESS=DEFLATE"),
has_xdeltapush: caps.has_str("XDELTAPUSH"),
is_chatmail: caps.has_str("XCHATMAIL"),
server_id,
};
Ok(capabilities)
+9
View File
@@ -99,6 +99,10 @@ impl Session {
self.capabilities.can_idle
}
pub fn can_move(&self) -> bool {
self.capabilities.can_move
}
pub fn can_check_quota(&self) -> bool {
self.capabilities.can_check_quota
}
@@ -107,6 +111,11 @@ impl Session {
self.capabilities.can_metadata
}
// Returns true if IMAP server has `XCHATMAIL` capability.
pub(crate) fn is_chatmail(&self) -> bool {
self.capabilities.is_chatmail
}
/// Prefetch `n_uids` messages starting from `uid_next`. Returns a list of fetch results in the
/// order of ascending UIDs.
#[expect(clippy::arithmetic_side_effects)]
+26 -16
View File
@@ -200,9 +200,6 @@ async fn import_backup(
backup_to_import: &Path,
passphrase: String,
) -> Result<()> {
if !passphrase.is_empty() {
bail!("Encrypted passphrase is not supported");
}
let backup_file = File::open(backup_to_import).await?;
let file_size = backup_file.metadata().await?.len();
info!(
@@ -213,7 +210,7 @@ async fn import_backup(
context.get_dbfile().display()
);
import_backup_stream(context, backup_file, file_size).await?;
import_backup_stream(context, backup_file, file_size, passphrase).await?;
Ok(())
}
@@ -234,6 +231,7 @@ pub(crate) async fn import_backup_stream<R: tokio::io::AsyncRead + Unpin>(
context: &Context,
backup_file: R,
file_size: u64,
passphrase: String,
) -> Result<()> {
ensure!(
!context.is_configured().await?,
@@ -244,7 +242,7 @@ pub(crate) async fn import_backup_stream<R: tokio::io::AsyncRead + Unpin>(
"Cannot import backup, IO is running"
);
import_backup_stream_inner(context, backup_file, file_size)
import_backup_stream_inner(context, backup_file, file_size, passphrase)
.await
.0
}
@@ -317,6 +315,7 @@ async fn import_backup_stream_inner<R: tokio::io::AsyncRead + Unpin>(
context: &Context,
backup_file: R,
file_size: u64,
passphrase: String,
) -> (Result<()>,) {
let backup_file = ProgressReader::new(backup_file, context.clone(), file_size);
let mut archive = Archive::new(backup_file);
@@ -363,7 +362,7 @@ async fn import_backup_stream_inner<R: tokio::io::AsyncRead + Unpin>(
if res.is_ok() {
res = context
.sql
.import(&unpacked_database)
.import(&unpacked_database, passphrase.clone())
.await
.context("cannot import unpacked database");
}
@@ -391,7 +390,7 @@ async fn import_backup_stream_inner<R: tokio::io::AsyncRead + Unpin>(
}
context
.sql
.open(context)
.open(context, "".to_string())
.await
.log_err(context)
.ok();
@@ -736,7 +735,7 @@ where
/// overwritten.
///
/// This also verifies that IO is not running during the export.
async fn export_database(context: &Context, dest: &Path, _passphrase: String) -> Result<()> {
async fn export_database(context: &Context, dest: &Path, passphrase: String) -> Result<()> {
ensure!(
!context.scheduler.is_running().await,
"cannot export backup, IO is running"
@@ -746,7 +745,6 @@ async fn export_database(context: &Context, dest: &Path, _passphrase: String) ->
let dest = dest
.to_str()
.with_context(|| format!("path {} is not valid unicode", dest.display()))?;
let mut dest_conn = rusqlite::Connection::open(dest)?;
context.set_config(Config::BccSelf, Some("1")).await?;
context
@@ -757,12 +755,22 @@ async fn export_database(context: &Context, dest: &Path, _passphrase: String) ->
context
.sql
.call_write(|conn| {
if let Err(err) = conn.execute("VACUUM", ()) {
warn!(context, "Vacuum failed, exporting anyway: {err:#}.");
}
let backup = rusqlite::backup::Backup::new(conn, &mut dest_conn)?;
backup.run_to_completion(5, std::time::Duration::ZERO, None)?;
conn.execute("VACUUM;", ())
.map_err(|err| warn!(context, "Vacuum failed, exporting anyway {err}"))
.ok();
conn.execute("ATTACH DATABASE ? AS backup KEY ?", (dest, passphrase))
.context("failed to attach backup database")?;
let res = conn
.query_row("SELECT sqlcipher_export('backup')", [], |_row| Ok(()))
.context("failed to export to attached backup database");
conn.execute(
"UPDATE backup.config SET value='0' WHERE keyname='verified_one_on_one_chats';",
[],
)
.ok(); // Deprecated 2025-07. If verified_one_on_one_chats was not set, this errors, which we ignore
conn.execute("DETACH DATABASE backup", [])
.context("failed to detach backup database")?;
res?;
Ok(())
})
.await
@@ -973,6 +981,7 @@ mod tests {
context1.get_config(Config::BccSelf).await?,
Some("0".to_string())
);
context1.set_config_bool(Config::IsChatmail, true).await?;
assert_eq!(context1.get_config_bool(Config::IsMuted).await?, false);
context1.set_config_bool(Config::IsMuted, true).await?;
@@ -992,6 +1001,7 @@ mod tests {
.get_matching(|evt| matches!(evt, EventType::ImexProgress(1000)))
.await;
assert!(context2.is_configured().await?);
assert!(context2.is_chatmail().await?);
for ctx in [context1, context2] {
// BccSelf should be enabled automatically when exporting a backup
assert_eq!(ctx.get_config_bool(Config::BccSelf).await?, true);
@@ -1032,7 +1042,7 @@ mod tests {
ar.unpack(&unpack_dir).await?;
let sql = sql::Sql::new(unpack_dir.path().join(DBFILE_BACKUP_NAME));
sql.open(&context2).await?;
sql.open(&context2, "".to_string()).await?;
assert_eq!(
sql.get_raw_config_int("backup_version").await?.unwrap(),
DCBACKUP_VERSION
+2 -1
View File
@@ -324,6 +324,7 @@ pub async fn get_backup2(
info!(context, "Sending backup authentication token.");
send_stream.write_all(auth_token.as_bytes()).await?;
let passphrase = String::new();
info!(context, "Starting to read backup from the stream.");
let mut file_size_buf = [0u8; 8];
@@ -333,7 +334,7 @@ pub async fn get_backup2(
// Emit a nonzero progress so that UIs can display smth like "Transferring...".
context.emit_event(EventType::ImexProgress(1));
import_backup_stream(context, recv_stream, file_size)
import_backup_stream(context, recv_stream, file_size, passphrase)
.await
.context("Failed to import backup from QUIC stream")?;
info!(context, "Finished importing backup from the stream.");
-4
View File
@@ -662,7 +662,6 @@ mod tests {
use crate::config::Config;
use crate::test_utils::{TestContext, TestContextManager, alice_keypair};
use crate::tools::SystemTime;
use crate::transport::add_pseudo_transport;
static KEYPAIR: LazyLock<SignedSecretKey> = LazyLock::new(alice_keypair);
@@ -812,7 +811,6 @@ i8pcjGO+IZffvyZJVRWfVooBJmWWbPB1pueo3tx8w3+fcuzpxz+RLFKaPyqXO+dD
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn test_load_self_generate_public() {
let t = TestContext::new().await;
add_pseudo_transport(&t, "alice@example.org").await.unwrap();
t.set_config(Config::ConfiguredAddr, Some("alice@example.org"))
.await
.unwrap();
@@ -823,7 +821,6 @@ i8pcjGO+IZffvyZJVRWfVooBJmWWbPB1pueo3tx8w3+fcuzpxz+RLFKaPyqXO+dD
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn test_load_self_generate_secret() {
let t = TestContext::new().await;
add_pseudo_transport(&t, "alice@example.org").await.unwrap();
t.set_config(Config::ConfiguredAddr, Some("alice@example.org"))
.await
.unwrap();
@@ -836,7 +833,6 @@ i8pcjGO+IZffvyZJVRWfVooBJmWWbPB1pueo3tx8w3+fcuzpxz+RLFKaPyqXO+dD
use std::thread;
let t = TestContext::new().await;
add_pseudo_transport(&t, "alice@example.org").await.unwrap();
t.set_config(Config::ConfiguredAddr, Some("alice@example.org"))
.await
.unwrap();
+2 -2
View File
@@ -52,7 +52,7 @@ const KEYUPDATE_CHUNK_CONTACTS: usize = 200;
const KEYUPDATE_MAX_SILENCE: i64 = 3 * 365 * 24 * 3600;
/// Upper bound on the contacts informed after a relay list change, keeping the freshest.
const KEYUPDATE_MAX_RECIPIENTS: u32 = 5000;
const KEYUPDATE_MAX_RECIPIENTS: usize = 5000;
/// A contact to inform: the relays to reach them at, and the key to encrypt to.
struct KeyupdateRecipient {
@@ -63,7 +63,7 @@ struct KeyupdateRecipient {
/// Returns at most `max_recipients` key-contacts to inform.
async fn keyupdate_recipients(
context: &Context,
max_recipients: u32,
max_recipients: usize,
) -> Result<Vec<KeyupdateRecipient>> {
// Single chat contacts only become keyupdate recipient candidates
// if we have a record of a sent message or `last_seen` is not 0.
+1 -1
View File
@@ -97,7 +97,7 @@ pub mod stock_str;
pub mod storage_usage;
mod sync;
mod token;
pub mod transport;
mod transport;
mod update_helper;
pub mod webxdc;
#[macro_use]
+95
View File
@@ -9,12 +9,14 @@
use std::fmt;
use anyhow::{Context as _, Result};
use num_traits::ToPrimitive as _;
use serde::{Deserialize, Serialize};
use crate::config::Config;
use crate::context::Context;
pub use crate::net::proxy::ProxyConfig;
pub use crate::provider::Socket;
use crate::tools::ToOption;
/// User-entered setting for certificate checks.
///
@@ -230,6 +232,62 @@ impl EnteredLoginParam {
oauth2: false,
})
}
/// Saves entered account settings,
/// so that they can be prefilled if the user wants to configure the server again.
///
/// This is needed in case a UI is not yet updated, and still uses `get_config("mail_pw")` etc.
/// in order to prefill the entered account settings.
pub(crate) async fn save_legacy(&self, context: &Context) -> Result<()> {
context.set_config(Config::Addr, Some(&self.addr)).await?;
context
.set_config(Config::MailServer, self.imap.server.to_option())
.await?;
context
.set_config(Config::MailPort, self.imap.port.to_option().as_deref())
.await?;
context
.set_config(
Config::MailSecurity,
self.imap.security.to_i32().to_option().as_deref(),
)
.await?;
context
.set_config(Config::MailUser, self.imap.user.to_option())
.await?;
context
.set_config(Config::MailPw, self.imap.password.to_option())
.await?;
context
.set_config(Config::SendServer, self.smtp.server.to_option())
.await?;
context
.set_config(Config::SendPort, self.smtp.port.to_option().as_deref())
.await?;
context
.set_config(
Config::SendSecurity,
self.smtp.security.to_i32().to_option().as_deref(),
)
.await?;
context
.set_config(Config::SendUser, self.smtp.user.to_option())
.await?;
context
.set_config(Config::SendPw, self.smtp.password.to_option())
.await?;
context
.set_config(
Config::ImapCertificateChecks,
self.certificate_checks.to_i32().to_option().as_deref(),
)
.await?;
Ok(())
}
}
impl fmt::Display for EnteredLoginParam {
@@ -311,4 +369,41 @@ mod tests {
Ok(())
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn test_save_entered_login_param() -> Result<()> {
let t = TestContext::new().await;
let param = EnteredLoginParam {
addr: "alice@example.org".to_string(),
imap: EnteredImapLoginParam {
server: "".to_string(),
port: 0,
folder: "".to_string(),
security: Socket::Starttls,
user: "".to_string(),
password: "foobar".to_string(),
},
smtp: EnteredSmtpLoginParam {
server: "".to_string(),
port: 2947,
security: Socket::default(),
user: "".to_string(),
password: "".to_string(),
},
certificate_checks: Default::default(),
oauth2: false,
};
param.save_legacy(&t).await?;
assert_eq!(
t.get_config(Config::Addr).await?.unwrap(),
"alice@example.org"
);
assert_eq!(t.get_config(Config::MailPw).await?.unwrap(), "foobar");
assert_eq!(t.get_config(Config::SendPw).await?, None);
assert_eq!(t.get_config_int(Config::SendPort).await?, 2947);
assert_eq!(EnteredLoginParam::load_legacy(&t).await?, param);
Ok(())
}
}
+13 -24
View File
@@ -1445,10 +1445,13 @@ impl std::fmt::Display for MessageState {
}
impl MessageState {
/// Returns true if the message can be marked as failed in the current state.
/// Returns true if the message can transition to `OutFailed` state from the current state.
pub fn can_fail(self) -> bool {
use MessageState::*;
matches!(self, OutPending | OutDelivered | OutFailed)
matches!(
self,
OutPending | OutDelivered | OutMdnRcvd // OutMdnRcvd can still fail because it could be a group message and only some recipients failed.
)
}
/// Returns true for any outgoing message states.
@@ -1641,15 +1644,11 @@ pub(crate) async fn delete_msgs_locally_done(
context: &Context,
msg_ids: &[MsgId],
modified_chat_ids: BTreeSet<ChatId>,
pinned_messages_changed_chat_ids: BTreeSet<ChatId>,
) -> Result<()> {
for modified_chat_id in modified_chat_ids {
context.emit_msgs_changed_without_msg_id(modified_chat_id);
chatlist_events::emit_chatlist_item_changed(context, modified_chat_id);
}
for chat_id in pinned_messages_changed_chat_ids {
context.emit_event(EventType::PinnedMessagesChanged { chat_id });
}
if !msg_ids.is_empty() {
context.emit_msgs_changed_without_ids();
chatlist_events::emit_chatlist_changed(context);
@@ -1675,7 +1674,6 @@ pub async fn delete_msgs_ext(
delete_for_all: bool,
) -> Result<()> {
let mut modified_chat_ids = BTreeSet::new();
let mut pinned_messages_changed_chat_ids = BTreeSet::new();
let mut deleted_rfc724_mid = Vec::new();
let mut res = Ok(());
@@ -1691,9 +1689,6 @@ pub async fn delete_msgs_ext(
);
modified_chat_ids.insert(msg.chat_id);
if msg.is_pinned() {
pinned_messages_changed_chat_ids.insert(msg.chat_id);
}
deleted_rfc724_mid.push(msg.rfc724_mid.clone());
let update_db = |trans: &mut rusqlite::Transaction| {
@@ -1751,13 +1746,7 @@ pub async fn delete_msgs_ext(
let msg = Message::load_from_db(context, msg_id).await?;
delete_msg_locally(context, &msg).await?;
}
delete_msgs_locally_done(
context,
msg_ids,
modified_chat_ids,
pinned_messages_changed_chat_ids,
)
.await?;
delete_msgs_locally_done(context, msg_ids, modified_chat_ids).await?;
// Interrupt Inbox loop to start message deletion, run housekeeping and call send_sync_msg().
context.scheduler.interrupt_inbox().await;
@@ -1986,15 +1975,15 @@ pub(crate) async fn set_msg_failed(
msg: &mut Message,
error: &str,
) -> Result<()> {
if !msg.state.can_fail() {
info!(
if msg.state.can_fail() {
msg.state = MessageState::OutFailed;
warn!(context, "{} failed: {}", msg.id, error);
} else {
warn!(
context,
"Ignoring failed {} in state {}: {}", msg.id, msg.state, error
);
return Ok(());
"{} seems to have failed ({}), but state is {}", msg.id, error, msg.state
)
}
msg.state = MessageState::OutFailed;
warn!(context, "{} failed: {}", msg.id, error);
msg.error = Some(error.to_string());
let exists = context
-46
View File
@@ -454,28 +454,6 @@ async fn test_get_state() -> Result<()> {
Ok(())
}
/// Tests that a failure reported after a read receipt leaves the message untouched.
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn test_set_msg_failed_after_mdn() -> Result<()> {
let mut tcm = TestContextManager::new();
let alice = &tcm.alice().await;
let bob = &tcm.bob().await;
let alice_chat = alice.create_chat(bob).await;
let sent = alice.send_text(alice_chat.id, "hi").await;
let bob_msg = bob.recv_msg(&sent).await;
alice.recv_mdn(bob, &bob_msg).await?;
let mut msg = sent.load_from_db().await;
assert_eq!(msg.state, MessageState::OutMdnRcvd);
set_msg_failed(alice, &mut msg, "relay bounced").await?;
let msg = sent.load_from_db().await;
assert_eq!(msg.state, MessageState::OutMdnRcvd);
assert_eq!(msg.error(), None);
let chats = Chatlist::try_load(alice, 0, None, None).await?;
assert_eq!(chats.get_msg_id(0)?, Some(msg.id));
Ok(())
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn test_is_bot() -> Result<()> {
let mut tcm = TestContextManager::new();
@@ -653,10 +631,6 @@ async fn test_delete_msgs_offline() -> Result<()> {
delete_msgs(alice, &[msg.id]).await?;
assert!(!alice.sql.exists(stmt, (msg.id,)).await?);
alice
.assert_warn("No SMTP connection candidates provided")
.await;
Ok(())
}
@@ -782,23 +756,3 @@ async fn test_get_existing_msg_ids() -> Result<()> {
Ok(())
}
#[test]
fn test_can_fail() -> Result<()> {
use MessageState::*;
// states that are not allowed to transition to OutFailed
assert!(!Undefined.can_fail());
assert!(!InFresh.can_fail());
assert!(!InNoticed.can_fail());
assert!(!InSeen.can_fail());
assert!(!OutDraft.can_fail());
assert!(!OutMdnRcvd.can_fail());
// states that are allowed to transition to OutFailed
assert!(OutPending.can_fail());
assert!(OutDelivered.can_fail());
assert!(OutFailed.can_fail());
Ok(())
}
+107 -10
View File
@@ -17,7 +17,7 @@ use tokio::fs;
use crate::aheader::{Aheader, EncryptPreference};
use crate::blob::BlobObject;
use crate::chat::{self, Chat, PARAM_BROADCAST_SECRET, load_broadcast_secret};
use crate::chat::{self, Chat, ChatId, PARAM_BROADCAST_SECRET, load_broadcast_secret};
use crate::config::Config;
use crate::constants::{Chattype, DC_FROM_HANDSHAKE};
use crate::contact::{Contact, ContactId, Origin};
@@ -36,9 +36,6 @@ use crate::param::Param;
use crate::peer_channels::{create_iroh_header, get_iroh_topic_for_msg};
use crate::pgp::{SeipdVersion, addresses_from_public_key, pubkey_supports_seipdv2, relay_addrs};
use crate::simplify::escape_message_footer_marks;
use crate::smtp::queue::{
Encryption as QueuedEncryption, QueuedMail, SideEffects as QueueSideEffects, ToBeQueuedMail,
};
use crate::stock_str;
use crate::tools::{IsNoneOrEmpty, create_outgoing_rfc724_mid, remove_subject_prefix, time};
use crate::webxdc::StatusUpdateSerial;
@@ -224,6 +221,110 @@ pub struct RenderedMessage {
sync_ids_to_delete: Option<String>,
}
#[derive(Debug, Clone)]
pub(crate) enum QueuedEncryption {
/// Unencrypted message.
No,
/// The message is encrypted asymmetrically to public keys.
Asymmetric {
/// OpenPGP keys to use for encryption.
///
/// The message is always encrypted to self,
/// no need to include own key here.
encryption_pubkeys: Vec<SignedPublicKey>,
},
/// Symmetrically encrypted message with a shared secret.
Symmetric { shared_secret: String },
}
impl QueuedEncryption {
pub(crate) fn is_encrypted(&self) -> bool {
match self {
Self::No => false,
Self::Asymmetric { .. } => true,
Self::Symmetric { .. } => true,
}
}
}
/// Email message queued, but not sent yet.
///
/// It is stored unencrypted to
/// make it possible to change protected headers
/// like the From address and Autocrypt header later.
#[derive(Debug, Clone)]
pub(crate) struct QueuedMail {
/// Unencrypted queued message.
///
/// This message has both the headers and the body,
/// but without the From, Autocrypt and Message-ID headers.
///
/// For encrypted messages this is the OpenPGP payload.
pub(crate) raw_message: Vec<u8>,
/// Display name to put in the `From:` field.
///
/// Email address is not determined yet here.
pub(crate) display_name: String,
/// Message-ID.
pub(crate) rfc724_mid: String,
/// Whether the message is encrypted and encryption keys.
pub(crate) encryption: QueuedEncryption,
/// If true, Autocrypt header should be added before sending.
pub(crate) should_attach_pubkey: bool,
/// If true, OpenPGP compression may be used.
pub(crate) should_compress: bool,
/// If true, encrypted message should be signed.
pub(crate) should_sign: bool,
/// Recipient addresses.
pub(crate) recipients: Vec<String>,
/// Addresses the messages was already sent to.
pub(crate) sent_to: Vec<String>,
/// If true, own addresses should be added to the list of recipients.
///
/// For unencrypted messages, only the sending addresses should be added.
/// For encrypted messages, all published addresses should be added.
pub(crate) bcc_self: bool,
}
/// Side effects that should be applied at the same time
/// as the message is persisted in the queue.
#[derive(Debug, Clone, Default)]
pub struct QueueSideEffects {
/// ID of the chat side effects should be applied to.
pub chat_id: ChatId,
/// Largest timestamp of the location sent in `location.kml` in this message.
pub last_added_location_timestamp: Option<i64>,
/// True if the message has the avatar attached.
///
/// Timestamp of the last time avatar was gossiped should be updated.
pub avatar_is_attached: bool,
/// A comma-separated string of sync-IDs that are used by the rendered email and must be deleted
/// from `multi_device_sync` once the message is actually queued for sending.
pub sync_ids_to_delete: Option<String>,
/// Subject that was rendered into the message.
///
/// Used to update the subject on the sent message object.
pub subject: String,
}
/// Email message ready to be queued with the side effects that should be applied at the same time.
pub(crate) type ToBeQueuedMail = (QueuedMail, Option<QueueSideEffects>);
/// Renders [`QueuedMail`].
///
/// Adds headers:
@@ -447,7 +548,6 @@ pub(crate) fn render_queued_mail(
}
/// Renders queued mail with the current sending address.
#[cfg(test)]
pub(crate) async fn render_queued_mail_with_context(
queued_mail: QueuedMail,
context: &Context,
@@ -1320,16 +1420,13 @@ impl MimeFactory {
///
/// Used for MDNs because they are fully rendered and sent in one go,
/// rather than first creating a [`QueuedMail`] and sending it later.
pub async fn render(self, context: &Context, from_addr: &str) -> Result<RenderedEmail> {
pub async fn render(self, context: &Context) -> Result<RenderedEmail> {
// Does not matter, we are not going to return the QueuedMail.
let bcc_self = false;
let public_key = key::load_self_public_key(context).await?;
let secret_key = key::load_self_secret_key(context).await?;
let (queued_mail, _side_effects) =
Box::pin(self.into_queued_mail(context, bcc_self)).await?;
let rendered_mail =
render_queued_mail(queued_mail, &public_key, &secret_key, from_addr.to_string())?;
let rendered_mail = render_queued_mail_with_context(queued_mail, context).await?;
Ok(rendered_mail)
}
+25 -149
View File
@@ -22,7 +22,7 @@ use crate::key::{load_self_secret_key, secret_key_to_public_key};
use crate::message;
use crate::mimeparser::MimeMessage;
use crate::receive_imf::receive_imf;
use crate::test_utils::{self, SentMessage};
use crate::test_utils;
use crate::test_utils::{TestContext, TestContextManager, get_chat_msg};
use crate::tools::SystemTime;
@@ -307,8 +307,7 @@ async fn test_mdn_create_encrypted() -> Result<()> {
let mimefactory =
MimeFactory::from_mdn(&bob, rcvd.from_id, rcvd.rfc724_mid.clone(), vec![]).await?;
assert!(!mimefactory.will_be_encrypted());
let bob_addr = bob.get_primary_self_addr().await?;
let rendered_msg = mimefactory.render(&bob, &bob_addr).await?;
let rendered_msg = mimefactory.render(&bob).await?;
assert!(!rendered_msg.message.contains("Bob Examplenet"));
assert!(!rendered_msg.message.contains("Alice Exampleorg"));
@@ -321,8 +320,7 @@ async fn test_mdn_create_encrypted() -> Result<()> {
let mimefactory = MimeFactory::from_mdn(&bob, rcvd.from_id, rcvd.rfc724_mid, vec![]).await?;
assert!(mimefactory.will_be_encrypted());
let bob_addr = bob.get_primary_self_addr().await?;
let rendered_msg = mimefactory.render(&bob, &bob_addr).await?;
let rendered_msg = mimefactory.render(&bob).await?;
assert!(!rendered_msg.message.contains("Bob Examplenet"));
assert!(!rendered_msg.message.contains("Alice Exampleorg"));
@@ -365,8 +363,7 @@ async fn test_mdn_autocrypt_throttle() -> Result<()> {
rcvd: &Message,
) -> Result<bool> {
let mf = MimeFactory::from_mdn(bob, rcvd.from_id, rcvd.rfc724_mid.clone(), vec![]).await?;
let addr = bob.get_primary_self_addr().await?;
let rendered_msg = mf.render(bob, &addr).await?;
let rendered_msg = mf.render(bob).await?;
let mime = MimeMessage::from_bytes(alice, rendered_msg.message.as_bytes()).await?;
Ok(mime.autocrypt_fingerprint.is_some())
}
@@ -647,8 +644,7 @@ async fn test_render_reply() {
let recipients = mimefactory.recipients();
assert_eq!(recipients, vec!["charlie@example.net"]);
let addr = t.get_primary_self_addr().await.unwrap();
let rendered_msg = mimefactory.render(t, &addr).await.unwrap();
let rendered_msg = mimefactory.render(t).await.unwrap();
let mail = mailparse::parse_mail(rendered_msg.message.as_bytes()).unwrap();
assert_eq!(
@@ -1022,7 +1018,7 @@ END:VCARD";
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn test_render_outer_headers_of_encrypted_msg() -> Result<()> {
async fn test_render_outer_headers() -> Result<()> {
let mut tcm = TestContextManager::new();
let alice = &tcm.alice().await;
let bob = &tcm.bob().await;
@@ -1030,9 +1026,26 @@ async fn test_render_outer_headers_of_encrypted_msg() -> Result<()> {
let chat_id = alice.create_chat_id(bob).await;
let sent = alice.send_text(chat_id, "Hello!").await;
let payload = normalized_payload(sent).await;
let (unencrypted, _encrypted) = sent
.payload()
.split_once("-----BEGIN PGP MESSAGE-----")
.unwrap();
let (unencrypted, _encrypted) = payload.split_once("-----BEGIN PGP MESSAGE-----").unwrap();
// Normalize the parts of the message that vary between runs
// (MIME boundary, Date, Message-ID)
let boundary = unencrypted
.split_once("boundary=\"")
.and_then(|(_, rest)| rest.split_once('"'))
.map(|(b, _)| b)
.unwrap_or_default();
let unencrypted = unencrypted.replace(boundary, "BOUNDARY");
let rfc724_mid = sent.load_from_db().await.rfc724_mid;
let unencrypted = unencrypted.replace(&rfc724_mid, "MESSAGE_ID@localhost");
let unencrypted = regex!(r"Date:[^\r\n]*")
.replace(&unencrypted, "Date: DATE")
.to_string();
let expected = r#"From: <alice@example.org>
Date: DATE
@@ -1068,140 +1081,3 @@ expected (debug print): {expected:?}"
Ok(())
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn test_render_unencrypted_msg_basic() -> Result<()> {
let alice = &TestContext::new_alice().await;
alice.allow_unencrypted().await?;
let chat = alice
.create_chat_with_contact("Bob", "bob@example.net")
.await;
let sent = alice.send_text(chat.id, "Hello!").await;
let unencrypted = normalized_payload(sent).await;
let expected = r#"From: <alice@example.org>
Message-ID: <MESSAGE_ID@localhost>
MIME-Version: 1.0
Autocrypt: addr=alice@example.org; prefer-encrypt=mutual; keydata=mDMEXlh13RYJKwYBBAHaRw8BAQdAzfVIAleCXMJrq8VeLlEVof6ITCviMktKjmcBKAu4m5
DCtAQfFggAZgUCXlh13RYhBC5vossjtTLXKGNLWGSwj2Gp7ZRDAhsDAh4JBAsJCAcFFQgJCgsDFgIB
AycJAgIZASwUgAAAAAASABFyZWxheXNAY2hhdG1haWwuYXRhbGljZUBleGFtcGxlLm9yZwAAb1QA/0
HbvPN3/Vn02Gk1dcQMEcyGyETld9dSsRo8uwHAyW35AQCrFJjAQFLTud7XK61uYt9BC/QHipCfIGbq
X1FjMbTUC80TPGFsaWNlQGV4YW1wbGUub3JnPsKRBBMWCAA5BQJeWHXdFiEELm+iyyO1MtcoY0tYZL
CPYantlEMCGwMCHgkECwkIBwUVCAkKCwMWAgEDJwkCAhkBAAoJEGSwj2Gp7ZRD1m4A/iOifEzIOiP8
wW0O8I/sg69gQtG8Czn4MsVV6Ea1EyIqAP4uByHaUJdy8MSQPfv/Usr09KsidNgy2Jh37yg82fKUBr
g4BF5Ydd0SCisGAQQBl1UBBQEBB0AG7cjWy2SFAU8KnltlubVW67rFiyfp01JrRe6Xqy22HQMBCAeI
eAQYFggAIBYhBC5vossjtTLXKGNLWGSwj2Gp7ZRDBQJeWHXdAhsMAAoJEGSwj2Gp7ZRDLo8BAObE8G
nsGVwKzNqCvHeWgJsqhjS3C6gvSlV3tEm9XmF6AQDXucIyVfoBwoyMh2h6cSn/ATn5QJb35pgo+ivp
3jsMAg==
Content-Type: text/plain; charset="utf-8"
Date: DATE
To: <bob@example.net>
Subject: Message from alice@example.org
References: <MESSAGE_ID@localhost>
Chat-Version: 1.0
Chat-Disposition-Notification-To: alice@example.org
Content-Transfer-Encoding: 7bit
Hello!"#
.replace("\n", "\r\n");
assert_eq!(
unencrypted, expected,
"---------------- Actual: ----------------
{unencrypted}
-----------------------------------------
actual (debug print): {unencrypted:?}
expected (debug print): {expected:?}"
);
Ok(())
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn test_render_unencrypted_msg_with_attachment() -> Result<()> {
let alice = &TestContext::new_alice().await;
alice.allow_unencrypted().await?;
let chat = alice
.create_chat_with_contact("Bob", "bob@example.net")
.await;
let mut msg = Message::new(Viewtype::File);
msg.set_text("Hello!".to_string());
msg.set_file_from_bytes(alice, "foo.bar", "content".as_bytes(), None)?;
let sent = alice.send_msg(chat.id, &mut msg).await;
let unencrypted = normalized_payload(sent).await;
let expected = r#"From: <alice@example.org>
Message-ID: <MESSAGE_ID@localhost>
MIME-Version: 1.0
Autocrypt: addr=alice@example.org; prefer-encrypt=mutual; keydata=mDMEXlh13RYJKwYBBAHaRw8BAQdAzfVIAleCXMJrq8VeLlEVof6ITCviMktKjmcBKAu4m5
DCtAQfFggAZgUCXlh13RYhBC5vossjtTLXKGNLWGSwj2Gp7ZRDAhsDAh4JBAsJCAcFFQgJCgsDFgIB
AycJAgIZASwUgAAAAAASABFyZWxheXNAY2hhdG1haWwuYXRhbGljZUBleGFtcGxlLm9yZwAAb1QA/0
HbvPN3/Vn02Gk1dcQMEcyGyETld9dSsRo8uwHAyW35AQCrFJjAQFLTud7XK61uYt9BC/QHipCfIGbq
X1FjMbTUC80TPGFsaWNlQGV4YW1wbGUub3JnPsKRBBMWCAA5BQJeWHXdFiEELm+iyyO1MtcoY0tYZL
CPYantlEMCGwMCHgkECwkIBwUVCAkKCwMWAgEDJwkCAhkBAAoJEGSwj2Gp7ZRD1m4A/iOifEzIOiP8
wW0O8I/sg69gQtG8Czn4MsVV6Ea1EyIqAP4uByHaUJdy8MSQPfv/Usr09KsidNgy2Jh37yg82fKUBr
g4BF5Ydd0SCisGAQQBl1UBBQEBB0AG7cjWy2SFAU8KnltlubVW67rFiyfp01JrRe6Xqy22HQMBCAeI
eAQYFggAIBYhBC5vossjtTLXKGNLWGSwj2Gp7ZRDBQJeWHXdAhsMAAoJEGSwj2Gp7ZRDLo8BAObE8G
nsGVwKzNqCvHeWgJsqhjS3C6gvSlV3tEm9XmF6AQDXucIyVfoBwoyMh2h6cSn/ATn5QJb35pgo+ivp
3jsMAg==
Content-Type: multipart/mixed;
boundary="BOUNDARY"
Date: DATE
To: <bob@example.net>
Subject: Message from alice@example.org
References: <MESSAGE_ID@localhost>
Chat-Version: 1.0
Chat-Disposition-Notification-To: alice@example.org
--BOUNDARY
Content-Type: text/plain; charset="utf-8"
Content-Transfer-Encoding: 7bit
Hello!
--BOUNDARY
Content-Type: application/octet-stream
Content-Disposition: attachment; filename="foo.bar"
Content-Transfer-Encoding: base64
Y29udGVudA==
--BOUNDARY--
"#
.replace("\n", "\r\n");
assert_eq!(
unencrypted, expected,
"---------------- Actual: ----------------
{unencrypted}
-----------------------------------------
actual (debug print): {unencrypted:?}
expected (debug print): {expected:?}"
);
Ok(())
}
/// Normalize the parts of the message that vary between runs
/// (MIME boundary, Date, Message-ID)
async fn normalized_payload(sent: SentMessage<'_>) -> String {
let rfc724_mid = sent.load_from_db().await.rfc724_mid;
let mut payload = sent.payload;
if let Some(boundary) = payload
.split_once("boundary=\"")
.and_then(|(_, rest)| rest.split_once('"'))
.map(|(b, _)| b)
{
payload = payload.replace(boundary, "BOUNDARY");
}
payload = payload.replace(&rfc724_mid, "MESSAGE_ID@localhost");
payload = regex!(r"Date:[^\r\n]*")
.replace(&payload, "Date: DATE")
.to_string();
payload
}
+12 -11
View File
@@ -2540,18 +2540,19 @@ async fn handle_ndn(
for msg_id in msg_ids {
let mut message = Message::load_from_db(context, msg_id).await?;
let chat = Chat::load_from_db(context, message.chat_id).await?;
if chat.typ == constants::Chattype::Single {
let aggregated_error = message
.error
.as_ref()
.map(|err| format!("{err}\n\n{err_msg}"));
set_msg_failed(
context,
&mut message,
aggregated_error.as_ref().unwrap_or(err_msg),
)
.await?;
if chat.typ == constants::Chattype::OutBroadcast {
continue;
}
let aggregated_error = message
.error
.as_ref()
.map(|err| format!("{err}\n\n{err_msg}"));
set_msg_failed(
context,
&mut message,
aggregated_error.as_ref().unwrap_or(err_msg),
)
.await?;
}
Ok(())
+1 -1
View File
@@ -33,7 +33,7 @@ use tls::wrap_tls;
pub(crate) const TIMEOUT: Duration = Duration::from_secs(60);
/// TTL for caches in seconds.
pub(crate) const CACHE_TTL: u32 = 30 * 24 * 60 * 60;
pub(crate) const CACHE_TTL: u64 = 30 * 24 * 60 * 60;
/// Removes connection history entries after `CACHE_TTL`.
pub(crate) async fn prune_connection_history(context: &Context) -> Result<()> {
-4
View File
@@ -2,7 +2,6 @@ use std::sync::LazyLock;
use tokio::sync::OnceCell;
use super::*;
use crate::transport::add_pseudo_transport;
use crate::{
config::Config,
decrypt,
@@ -20,9 +19,6 @@ async fn decrypt_bytes(
auth_tokens_for_decryption: &[String],
) -> Result<pgp::composed::Message<'static>> {
let t = &TestContext::new().await;
add_pseudo_transport(t, "alice@example.org")
.await
.expect("Failed to add pseudo transport");
t.set_config(Config::ConfiguredAddr, Some("alice@example.org"))
.await
.expect("Failed to configure address");
+223 -6
View File
@@ -13,7 +13,6 @@ use anyhow::{Result, ensure};
use crate::chat::{ChatId, send_msg};
use crate::contact::ContactId;
use crate::context::Context;
use crate::events::EventType;
use crate::log::warn;
use crate::message::{Message, MessageState, MsgId, Viewtype};
use crate::mimeparser::SystemMessage;
@@ -79,11 +78,7 @@ async fn update_pinned_state_in_db(
(new_pinned_state, msg.id),
)
.await?;
context.emit_msgs_changed(msg.chat_id, msg.id);
context.emit_event(EventType::PinnedMessagesChanged {
chat_id: msg.chat_id,
});
Ok(())
}
@@ -145,4 +140,226 @@ pub(crate) async fn handle_pinned_state_from_wire(
}
#[cfg(test)]
mod pinned_messages_tests;
mod tests {
use super::*;
use crate::chat::{ChatItem, add_info_msg, create_broadcast, get_chat_msgs};
use crate::config::Config;
use crate::securejoin::get_securejoin_qr;
use crate::test_utils::{TestContextManager, sync};
use std::time::Duration;
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn test_pinned_messages() -> Result<()> {
let mut tcm = TestContextManager::new();
let alice = &tcm.alice().await;
let alice2 = &tcm.alice().await; // Alice's second device
let bob = &tcm.bob().await;
alice.set_config_bool(Config::SyncMsgs, true).await?;
alice2.set_config_bool(Config::SyncMsgs, true).await?;
// Alice creates all chat types upfront, with Bob as member if possible
let single_chat_id = alice.create_chat(bob).await.id;
let group_chat_id = alice.create_group_with_members("Group", &[bob]).await;
let broadcast_chat_id = create_broadcast(alice, "Channel".to_string()).await?;
let qr = get_securejoin_qr(alice, Some(broadcast_chat_id)).await?;
tcm.exec_securejoin_qr(bob, alice, &qr).await;
let self_chat_id = alice.get_self_chat().await.id;
sync(alice, alice2).await;
for alice_chat_id in [
single_chat_id,
group_chat_id,
broadcast_chat_id,
self_chat_id,
] {
let pinned = get_pinned_messages(alice, alice_chat_id).await?;
assert!(pinned.is_empty());
// Alice sends message "Foo" and pins it
let sent1 = alice.send_text(alice_chat_id, "Foo").await;
let msg1 = sent1.load_from_db().await;
assert!(!msg1.is_pinned());
set_pinned_state(alice, msg1.id, true).await?;
let sent2 = alice.pop_sent_msg().await;
assert!(sent1.load_from_db().await.is_pinned());
let info_msg = sent2.load_from_db().await;
assert!(info_msg.is_info());
assert!(!info_msg.hidden);
assert_eq!(info_msg.get_info_type(), SystemMessage::MessagePinned);
assert!(info_msg.get_info_contact_id(alice).await?.is_none()); // contact not needed, tapping shall jump to message
let pinned = get_pinned_messages(alice, alice_chat_id).await?;
assert_eq!(pinned.len(), 1);
assert_eq!(pinned[0], msg1.id);
// Pinning an info message does not work
assert!(set_pinned_state(alice, info_msg.id, true).await.is_err());
// Unpin the initially pinned message.
// Before, send another message "Bar". To test, no visible info message is added this time,
let sent3 = alice.send_text(alice_chat_id, "Bar").await;
set_pinned_state(alice, msg1.id, false).await?;
let sent4 = alice.pop_sent_msg().await;
assert!(!sent1.load_from_db().await.is_pinned());
let pinned = get_pinned_messages(alice, alice_chat_id).await?;
assert!(pinned.is_empty());
let msg3 = sent3.load_from_db().await;
assert!(!msg3.is_info());
assert!(!msg3.is_pinned());
assert_eq!(alice.get_last_msg_id_in(msg3.chat_id).await, msg3.id); // last message is still "Bar", not an info message
if alice_chat_id != self_chat_id {
// Bob receives message "Foo"
let msg1 = bob.recv_msg(&sent1).await;
assert!(!msg1.is_pinned());
let pinned = get_pinned_messages(bob, msg1.chat_id).await?;
assert!(pinned.is_empty());
// Bob receives info message to pin "Foo"
bob.recv_msg(&sent2).await;
assert!(Message::load_from_db(bob, msg1.id).await?.is_pinned());
let pinned = get_pinned_messages(bob, msg1.chat_id).await?;
assert_eq!(pinned.len(), 1);
assert_eq!(pinned[0], msg1.id);
let info_msg =
Message::load_from_db(bob, bob.get_last_msg_id_in(msg1.chat_id).await).await?;
assert!(info_msg.is_info());
assert!(!info_msg.hidden);
assert_eq!(info_msg.get_info_type(), SystemMessage::MessagePinned);
assert!(info_msg.get_info_contact_id(bob).await?.is_none());
// Bob receives message "Bar" and hidden message to unpin message "Foo"
bob.recv_msg(&sent3).await;
bob.recv_msg_trash(&sent4).await;
assert!(!Message::load_from_db(bob, msg1.id).await?.is_pinned());
let pinned = get_pinned_messages(bob, msg1.chat_id).await?;
assert!(pinned.is_empty());
let no_info_msg =
Message::load_from_db(bob, bob.get_last_msg_id_in(msg1.chat_id).await).await?;
assert!(!no_info_msg.is_info());
assert_eq!(no_info_msg.text, "Bar");
}
// Alice's second device receives all four messages and ends up in the same state
let msg1 = alice2.recv_msg(&sent1).await;
alice2.recv_msg(&sent2).await;
assert!(Message::load_from_db(alice2, msg1.id).await?.is_pinned());
alice2.recv_msg(&sent3).await;
alice2.recv_msg_trash(&sent4).await;
assert!(!Message::load_from_db(alice2, msg1.id).await?.is_pinned());
let no_info_msg =
Message::load_from_db(alice2, alice2.get_last_msg_id_in(msg1.chat_id).await)
.await?;
assert!(!no_info_msg.is_info());
assert_eq!(no_info_msg.text, "Bar");
}
Ok(())
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn test_get_pinned_messages_order() -> Result<()> {
let mut tcm = TestContextManager::new();
let alice = &tcm.alice().await;
let chat_id = alice.get_self_chat().await.id;
let boilerplate_msg_count = get_chat_msgs(alice, chat_id).await?.len();
// create three messages, sent1 and sent2 have different timestamp, sent2 and sent3 may differ by ID only
let sent1 = alice.send_text(chat_id, "1").await;
tokio::time::sleep(Duration::from_millis(1100)).await;
let sent2 = alice.send_text(chat_id, "2").await;
let sent3 = alice.send_text(chat_id, "3").await;
// get_chat_msgs() start with the oldest message
let chat_msgs = get_chat_msgs(alice, chat_id).await?;
let msg_ids: Vec<_> = chat_msgs
.into_iter()
.filter_map(|item| match item {
ChatItem::Message { msg_id } => Some(msg_id),
ChatItem::DayMarker { .. } => None,
})
.collect();
assert_eq!(
&msg_ids[boilerplate_msg_count..],
&[
sent1.sender_msg_id,
sent2.sender_msg_id,
sent3.sender_msg_id
]
);
// get_pinned_messages() has the same order, also starting with the oldest message
set_pinned_state(alice, sent1.sender_msg_id, true).await?;
set_pinned_state(alice, sent2.sender_msg_id, true).await?;
set_pinned_state(alice, sent3.sender_msg_id, true).await?;
let pinned = get_pinned_messages(alice, chat_id).await?;
assert_eq!(pinned.len(), 3);
assert_eq!(pinned[0], sent1.sender_msg_id);
assert_eq!(pinned[1], sent2.sender_msg_id);
assert_eq!(pinned[2], sent3.sender_msg_id);
// order of pinning does not affect the order of pinned messages.
// this is to keep scrolling direction of chat bubbles and pinned banner scrollbar in sync,
// and not jumping wildly around.
// this is also what most other messengers are doing.
set_pinned_state(alice, sent1.sender_msg_id, false).await?;
set_pinned_state(alice, sent2.sender_msg_id, false).await?;
set_pinned_state(alice, sent3.sender_msg_id, false).await?;
let pinned = get_pinned_messages(alice, chat_id).await?;
assert_eq!(pinned.len(), 0);
set_pinned_state(alice, sent3.sender_msg_id, true).await?;
set_pinned_state(alice, sent2.sender_msg_id, true).await?;
set_pinned_state(alice, sent1.sender_msg_id, true).await?;
let pinned = get_pinned_messages(alice, chat_id).await?;
assert_eq!(pinned.len(), 3);
assert_eq!(pinned[0], sent1.sender_msg_id);
assert_eq!(pinned[1], sent2.sender_msg_id);
assert_eq!(pinned[2], sent3.sender_msg_id);
Ok(())
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn test_handle_pinned_state_from_wire() -> Result<()> {
let mut tcm = TestContextManager::new();
let alice = &tcm.alice().await;
let chat_id = alice.get_self_chat().await.id;
let sent1 = alice.send_text(chat_id, "pinnable").await;
let msg1 = sent1.load_from_db().await;
assert!(is_pinnable(&msg1));
assert!(
handle_pinned_state_from_wire(alice, &msg1, true)
.await
.is_ok()
);
// For not-pinnable messages, handle_pinned_state_from_wire() logs a warning and returns "ok".
// otherwise if there is an incompatibility in which messages are treated as "pinnable",
// this error will bubble up and user will get a device message saying "please report a bug".
let msg2_id = add_info_msg(alice, chat_id, "not pinnable").await?;
let msg2 = Message::load_from_db(alice, msg2_id).await?;
assert!(!is_pinnable(&msg2));
assert!(
handle_pinned_state_from_wire(alice, &msg2, true)
.await
.is_ok()
);
alice.assert_warn("Message is not pinnable").await;
Ok(())
}
}
@@ -1,309 +0,0 @@
use super::*;
use crate::chat::{ChatItem, add_info_msg, create_broadcast, get_chat_msgs};
use crate::config::Config;
use crate::ephemeral;
use crate::message;
use crate::securejoin::get_securejoin_qr;
use crate::test_utils::{TestContext, TestContextManager, sync};
use crate::tools::{SystemTime, time};
use std::time::Duration;
/// Waits for a PinnedMessagesChanged for a given `chat_id`.
///
/// Panics if event arrives for the wrong `chat_id`.
async fn expect_pinned_message_event(context: &TestContext, chat_id: ChatId) {
let EventType::PinnedMessagesChanged {
chat_id: event_chat_id,
} = context
.evtracker
.get_matching(|evt| matches!(evt, EventType::PinnedMessagesChanged { .. }))
.await
else {
unreachable!();
};
assert_eq!(event_chat_id, chat_id);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn test_pinned_messages() -> Result<()> {
let mut tcm = TestContextManager::new();
let alice = &tcm.alice().await;
let alice2 = &tcm.alice().await; // Alice's second device
let bob = &tcm.bob().await;
alice.set_config_bool(Config::SyncMsgs, true).await?;
alice2.set_config_bool(Config::SyncMsgs, true).await?;
// Alice creates all chat types upfront, with Bob as member if possible
let single_chat_id = alice.create_chat(bob).await.id;
let group_chat_id = alice.create_group_with_members("Group", &[bob]).await;
let broadcast_chat_id = create_broadcast(alice, "Channel".to_string()).await?;
let qr = get_securejoin_qr(alice, Some(broadcast_chat_id)).await?;
tcm.exec_securejoin_qr(bob, alice, &qr).await;
let self_chat_id = alice.get_self_chat().await.id;
sync(alice, alice2).await;
for alice_chat_id in [
single_chat_id,
group_chat_id,
broadcast_chat_id,
self_chat_id,
] {
let pinned = get_pinned_messages(alice, alice_chat_id).await?;
assert!(pinned.is_empty());
// Alice sends message "Foo" and pins it
let sent1 = alice.send_text(alice_chat_id, "Foo").await;
let msg1 = sent1.load_from_db().await;
assert!(!msg1.is_pinned());
set_pinned_state(alice, msg1.id, true).await?;
let sent2 = alice.pop_sent_msg().await;
expect_pinned_message_event(alice, msg1.chat_id).await;
assert!(sent1.load_from_db().await.is_pinned());
let info_msg = sent2.load_from_db().await;
assert!(info_msg.is_info());
assert!(!info_msg.hidden);
assert_eq!(info_msg.get_info_type(), SystemMessage::MessagePinned);
assert!(info_msg.get_info_contact_id(alice).await?.is_none()); // contact not needed, tapping shall jump to message
let pinned = get_pinned_messages(alice, alice_chat_id).await?;
assert_eq!(pinned.len(), 1);
assert_eq!(pinned[0], msg1.id);
// Pinning an info message does not work
assert!(set_pinned_state(alice, info_msg.id, true).await.is_err());
// Unpin the initially pinned message.
// Before, send another message "Bar". To test, no visible info message is added this time,
let sent3 = alice.send_text(alice_chat_id, "Bar").await;
set_pinned_state(alice, msg1.id, false).await?;
let sent4 = alice.pop_sent_msg().await;
assert!(!sent1.load_from_db().await.is_pinned());
expect_pinned_message_event(alice, msg1.chat_id).await;
let pinned = get_pinned_messages(alice, alice_chat_id).await?;
assert!(pinned.is_empty());
let msg3 = sent3.load_from_db().await;
assert!(!msg3.is_info());
assert!(!msg3.is_pinned());
assert_eq!(alice.get_last_msg_id_in(msg3.chat_id).await, msg3.id); // last message is still "Bar", not an info message
if alice_chat_id != self_chat_id {
// Bob receives message "Foo"
let msg1 = bob.recv_msg(&sent1).await;
assert!(!msg1.is_pinned());
let pinned = get_pinned_messages(bob, msg1.chat_id).await?;
assert!(pinned.is_empty());
// Bob receives info message to pin "Foo"
bob.recv_msg(&sent2).await;
expect_pinned_message_event(bob, msg1.chat_id).await;
assert!(Message::load_from_db(bob, msg1.id).await?.is_pinned());
let pinned = get_pinned_messages(bob, msg1.chat_id).await?;
assert_eq!(pinned.len(), 1);
assert_eq!(pinned[0], msg1.id);
let info_msg =
Message::load_from_db(bob, bob.get_last_msg_id_in(msg1.chat_id).await).await?;
assert!(info_msg.is_info());
assert!(!info_msg.hidden);
assert_eq!(info_msg.get_info_type(), SystemMessage::MessagePinned);
assert!(info_msg.get_info_contact_id(bob).await?.is_none());
// Bob receives message "Bar" and hidden message to unpin message "Foo"
bob.recv_msg(&sent3).await;
bob.recv_msg_trash(&sent4).await;
expect_pinned_message_event(bob, msg1.chat_id).await;
assert!(!Message::load_from_db(bob, msg1.id).await?.is_pinned());
let pinned = get_pinned_messages(bob, msg1.chat_id).await?;
assert!(pinned.is_empty());
let no_info_msg =
Message::load_from_db(bob, bob.get_last_msg_id_in(msg1.chat_id).await).await?;
assert!(!no_info_msg.is_info());
assert_eq!(no_info_msg.text, "Bar");
}
// Alice's second device receives all four messages and ends up in the same state
let msg1 = alice2.recv_msg(&sent1).await;
alice2.recv_msg(&sent2).await;
expect_pinned_message_event(alice2, msg1.chat_id).await;
assert!(Message::load_from_db(alice2, msg1.id).await?.is_pinned());
alice2.recv_msg(&sent3).await;
alice2.recv_msg_trash(&sent4).await;
expect_pinned_message_event(alice2, msg1.chat_id).await;
assert!(!Message::load_from_db(alice2, msg1.id).await?.is_pinned());
let no_info_msg =
Message::load_from_db(alice2, alice2.get_last_msg_id_in(msg1.chat_id).await).await?;
assert!(!no_info_msg.is_info());
assert_eq!(no_info_msg.text, "Bar");
}
Ok(())
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn test_get_pinned_messages_order() -> Result<()> {
let mut tcm = TestContextManager::new();
let alice = &tcm.alice().await;
let chat_id = alice.get_self_chat().await.id;
let boilerplate_msg_count = get_chat_msgs(alice, chat_id).await?.len();
// create three messages, sent1 and sent2 have different timestamp, sent2 and sent3 may differ by ID only
let sent1 = alice.send_text(chat_id, "1").await;
tokio::time::sleep(Duration::from_millis(1100)).await;
let sent2 = alice.send_text(chat_id, "2").await;
let sent3 = alice.send_text(chat_id, "3").await;
// get_chat_msgs() start with the oldest message
let chat_msgs = get_chat_msgs(alice, chat_id).await?;
let msg_ids: Vec<_> = chat_msgs
.into_iter()
.filter_map(|item| match item {
ChatItem::Message { msg_id } => Some(msg_id),
ChatItem::DayMarker { .. } => None,
})
.collect();
assert_eq!(
&msg_ids[boilerplate_msg_count..],
&[
sent1.sender_msg_id,
sent2.sender_msg_id,
sent3.sender_msg_id
]
);
// get_pinned_messages() has the same order, also starting with the oldest message
set_pinned_state(alice, sent1.sender_msg_id, true).await?;
set_pinned_state(alice, sent2.sender_msg_id, true).await?;
set_pinned_state(alice, sent3.sender_msg_id, true).await?;
let pinned = get_pinned_messages(alice, chat_id).await?;
assert_eq!(pinned.len(), 3);
assert_eq!(pinned[0], sent1.sender_msg_id);
assert_eq!(pinned[1], sent2.sender_msg_id);
assert_eq!(pinned[2], sent3.sender_msg_id);
// order of pinning does not affect the order of pinned messages.
// this is to keep scrolling direction of chat bubbles and pinned banner scrollbar in sync,
// and not jumping wildly around.
// this is also what most other messengers are doing.
set_pinned_state(alice, sent1.sender_msg_id, false).await?;
set_pinned_state(alice, sent2.sender_msg_id, false).await?;
set_pinned_state(alice, sent3.sender_msg_id, false).await?;
let pinned = get_pinned_messages(alice, chat_id).await?;
assert_eq!(pinned.len(), 0);
set_pinned_state(alice, sent3.sender_msg_id, true).await?;
set_pinned_state(alice, sent2.sender_msg_id, true).await?;
set_pinned_state(alice, sent1.sender_msg_id, true).await?;
let pinned = get_pinned_messages(alice, chat_id).await?;
assert_eq!(pinned.len(), 3);
assert_eq!(pinned[0], sent1.sender_msg_id);
assert_eq!(pinned[1], sent2.sender_msg_id);
assert_eq!(pinned[2], sent3.sender_msg_id);
Ok(())
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn test_handle_pinned_state_from_wire() -> Result<()> {
let mut tcm = TestContextManager::new();
let alice = &tcm.alice().await;
let chat_id = alice.get_self_chat().await.id;
let sent1 = alice.send_text(chat_id, "pinnable").await;
let msg1 = sent1.load_from_db().await;
assert!(is_pinnable(&msg1));
assert!(
handle_pinned_state_from_wire(alice, &msg1, true)
.await
.is_ok()
);
// For not-pinnable messages, handle_pinned_state_from_wire() logs a warning and returns "ok".
// otherwise if there is an incompatibility in which messages are treated as "pinnable",
// this error will bubble up and user will get a device message saying "please report a bug".
let msg2_id = add_info_msg(alice, chat_id, "not pinnable").await?;
let msg2 = Message::load_from_db(alice, msg2_id).await?;
assert!(!is_pinnable(&msg2));
assert!(
handle_pinned_state_from_wire(alice, &msg2, true)
.await
.is_ok()
);
alice.assert_warn("Message is not pinnable").await;
Ok(())
}
/// Tests that disappearing pinned message expires and emits `PinnedMessagesChanged` event.
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn test_ephemeral_pinned_message() -> Result<()> {
let mut tcm = TestContextManager::new();
let alice = &tcm.alice().await;
let bob = &tcm.bob().await;
let alice_chat_id = alice.create_chat_id(bob).await;
// Alice sends ephemeral timer in single chat with Bob.
alice_chat_id
.set_ephemeral_timer(alice, ephemeral::Timer::from_u32(60))
.await?;
let sent = alice.pop_sent_msg().await;
bob.recv_msg(&sent).await;
// Alice sends "Hello!" message to Bob.
let bob_msg = tcm.send_recv_accept(alice, bob, "Hello!").await;
let bob_chat_id = bob_msg.chat_id;
// Bob reads the message, so the timer starts.
message::markseen_msgs(bob, vec![bob_msg.id]).await?;
// Bob pins "Hello!" message received from Alice.
set_pinned_state(bob, bob_msg.id, true).await?;
expect_pinned_message_event(bob, bob_chat_id).await;
let pinned = get_pinned_messages(bob, bob_chat_id).await?;
assert_eq!(pinned.len(), 1);
// Wait until the message expires.
SystemTime::shift(Duration::from_secs(100));
ephemeral::delete_expired_messages(bob, time()).await?;
expect_pinned_message_event(bob, bob_chat_id).await;
let pinned = get_pinned_messages(bob, bob_chat_id).await?;
assert!(pinned.is_empty());
Ok(())
}
/// Tests that `PinnedMessagesChanged` event is emitted when pinned message is deleted.
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn test_delete_pinned_message() -> Result<()> {
let mut tcm = TestContextManager::new();
let alice = &tcm.alice().await;
let bob = &tcm.bob().await;
let bob_msg = tcm.send_recv_accept(alice, bob, "Hello!").await;
let bob_chat_id = bob_msg.chat_id;
// Bob pins "Hello!" message received from Alice.
set_pinned_state(bob, bob_msg.id, true).await?;
expect_pinned_message_event(bob, bob_chat_id).await;
let pinned = get_pinned_messages(bob, bob_chat_id).await?;
assert_eq!(pinned.len(), 1);
message::delete_msgs(bob, &[bob_msg.id]).await?;
expect_pinned_message_event(bob, bob_chat_id).await;
let pinned = get_pinned_messages(bob, bob_chat_id).await?;
assert!(pinned.is_empty());
Ok(())
}
+30 -18
View File
@@ -13,7 +13,6 @@ use serde::Deserialize;
use crate::autorelay::login_param_from_host;
use crate::config::Config;
use crate::configure::MAX_RELAYS;
use crate::contact::{Contact, ContactId, Origin};
use crate::context::Context;
use crate::key::Fingerprint;
@@ -134,6 +133,18 @@ pub enum Qr {
contact_id: ContactId,
},
/// Scanned fingerprint does not match the last seen fingerprint.
FprMismatch {
/// Contact ID.
contact_id: Option<ContactId>,
},
/// The scanned QR code contains a fingerprint but no e-mail address.
FprWithoutAddr {
/// Key fingerprint.
fingerprint: String,
},
/// Ask the user if they want to create an account on the given domain.
Account {
/// Server domain name.
@@ -489,8 +500,7 @@ async fn decode_openpgp(context: &Context, qr: &str) -> Result<Qr> {
addrs.push(normalize_address(primary_addr)?);
};
if let Some(secondary_addrs_raw) = param.get("r") {
let max_secondary = MAX_RELAYS.saturating_sub(addrs.len());
for secondary_address in secondary_addrs_raw.split(',').take(max_secondary) {
for secondary_address in secondary_addrs_raw.split(',') {
addrs.push(normalize_address(secondary_address)?)
}
}
@@ -641,21 +651,24 @@ async fn decode_openpgp(context: &Context, qr: &str) -> Result<Qr> {
is_v3,
})
}
} else {
} else if let Some(addr) = addrs.first() {
let fingerprint = fingerprint.hex();
let contact_id: Option<ContactId> = context
.sql
.query_get_value(
"SELECT id FROM contacts WHERE fingerprint=?",
(fingerprint,),
)
.await?;
let (contact_id, _) =
Contact::add_or_lookup_ext(context, "", addr, &fingerprint, Origin::UnhandledQrScan)
.await?;
let contact = Contact::get_by_id(context, contact_id).await?;
let Some(contact_id) = contact_id else {
bail!("Contact matching the fingerprint is not found");
};
Ok(Qr::FprOk { contact_id })
if contact.public_key(context).await?.is_some() {
Ok(Qr::FprOk { contact_id })
} else {
Ok(Qr::FprMismatch {
contact_id: Some(contact_id),
})
}
} else {
Ok(Qr::FprWithoutAddr {
fingerprint: fingerprint.human_readable(),
})
}
}
@@ -825,8 +838,7 @@ pub(crate) async fn login_param_from_account_qr(
.context("Invalid DCACCOUNT scheme")?;
if !payload.starts_with(HTTPS_SCHEME) {
let mark_as_autorelay = false;
let param = login_param_from_host(payload, mark_as_autorelay);
let param = login_param_from_host(payload);
return Ok(param);
}
+38 -52
View File
@@ -305,10 +305,12 @@ async fn test_decode_openpgp_invalid_token() -> Result<()> {
let ctx = TestContext::new_alice().await;
// Token cannot contain "/"
assert!(check_qr(
let qr = check_qr(
&ctx.ctx,
"OPENPGP4FPR:79252762C34C5096AF57958F4FC3D21A81B0F0A7#a=cli%40deltachat.de&g=test%20%3F+test%20%21&x=h-0oKQf2CDK&i=9JEXlxAqGM0&s=0V7LzL/cxRL"
).await.is_err());
).await?;
assert!(matches!(qr, Qr::FprMismatch { .. }));
Ok(())
}
@@ -359,23 +361,6 @@ async fn test_decode_openpgp_secure_join() -> Result<()> {
bail!("Wrong QR code type");
}
// A bad invite code must not be able to steer us onto arbitrarily many relays.
let relays: Vec<String> = (0..20).map(|i| format!("cli%40r{i}.example.org")).collect();
let qr = check_qr(
&ctx.ctx,
&format!(
"openpgp4fpr:79252762C34C5096AF57958F4FC3D21A81B0F0A7#a=cli%40deltachat.de&r={}&i=TbnwJ6lSvD5&s=0ejvbdFSQxB",
relays.join(",")
),
)
.await?;
if let Qr::AskVerifyContact { addrs, .. } = qr {
assert_eq!(addrs.len(), MAX_RELAYS);
} else {
bail!("Wrong QR code type");
}
Ok(())
}
@@ -388,17 +373,16 @@ async fn test_decode_openpgp_fingerprint() -> Result<()> {
let alice_contact = bob.add_or_lookup_contact(alice).await;
let alice_contact_id = alice_contact.id;
// OPENPGP4FPR may have an address,
// but it is not used anymore since key contacts are introduced.
// We lookup the contact only by fingerprint and ignore the address.
assert!(
check_qr(
bob,
"OPENPGP4FPR:1234567890123456789012345678901234567890#a=alice@example.org",
)
.await
.is_err()
);
let qr = check_qr(
bob,
"OPENPGP4FPR:1234567890123456789012345678901234567890#a=alice@example.org",
)
.await?;
if let Qr::FprMismatch { contact_id, .. } = qr {
assert_ne!(contact_id.unwrap(), alice_contact_id);
} else {
bail!("Wrong QR code type");
}
let qr = check_qr(
bob,
@@ -414,44 +398,46 @@ async fn test_decode_openpgp_fingerprint() -> Result<()> {
bail!("Wrong QR code type");
}
assert!(
assert!(matches!(
check_qr(
bob,
"OPENPGP4FPR:1234567890123456789012345678901234567890#a=bob@example.org",
)
.await
.is_err()
);
.await?,
Qr::FprMismatch { .. }
));
Ok(())
}
/// Tests OPENPGP4FPR QR codes without an email address.
///
/// Email address in OPENPGP4FPR QR codes was an extension
/// used before switch to identifying contacts by fingerprint.
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn test_decode_openpgp_without_addr() -> Result<()> {
let ctx = TestContext::new().await;
assert!(
check_qr(
&ctx.ctx,
"OPENPGP4FPR:1234567890123456789012345678901234567890",
)
.await
.is_err()
let qr = check_qr(
&ctx.ctx,
"OPENPGP4FPR:1234567890123456789012345678901234567890",
)
.await?;
assert_eq!(
qr,
Qr::FprWithoutAddr {
fingerprint: "1234 5678 9012 3456 7890\n1234 5678 9012 3456 7890".to_string()
}
);
// Test it again with lowercased "openpgp4fpr:" uri scheme
assert!(
check_qr(
&ctx.ctx,
"openpgp4fpr:1234567890123456789012345678901234567890",
)
.await
.is_err()
let qr = check_qr(
&ctx.ctx,
"openpgp4fpr:1234567890123456789012345678901234567890",
)
.await?;
assert_eq!(
qr,
Qr::FprWithoutAddr {
fingerprint: "1234 5678 9012 3456 7890\n1234 5678 9012 3456 7890".to_string()
}
);
let res = check_qr(&ctx.ctx, "OPENPGP4FPR:12345678901234567890").await;
+2 -2
View File
@@ -84,7 +84,7 @@ pub struct ReactionFrequency {
pub reaction: Reaction,
/// Number of contacts that reacted with this emoji.
pub count: u32,
pub count: usize,
/// True if `ContactId::SELF` is among the contacts that reacted with this emoji.
pub is_from_self: bool,
@@ -420,7 +420,7 @@ pub(crate) async fn apply_pending_reactions(
/// sorted in descending order of frequency.
fn calc_frequencies(by_contact: &BTreeMap<ContactId, Reaction>) -> Vec<ReactionFrequency> {
let mut self_reaction = Reaction::new("");
let mut counts: BTreeMap<&str, u32> = BTreeMap::new();
let mut counts: BTreeMap<&str, usize> = BTreeMap::new();
for (contact_id, reaction) in by_contact {
let count = counts.entry(reaction.as_str()).or_insert(0);
*count = count.saturating_add(1);
+4 -4
View File
@@ -42,7 +42,7 @@ struct WireMessage {
#[derive(Debug, Serialize, Deserialize)]
struct WireEntry {
emoji: String,
count: u32,
count: usize,
}
/// Renders one or more message's states as a JSON string, ready to be sent in `Chat-Broadcast-States:` header.
@@ -235,10 +235,10 @@ pub(crate) async fn load_broadcast_reactions(
(msg_id,),
|row| {
let reaction: String = row.get(0)?;
let count: u32 = row.get(1)?;
let count: i64 = row.get(1)?;
Ok(ReactionFrequency {
reaction: Reaction::new(&reaction),
count,
count: count as usize,
is_from_self: false,
})
},
@@ -405,7 +405,7 @@ mod tests {
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn test_modify_frequencies() {
// Helper to create a ReactionFrequency entry
let freq = |emoji: &str, count: u32, is_from_self: bool| -> ReactionFrequency {
let freq = |emoji: &str, count: usize, is_from_self: bool| -> ReactionFrequency {
ReactionFrequency {
reaction: Reaction::new(emoji),
count,
+12 -23
View File
@@ -2107,8 +2107,7 @@ async fn add_parts(
let hidden = part.is_reaction;
if part.is_reaction {
let reaction_str = simplify::remove_footers(part.msg.as_str());
let is_incoming_fresh =
mime_parser.incoming && !seen && chat_id_blocked == Blocked::Not;
let is_incoming_fresh = mime_parser.incoming && !seen;
set_msg_reaction(
context,
mime_in_reply_to,
@@ -2392,7 +2391,6 @@ async fn handle_edit_delete(
}
let mut modified_chat_ids = BTreeSet::new();
let mut pinned_messages_changed_chat_ids = BTreeSet::new();
let mut msg_ids = Vec::new();
let rfc724_mid_vec: Vec<&str> = rfc724_mid_list.split_whitespace().collect();
@@ -2417,17 +2415,8 @@ async fn handle_edit_delete(
message::delete_msg_locally(context, &msg).await?;
msg_ids.push(msg.id);
modified_chat_ids.insert(msg.chat_id);
if msg.is_pinned() {
pinned_messages_changed_chat_ids.insert(msg.chat_id);
}
}
message::delete_msgs_locally_done(
context,
&msg_ids,
modified_chat_ids,
pinned_messages_changed_chat_ids,
)
.await?;
message::delete_msgs_locally_done(context, &msg_ids, modified_chat_ids).await?;
}
Ok(())
}
@@ -2716,12 +2705,6 @@ async fn lookup_or_create_adhoc_group(
for &id in &contact_ids {
stmt.execute((id,)).context("INSERT INTO temp.contacts")?;
}
// Contact IDs are 32-bit internally,
// so this conversion of contact ID set size
// to u32 should never fail.
let contact_ids_len = u32::try_from(contact_ids.len())?;
let val = t
.query_row(
"SELECT c.id, c.blocked
@@ -2735,7 +2718,7 @@ async fn lookup_or_create_adhoc_group(
AND contact_id NOT IN (SELECT id FROM temp.contacts)
AND add_timestamp >= remove_timestamp)=0
ORDER BY m.timestamp DESC",
(&grpname, contact_ids_len),
(&grpname, contact_ids.len()),
|row| {
let id: ChatId = row.get(0)?;
let blocked: Blocked = row.get(1)?;
@@ -3749,9 +3732,10 @@ async fn apply_out_broadcast_changes(
if from_id == ContactId::SELF {
let added_id = lookup_key_contact_by_fingerprint(context, added_fpr).await?;
if let Some(added_id) = added_id {
info!(context, "Broadcast addition (TRASH)");
better_msg.get_or_insert("".to_string());
if !chat::is_contact_in_chat(context, chat.id, added_id).await? {
if chat::is_contact_in_chat(context, chat.id, added_id).await? {
info!(context, "No-op broadcast addition (TRASH)");
better_msg.get_or_insert("".to_string());
} else {
chat::add_to_chat_contacts_table(
context,
mime_parser.timestamp_sent,
@@ -3759,6 +3743,11 @@ async fn apply_out_broadcast_changes(
&[added_id],
)
.await?;
let msg =
stock_str::msg_add_member_local(context, added_id, ContactId::UNDEFINED)
.await;
better_msg.get_or_insert(msg);
added_removed_id = Some(added_id);
send_event_chat_modified = true;
}
} else {
+31 -34
View File
@@ -726,7 +726,6 @@ async fn test_resend_after_ndn() -> Result<()> {
Ok(())
}
// an NDN in a group does not make the whole message as failed
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn test_parse_ndn_group_msg() -> Result<()> {
let t = TestContext::new().await;
@@ -759,7 +758,7 @@ async fn test_parse_ndn_group_msg() -> Result<()> {
let msg = Message::load_from_db(&t, msg_id).await?;
assert_eq!(msg.state, MessageState::OutDelivered);
assert_eq!(msg.state, MessageState::OutFailed);
let msgs = chat::get_chat_msgs(&t, msg.chat_id).await?;
assert!(matches!(
@@ -767,6 +766,9 @@ async fn test_parse_ndn_group_msg() -> Result<()> {
ChatItem::Message { msg_id } if msg_id == msg.id
));
t.assert_warn("Delivery Status Notification (Failure)")
.await;
Ok(())
}
@@ -2723,7 +2725,19 @@ async fn test_read_receipts_dont_create_chats() -> Result<()> {
let chats = Chatlist::try_load(&alice, 0, None, None).await?;
assert_eq!(chats.len(), 0);
alice.recv_mdn(&bob, &received_msg).await?;
// Bob sends a read receipt.
let mdn_mimefactory = crate::mimefactory::MimeFactory::from_mdn(
&bob,
received_msg.from_id,
received_msg.rfc724_mid,
vec![],
)
.await?;
let rendered_mdn = mdn_mimefactory.render(&bob).await?;
let mdn_body = rendered_mdn.message;
// Alice receives the read receipt.
receive_imf(&alice, mdn_body.as_bytes(), false).await?;
// Chat should not pop up in the chatlist.
let chats = Chatlist::try_load(&alice, 0, None, None).await?;
@@ -2746,7 +2760,19 @@ async fn test_read_receipts_dont_unmark_bots() -> Result<()> {
.await;
let received_msg = bob.get_last_msg().await;
alice.recv_mdn(bob, &received_msg).await?;
// Bob sends a read receipt.
let mdn_mimefactory = crate::mimefactory::MimeFactory::from_mdn(
bob,
received_msg.from_id,
received_msg.rfc724_mid,
vec![],
)
.await?;
let rendered_mdn = mdn_mimefactory.render(bob).await?;
let mdn_body = rendered_mdn.message;
// Alice receives the read receipt.
receive_imf(alice, mdn_body.as_bytes(), false).await?;
let msg = alice.get_last_msg_in(alice_chat.id).await;
assert_eq!(msg.state, MessageState::OutMdnRcvd);
let ab_contact = alice.add_or_lookup_contact(bob).await;
@@ -3338,35 +3364,6 @@ async fn test_blocked_contact_creates_group() -> Result<()> {
Ok(())
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn test_blocked_contact_sends_reaction() -> Result<()> {
let mut tcm = TestContextManager::new();
let alice = &tcm.alice().await;
let bob = &tcm.bob().await;
let bob_msg_id = tcm.send_recv_accept(alice, bob, "Hi!").await.id;
let chat = alice.get_chat(bob).await;
chat.id.block(alice).await?;
crate::reaction::send_reaction(bob, bob_msg_id, "👍").await?;
let sent = bob.pop_sent_msg().await;
alice.recv_msg_hidden(&sent).await;
alice.emit_event(EventType::Test);
while let Some(ev) = alice.evtracker.recv().await {
match ev.typ {
EventType::IncomingReaction { .. } => {
panic!("Alice is not supposed to receive a notification, since she blocked Bob")
}
EventType::Test => break,
_ => {}
}
}
Ok(())
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn test_outgoing_undecryptable() -> Result<()> {
let alice = &TestContext::new().await;
@@ -5766,7 +5763,7 @@ pub(crate) async fn first_row_in_smtp_queue(context: &TestContext) -> (MsgId, St
let queued_mail = context
.sql
.transaction_ext(query_only, |transaction| {
smtp::queue::load_queued_mail(transaction, rowid)
smtp::load_queued_mail(transaction, rowid)
})
.await
.unwrap();
+12 -13
View File
@@ -215,17 +215,10 @@ impl SchedulerState {
/// Indicate that the network likely has come back.
pub(crate) async fn maybe_network(&self) {
self.interrupt_inbox_idle().await;
self.interrupt_smtp().await;
}
/// Interrupts IDLE on all transports so that they fetch,
/// and marks them as having work to do.
pub(crate) async fn interrupt_inbox_idle(&self) {
let inner = self.inner.read().await;
let inboxes = match *inner {
InnerSchedulerState::Started(ref scheduler) => {
scheduler.interrupt_inbox();
scheduler.maybe_network();
scheduler
.inboxes
.iter()
@@ -341,7 +334,7 @@ async fn background_fetch_from_transport(
let folder = connection.folder.clone();
connection
.fetch_delete(context, &mut session, &folder)
.fetch_move_delete(context, &mut session, &folder)
.await
}
@@ -521,6 +514,7 @@ async fn inbox_fetch_idle(ctx: &Context, imap: &mut Imap, mut session: Session)
};
maybe_broadcast_reactions(ctx).await.log_err(ctx).ok();
maybe_send_stats(ctx).await.log_err(ctx).ok();
session
.update_metadata(ctx)
@@ -533,8 +527,6 @@ async fn inbox_fetch_idle(ctx: &Context, imap: &mut Imap, mut session: Session)
);
}
maybe_send_stats(ctx).await.log_err(ctx).ok();
let session = fetch_idle(ctx, imap, session).await?;
Ok(session)
}
@@ -556,9 +548,9 @@ async fn fetch_idle(ctx: &Context, connection: &mut Imap, mut session: Session)
// Fetch the watched folder.
connection
.fetch_delete(ctx, &mut session, &watch_folder)
.fetch_move_delete(ctx, &mut session, &watch_folder)
.await
.context("fetch_delete")?;
.context("fetch_move_delete")?;
download_known_post_messages_without_pre_message(ctx, &mut session).await?;
download_msgs(ctx, &mut session)
@@ -807,6 +799,13 @@ impl Scheduler {
self.inboxes.iter()
}
fn maybe_network(&self) {
for b in self.boxes() {
b.conn_state.interrupt();
}
self.interrupt_smtp();
}
fn maybe_network_lost(&self) {
for b in self.boxes() {
b.conn_state.interrupt();
+6 -42
View File
@@ -1,6 +1,6 @@
use core::fmt;
use std::cmp::min;
use std::{ops::Deref, sync::Arc};
use std::{iter::once, ops::Deref, sync::Arc};
use anyhow::Result;
use humansize::{BINARY, format_size};
@@ -531,15 +531,14 @@ impl Context {
Ok(ret)
}
/// Returns true if all background work is done,
/// checking the outgoing message queue only if `include_smtp` is set.
async fn work_done(&self, include_smtp: bool) -> bool {
/// Returns true if all background work is done.
async fn all_work_done(&self) -> bool {
let lock = self.scheduler.inner.read().await;
let stores: Vec<_> = match *lock {
InnerSchedulerState::Started(ref sched) => sched
.boxes()
.map(|b| &b.conn_state.state)
.chain(include_smtp.then_some(&sched.smtp.state))
.chain(once(&sched.smtp.state))
.map(|state| state.connectivity.clone())
.collect(),
_ => return false,
@@ -556,26 +555,19 @@ impl Context {
/// Waits until background work is finished.
pub async fn wait_for_all_work_done(&self) {
let include_smtp = true;
self.wait_for_work_done(include_smtp).await
}
/// Waits until background work is finished,
/// checking the outgoing message queue only if `include_smtp` is set.
pub(crate) async fn wait_for_work_done(&self, include_smtp: bool) {
// Ideally we could wait for connectivity change events,
// but sleep loop is good enough.
// First 100 ms sleep in chunks of 10 ms.
for _ in 0..10 {
if self.work_done(include_smtp).await {
if self.all_work_done().await {
break;
}
tokio::time::sleep(std::time::Duration::from_millis(10)).await;
}
// If we are not finished in 100 ms, keep waking up every 100 ms.
while !self.work_done(include_smtp).await {
while !self.all_work_done().await {
tokio::time::sleep(std::time::Duration::from_millis(100)).await;
}
}
@@ -584,34 +576,6 @@ impl Context {
#[cfg(test)]
mod tests {
use super::*;
use crate::test_utils::TestContext;
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn test_background_fetch_leaves_smtp_alone() -> Result<()> {
let alice = TestContext::new_alice().await;
alice.start_io().await;
alice.wait_for_all_work_done().await;
let smtp = match *alice.scheduler.inner.read().await {
InnerSchedulerState::Started(ref scheduler) => {
scheduler.smtp.state.connectivity.clone()
}
_ => panic!("scheduler is not running"),
};
smtp.set_working(&alice);
alice.background_fetch().await?;
assert!(!smtp.get_all_work_done());
assert!(alice.scheduler.is_running().await);
alice
.assert_warns_or_errors(&[
"No IMAP connection candidates provided",
"IMAP got rate limited",
])
.await;
Ok(())
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn test_combine_connectivities() {
+158 -85
View File
@@ -1,7 +1,6 @@
//! # SMTP transport module.
mod connect;
pub(crate) mod queue;
pub mod send;
use std::collections::BTreeSet;
@@ -9,24 +8,27 @@ use std::collections::BTreeSet;
use anyhow::{Context as _, Error, Result, bail, format_err};
use async_smtp::response::{Category, Code, Detail};
use async_smtp::{EmailAddress, SmtpTransport};
use pgp::composed::SignedPublicKey;
use rusqlite::OptionalExtension as _;
use tokio::task;
use crate::chat;
use crate::chat::{ChatId, add_info_msg_with_cmd};
use crate::config::Config;
use crate::contact::{Contact, ContactId};
use crate::context::Context;
use crate::events::EventType;
use crate::key;
use crate::key::DcKey;
use crate::log::{LogExt, warn};
use crate::message::Message;
use crate::message::{self, MsgId};
use crate::mimefactory;
use crate::mimefactory::MimeFactory;
use crate::mimefactory::QueuedMail;
use crate::net::proxy::ProxyConfig;
use crate::net::session::SessionBufStream;
use crate::scheduler::connectivity::ConnectivityStore;
use crate::smtp::queue::QueuedMail;
use crate::stock_str::unencrypted_email;
use crate::tools::{self, time_elapsed};
use crate::transport::{
@@ -41,9 +43,6 @@ pub(crate) struct Smtp {
/// Email address we are sending from.
from: Option<EmailAddress>,
/// Transport we are connected to.
transport_id: Option<u32>,
/// Timestamp of last successful send/receive network interaction
/// (eg connect or send succeeded). On initialization and disconnect
/// it is set to None.
@@ -99,39 +98,19 @@ impl Smtp {
}
self.connectivity.set_connecting(context);
let (_transport_id, lp) = ConfiguredLoginParam::load(context)
.await?
.context("Not configured")?;
let proxy_config = ProxyConfig::load(context).await?;
let transports = ConfiguredLoginParam::load_all(context).await?;
// Try to connect to the newest transport first. If sending is unreliable,
// user can configure a new transport and it will be the one used.
// Conversely, if user just added a new transport and sending got less reliable,
// user can restore old state by removing the just added transport.
for (transport_id, lp) in transports.into_iter().rev() {
info!(context, "Trying to connect to transport {transport_id}.");
match self
.connect(
context,
&lp.smtp,
&lp.smtp_password,
&proxy_config,
&lp.addr,
lp.strict_tls(proxy_config.is_some())?,
)
.await
{
Ok(()) => {
self.transport_id = Some(transport_id);
return Ok(());
}
Err(err) => {
warn!(
context,
"Failed to connect to SMTP transport {transport_id}: {err:#}."
);
}
}
}
bail!("Failed to connect to any SMTP server");
self.connect(
context,
&lp.smtp,
&lp.smtp_password,
&proxy_config,
&lp.addr,
lp.strict_tls(proxy_config.is_some())?,
)
.await
}
/// Connect using the provided login params.
@@ -218,6 +197,15 @@ pub(crate) async fn smtp_send(
smtp.connectivity.set_working(context);
if let Err(err) = smtp
.connect_configured(context)
.await
.context("Failed to open SMTP connection")
{
smtp.last_send_error = Some(format!("{err:#}"));
return SendResult::Retry;
}
let send_result = smtp.send(context, recipients, message.as_bytes()).await;
smtp.last_send_error = send_result.as_ref().err().map(|e| e.to_string());
@@ -359,7 +347,7 @@ pub(crate) async fn insert_into_smtp(
let msg_id = message::insert_tombstone(context, rfc724_mid).await?;
context
.sql
.transaction(|transaction| queue::enqueue_mail(transaction, now, msg_id, queued_msg, None))
.transaction(|transaction| chat::enqueue_mail(transaction, now, msg_id, queued_msg, None))
.await?;
Ok(())
}
@@ -408,7 +396,7 @@ pub(crate) async fn send_msg_to_smtp(
else {
return Ok(None);
};
let queued_mail = queue::load_queued_mail(transaction, rowid)
let queued_mail = load_queued_mail(transaction, rowid)
.with_context(|| format!("Failed to load queued mail for {rowid}"))?;
Ok(Some((queued_mail, msg_id, retries)))
})
@@ -431,16 +419,10 @@ pub(crate) async fn send_msg_to_smtp(
let mut recipients = queued_mail.recipients.clone();
if queued_mail.bcc_self {
let from_addr = smtp
.from
.as_ref()
.context("No From address available, likely not connected")?
.to_string();
add_self_recipients(
context,
&mut recipients,
queued_mail.encryption.is_encrypted(),
from_addr,
)
.await
.context("Failed to add self recipients")?;
@@ -469,12 +451,6 @@ pub(crate) async fn send_msg_to_smtp(
.context("No From address available, likely not connected")?
.to_string();
let transport_id = smtp.transport_id.context("Not connected")?;
let chunk_size = context
.get_max_smtp_rcpt_to(transport_id, &from_addr)
.await?
.max(1);
let rendered_mail =
mimefactory::render_queued_mail(queued_mail, &public_key, &secret_key, from_addr)?;
let body = rendered_mail.message;
@@ -483,22 +459,24 @@ pub(crate) async fn send_msg_to_smtp(
context,
"Try number {retries} to send message {msg_id} (entry {rowid}) over SMTP."
);
let chunk_size = context.get_max_smtp_rcpt_to().await?.max(1);
let mut unsent = recipients_list.as_slice();
let status = loop {
let unsent_len = u32::try_from(unsent.len()).context("Too many SMTP recipients")?;
let split_index = usize::try_from(chunk_size.min(unsent_len))
.context("Failed to convert SMTP chunk size")?;
let (chunk, rest) = unsent.split_at(split_index);
let status = smtp_send(context, chunk, body.as_str(), smtp, Some(msg_id)).await;
if !matches!(status, SendResult::Success) || rest.is_empty() {
break status;
}
for sent_to_addr in chunk {
let sent_to_addr_string = sent_to_addr.to_string();
if !sent_to_set.insert(sent_to_addr_string) {
error!(context, "Attempted to send to {sent_to_addr} twice.");
}
}
let status = smtp_send(context, chunk, body.as_str(), smtp, Some(msg_id)).await;
if !matches!(status, SendResult::Success) || rest.is_empty() {
break status;
}
let sent_to_str = sent_to_set
.iter()
.map(|a| a.as_ref())
@@ -521,15 +499,6 @@ pub(crate) async fn send_msg_to_smtp(
.sql
.execute("DELETE FROM smtp2 WHERE id=?", (rowid,))
.await?;
let sent_to_len = sent_to_set.len();
debug_assert!(sent_to_len > 0);
let info_msg = format!(
"Message len={} was SMTP-sent to {sent_to_len} recipients.",
body.len()
);
info!(context, "{info_msg}.");
context.emit_event(EventType::SmtpMessageSent(info_msg));
}
SendResult::Failure(ref err) => {
if err
@@ -700,16 +669,11 @@ async fn send_mdn_rfc724_mid(
} else {
mimefactory.recipients()
};
let from = smtp
.from
.as_ref()
.context("No From address, not connected")?
.to_string();
let rendered_msg = Box::pin(mimefactory.render(context, &from)).await?;
let rendered_msg = Box::pin(mimefactory.render(context)).await?;
let body = rendered_msg.message;
if context.get_config_bool(Config::BccSelf).await? {
add_self_recipients(context, &mut recipients, encrypted, from).await?;
add_self_recipients(context, &mut recipients, encrypted).await?;
}
let recipients: Vec<_> = recipients
.into_iter()
@@ -774,19 +738,6 @@ async fn send_mdn(context: &Context, smtp: &mut Smtp) -> Result<bool> {
};
let (rfc724_mid, contact_id) = msg_row;
// Connect SMTP after checking that we have an MDN to send,
// but before increasing MDN retry counter.
// If we don't have an MDN, no need to connect.
// If we are offline and cannot connect, it is not a failure of an MDN.
if let Err(err) = smtp
.connect_configured(context)
.await
.context("SMTP connection failure while preparing to send MDNs")
{
smtp.last_send_error = Some(format!("{err:#}"));
return Err(err);
}
context
.sql
.execute(
@@ -826,11 +777,11 @@ pub(crate) async fn add_self_recipients(
context: &Context,
recipients: &mut Vec<String>,
encrypted: bool,
from: String,
) -> Result<()> {
// Avoid sending unencrypted messages to all transports,
// chatmail relays won't accept them. Normally the user should have
// a non-chatmail sending transport to send unencrypted messages.
let from = context.get_primary_self_addr().await?;
if encrypted {
for addr in context.get_self_addrs().await? {
if addr != from {
@@ -848,3 +799,125 @@ pub(crate) async fn add_self_recipients(
Ok(())
}
/// Returns true if SMTP queue is empty.
pub(crate) async fn is_queue_empty(context: &Context) -> Result<bool> {
let sending_finished = !context.sql.exists("SELECT COUNT(*) FROM smtp2", ()).await?;
Ok(sending_finished)
}
/// Loads the queued mail from `smtp2` table and the list of recipients.
pub(crate) fn load_queued_mail(
transaction: &mut rusqlite::Transaction<'_>,
row_id: i64,
) -> Result<QueuedMail> {
let (mut queued_mail, encryption_fingerprints) = transaction
.query_row_and_then(
"
SELECT display_name,
rfc724_mid,
mime,
should_attach_pubkey,
should_compress,
should_sign,
is_encrypted,
shared_secret,
encryption_fingerprints,
recipients,
sent_to,
bcc_self
FROM smtp2 WHERE id = ?
",
(row_id,),
|row| {
let display_name: String = row.get(0)?;
let rfc724_mid: String = row.get(1)?;
let raw_message: Vec<u8> = row.get(2)?;
let should_attach_pubkey: bool = row.get(3)?;
let should_compress: bool = row.get(4)?;
let should_sign: bool = row.get(5)?;
let is_encrypted: bool = row.get(6)?;
let shared_secret: String = row.get(7)?;
let encryption_fingerprints: String = row.get(8)?;
let encryption_fingerprints: Vec<String> = if encryption_fingerprints.is_empty() {
Vec::new()
} else {
encryption_fingerprints
.split(' ')
.map(|s| s.to_string())
.collect()
};
let recipients: String = row.get(9)?;
let recipients: Vec<String> = if recipients.is_empty() {
Vec::new()
} else {
recipients.split(' ').map(|s| s.to_string()).collect()
};
debug_assert!(!recipients.iter().any(|s| s.is_empty()));
let sent_to: String = row.get(10)?;
let sent_to: Vec<String> = if sent_to.is_empty() {
Vec::new()
} else {
sent_to.split(' ').map(|s| s.to_string()).collect()
};
let bcc_self: bool = row.get(11)?;
let encryption = match (
is_encrypted,
shared_secret.is_empty(),
encryption_fingerprints.is_empty(),
) {
(false, true, true) => mimefactory::QueuedEncryption::No,
(true, false, true) => {
mimefactory::QueuedEncryption::Symmetric { shared_secret }
}
(true, true, _) => mimefactory::QueuedEncryption::Asymmetric {
// Public keys are loaded below based on the encryption fingerprints.
encryption_pubkeys: Vec::new(),
},
_ => bail!("Invalid encryption in smtp2 row"),
};
Ok::<_, anyhow::Error>((
QueuedMail {
raw_message,
display_name,
rfc724_mid,
encryption,
should_attach_pubkey,
should_compress,
should_sign,
recipients,
sent_to,
bcc_self,
},
encryption_fingerprints,
))
},
)
.with_context(|| format!("Failed to select row {row_id} from smtp2 table"))?;
if let mimefactory::QueuedEncryption::Asymmetric {
ref mut encryption_pubkeys,
} = queued_mail.encryption
{
for fingerprint in encryption_fingerprints {
let public_key_bytes: Option<Vec<u8>> = transaction
.query_row(
"SELECT public_key FROM public_keys WHERE fingerprint=?",
(fingerprint,),
|row| {
let bytes: Vec<u8> = row.get(0)?;
Ok(bytes)
},
)
.optional()
.context("Failed to select public key by fingerprint")?;
if let Some(public_key_bytes) = public_key_bytes {
let public_key = SignedPublicKey::from_slice(&public_key_bytes)?;
encryption_pubkeys.push(public_key);
}
}
}
Ok(queued_mail)
}
-323
View File
@@ -1,323 +0,0 @@
//! # SMTP queue module.
use anyhow::{Context as _, Result, bail};
use crate::chat::ChatId;
use crate::context::Context;
use crate::key::{DcKey, SignedPublicKey};
use crate::message::MsgId;
use rusqlite::OptionalExtension as _;
#[derive(Debug, Clone)]
pub(crate) enum Encryption {
/// Unencrypted message.
No,
/// The message is encrypted asymmetrically to public keys.
Asymmetric {
/// OpenPGP keys to use for encryption.
///
/// The message is always encrypted to self,
/// no need to include own key here.
encryption_pubkeys: Vec<SignedPublicKey>,
},
/// Symmetrically encrypted message with a shared secret.
Symmetric { shared_secret: String },
}
impl Encryption {
pub(crate) fn is_encrypted(&self) -> bool {
match self {
Self::No => false,
Self::Asymmetric { .. } => true,
Self::Symmetric { .. } => true,
}
}
}
/// Email message queued, but not sent yet.
///
/// It is stored unencrypted to
/// make it possible to change protected headers
/// like the From address and Autocrypt header later.
#[derive(Debug, Clone)]
pub(crate) struct QueuedMail {
/// Unencrypted queued message.
///
/// This message has both the headers and the body,
/// but without the From, Autocrypt and Message-ID headers.
///
/// For encrypted messages this is the OpenPGP payload.
pub(crate) raw_message: Vec<u8>,
/// Display name to put in the `From:` field.
///
/// Email address is not determined yet here.
pub(crate) display_name: String,
/// Message-ID.
pub(crate) rfc724_mid: String,
/// Whether the message is encrypted and encryption keys.
pub(crate) encryption: Encryption,
/// If true, Autocrypt header should be added before sending.
pub(crate) should_attach_pubkey: bool,
/// If true, OpenPGP compression may be used.
pub(crate) should_compress: bool,
/// If true, encrypted message should be signed.
pub(crate) should_sign: bool,
/// Recipient addresses.
pub(crate) recipients: Vec<String>,
/// Addresses the messages was already sent to.
pub(crate) sent_to: Vec<String>,
/// If true, own addresses should be added to the list of recipients.
///
/// For unencrypted messages, only the sending addresses should be added.
/// For encrypted messages, all published addresses should be added.
pub(crate) bcc_self: bool,
}
/// Side effects that should be applied at the same time
/// as the message is persisted in the queue.
#[derive(Debug, Clone, Default)]
pub struct SideEffects {
/// ID of the chat side effects should be applied to.
pub chat_id: ChatId,
/// Largest timestamp of the location sent in `location.kml` in this message.
pub last_added_location_timestamp: Option<i64>,
/// True if the message has the avatar attached.
///
/// Timestamp of the last time avatar was gossiped should be updated.
pub avatar_is_attached: bool,
/// A comma-separated string of sync-IDs that are used by the rendered email and must be deleted
/// from `multi_device_sync` once the message is actually queued for sending.
pub sync_ids_to_delete: Option<String>,
/// Subject that was rendered into the message.
///
/// Used to update the subject on the sent message object.
pub subject: String,
}
/// Email message ready to be queued with the side effects that should be applied at the same time.
pub(crate) type ToBeQueuedMail = (QueuedMail, Option<SideEffects>);
/// Process side effects and store queued mail.
pub(crate) fn enqueue_mail(
transaction: &mut rusqlite::Transaction<'_>,
now: i64,
msg_id: MsgId,
queued_mail: &QueuedMail,
side_effects: Option<&SideEffects>,
) -> Result<i64> {
if let Some(side_effects) = side_effects {
if let Some(last_added_location_timestamp) = side_effects.last_added_location_timestamp {
transaction.execute(
"UPDATE chats SET locations_last_sent=? WHERE id=?;",
(last_added_location_timestamp, side_effects.chat_id),
)?;
}
if side_effects.avatar_is_attached {
side_effects
.chat_id
.set_selfavatar_timestamp(transaction, now)
.context("Failed to set selfavatar timestamp")?;
}
if let Some(ref sync_ids) = side_effects.sync_ids_to_delete {
transaction.execute(
&format!("DELETE FROM multi_device_sync WHERE id IN ({sync_ids})"),
(),
)?;
}
}
// Store mail into queue.
let all_recipients = queued_mail.recipients.join(" ");
let is_encrypted = queued_mail.encryption.is_encrypted();
transaction
.execute(
"
INSERT INTO smtp2 (
display_name,
rfc724_mid,
mime,
should_attach_pubkey,
should_compress,
should_sign,
msg_id,
recipients,
bcc_self,
is_encrypted,
shared_secret,
encryption_fingerprints
)
VALUES (
?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?
)
",
(
&queued_mail.display_name,
&queued_mail.rfc724_mid,
&queued_mail.raw_message,
queued_mail.should_attach_pubkey,
queued_mail.should_compress,
queued_mail.should_sign,
msg_id,
&all_recipients,
queued_mail.bcc_self,
is_encrypted,
if let Encryption::Symmetric { ref shared_secret } = queued_mail.encryption {
shared_secret
} else {
""
},
if let Encryption::Asymmetric {
ref encryption_pubkeys,
} = queued_mail.encryption
{
let res: Vec<String> = encryption_pubkeys
.iter()
.map(|pubkey| pubkey.dc_fingerprint().hex())
.collect();
res.join(" ")
} else {
"".to_string()
},
),
)
.context("Failed to insert a row into smtp2 table")?;
let row_id = transaction.last_insert_rowid();
Ok(row_id)
}
/// Loads the queued mail from `smtp2` table and the list of recipients.
pub(crate) fn load_queued_mail(
transaction: &mut rusqlite::Transaction<'_>,
row_id: i64,
) -> Result<QueuedMail> {
let (mut queued_mail, encryption_fingerprints) = transaction
.query_row_and_then(
"
SELECT display_name,
rfc724_mid,
mime,
should_attach_pubkey,
should_compress,
should_sign,
is_encrypted,
shared_secret,
encryption_fingerprints,
recipients,
sent_to,
bcc_self
FROM smtp2 WHERE id = ?
",
(row_id,),
|row| {
let display_name: String = row.get(0)?;
let rfc724_mid: String = row.get(1)?;
let raw_message: Vec<u8> = row.get(2)?;
let should_attach_pubkey: bool = row.get(3)?;
let should_compress: bool = row.get(4)?;
let should_sign: bool = row.get(5)?;
let is_encrypted: bool = row.get(6)?;
let shared_secret: String = row.get(7)?;
let encryption_fingerprints: String = row.get(8)?;
let encryption_fingerprints: Vec<String> = if encryption_fingerprints.is_empty() {
Vec::new()
} else {
encryption_fingerprints
.split(' ')
.map(|s| s.to_string())
.collect()
};
let recipients: String = row.get(9)?;
let recipients: Vec<String> = if recipients.is_empty() {
Vec::new()
} else {
recipients.split(' ').map(|s| s.to_string()).collect()
};
debug_assert!(!recipients.iter().any(|s| s.is_empty()));
let sent_to: String = row.get(10)?;
let sent_to: Vec<String> = if sent_to.is_empty() {
Vec::new()
} else {
sent_to.split(' ').map(|s| s.to_string()).collect()
};
let bcc_self: bool = row.get(11)?;
let encryption = match (
is_encrypted,
shared_secret.is_empty(),
encryption_fingerprints.is_empty(),
) {
(false, true, true) => Encryption::No,
(true, false, true) => Encryption::Symmetric { shared_secret },
(true, true, _) => Encryption::Asymmetric {
// Public keys are loaded below based on the encryption fingerprints.
encryption_pubkeys: Vec::new(),
},
_ => bail!("Invalid encryption in smtp2 row"),
};
Ok::<_, anyhow::Error>((
QueuedMail {
raw_message,
display_name,
rfc724_mid,
encryption,
should_attach_pubkey,
should_compress,
should_sign,
recipients,
sent_to,
bcc_self,
},
encryption_fingerprints,
))
},
)
.with_context(|| format!("Failed to select row {row_id} from smtp2 table"))?;
if let Encryption::Asymmetric {
ref mut encryption_pubkeys,
} = queued_mail.encryption
{
for fingerprint in encryption_fingerprints {
let public_key_bytes: Option<Vec<u8>> = transaction
.query_row(
"SELECT public_key FROM public_keys WHERE fingerprint=?",
(fingerprint,),
|row| {
let bytes: Vec<u8> = row.get(0)?;
Ok(bytes)
},
)
.optional()
.context("Failed to select public key by fingerprint")?;
if let Some(public_key_bytes) = public_key_bytes {
let public_key = SignedPublicKey::from_slice(&public_key_bytes)?;
encryption_pubkeys.push(public_key);
}
}
}
Ok(queued_mail)
}
/// Returns true if SMTP queue is empty.
pub(crate) async fn is_empty(context: &Context) -> Result<bool> {
let sending_finished = !context.sql.exists("SELECT COUNT(*) FROM smtp2", ()).await?;
Ok(sending_finished)
}
+9
View File
@@ -5,6 +5,7 @@ use async_smtp::{EmailAddress, Envelope, SendableEmail};
use super::Smtp;
use crate::config::Config;
use crate::context::Context;
use crate::events::EventType;
use crate::log::warn;
use crate::tools;
@@ -38,6 +39,8 @@ impl Smtp {
context.ratelimit.write().await.send();
}
let message_len_bytes = message.len();
let envelope =
Envelope::new(self.from.clone(), recipients.to_vec()).map_err(Error::Envelope)?;
let mail = SendableEmail::new(envelope, message);
@@ -52,6 +55,12 @@ impl Smtp {
transport.send(mail).await.map_err(Error::SmtpSend)?;
let info_msg = format!(
"Message len={message_len_bytes} was SMTP-sent to {} recipients.",
recipients.len()
);
info!(context, "{info_msg}.");
context.emit_event(EventType::SmtpMessageSent(info_msg));
self.last_success = Some(tools::Time::now());
Ok(())
}
+147 -15
View File
@@ -56,6 +56,10 @@ pub struct Sql {
/// SQL connection pool.
pool: RwLock<Option<Pool>>,
/// None if the database is not open, true if it is open with passphrase and false if it is
/// open without a passphrase.
is_encrypted: RwLock<Option<bool>>,
/// Cache of `config` table.
pub(crate) config_cache: RwLock<HashMap<String, Option<String>>>,
}
@@ -66,15 +70,52 @@ impl Sql {
Self {
dbfile,
pool: Default::default(),
is_encrypted: Default::default(),
config_cache: Default::default(),
}
}
/// Tests SQLCipher passphrase.
///
/// Returns true if passphrase is correct, i.e. the database is new or can be unlocked with
/// this passphrase, and false if the database is already encrypted with another passphrase or
/// corrupted.
///
/// Fails if database is already open.
pub async fn check_passphrase(&self, passphrase: String) -> Result<bool> {
if self.is_open().await {
bail!("Database is already opened.");
}
// Hold the lock to prevent other thread from opening the database.
let _lock = self.pool.write().await;
// Test that the key is correct using a single connection.
let connection = Connection::open(&self.dbfile)?;
if !passphrase.is_empty() {
connection
.pragma_update(None, "key", &passphrase)
.context("Failed to set PRAGMA key")?;
}
let key_is_correct = connection
.query_row("SELECT count(*) FROM sqlite_master", [], |_row| Ok(()))
.is_ok();
Ok(key_is_correct)
}
/// Checks if there is currently a connection to the underlying Sqlite database.
pub async fn is_open(&self) -> bool {
self.pool.read().await.is_some()
}
/// Returns true if the database is encrypted.
///
/// If database is not open, returns `None`.
pub(crate) async fn is_encrypted(&self) -> Option<bool> {
*self.is_encrypted.read().await
}
/// Closes all underlying Sqlite connections.
pub(crate) async fn close(&self) {
let _ = self.pool.write().await.take();
@@ -82,22 +123,52 @@ impl Sql {
}
/// Imports the database from a separate file with the given passphrase.
pub(crate) async fn import(&self, path: &Path) -> Result<()> {
pub(crate) async fn import(&self, path: &Path, passphrase: String) -> Result<()> {
let path_str = path
.to_str()
.with_context(|| format!("path {path:?} is not valid unicode"))?
.to_string();
// Keep `config_cache` locked all the time the db is imported so that nobody can use invalid
// values from there. And clear it immediately so as not to forget in case of errors.
let mut config_cache = self.config_cache.write().await;
config_cache.clear();
let src_conn =
rusqlite::Connection::open(path).context("Failed to open source database")?;
let query_only = false;
self.call(query_only, move |conn| {
let backup = rusqlite::backup::Backup::new(&src_conn, &mut *conn)?;
backup.run_to_completion(5, std::time::Duration::ZERO, None)?;
drop(backup);
// Check that backup passphrase is correct before resetting our database.
conn.execute("ATTACH DATABASE ? AS backup KEY ?", (path_str, passphrase))
.context("failed to attach backup database")?;
let res = conn
.query_row("SELECT count(*) FROM sqlite_master", [], |_row| Ok(()))
.context("backup passphrase is not correct");
conn.execute("VACUUM", [])
.context("failed to vacuum the database")?;
// Reset the database without reopening it. We don't want to reopen the database because we
// don't have main database passphrase at this point.
// See <https://sqlite.org/c3ref/c_dbconfig_enable_fkey.html> for documentation.
// Without resetting import may fail due to existing tables.
res.and_then(|_| {
conn.set_db_config(DbConfig::SQLITE_DBCONFIG_RESET_DATABASE, true)
.context("failed to set SQLITE_DBCONFIG_RESET_DATABASE")
})
.and_then(|_| {
conn.execute("VACUUM", [])
.context("failed to vacuum the database")
})
.and(
conn.set_db_config(DbConfig::SQLITE_DBCONFIG_RESET_DATABASE, false)
.context("failed to unset SQLITE_DBCONFIG_RESET_DATABASE"),
)
.and_then(|_| {
conn.query_row("SELECT sqlcipher_export('main', 'backup')", [], |_row| {
Ok(())
})
.context("failed to import from attached backup database")
})
.and(
conn.execute("DETACH DATABASE backup", [])
.context("failed to detach backup database"),
)?;
Ok(())
})
.await
@@ -106,10 +177,10 @@ impl Sql {
const N_DB_CONNECTIONS: usize = 3;
/// Creates a new connection pool.
fn new_pool(dbfile: &Path) -> Result<Pool> {
fn new_pool(dbfile: &Path, passphrase: String) -> Result<Pool> {
let mut connections = Vec::with_capacity(Self::N_DB_CONNECTIONS);
for _ in 0..Self::N_DB_CONNECTIONS {
let connection = new_connection(dbfile)?;
let connection = new_connection(dbfile, &passphrase)?;
connections.push(connection);
}
@@ -117,8 +188,8 @@ impl Sql {
Ok(pool)
}
async fn try_open(&self, context: &Context, dbfile: &Path) -> Result<()> {
*self.pool.write().await = Some(Self::new_pool(dbfile)?);
async fn try_open(&self, context: &Context, dbfile: &Path, passphrase: String) -> Result<()> {
*self.pool.write().await = Some(Self::new_pool(dbfile, passphrase.to_string())?);
if let Err(e) = self.run_migrations(context).await {
error!(context, "Running migrations failed: {e:#}");
@@ -176,7 +247,7 @@ impl Sql {
/// Opens the provided database and runs any necessary migrations.
/// If a database is already open, this will return an error.
pub async fn open(&self, context: &Context) -> Result<()> {
pub async fn open(&self, context: &Context, passphrase: String) -> Result<()> {
if self.is_open().await {
error!(
context,
@@ -185,8 +256,10 @@ impl Sql {
bail!("SQL database is already opened.");
}
self.try_open(context, &self.dbfile).await?;
let passphrase_nonempty = !passphrase.is_empty();
self.try_open(context, &self.dbfile, passphrase).await?;
info!(context, "Opened database {:?}.", self.dbfile);
*self.is_encrypted.write().await = Some(passphrase_nonempty);
// setup debug logging if there is an entry containing its id
if let Some(xdc_id) = self
@@ -198,6 +271,28 @@ impl Sql {
Ok(())
}
/// Changes the passphrase of encrypted database.
///
/// The database must already be encrypted and the passphrase cannot be empty.
/// It is impossible to turn encrypted database into unencrypted
/// and vice versa this way, use import/export for this.
pub async fn change_passphrase(&self, passphrase: String) -> Result<()> {
let mut lock = self.pool.write().await;
let pool = lock.take().context("SQL connection pool is not open")?;
let query_only = false;
let conn = pool.get(query_only).await?;
if !passphrase.is_empty() {
conn.pragma_update(None, "rekey", passphrase.clone())
.context("Failed to set PRAGMA rekey")?;
}
drop(pool);
*lock = Some(Self::new_pool(&self.dbfile, passphrase.to_string())?);
Ok(())
}
/// Allocates a connection and calls `function` with the connection.
///
/// If `query_only` is true, allocates read-only connection,
@@ -589,13 +684,47 @@ impl Sql {
}
}
pub(crate) trait TransactionExt {
/// Prepares and executes the statement and maps a function over the resulting rows.
///
/// Collects the resulting rows into a collection.
fn query_map_collect<T, C, F>(
&self,
sql: &str,
params: impl rusqlite::Params + Send,
f: F,
) -> Result<C>
where
T: Send + 'static,
C: Send + 'static + std::iter::FromIterator<T>,
F: Send + FnMut(&rusqlite::Row) -> Result<T>;
}
impl TransactionExt for rusqlite::Transaction<'_> {
fn query_map_collect<T, C, F>(
&self,
sql: &str,
params: impl rusqlite::Params + Send,
f: F,
) -> Result<C>
where
T: Send + 'static,
C: Send + 'static + std::iter::FromIterator<T>,
F: Send + FnMut(&rusqlite::Row) -> Result<T>,
{
let mut stmt = self.prepare(sql)?;
let res = stmt.query_and_then(params, f)?;
res.collect()
}
}
/// Creates a new SQLite connection.
///
/// `path` is the database path.
///
/// `passphrase` is the SQLCipher database passphrase.
/// Empty string if database is not encrypted.
fn new_connection(path: &Path) -> Result<Connection> {
fn new_connection(path: &Path, passphrase: &str) -> Result<Connection> {
let flags = OpenFlags::SQLITE_OPEN_NO_MUTEX
| OpenFlags::SQLITE_OPEN_READ_WRITE
| OpenFlags::SQLITE_OPEN_CREATE;
@@ -631,6 +760,9 @@ fn new_connection(path: &Path) -> Result<Connection> {
conn.busy_timeout(Duration::ZERO)?;
}
if !passphrase.is_empty() {
conn.pragma_update(None, "key", passphrase)?;
}
// Try to enable auto_vacuum. This will only be
// applied if the database is new or after successful
// VACUUM, which usually happens before backup export.
+12 -6
View File
@@ -657,7 +657,7 @@ pub(crate) async fn msgs_to_key_contacts(context: &Context) -> Result<()> {
return Ok(());
}
let trans_fn = |t: &mut rusqlite::Transaction| {
let mut first_key_contacts_msg_id: u32 = t
let mut first_key_contacts_msg_id: u64 = t
.query_one(
"SELECT CAST(value AS INTEGER) FROM config WHERE keyname='first_key_contacts_msg_id'",
(),
@@ -681,10 +681,10 @@ pub(crate) async fn msgs_to_key_contacts(context: &Context) -> Result<()> {
)
.context("Prepare stmt")?;
let msgs_to_migrate = 1000;
let mut msgs_migrated: u32 = 0;
let mut msgs_migrated: u64 = 0;
while first_key_contacts_msg_id > 0 && msgs_migrated < msgs_to_migrate {
let start_msg_id = first_key_contacts_msg_id.saturating_sub(msgs_to_migrate);
let cnt: u32 = stmt
let cnt: u64 = stmt
.execute((start_msg_id, first_key_contacts_msg_id))
.context("UPDATE msgs")?
.try_into()?;
@@ -2675,9 +2675,15 @@ CREATE TABLE smtp2 (
inc_and_check(&mut migration_version, 167)?;
if dbversion < migration_version {
// The default relay candidates are no longer seeded into the table.
sql.execute_migration("DELETE FROM relay_candidates", migration_version)
.await?;
sql.execute_migration(
"DELETE FROM relay_candidates;
CREATE TABLE relay_candidates_last_tried(
host TEXT PRIMARY KEY NOT NULL,
last_tried INTEGER NOT NULL DEFAULT 0
) STRICT",
migration_version,
)
.await?;
}
let new_version = sql
+3 -3
View File
@@ -16,15 +16,15 @@ async fn test_clear_config_cache() -> anyhow::Result<()> {
// This test checks that the config cache is invalidated in `execute_migration()`.
let t = TestContext::new().await;
assert_eq!(t.sql.get_raw_config_bool("cached_key").await?, false);
assert_eq!(t.get_config_bool(Config::IsChatmail).await?, false);
t.sql
.execute_migration(
"INSERT INTO config (keyname, value) VALUES ('cached_key', '1')",
"INSERT INTO config (keyname, value) VALUES ('is_chatmail', '1')",
1000,
)
.await?;
assert_eq!(t.sql.get_raw_config_bool("cached_key").await?, true);
assert_eq!(t.get_config_bool(Config::IsChatmail).await?, true);
assert_eq!(t.sql.get_raw_config_int(VERSION_CFG).await?.unwrap(), 1000);
Ok(())
+96 -3
View File
@@ -83,7 +83,7 @@ async fn test_housekeeping_db_closed() {
t.sql.close().await;
housekeeping(&t).await.unwrap(); // housekeeping should emit warnings but not fail
t.sql.open(&t).await.unwrap();
t.sql.open(&t, "".to_string()).await.unwrap();
let a = t.get_config(Config::Selfavatar).await.unwrap().unwrap();
assert_eq!(avatar_bytes, &tokio::fs::read(&a).await.unwrap()[..]);
@@ -155,11 +155,11 @@ async fn test_db_reopen() -> Result<()> {
let sql = Sql::new(dbfile);
// Create database with all the tables.
sql.open(&t).await.unwrap();
sql.open(&t, "".to_string()).await.unwrap();
sql.close().await;
// Reopen the database
sql.open(&t).await?;
sql.open(&t, "".to_string()).await?;
sql.execute(
"INSERT INTO config (keyname, value) VALUES (?, ?);",
("foo", "bar"),
@@ -209,6 +209,99 @@ async fn test_migration_flags() -> Result<()> {
Ok(())
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn test_check_passphrase() -> Result<()> {
use tempfile::tempdir;
// The context is used only for logging.
let t = TestContext::new().await;
// Create a separate empty database for testing.
let dir = tempdir()?;
let dbfile = dir.path().join("testdb.sqlite");
let sql = Sql::new(dbfile.clone());
sql.check_passphrase("foo".to_string()).await?;
sql.open(&t, "foo".to_string())
.await
.context("failed to open the database first time")?;
sql.close().await;
// Reopen the database
let sql = Sql::new(dbfile);
// Test that we can't open encrypted database without a passphrase.
assert!(sql.open(&t, "".to_string()).await.is_err());
// Now open the database with passpharse, it should succeed.
sql.check_passphrase("foo".to_string()).await?;
sql.open(&t, "foo".to_string())
.await
.context("failed to open the database second time")?;
Ok(())
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn test_sql_change_passphrase() -> Result<()> {
use tempfile::tempdir;
// The context is used only for logging.
let t = TestContext::new().await;
// Create a separate empty database for testing.
let dir = tempdir()?;
let dbfile = dir.path().join("testdb.sqlite");
let sql = Sql::new(dbfile.clone());
sql.open(&t, "foo".to_string())
.await
.context("failed to open the database first time")?;
sql.close().await;
// Change the passphrase from "foo" to "bar".
let sql = Sql::new(dbfile.clone());
sql.open(&t, "foo".to_string())
.await
.context("failed to open the database second time")?;
sql.change_passphrase("bar".to_string())
.await
.context("failed to change passphrase")?;
// Test that at least two connections are still working.
// This ensures that not only the connection which changed the password is working,
// but other connections as well.
{
let lock = sql.pool.read().await;
let pool = lock.as_ref().unwrap();
let query_only = true;
let conn1 = pool.get(query_only).await?;
let conn2 = pool.get(query_only).await?;
conn1
.query_row("SELECT count(*) FROM sqlite_master", [], |_row| Ok(()))
.unwrap();
conn2
.query_row("SELECT count(*) FROM sqlite_master", [], |_row| Ok(()))
.unwrap();
}
sql.close().await;
let sql = Sql::new(dbfile);
// Test that old passphrase is not working.
assert!(sql.open(&t, "foo".to_string()).await.is_err());
// Open the database with the new passphrase.
sql.check_passphrase("bar".to_string()).await?;
sql.open(&t, "bar".to_string())
.await
.context("failed to open the database third time")?;
sql.close().await;
Ok(())
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn test_query_only() -> Result<()> {
let t = TestContext::new().await;
+6 -13
View File
@@ -42,8 +42,7 @@ struct Statistics {
/// Size of the public key in bytes (encoded in binary, not base64).
pubkey_size: usize,
stats_id: String,
/// Whether all transports are chatmail relays, `None` if not known for all of them.
is_chatmail: Option<bool>,
is_chatmail: bool,
contact_stats: Vec<ContactStat>,
message_stats: BTreeMap<Chattype, MessageStats>,
securejoin_sources: SecurejoinSources,
@@ -69,7 +68,7 @@ struct ContactStat {
#[serde(skip_serializing_if = "is_false", rename = "direct_chat")]
single_chat: bool,
last_seen: i64,
last_seen: u64,
/// Whether the contact was established after stats-sending was enabled
#[serde(skip_serializing_if = "is_false")]
@@ -312,7 +311,7 @@ async fn ensure_last_old_contact_id(context: &Context) -> Result<()> {
return Ok(());
}
let last_contact_id: u32 = context
let last_contact_id: u64 = context
.sql
.query_get_value("SELECT MAX(id) FROM contacts", ())
.await?
@@ -351,22 +350,16 @@ async fn get_stats(context: &Context) -> Result<String> {
let sending_disabled_timestamps =
get_timestamps(context, "stats_sending_disabled_events").await?;
let number_of_transports = context.count_transports().await?;
let metadata = context.metadata.read().await;
let is_chatmail = (!metadata.is_empty() && metadata.len() == number_of_transports)
.then(|| metadata.values().all(|m| m.supports_push));
drop(metadata);
let stats = Statistics {
core_version: DC_VERSION_STR.to_string(),
number_of_transports,
number_of_transports: context.count_transports().await?,
key_create_timestamps,
number_of_keys,
key_version: self_public_key.primary_key.version().into(),
key_algorithm: format!("{:?}", self_public_key.algorithm()),
pubkey_size: DcKey::to_bytes(&self_public_key).len(),
stats_id: stats_id(context).await?,
is_chatmail,
is_chatmail: context.is_chatmail().await?,
contact_stats: get_contact_stats(context, last_old_contact).await?,
message_stats: get_message_stats(context).await?,
securejoin_sources: get_securejoin_source_stats(context).await?,
@@ -436,7 +429,7 @@ async fn get_contact_stats(context: &Context, last_old_contact: u32) -> Result<V
|row| {
let id = row.get(0)?;
let encrypted: bool = row.get(1)?;
let last_seen: i64 = row.get(2)?;
let last_seen: u64 = row.get(2)?;
let bot: bool = row.get(3)?;
Ok(ContactStat {
+7 -24
View File
@@ -4,13 +4,11 @@ use super::*;
use crate::chat::{
Chat, create_broadcast, create_group, create_group_unencrypted, get_chat_contacts,
};
use crate::imap::ServerMetadata;
use crate::mimeparser::SystemMessage;
use crate::qr::check_qr;
use crate::securejoin::{get_securejoin_qr, join_securejoin, join_securejoin_with_ux_info};
use crate::test_utils::{TestContext, TestContextManager, get_chat_msg};
use crate::tools::SystemTime;
use crate::transport::add_pseudo_transport;
use pretty_assertions::assert_eq;
use serde_json::{Number, Value};
@@ -465,33 +463,18 @@ async fn test_stats_securejoin_invites() -> Result<()> {
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn test_stats_is_chatmail() -> Result<()> {
async fn is_chatmail(context: &TestContext) -> Result<Value> {
let stats: Value = serde_json::from_str(&get_stats(context).await?)?;
Ok(stats.get("is_chatmail").unwrap().clone())
}
let alice = &TestContext::new_alice().await;
alice.set_config_bool(Config::StatsSending, true).await?;
assert!(is_chatmail(alice).await?.is_null());
alice.metadata.write().await.insert(
0,
ServerMetadata {
supports_push: true,
..Default::default()
},
);
assert_eq!(is_chatmail(alice).await?, Value::Bool(true));
let r = get_stats(alice).await?;
let r: serde_json::Value = serde_json::from_str(&r)?;
assert_eq!(r.get("is_chatmail").unwrap().as_bool().unwrap(), false);
add_pseudo_transport(alice, "alice@example.net").await?;
assert!(is_chatmail(alice).await?.is_null());
alice.set_config_bool(Config::IsChatmail, true).await?;
alice
.metadata
.write()
.await
.insert(1, ServerMetadata::default());
assert_eq!(is_chatmail(alice).await?, Value::Bool(false));
let r = get_stats(alice).await?;
let r: serde_json::Value = serde_json::from_str(&r)?;
assert_eq!(r.get("is_chatmail").unwrap().as_bool().unwrap(), true);
Ok(())
}
+8 -16
View File
@@ -56,18 +56,16 @@ pub async fn get_storage_usage(ctx: &Context) -> Result<StorageUsage> {
let blobdir_size =
tokio::task::spawn_blocking(move || get_blobdir_storage_usage(&context_clone));
let page_size: i64 = ctx
let page_size: u64 = ctx
.sql
.query_get_value("PRAGMA page_size", ())
.await?
.unwrap_or_default();
let page_size = u64::try_from(page_size)?;
let page_count: i64 = ctx
let page_count: u64 = ctx
.sql
.query_get_value("PRAGMA page_count", ())
.await?
.unwrap_or_default();
let page_count = u64::try_from(page_count)?;
let mut largest_tables = ctx
.sql
@@ -80,8 +78,7 @@ pub async fn get_storage_usage(ctx: &Context) -> Result<StorageUsage> {
(),
|row| {
let name: String = row.get(0)?;
let size: i64 = row.get(1)?;
let size: u64 = u64::try_from(size)?;
let size: u64 = row.get(1)?;
Ok((name, size, None))
},
)
@@ -89,13 +86,12 @@ pub async fn get_storage_usage(ctx: &Context) -> Result<StorageUsage> {
for row in &mut largest_tables {
let name = &row.0;
let row_count: Option<i64> = ctx
let row_count: Result<Option<u64>> = ctx
.sql
// SECURITY: the table name comes from the db, not from the user
.query_get_value(&format!("SELECT COUNT(*) FROM {name}"), ())
.await
.unwrap_or_default();
row.2 = row_count.map(|count| u64::try_from(count).unwrap_or_default());
.await;
row.2 = row_count.unwrap_or_default();
}
let largest_webxdc_data = ctx
@@ -107,12 +103,8 @@ pub async fn get_storage_usage(ctx: &Context) -> Result<StorageUsage> {
(),
|row| {
let msg_id: MsgId = row.get(0)?;
let size: i64 = row.get(1)?;
let count: i64 = row.get(2)?;
// This should never fail as the count cannot be negative.
let size: u64 = u64::try_from(size)?;
let count: u64 = u64::try_from(count)?;
let size: u64 = row.get(1)?;
let count: u64 = row.get(2)?;
Ok((msg_id, size, count))
},
+8 -13
View File
@@ -384,7 +384,6 @@ impl Context {
async fn sync_message_deletion(&self, msgs: &[String]) -> Result<()> {
let mut modified_chat_ids = BTreeSet::new();
let mut pinned_messages_changed_chat_ids = BTreeSet::new();
let mut msg_ids = Vec::new();
for rfc724_mid in msgs {
if let Some(msg_id) = message::rfc724_mid_exists(self, rfc724_mid).await? {
@@ -392,9 +391,6 @@ impl Context {
message::delete_msg_locally(self, &msg).await?;
msg_ids.push(msg.id);
modified_chat_ids.insert(msg.chat_id);
if msg.is_pinned() {
pinned_messages_changed_chat_ids.insert(msg.chat_id);
}
} else {
warn!(self, "Sync message delete: Database entry does not exist.");
}
@@ -402,13 +398,7 @@ impl Context {
warn!(self, "Sync message delete: {rfc724_mid:?} not found.");
}
}
message::delete_msgs_locally_done(
self,
&msg_ids,
modified_chat_ids,
pinned_messages_changed_chat_ids,
)
.await?;
message::delete_msgs_locally_done(self, &msg_ids, modified_chat_ids).await?;
Ok(())
}
}
@@ -713,7 +703,9 @@ mod tests {
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn test_send_sync_msg_enables_bccself() -> Result<()> {
for sync_message_sent in [false, true] {
for (chatmail, sync_message_sent) in
[(false, false), (false, true), (true, false), (true, true)]
{
let alice1 = TestContext::new_alice().await;
let alice2 = TestContext::new_alice().await;
@@ -722,6 +714,9 @@ mod tests {
alice1.set_config_bool(Config::SyncMsgs, true).await?;
alice2.set_config_bool(Config::SyncMsgs, true).await?;
alice1.set_config_bool(Config::IsChatmail, chatmail).await?;
alice2.set_config_bool(Config::IsChatmail, chatmail).await?;
alice1.set_config_bool(Config::BccSelf, true).await?;
alice2.set_config_bool(Config::BccSelf, false).await?;
@@ -740,7 +735,7 @@ mod tests {
alice1.send_text(chat.id, "Hi").await
};
// BccSelf defaults to false.
// On chatmail accounts, BccSelf defaults to false.
// When receiving a sync message from another device,
// there obviously is a multi-device-setup, and BccSelf
// should be enabled.
+3 -27
View File
@@ -35,7 +35,7 @@ use crate::context::Context;
use crate::events::{Event, EventEmitter, EventType, Events};
use crate::key::{self, DcKey, self_fingerprint};
use crate::message::{Message, MessageState, MsgId};
use crate::mimefactory::{self, MimeFactory};
use crate::mimefactory;
use crate::mimeparser::{MimeMessage, SystemMessage};
use crate::pgp::SeipdVersion;
use crate::receive_imf::{ReceivedMsg, receive_imf};
@@ -557,9 +557,6 @@ impl TestContext {
/// The context will be configured but the key will not be pre-generated so if a key is
/// used the fingerprint will be different every time.
pub async fn configure_addr(&self, addr: &str) {
add_pseudo_transport(&self.ctx, addr)
.await
.expect("Failed to add pseudo transport");
self.ctx
.set_config(Config::ConfiguredAddr, Some(addr))
.await
@@ -627,20 +624,15 @@ ORDER BY id"
.ctx
.sql
.transaction_ext(query_only, |transaction| {
smtp::queue::load_queued_mail(transaction, rowid)
smtp::load_queued_mail(transaction, rowid)
})
.await
.expect("Failed to load queued mail");
if queued_mail.bcc_self {
let from = self
.get_primary_self_addr()
.await
.expect("Cannot get From address");
smtp::add_self_recipients(
&self.ctx,
&mut queued_mail.recipients,
queued_mail.encryption.is_encrypted(),
from,
)
.await
.expect("Failed to add self recipients");
@@ -729,20 +721,15 @@ ORDER BY id"
.ctx
.sql
.transaction_ext(query_only, |transaction| {
smtp::queue::load_queued_mail(transaction, rowid)
smtp::load_queued_mail(transaction, rowid)
})
.await
.expect("Failed to load queued mail");
if queued_mail.bcc_self {
let from = self
.get_primary_self_addr()
.await
.expect("Cannot get self address");
smtp::add_self_recipients(
&self.ctx,
&mut queued_mail.recipients,
queued_mail.encryption.is_encrypted(),
from,
)
.await
.expect("Failed to add self recipients");
@@ -850,17 +837,6 @@ ORDER BY id"
assert_eq!(received.chat_id, ChatId::TRASH);
}
/// Receives a read receipt from `reader`, who received `msg`.
pub async fn recv_mdn(&self, reader: &TestContext, msg: &Message) -> Result<()> {
let mdn = MimeFactory::from_mdn(reader, msg.from_id, msg.rfc724_mid.clone(), vec![])
.await?
.render(reader, &reader.get_primary_self_addr().await?)
.await?
.message;
receive_imf(self, mdn.as_bytes(), false).await?;
Ok(())
}
/// Gets the most recent message ID of a chat.
///
/// Panics on errors or if the most recent message is a marker.
+26
View File
@@ -520,6 +520,32 @@ where
}
}
pub(crate) trait ToOption<T> {
fn to_option(self) -> Option<T>;
}
impl<'a> ToOption<&'a str> for &'a String {
fn to_option(self) -> Option<&'a str> {
if self.is_empty() { None } else { Some(self) }
}
}
impl ToOption<String> for u16 {
fn to_option(self) -> Option<String> {
if self == 0 {
None
} else {
Some(self.to_string())
}
}
}
impl ToOption<String> for Option<i32> {
fn to_option(self) -> Option<String> {
match self {
None | Some(0) => None,
Some(v) => Some(v.to_string()),
}
}
}
#[expect(clippy::arithmetic_side_effects)]
pub(crate) fn remove_subject_prefix(last_subject: &str) -> String {
let subject_start = if last_subject.starts_with("Chat:") {
+30 -2
View File
@@ -11,7 +11,7 @@
use std::fmt;
use std::sync::atomic::Ordering;
use anyhow::{Context as _, Result, format_err};
use anyhow::{Context as _, Result, bail, format_err};
use deltachat_contact_tools::{EmailAddress, addr_normalize};
use rusqlite::OptionalExtension;
use serde::{Deserialize, Serialize};
@@ -260,6 +260,34 @@ impl fmt::Display for ConfiguredLoginParam {
}
impl ConfiguredLoginParam {
/// Load configured account settings from the database.
///
/// Returns transport ID and configured parameters
/// of the transport currently used for sending.
/// Returns `None` if account is not configured.
pub(crate) async fn load(context: &Context) -> Result<Option<(u32, Self)>> {
let Some(self_addr) = context.get_config(Config::ConfiguredAddr).await? else {
return Ok(None);
};
let Some((id, json)) = context
.sql
.query_row_optional(
"SELECT id, configured_param FROM transports WHERE addr=?",
(&self_addr,),
|row| {
let id: u32 = row.get(0)?;
let json: String = row.get(1)?;
Ok((id, json))
},
)
.await?
else {
bail!("Self address {self_addr} doesn't have a corresponding transport");
};
Ok(Some((id, Self::from_json(&json)?)))
}
/// Loads configured login parameters for all transports.
///
/// Returns a vector of all transport IDs
@@ -737,7 +765,7 @@ pub(crate) fn maybe_update_sending_transport(
}
/// Adds transport entry to the `transports` table with empty configuration.
pub async fn add_pseudo_transport(context: &Context, addr: &str) -> Result<()> {
pub(crate) async fn add_pseudo_transport(context: &Context, addr: &str) -> Result<()> {
context.sql
.execute(
"INSERT OR IGNORE INTO transports (addr, entered_param, configured_param) VALUES (?, ?, ?)",
+39 -34
View File
@@ -62,7 +62,7 @@ async fn test_save_load_login_param() -> Result<()> {
expected_param
);
assert_eq!(t.is_configured().await?, true);
let (_transport_id, loaded) = ConfiguredLoginParam::load_all(&t).await?.remove(0);
let (_transport_id, loaded) = ConfiguredLoginParam::load(&t).await?.unwrap();
assert_eq!(param, loaded);
let formatted = format!(" {loaded}");
@@ -75,7 +75,7 @@ async fn test_save_load_login_param() -> Result<()> {
// Legacy ConfiguredImapCertificateChecks config is ignored
t.set_config(Config::ConfiguredImapCertificateChecks, Some("999"))
.await?;
assert!(ConfiguredLoginParam::load_all(&t).await.is_ok());
assert!(ConfiguredLoginParam::load(&t).await.is_ok());
// Test that we don't panic on unknown ConfiguredImapCertificateChecks values.
let wrong_param = expected_param.replace("Strict", "Stricct");
@@ -83,7 +83,7 @@ async fn test_save_load_login_param() -> Result<()> {
t.sql
.execute("UPDATE transports SET configured_param=?", (wrong_param,))
.await?;
assert!(ConfiguredLoginParam::load_all(&t).await.is_err());
assert!(ConfiguredLoginParam::load(&t).await.is_err());
Ok(())
}
@@ -226,11 +226,10 @@ async fn test_delete_transport() -> Result<()> {
Ok(())
}
/// Tests that selecting sending transport by setting "configured_addr" does not
/// send the sync message, is not synchronized between devices even if sync message is forced,
/// and does not bump sending transport `add_timestamp`.
/// Tests that promoting a transport bumps its `add_timestamp` on other devices
/// even if it was added within the same second.
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn test_no_configured_addr_synchronization() -> Result<()> {
async fn test_promote_transport_same_second() -> Result<()> {
let mut tcm = TestContextManager::new();
let alice = &tcm.alice().await;
let alice2 = &tcm.alice().await;
@@ -239,35 +238,11 @@ async fn test_no_configured_addr_synchronization() -> Result<()> {
a.set_config_bool(Config::BccSelf, true).await?;
}
let addr = "alice@otherprovider.com";
add_dummy_transport(alice, addr).await?;
add_dummy_transport(alice, "alice@otherprovider.com").await?;
send_sync_transports(alice).await?;
sync_and_check_recipients(alice, alice2, &format!("{addr} alice@example.org")).await;
sync_and_check_recipients(alice, alice2, "alice@otherprovider.com alice@example.org").await;
// Selects `addr` on `alice` as the sending transport
// and syncs the transport update to `alice2`,
// whose own sending transport must stay unchanged.
let old_timestamp = add_timestamp(alice2, addr).await;
let alice2_primary = alice2.get_config(Config::ConfiguredAddr).await?;
alice.set_config(Config::ConfiguredAddr, Some(addr)).await?;
assert_eq!(add_timestamp(alice, addr).await, old_timestamp);
send_sync_transports(alice).await?;
alice.send_sync_msg().await?.unwrap();
let sync_msg = alice.pop_sent_msg().await;
assert_eq!(sync_msg.recipients, format!("alice@example.org {addr}"));
// The sync message comes from the new sending address,
// which must not make `alice2` adopt it as its own sending address.
assert!(sync_msg.payload.contains(&format!("From: <{addr}>")));
alice2.recv_msg_trash(&sync_msg).await;
// add_timestamp must not change.
assert_eq!(add_timestamp(alice2, addr).await, old_timestamp);
assert_eq!(
alice2.get_config(Config::ConfiguredAddr).await?,
alice2_primary
);
Ok(())
promote_transport_and_sync(alice, alice2, "alice@otherprovider.com").await
}
/// Tests that `sync_transports()` requests an IO restart
@@ -288,6 +263,36 @@ async fn test_sync_transports_requests_io_restart() -> Result<()> {
Ok(())
}
/// Promotes `addr` on `alice` and syncs the transport update to `alice2`,
/// whose own primary transport must stay unchanged.
async fn promote_transport_and_sync(
alice: &TestContext,
alice2: &TestContext,
addr: &str,
) -> Result<()> {
let old_timestamp = add_timestamp(alice2, addr).await;
let alice2_primary = alice2.get_config(Config::ConfiguredAddr).await?;
alice.set_config(Config::ConfiguredAddr, Some(addr)).await?;
assert!(add_timestamp(alice, addr).await > old_timestamp);
alice.send_sync_msg().await?.unwrap();
let sync_msg = alice.pop_sent_msg().await;
assert_eq!(sync_msg.recipients, format!("alice@example.org {addr}"));
// The sync message comes from the new primary,
// which must not make `alice2` adopt it as its own primary.
assert!(sync_msg.payload.contains(&format!("From: <{addr}>")));
alice2.recv_msg_trash(&sync_msg).await;
// add_timestamp must monotonically increase because
// other devices ignore the change otherwise.
assert!(add_timestamp(alice2, addr).await > old_timestamp);
assert_eq!(
alice2.get_config(Config::ConfiguredAddr).await?,
alice2_primary
);
Ok(())
}
async fn add_timestamp(t: &TestContext, addr: &str) -> i64 {
t.sql
.query_get_value("SELECT add_timestamp FROM transports WHERE addr=?", (addr,))
@@ -2,4 +2,5 @@ OutBroadcast#Chat#1001: My Channel [1 member(s)]🔇 Icon: e9b6c7a78aa2e4f415644
--------------------------------------------------------------------------------
Msg#1001: info (Contact#Contact#Info): Messages are end-to-end encrypted. [NOTICED][INFO]
Msg#1002🔒: Me (Contact#Contact#Self): Channel image changed. [INFO] √
Msg#1005🔒: Me (Contact#Contact#Self): Member bob@example.net added. [INFO] √
--------------------------------------------------------------------------------
@@ -1,6 +1,7 @@
OutBroadcast#Chat#1001: Channel [0 member(s)]🔇
--------------------------------------------------------------------------------
Msg#1001: info (Contact#Contact#Info): Messages are end-to-end encrypted. [NOTICED][INFO]
Msg#1006🔒: Me (Contact#Contact#Self): Member bob@example.net added. [INFO] √
Msg#1009🔒: Me (Contact#Contact#Self): hi √
Msg#1010🔒: Me (Contact#Contact#Self): You removed member bob@example.net. [INFO] √
--------------------------------------------------------------------------------
@@ -1,6 +1,7 @@
OutBroadcast#Chat#1001: Channel [0 member(s)]🔇
--------------------------------------------------------------------------------
Msg#1002: info (Contact#Contact#Info): Messages are end-to-end encrypted. [NOTICED][INFO]
Msg#1006🔒: Me (Contact#Contact#Self): Member bob@example.net added. [INFO] √
Msg#1009🔒: Me (Contact#Contact#Self): hi √
Msg#1010🔒: Me (Contact#Contact#Self): You removed member bob@example.net. [INFO] √
--------------------------------------------------------------------------------