Compare commits

..
Author SHA1 Message Date
holger krekel d9b0470da8 fix: (linux-only) make sure large attachments return memory to kernel
The test fails on main with ~150MB memory remaining allocated
after sending/receiving two 20MB messages,
and stays well below when explicitely setting libc's M_MMAP_THRESHOLD.
2026-10-04 20:40:36 +02:00
24 changed files with 219 additions and 329 deletions
+4 -5
View File
@@ -30,9 +30,6 @@ jobs:
name: Lint Rust
runs-on: ubuntu-latest
timeout-minutes: 60
env:
# Tests always unwind: match it so dependencies are only checked once.
CARGO_PROFILE_DEV_PANIC: unwind
steps:
- uses: actions/checkout@v7
with:
@@ -51,6 +48,8 @@ jobs:
run: cargo fmt --all -- --check
- name: Run clippy
run: scripts/clippy.sh
- name: Check with all features
run: cargo check --workspace --all-targets --all-features
- name: Check with only default features
run: cargo check --all-targets
@@ -140,12 +139,12 @@ jobs:
- name: Tests
env:
RUST_BACKTRACE: 1
run: cargo nextest run --workspace --exclude deltachat-jsonrpc-bindings --locked
run: cargo nextest run --workspace --locked
- name: Doc-Tests
env:
RUST_BACKTRACE: 1
run: cargo test --workspace --exclude deltachat-jsonrpc-bindings --locked --doc
run: cargo test --workspace --locked --doc
- name: Test cargo vendor
run: cargo vendor
Generated
+4 -2
View File
@@ -1406,6 +1406,7 @@ dependencies = [
"tokio-stream",
"tokio-util",
"toml",
"tracing",
"url",
"uuid",
"walkdir",
@@ -1484,6 +1485,7 @@ dependencies = [
"deltachat",
"deltachat-jsonrpc",
"futures-lite",
"libc",
"log",
"serde",
"serde_json",
@@ -4246,9 +4248,9 @@ dependencies = [
[[package]]
name = "pgp"
version = "0.21.0"
version = "0.20.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "ae70f4d9325a391db30d115d6191a7bb67cec856aa0bba131290f0c4e09a532a"
checksum = "1cfa4743b28656065ff4c0ba09e46b357a65e8c00fc2341e89084b82f87cbdf1"
dependencies = [
"aead",
"aes",
+2 -1
View File
@@ -77,7 +77,7 @@ num-derive = "0.4"
num-traits = { workspace = true }
parking_lot = "0.12.4"
percent-encoding = "2.3"
pgp = { version = "0.21.0", features = ["pqc"], default-features = false }
pgp = { version = "0.20.0", features = ["draft-pqc"], default-features = false }
pin-project = "1"
qrcodegen = "1.7.0"
quick-xml = { version = "0.41", features = ["escape-html"] }
@@ -104,6 +104,7 @@ astral-tokio-tar = { version = "0.7.0", default-features = false }
tokio-util = { workspace = true }
tokio = { workspace = true, features = ["fs", "rt-multi-thread", "macros"] }
toml = "0.9"
tracing = "0.1.41"
url = "2"
uuid = { version = "1", features = ["serde", "v4"] }
walkdir = "2.5.0"
+22
View File
@@ -3543,6 +3543,28 @@ uint32_t dc_chatlist_get_msg_id (const dc_chatlist_t* chatlist, siz
dc_lot_t* dc_chatlist_get_summary (const dc_chatlist_t* chatlist, size_t index, dc_chat_t* chat);
/**
* Create a chatlist summary item when the chatlist object is already unref()'d.
*
* This function is similar to dc_chatlist_get_summary(), however,
* it takes the chat ID and the message ID as returned by dc_chatlist_get_chat_id() and dc_chatlist_get_msg_id()
* as arguments. The chatlist object itself is not needed directly.
*
* This maybe useful if you convert the complete object into a different representation
* as done e.g. in the node-bindings.
* If you have access to the chatlist object in some way, using this function is not recommended,
* use dc_chatlist_get_summary() in this case instead.
*
* @memberof dc_context_t
* @param context The context object.
* @param chat_id The chat ID to get a summary for.
* @param msg_id The message ID to get a summary for.
* @return The summary as an dc_lot_t object, see dc_chatlist_get_summary() for details.
* Must be freed using dc_lot_unref(). NULL is never returned.
*/
dc_lot_t* dc_chatlist_get_summary2 (dc_context_t* context, uint32_t chat_id, uint32_t msg_id);
/**
* Get info summary for a chat, in JSON format.
*
+29
View File
@@ -47,6 +47,7 @@ mod dc_array;
mod lot;
mod string;
use deltachat::chatlist::Chatlist;
use self::string::*;
@@ -2843,6 +2844,34 @@ pub unsafe extern "C" fn dc_chatlist_get_summary(
Box::into_raw(Box::new(summary.into()))
}
#[unsafe(no_mangle)]
pub unsafe extern "C" fn dc_chatlist_get_summary2(
context: *mut dc_context_t,
chat_id: u32,
msg_id: u32,
) -> *mut dc_lot_t {
if context.is_null() {
eprintln!("ignoring careless call to dc_chatlist_get_summary2()");
return ptr::null_mut();
}
let ctx = unsafe { &*context };
let msg_id = if msg_id == 0 {
None
} else {
Some(MsgId::new(msg_id))
};
let summary = block_on(Chatlist::get_summary2(
ctx,
ChatId::new(chat_id),
msg_id,
None,
))
.context("get_summary2 failed")
.log_err(ctx)
.unwrap_or_default();
Box::into_raw(Box::new(summary.into()))
}
// dc_chat_t
/// FFI struct for [dc_chat_t]
@@ -31,6 +31,8 @@ pub enum ChatListItemFetchResult {
summary_text1: String,
summary_text2: String,
summary_status: u32,
/// showing preview if last chat message is image
summary_preview_image: Option<String>,
/// True if the chat is encrypted.
/// This means that all messages in the chat are encrypted,
@@ -101,6 +103,8 @@ pub(crate) async fn get_chat_list_item_by_id(
let summary_text1 = summary.prefix.map_or_else(String::new, |s| s.to_string());
let summary_text2 = summary.text.to_owned();
let summary_preview_image = summary.thumbnail_path;
let visibility = chat.get_visibility();
let avatar_path = chat
@@ -153,6 +157,7 @@ pub(crate) async fn get_chat_list_item_by_id(
summary_text1,
summary_text2,
summary_status: summary.state.to_u32().expect("impossible"), // idea and a function to transform the constant to strings? or return string enum
summary_preview_image,
is_encrypted: chat.is_encrypted(ctx).await?,
is_group: chat.get_type() == Chattype::Group,
fresh_message_counter,
+33
View File
@@ -0,0 +1,33 @@
import os
import sys
import pytest
def anonymous_mib(pid):
with open(f"/proc/{pid}/smaps_rollup") as f:
for line in f:
if line.startswith("Anonymous:"):
return int(line.split()[1]) // 1024
raise LookupError("Anonymous")
@pytest.mark.skipif(sys.platform != "linux", reason="reads /proc")
def test_attachment_memory_is_returned(acf, rpc, tmp_path):
# See also comments for `tune_malloc` in `deltachat-rpc-server/src/main.rs`
ac1, ac2 = acf.get_online_accounts(2)
chat1 = acf.get_accepted_chat(ac1, ac2)
chat2 = ac2.create_chat(ac1)
blob = tmp_path / "blob.bin"
blob.write_bytes(os.urandom(20 << 20))
before = anonymous_mib(rpc.process.pid)
for sender_chat, receiver in ((chat1, ac2), (chat2, ac1)):
sender_chat.send_file(str(blob))
event = receiver.wait_for_incoming_msg_event()
assert receiver.get_message_by_id(event.msg_id).get_snapshot().file_bytes == 20 << 20
for ac in (ac1, ac2):
rpc.wait_for_all_work_done(ac.id)
grown = anonymous_mib(rpc.process.pid) - before
assert grown < 64, f"the server kept {grown} MiB after two 20 MiB attachments"
+1
View File
@@ -14,6 +14,7 @@ deltachat = { workspace = true }
anyhow = { workspace = true }
futures-lite = { workspace = true }
libc = { workspace = true }
log = { workspace = true }
serde_json = { workspace = true }
serde = { workspace = true, features = ["derive"] }
+24
View File
@@ -22,8 +22,32 @@ use tokio::task::JoinHandle;
use tokio_util::sync::CancellationToken;
use yerpc::{RpcClient, RpcSession};
/// Pins Linux glibc's mmap threshold so that freed message buffers go back to the kernel.
///
/// See M_MMAP_THRESHOLD in <https://man7.org/linux/man-pages/man3/mallopt.3.html>:
/// glibc by default starts with a M_MMAP_THRESHOLD threshold of 128 KiB
/// but raises it to the size of every freed block that exceeds it,
/// up to 32 MiB on 64-bit systems,
/// and trims the heap only from its top end once twice that much is free.
/// Fixating the threshold disables the adjustment: allocations at or above it
/// that the free list cannot satisfy are mmapped and unmapped on free,
/// at the price of the kernel zeroing each such buffer after unmap.
/// Large message processing (allocations above 128KiB) very slightly slows
/// down to the kernel zeroing the buffers, but it's hardly measurable,
/// while overall process memory allocation behaviour significantly improves.
#[cfg(all(target_os = "linux", target_env = "gnu"))]
fn tune_malloc() {
unsafe {
libc::mallopt(libc::M_MMAP_THRESHOLD, 128 * 1024);
}
}
#[cfg(not(all(target_os = "linux", target_env = "gnu")))]
fn tune_malloc() {}
#[tokio::main(flavor = "multi_thread")]
async fn main() {
tune_malloc();
// Logs from `log` crate and traces from `tracing` crate
// are configurable with `RUST_LOG` environment variable
// and go to stderr to avoid interfering with JSON-RPC using stdout.
-15
View File
@@ -456,21 +456,6 @@ CREATE TABLE smtp_status_updates (
descr TEXT NOT NULL -- text to send along with the updates
);
-- Table to record the successful usage transports for sending.
-- Sorting the table by rowid in descending order
-- returns most recently successfully used transport first.
CREATE TABLE smtp_success (
-- Sequentially increasing ID of the success.
-- Transport with the highest ID is to be used first.
id INTEGER PRIMARY KEY AUTOINCREMENT NOT NULL,
-- ID of the transport that was used to send a message.
transport_id INTEGER UNIQUE NOT NULL,
-- Delete `smtp_success` rows when the transport is deleted.
FOREIGN KEY(transport_id) REFERENCES transports(id) ON DELETE CASCADE
) STRICT;
-- Table of "sync items" to be grouped into sync messages
-- and sent to own devices.
CREATE TABLE multi_device_sync (
+1 -1
View File
@@ -6,4 +6,4 @@
#
# To automatically fix warnings, run
# scripts/clippy.sh --fix --allow-dirty
cargo clippy --locked --workspace --exclude deltachat-jsonrpc-bindings --all-targets --all-features "$@" -- -D warnings
cargo clippy --locked --workspace --all-targets --all-features "$@" -- -D warnings
+22 -3
View File
@@ -76,8 +76,12 @@ impl Accounts {
Accounts::open(events, dir, writable).await
}
fn log_info(&self, file: &str, line: u32, msg: String) {
self.emit_event(EventType::Info(format!("{file}:{line}: {msg}")));
/// Get the ID used to log events.
///
/// Account manager logs events with ID 0
/// which is not used by any accounts.
fn get_id(&self) -> u32 {
0
}
/// Ensures the accounts directory and config file exist.
@@ -391,6 +395,11 @@ impl Accounts {
"Starting background fetch for {n_accounts} accounts."
)),
});
::tracing::event!(
::tracing::Level::INFO,
account_id = 0,
"Starting background fetch for {n_accounts} accounts."
);
let mut set = JoinSet::new();
for account in accounts {
set.spawn(async move {
@@ -406,6 +415,11 @@ impl Accounts {
"Finished background fetch for {n_accounts} accounts."
)),
});
::tracing::event!(
::tracing::Level::INFO,
account_id = 0,
"Finished background fetch for {n_accounts} accounts."
);
}
/// Auxiliary function for [Accounts::background_fetch].
@@ -448,6 +462,11 @@ impl Accounts {
id: 0,
typ: EventType::Warning("Background fetch timed out.".to_string()),
});
::tracing::event!(
::tracing::Level::WARN,
account_id = 0,
"Background fetch timed out."
);
}
events.emit(Event {
id: 0,
@@ -530,7 +549,7 @@ impl Accounts {
}
}
/// Emits a single event with ID 0, which is not used by any accounts.
/// Emits a single event.
pub fn emit_event(&self, event: EventType) {
self.events.emit(Event { id: 0, typ: event })
}
+4 -1
View File
@@ -125,11 +125,14 @@ pub trait DcKey: Serialize + Deserializable + Clone {
/// Converts secret key to public key.
pub(crate) fn secret_key_to_public_key(
context: &Context,
mut signed_secret_key: SignedSecretKey,
timestamp: u32,
addr: &str,
relay_addrs: &str,
) -> Result<SignedPublicKey> {
info!(context, "Converting secret key to public key.");
// Make sure timestamp of created signatures
// is not in the past compared to the primary key timestamp.
let timestamp = std::cmp::max(
@@ -302,7 +305,7 @@ pub(crate) async fn load_self_public_key_opt(context: &Context) -> Result<Option
let addr = context.get_primary_self_addr().await?;
let all_addrs = context.get_self_addrs().await?.join(",");
let signed_public_key =
secret_key_to_public_key(signed_secret_key, timestamp, &addr, &all_addrs)?;
secret_key_to_public_key(context, signed_secret_key, timestamp, &addr, &all_addrs)?;
*lock = Some(signed_public_key.clone());
Ok(Some(signed_public_key))
+30 -23
View File
@@ -3,7 +3,6 @@
#![allow(missing_docs)]
use crate::context::Context;
use crate::events::EventType;
mod stream;
@@ -13,9 +12,15 @@ macro_rules! info {
($ctx:expr, $msg:expr) => {
info!($ctx, $msg,)
};
($ctx:expr, $msg:expr, $($args:expr),* $(,)?) => {
$ctx.log_info(file!(), line!(), format!($msg, $($args),*))
};
($ctx:expr, $msg:expr, $($args:expr),* $(,)?) => {{
let formatted = format!($msg, $($args),*);
let full = format!("{file}:{line}: {msg}",
file = file!(),
line = line!(),
msg = &formatted);
::tracing::event!(::tracing::Level::INFO, account_id = $ctx.get_id(), "{}", &formatted);
$ctx.emit_event($crate::EventType::Info(full));
}};
}
// Workaround for <https://github.com/rust-lang/rust/issues/133708>.
@@ -25,9 +30,15 @@ mod warn_macro_mod {
($ctx:expr, $msg:expr) => {
warn_macro!($ctx, $msg,)
};
($ctx:expr, $msg:expr, $($args:expr),* $(,)?) => {
$ctx.log_warn(file!(), line!(), format!($msg, $($args),*))
};
($ctx:expr, $msg:expr, $($args:expr),* $(,)?) => {{
let formatted = format!($msg, $($args),*);
let full = format!("{file}:{line}: {msg}",
file = file!(),
line = line!(),
msg = &formatted);
::tracing::event!(::tracing::Level::WARN, account_id = $ctx.get_id(), "{}", &formatted);
$ctx.emit_event($crate::EventType::Warning(full));
}};
}
pub(crate) use warn_macro;
@@ -39,25 +50,15 @@ macro_rules! error {
($ctx:expr, $msg:expr) => {
error!($ctx, $msg,)
};
($ctx:expr, $msg:expr, $($args:expr),* $(,)?) => {
$ctx.log_error(format!($msg, $($args),*))
};
($ctx:expr, $msg:expr, $($args:expr),* $(,)?) => {{
let formatted = format!($msg, $($args),*);
::tracing::event!(::tracing::Level::ERROR, account_id = $ctx.get_id(), "{}", &formatted);
$ctx.set_last_error(&formatted);
$ctx.emit_event($crate::EventType::Error(formatted));
}};
}
impl Context {
pub(crate) fn log_info(&self, file: &str, line: u32, msg: String) {
self.emit_event(EventType::Info(format!("{file}:{line}: {msg}")));
}
pub(crate) fn log_warn(&self, file: &str, line: u32, msg: String) {
self.emit_event(EventType::Warning(format!("{file}:{line}: {msg}")));
}
pub(crate) fn log_error(&self, msg: String) {
self.set_last_error(&msg);
self.emit_event(EventType::Error(msg));
}
/// Set last error string.
/// Implemented as blocking as used from macros in different, not always async blocks.
pub fn set_last_error(&self, error: &str) {
@@ -115,6 +116,12 @@ impl<T, E: std::fmt::Display> LogExt<T, E> for Result<T, E> {
);
// We can't use the warn!() macro here as the file!() and line!() macros
// don't work with #[track_caller]
tracing::event!(
::tracing::Level::WARN,
account_id = context.get_id(),
"{}",
&full
);
context.emit_event(crate::EventType::Warning(full));
};
self
+5
View File
@@ -93,6 +93,11 @@ impl<S: SessionStream> AsyncRead for LoggingStream<S> {
"Read error on stream {peer_addr:?} after reading {} and writing {} bytes: {err}.",
this.metrics.total_read, this.metrics.total_written
);
tracing::event!(
::tracing::Level::WARN,
account_id = *this.account_id,
log_message
);
this.events.emit(Event {
id: *this.account_id,
typ: EventType::Warning(log_message),
+1
View File
@@ -366,6 +366,7 @@ async fn test_mdn_sent_to_all_relays() -> Result<()> {
// Bob's key gets a second relay address and Alice merges the newer key.
let bob_secret_key = load_self_secret_key(bob).await?;
let bob_public_key = secret_key_to_public_key(
bob,
bob_secret_key,
u32::try_from(time())? + 100,
"bob@example.net",
+6 -58
View File
@@ -1,6 +1,5 @@
//! OpenPGP helper module using [rPGP facilities](https://github.com/rpgp/rpgp).
use std::cmp::Ordering;
use std::collections::{HashMap, HashSet};
use std::io::Cursor;
@@ -84,62 +83,13 @@ pub(crate) fn create_keypair(addr: EmailAddress) -> Result<SignedSecretKey> {
/// Selects a subkey of the public key to use for encryption.
///
/// The key is selected according to
/// <https://www.ietf.org/archive/id/draft-autocrypt-openpgp-v2-cert-03.html#section-4.3-4>.
/// If multiple keys are available, the one that will expire sooner is selected.
///
/// Returns `None` if the public key cannot be used for encryption.
fn select_pk_for_encryption(now: u32, key: &SignedPublicKey) -> Option<&SignedPublicSubKey> {
///
/// TODO: take key flags and expiration dates into account
fn select_pk_for_encryption(key: &SignedPublicKey) -> Option<&SignedPublicSubKey> {
key.public_subkeys
.iter()
.filter(|subkey| subkey.algorithm().can_encrypt())
.filter_map(|subkey| {
let signature = subkey.signatures.first()?;
let key_flags = signature.key_flags();
if !key_flags.encrypt_comms() {
return None;
}
if let Some(expiration_duration) = signature
.key_expiration_time()
.filter(|duration| duration.as_secs() != 0)
&& now
> subkey
.created_at()
.as_secs()
.saturating_add(expiration_duration.as_secs())
{
// Key is expired.
return None;
}
Some((subkey, signature))
})
.min_by(|(subkey1, signature1), (subkey2, signature2)| {
match (
signature1
.key_expiration_time()
.filter(|duration| duration.as_secs() != 0),
signature2
.key_expiration_time()
.filter(|duration| duration.as_secs() != 0),
) {
(None, None) => Ordering::Equal,
(None, Some(_)) => Ordering::Greater,
(Some(_), None) => Ordering::Less,
(Some(expiration1), Some(expiration2)) => (subkey1
.created_at()
.as_secs()
.saturating_add(expiration1.as_secs()))
.cmp(
&(subkey2
.created_at()
.as_secs()
.saturating_add(expiration2.as_secs())),
),
}
})
.map(|(subkey, _signature)| subkey)
.find(|subkey| subkey.algorithm().can_encrypt())
}
/// Version of SEIPD packet to use.
@@ -201,11 +151,10 @@ pub fn pk_encrypt(
) -> Result<String> {
tokio::task::block_in_place(|| {
let mut rng = thread_rng();
let now = pgp::types::Timestamp::now();
let pkeys = public_keys_for_encryption
.iter()
.filter_map(|key| select_pk_for_encryption(now.as_secs(), key));
.filter_map(select_pk_for_encryption);
let msg = MessageBuilder::from_bytes("", plain);
let encoded_msg = match seipd_version {
@@ -536,8 +485,7 @@ pub(crate) fn relay_addrs(public_key: &SignedPublicKey, addr: &str) -> Vec<Strin
/// Returns true if the key can be encrypted to, i.e. has an encryption subkey.
pub(crate) fn pubkey_can_encrypt(public_key: &SignedPublicKey) -> bool {
let now = pgp::types::Timestamp::now();
select_pk_for_encryption(now.as_secs(), public_key).is_some()
select_pk_for_encryption(public_key).is_some()
}
/// Returns true if public key advertises SEIPDv2 feature.
-107
View File
@@ -423,110 +423,3 @@ async fn test_securejoin_pqc_joiner() {
tcm.execute_securejoin(bob, pqc).await;
}
/// Tests that public subkey selection for encryption follows Autocrypt 2 rules.
///
/// If there is an expiring subkey, it is preferred, otherwise non-expiring subkey is selected.
/// Non-encryption subkeys such as RSA subkey for authentication are ignored.
#[test]
fn test_select_pk_for_encryption() {
// Public key generated with GnuPG 2.4.9 with the following subkeys:
// 1. Auth-only RSA subkey (92E762B9084CA740).
// 2. Expired Curve25519 encryption subkey with 1-day expiration (C8F382BD0F35C49E)
// 3. Ed25519 signing subkey (F177AC3118F923CC).
// 4. Curve25519 encryption subkey with fingerprint (36188C6FFC8E267B)
// 5. Curve25519 encryption subkey with 1 year expiration, valid in the beginning of 2008, with key ID 9223FCEE7546CDE7
// 6. Curve25519 encryption subkey with no expiration (FD2C0567967223D8).
// Primary key is an Ed25519 not expiring key.
// Key 4 is the one that should be selected.
// Subkey 4 fingerprint.
let expected_fallback_fingerprint = "cdeb3ba3999bf7880f0ee1f536188c6ffc8e267b";
// Subkey 5 fingerprint, should be preferred to fallback when not expired.
let expected_expiring_fingerprint = "5133fab157c4a46ca41f6dc39223fcee7546cde7";
let alice_tpk_asc = "
-----BEGIN PGP PUBLIC KEY BLOCK-----
mDMERvfcPBYJKwYBBAHaRw8BAQdAimvPsr7NdJ4dBoFPySwhpTQqoYOoHL3AzfE7
mGWQOOC0GUFsaWNlIDxhbGljZUBleGFtcGxlLm9yZz6IkAQTFgoAOBYhBCi19Yqv
ugVVkhyHgicqAms0FFoiBQJG99w8AhsDBQsJCAcCBhUKCQgLAgQWAgMBAh4BAheA
AAoJECcqAms0FFoiUGEA/3VZMBCoRq0ZpHarzmvzgdZCoL3r3m9en/eZScFzxITx
AP42Mn27r0SOKwIln0VcPTdAQCk49mBW/EX3CMOlLPU3DLkBjQRG99w8AQwArsYe
Jkdl6sSM/hfoEw0vgx/RdUBQ6QRYi1uc0UUNlIGy8mlczLFdkD3JF/hGocjPvt45
XAQoK110zAkZlfpFRqNT1M/IC68Er8rLkYPC4OeFh6W4Iyn17fcUanP0lf8em/jh
Vffvgy8sFOMdO235lvFA3txNA98s4fHdmU3PScyd1hc3C4M0yP83LnYyWt4X59Xc
E/Om5Dm458eKCSeYkLI6752W0mXsBxSi3/dLn0XeuNRpgmKxSkm562FHOFaLbKtR
Y5hobAI9PkNcVgRxZZvWQls5PHTZWjqngF21lKlaLfdqZ/Uae1i2hzZOihgitL43
Le00qwNBi5hKYMeuDnHaQrLmb+A+0/IAEE4Ub+TkZhzI+2ZP2k1cAT0qZmZi0w+q
xqP9INk0hY/oZFCnV2wkHN7zvQmVlUIcQ2rmfbafK1yiEL1qeGT96zyjbdXeGPqJ
3O9h+YIG38JMRJBijsFujUN34Z546zS/kzOPXsz/WlGUMjwu8n5s5uf2TFyFABEB
AAGIeAQYFgoAIBYhBCi19YqvugVVkhyHgicqAms0FFoiBQJG99w8AhsgAAoJECcq
Ams0FFoiCOUBAPafRLDpWN9iT4hcCXjESf1Hw5KNVkJpwfzPfu2H9BkMAQCaHhKg
pq9ywH4pyOHZCPV8P2ywkyn+EsjBC3fG+GBBBbg4BEb33DwSCisGAQQBl1UBBQEB
B0DfI8AJFT3nWa6ZXLkHSf7W8W7S6AWIO7LAcjoyHwb8CwMBCAeIfgQYFgoAJhYh
BCi19YqvugVVkhyHgicqAms0FFoiBQJG99w8AhsMBQkAAVGAAAoJECcqAms0FFoi
U5EA/3G74HRwIMJlNOEW5gkYYV5KJW2qgtMfxHCUjoHvNWU1AQCHt/bLU2aviAiS
of1R43qojxKUqzzoi8lYRQ+1sYhvB7gzBEb33DwWCSsGAQQB2kcPAQEHQKZXUJ7s
xqH3kVcMnhasw6DrFMwCxHDdj+qvkg8r/DvtiO8EGBYKACAWIQQotfWKr7oFVZIc
h4InKgJrNBRaIgUCRvfcPAIbAgCBCRAnKgJrNBRaInYgBBkWCgAdFiEEI8hstnVQ
9sgylIYg8XesMRj5I8wFAkb33DwACgkQ8XesMRj5I8xo6QEAu4o/TyEZwFcyqZpw
LEo9vTLCsc7fo0nx0ssiP6FyV5cBAOWal1DznDhsXWCNt+U8UaafXsU2DTV51KaD
VBOVFo4CyeQBAMirIjXV5PbUV674TNLhYl2s0jTtNz+GKtOjSdZuRm1mAPwPG6ya
K1b7iMRdBT92gNZMw30LbtcXmttCxpZAwr9lBLg4BEb33DwSCisGAQQBl1UBBQEB
B0C5fFb4WTHoIoI6ou/31+1N1wn8ghsSkUVzpbtv/aTkegMBCAeIeAQYFgoAIBYh
BCi19YqvugVVkhyHgicqAms0FFoiBQJG99w8AhsMAAoJECcqAms0FFoi8z8BALPL
7V0ICLEY5YSUa4lQ2rjiXOcVTlWkG3h4TATPrr08AP9tIAQIE0o50IGdQAcKJoTn
Lyxnf2wfjZ16vL3JLLjSBrg4BEb33DwSCisGAQQBl1UBBQEBB0DXDrcGrnuLjAUO
eo/t8MQNOe+ZKYSDPGTkO7iM5IloSAMBCAeIfgQYFgoAJhYhBCi19YqvugVVkhyH
gicqAms0FFoiBQJG99w8AhsMBQkB4TOAAAoJECcqAms0FFoivQIA/2PcZ1vcImAa
7ldPY00JkcW6WlSSd6yOIZsVa4TdA1FiAP42gOji+4RrLps2+NX6L1znSc8EJBXo
RMbND/CZQWXfA7g4BEb33DwSCisGAQQBl1UBBQEBB0BHxCvo5zuygw2XiluYNobx
7iFJqlmCkjekKyoVFquKHQMBCAeIeAQYFgoAIBYhBCi19YqvugVVkhyHgicqAms0
FFoiBQJG99w8AhsMAAoJECcqAms0FFoirSMA/0gVP98sPFga+UhQ3uJxJw5bO2Rs
7hxVk6aPREWgBYg1AQCT7AE8m7j17SP/1fl8OjpxsQQmCJyv2wNcP48OfKGOCA==
=AXAi
-----END PGP PUBLIC KEY BLOCK-----
";
let alice_tpk = SignedPublicKey::from_asc(alice_tpk_asc).unwrap();
let now = pgp::types::Timestamp::now();
let encryption_subkey = select_pk_for_encryption(now.as_secs(), &alice_tpk).unwrap();
assert_eq!(
encryption_subkey.fingerprint().to_string().as_str(),
expected_fallback_fingerprint
);
// In the beginning of 2008 expiring subkey is not expired yet and should be used.
let encryption_subkey = select_pk_for_encryption(1199149200, &alice_tpk).unwrap();
assert_eq!(
encryption_subkey.fingerprint().to_string().as_str(),
expected_expiring_fingerprint
);
}
/// Tests that key selection fails if there is only an expired subkey.
#[test]
fn test_select_only_expired_subkey() {
let alice_tpk_asc = "
-----BEGIN PGP PUBLIC KEY BLOCK-----
mDMERvfcPBYJKwYBBAHaRw8BAQdAUz6AZXIRE8T04Vh8RReFiP3tEV8UfSs2EiYs
8b8i4ce0GUFsaWNlIDxhbGljZUBleGFtcGxlLm9yZz6IkAQTFgoAOBYhBNRG3w/A
qWyPitBaqkGLyqDx3CaUBQJG99w8AhsDBQsJCAcCBhUKCQgLAgQWAgMBAh4BAheA
AAoJEEGLyqDx3CaUMjIA/Rq9/iORLP360s6EsIDe9qSmlmSCggtivafH+uVBWa0X
AP0Wb+kailDwISq2O9Ef/jceZw7ozyzDLeDskKpCYiL3A7g4BEb33DwSCisGAQQB
l1UBBQEBB0BXhbrks9iskW1GfT3B022W3KhJCgz8gw81z9lAWv7OZgMBCAeIfgQY
FgoAJhYhBNRG3w/AqWyPitBaqkGLyqDx3CaUBQJG99w8AhsMBQkB4TOAAAoJEEGL
yqDx3CaU6/YA/jaYJsZMXvu5grrMq1wq3Z/yzGXd15zeWp+alY/6UYxoAQCmSrBC
SOtE3ODLZN9tC3F7k1N9clme1cHXyUiH3EliBQ==
=zWPq
-----END PGP PUBLIC KEY BLOCK-----
";
let alice_tpk = SignedPublicKey::from_asc(alice_tpk_asc).unwrap();
let now = pgp::types::Timestamp::now();
assert_eq!(select_pk_for_encryption(now.as_secs(), &alice_tpk), None);
}
+1
View File
@@ -1032,6 +1032,7 @@ Content-Disposition: reaction\n\
assert_eq!(summary.timestamp, bob_msg1.get_timestamp()); // time refers to message, not to reaction
assert_eq!(summary.state, MessageState::InFresh); // state refers to message, not to reaction
assert!(summary.prefix.is_none());
assert!(summary.thumbnail_path.is_none());
assert_summary(&alice, "BOB reacted 👍 to \"Party?\"").await;
// Alice reacts to own message as well
+7 -49
View File
@@ -54,39 +54,6 @@ pub(crate) struct Smtp {
pub(crate) last_send_error: Option<String>,
}
/// Returns transports with their IDs in the order in which they should be tried.
async fn sorted_transports(context: &Context) -> Result<Vec<(u32, ConfiguredLoginParam)>> {
context
.sql
.query_map_vec(
"SELECT transports.id, configured_param FROM transports
LEFT JOIN smtp_success ON smtp_success.transport_id=transports.id
ORDER BY IFNULL(smtp_success.id, 0) DESC, transports.id ASC",
(),
|row| {
let id: u32 = row.get(0)?;
let json: String = row.get(1)?;
let param = ConfiguredLoginParam::from_json(&json)?;
Ok((id, param))
},
)
.await
}
/// Records successful use of SMTP transport so it is tried first next time we connect to SMTP.
async fn record_success(context: &Context, transport_id: u32) -> Result<()> {
// INSERT OR REPLACE essentially replaces rowid of the row
// if the row exists already, so it becomes the highest rowid in the table.
context
.sql
.execute(
"INSERT OR REPLACE INTO smtp_success (transport_id) VALUES (?)",
(transport_id,),
)
.await?;
Ok(())
}
impl Smtp {
/// Create a new Smtp instances.
pub fn new() -> Self {
@@ -134,7 +101,13 @@ impl Smtp {
self.connectivity.set_connecting(context);
let proxy_config = ProxyConfig::load(context).await?;
for (transport_id, lp) in sorted_transports(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(
@@ -354,18 +327,6 @@ pub(crate) async fn smtp_send(
Ok(()) => SendResult::Success,
};
if matches!(status, SendResult::Success) {
debug_assert!(smtp.transport_id.is_some());
if let Some(transport_id) = smtp.transport_id
&& let Err(err) = record_success(context, transport_id).await
{
warn!(
context,
"Failed to record successful use of transport {transport_id} in smtp_success table: {err:#}."
);
}
}
if let SendResult::Failure(err) = &status
&& let Some(msg_id) = msg_id
{
@@ -897,6 +858,3 @@ pub(crate) async fn add_self_recipients(
Ok(())
}
#[cfg(test)]
mod smtp_tests;
-49
View File
@@ -1,49 +0,0 @@
use anyhow::Result;
use crate::test_utils::TestContextManager;
use crate::transport;
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn test_smtp_candidates() -> Result<()> {
let mut tcm = TestContextManager::new();
let t = &tcm.unconfigured().await;
transport::add_pseudo_transport(t, "foo@example.net").await?;
transport::add_pseudo_transport(t, "bar@example.net").await?;
transport::add_pseudo_transport(t, "baz@example.net").await?;
let transports = super::sorted_transports(t).await?;
let [
(transport_id1, ref transport1),
(transport_id2, ref transport2),
(transport_id3, ref transport3),
] = transports[..]
else {
panic!("Unexpected number of transports");
};
// By default first added transport is used first.
assert_eq!(transport1.addr, "foo@example.net");
assert_eq!(transport2.addr, "bar@example.net");
assert_eq!(transport3.addr, "baz@example.net");
super::record_success(t, transport_id3).await?;
let transports2 = super::sorted_transports(t).await?;
assert_eq!(transports2[0].0, transport_id3);
assert_eq!(transports2[1].0, transport_id1);
assert_eq!(transports2[2].0, transport_id2);
super::record_success(t, transport_id2).await?;
let transports3 = super::sorted_transports(t).await?;
assert_eq!(transports3[0].0, transport_id2);
assert_eq!(transports3[1].0, transport_id3);
assert_eq!(transports3[2].0, transport_id1);
super::record_success(t, transport_id3).await?;
let transports4 = super::sorted_transports(t).await?;
assert_eq!(transports4[0].0, transport_id3);
assert_eq!(transports4[1].0, transport_id2);
assert_eq!(transports4[2].0, transport_id1);
Ok(())
}
-15
View File
@@ -2672,21 +2672,6 @@ CREATE TABLE smtp2 (
.await?;
}
inc_and_check(&mut migration_version, 168)?;
if dbversion < migration_version {
sql.execute_migration(
"
CREATE TABLE smtp_success (
id INTEGER PRIMARY KEY AUTOINCREMENT NOT NULL,
transport_id INTEGER UNIQUE NOT NULL,
FOREIGN KEY(transport_id) REFERENCES transports(id) ON DELETE CASCADE
) STRICT;
",
migration_version,
)
.await?;
}
let new_version = sql
.get_raw_config_int(VERSION_CFG)
.await?
+17
View File
@@ -56,6 +56,9 @@ pub struct Summary {
/// Message state.
pub state: MessageState,
/// Message preview image path
pub thumbnail_path: Option<String>,
}
impl Summary {
@@ -80,6 +83,7 @@ impl Summary {
text: msg_reacted(context, reaction_contact_id, &reaction, &summary).await,
timestamp: msg.get_timestamp(), // message timestamp (not reaction) to make timestamps more consistent with chats ordering
state: msg.state, // message state (not reaction) - indicating if it was me sending the last message
thumbnail_path: None,
});
}
Self::new(context, msg, chat, contact).await
@@ -123,11 +127,24 @@ impl Summary {
text = stock_str::reply_noun(context)
}
let thumbnail_path = if msg.viewtype == Viewtype::Image
|| msg.viewtype == Viewtype::Gif
|| msg.viewtype == Viewtype::Sticker
{
msg.get_file(context)
.and_then(|path| path.to_str().map(|p| p.to_owned()))
} else if msg.viewtype == Viewtype::Webxdc {
Some("webxdc-icon://last-msg-id".to_string())
} else {
None
};
Ok(Summary {
prefix,
text,
timestamp: msg.get_timestamp(),
state: msg.state,
thumbnail_path,
})
}
+1
View File
@@ -1664,6 +1664,7 @@ async fn test_webxdc_chatlist_summary() -> Result<()> {
assert_eq!(chatlist.len(), 1);
let summary = chatlist.get_summary(&t, 0, None).await?;
assert_eq!(summary.text, "📱 nice app!".to_string());
assert_eq!(summary.thumbnail_path.unwrap(), "webxdc-icon://last-msg-id");
Ok(())
}