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
20 changed files with 208 additions and 163 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 })
}
+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
@@ -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(())
}