From 6c31d9ef69fd15c633c5cef6ce70dd9be506d761 Mon Sep 17 00:00:00 2001 From: holger krekel Date: Fri, 14 Aug 2026 17:20:33 +0200 Subject: [PATCH] feat: perform background fetch from all transports With I/O stopped, `background_fetch()` connected only to the transport of `configured_addr` and we now instead fan out to all transports in a controlled loop. Also drop the quota check from this background fetch path: its result is in-memory only, discarded when the iOS notification service exits, and the regular scheduler fetching refreshes it every 60s anyway. Moreover, quota errors/running full is pretty rare since relays generally automatically stay under quota these days. It's another round trip for each transport of each profile and simply not neccessary. Also adds previosly missing online tests and adds `transport_id` to `ImapInboxIdle` Event --- deltachat-ffi/deltachat.h | 2 +- deltachat-jsonrpc/src/api.rs | 8 +++ .../src/deltachat_rpc_client/account.py | 5 +- deltachat-rpc-client/tests/conftest.py | 8 ++- .../tests/test_multitransport.py | 66 +++++++++++++++---- src/constants.rs | 3 - src/context.rs | 42 +++--------- src/imap.rs | 12 ---- src/scheduler.rs | 34 ++++++++++ 9 files changed, 113 insertions(+), 67 deletions(-) diff --git a/deltachat-ffi/deltachat.h b/deltachat-ffi/deltachat.h index 721882842f..232626dd97 100644 --- a/deltachat-ffi/deltachat.h +++ b/deltachat-ffi/deltachat.h @@ -3187,7 +3187,7 @@ 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 messages from imap and then resumes the scheduler. + * Pauses the scheduler, fetches messages from all transports and then resumes the scheduler. * * dc_accounts_background_fetch() was created for the iOS Background fetch. * diff --git a/deltachat-jsonrpc/src/api.rs b/deltachat-jsonrpc/src/api.rs index a44bde8314..8bca552f35 100644 --- a/deltachat-jsonrpc/src/api.rs +++ b/deltachat-jsonrpc/src/api.rs @@ -2089,6 +2089,14 @@ impl CommandApi { Ok(()) } + /// Waits until all transports fetched their inbox and no background work is left. + /// Never returns unless I/O is started. Must ONLY be used by tests. + async fn wait_for_transports_idle_and_all_work_done(&self, account_id: u32) -> Result<()> { + let ctx = self.get_context(account_id).await?; + ctx.wait_for_all_work_done().await; + Ok(()) + } + /// Get the current connectivity, i.e. whether the device is connected to the IMAP server. /// One of: /// - DC_CONNECTIVITY_NOT_CONNECTED (1000): Show e.g. the string "Not connected" or a red dot diff --git a/deltachat-rpc-client/src/deltachat_rpc_client/account.py b/deltachat-rpc-client/src/deltachat_rpc_client/account.py index 9e2a7b06c3..725f2e0361 100644 --- a/deltachat-rpc-client/src/deltachat_rpc_client/account.py +++ b/deltachat-rpc-client/src/deltachat_rpc_client/account.py @@ -154,9 +154,10 @@ def list_transports(self): return transports def bring_online(self): - """Start I/O and wait until IMAP becomes IDLE.""" + """Start I/O, wait until all transports became IDLE and drop the events seen so far.""" self.start_io() - self.wait_for_event(EventType.IMAP_INBOX_IDLE) + self._rpc.wait_for_transports_idle_and_all_work_done(self.id) + self.clear_all_events() def create_contact(self, obj: Union[int, str, Contact, "Account"], name: Optional[str] = None) -> Contact: """Create a new Contact or return an existing one. diff --git a/deltachat-rpc-client/tests/conftest.py b/deltachat-rpc-client/tests/conftest.py index b8f1600457..9bf4a1f75e 100644 --- a/deltachat-rpc-client/tests/conftest.py +++ b/deltachat-rpc-client/tests/conftest.py @@ -22,8 +22,10 @@ class DirectImap: """Internal Python-level IMAP handling.""" - def __init__(self, account: Account) -> None: + def __init__(self, account: Account, addr=None, password=None) -> None: self.account = account + 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() @@ -33,9 +35,9 @@ def connect(self): host = self.account.get_config("configured_mail_server") port = 993 - user = self.account.get_config("addr") + user = self.addr host = user.rsplit("@")[-1] - pw = self.account.get_config("mail_pw") + pw = self.password ssl_context = ssl.create_default_context() if host.startswith("_"): diff --git a/deltachat-rpc-client/tests/test_multitransport.py b/deltachat-rpc-client/tests/test_multitransport.py index 4a52e413a3..ce67e25fe7 100644 --- a/deltachat-rpc-client/tests/test_multitransport.py +++ b/deltachat-rpc-client/tests/test_multitransport.py @@ -1,3 +1,5 @@ +import time + import pytest from deltachat_rpc_client import EventType @@ -5,6 +7,24 @@ from deltachat_rpc_client.rpc import JsonRpcError +def alice_with_two_transports_and_bob(acf): + alice, bob = acf.get_online_accounts(2) + alice.add_transport_from_qr(acf.get_account_qr()) + alice.bring_online() + return alice, alice.create_chat(bob), bob.create_chat(alice) + + +def messages_with_text(chat, text): + return [msg for msg in chat.get_messages() if msg.get_snapshot().text == text] + + +def wait_for_imap_message(imap, timeout=60): + deadline = time.time() + timeout + while not imap.get_all_messages(): + assert time.time() < deadline, f"no message arrived in {imap.addr}" + time.sleep(1) + + def test_add_second_address(acf) -> None: account = acf.new_configured_account() assert len(account.list_transports()) == 1 @@ -256,11 +276,10 @@ def test_message_info_imap_urls(acf) -> None: alice, bob = acf.get_online_accounts(2) qr = acf.get_account_qr() - for i in range(3): + for _ in range(3): alice.add_transport_from_qr(qr) # Wait for all transports to go IDLE after adding each one. - for _ in range(i + 1): - alice.bring_online() + alice.bring_online() # Enable multi-device mode so messages are not deleted immediately. alice.set_config("bcc_self", "1") @@ -292,14 +311,7 @@ def test_message_info_imap_urls(acf) -> None: def test_remove_primary_transport(acf, log) -> None: """Test that after removing the primary relay, Alice can still receive messages.""" - alice, bob = acf.get_online_accounts(2) - qr = acf.get_account_qr() - - alice.add_transport_from_qr(qr) - alice.bring_online() - - bob_chat = bob.create_chat(alice) - alice.create_chat(bob) + alice, alice_chat, bob_chat = alice_with_two_transports_and_bob(acf) log.section("Alice sets up second transport") [transport1, transport2] = alice.list_transports() @@ -318,4 +330,34 @@ def test_remove_primary_transport(acf, log) -> None: msg2 = alice.wait_for_incoming_msg().get_snapshot() assert msg2.text == "Hello again!" assert msg2.chat.get_basic_snapshot().chat_type == ChatType.SINGLE - assert msg2.chat == alice.create_chat(bob) + assert msg2.chat == alice_chat + + +def test_background_fetch_from_second_transport(acf, direct_imap, dc): + alice, alice_chat, bob_chat = alice_with_two_transports_and_bob(acf) + [transport1, transport2] = alice.list_transports() + assert alice.get_config("configured_addr") == transport1["addr"] + + alice.stop_io() + bob_chat.send_text("hello") + imap1 = direct_imap(alice, transport1["addr"], transport1["password"]) + wait_for_imap_message(direct_imap(alice, transport2["addr"], transport2["password"])) + wait_for_imap_message(imap1) + + # Leave the message on the second transport only. + imap1.delete("1:*") + + dc.background_fetch(30) + assert len(messages_with_text(alice_chat, "hello")) == 1 + + +def test_background_fetch_no_duplicates(acf, direct_imap, dc): + alice, alice_chat, bob_chat = alice_with_two_transports_and_bob(acf) + + alice.stop_io() + bob_chat.send_text("hello") + for transport in alice.list_transports(): + wait_for_imap_message(direct_imap(alice, transport["addr"], transport["password"])) + + dc.background_fetch(30) + assert len(messages_with_text(alice_chat, "hello")) == 1 diff --git a/src/constants.rs b/src/constants.rs index 9a0aefd4d4..795c9a2b34 100644 --- a/src/constants.rs +++ b/src/constants.rs @@ -185,9 +185,6 @@ pub const MAX_RCVD_IMAGE_PIXELS: u32 = 50_000_000; // Relays typically advertise their limit via IMAP METADATA. pub(crate) const DEFAULT_MAX_SMTP_RCPT_TO: u32 = 50; -/// How far the last quota check needs to be in the past to be checked by the background function (in seconds). -pub(crate) const DC_BACKGROUND_FETCH_QUOTA_CHECK_RATELIMIT: u64 = 12 * 60 * 60; // 12 hours - /// How far in the future the sender timestamp of a message is allowed to be, in seconds. Also used /// in the group membership consistency algo to reject outdated membership changes. pub(crate) const TIMESTAMP_SENT_TOLERANCE: i64 = 60; diff --git a/src/context.rs b/src/context.rs index fc50fa2c06..a39cd81c8c 100644 --- a/src/context.rs +++ b/src/context.rs @@ -16,11 +16,11 @@ use tokio::sync::{Mutex, Notify, RwLock}; use crate::chat::{ChatId, get_chat_cnt}; use crate::config::Config; -use crate::constants::{self, DC_BACKGROUND_FETCH_QUOTA_CHECK_RATELIMIT, DC_VERSION_STR}; +use crate::constants::{self, DC_VERSION_STR}; use crate::contact::{Contact, ContactId}; use crate::debug_logging::DebugLogging; use crate::events::{Event, EventEmitter, EventType, Events}; -use crate::imap::{Imap, ServerMetadata}; +use crate::imap::ServerMetadata; use crate::log::warn; use crate::logged_debug_assert; use crate::message::{self, MessageState, MsgId}; @@ -598,56 +598,30 @@ impl Context { Ok(constants::DEFAULT_MAX_SMTP_RCPT_TO) } - /// Does a single round of fetching from IMAP and returns. + /// Does a single round of fetching messages from all transports and returns. /// /// Can be used even if I/O is currently stopped. - /// If I/O is currently stopped, starts a new IMAP connection - /// and fetches from Inbox and DeltaChat folders. + /// If I/O is stopped, starts a new IMAP connection for each transport. pub async fn background_fetch(&self) -> Result<()> { if !(self.is_configured().await?) { return Ok(()); } - let address = self.get_primary_self_addr().await?; let time_start = tools::Time::now(); - info!(self, "background_fetch started fetching {address}."); + info!(self, "background_fetch started."); if self.scheduler.is_running().await { self.scheduler.maybe_network().await; self.wait_for_all_work_done().await; } else { - // Pause the scheduler to ensure another connection does not start - // while we are fetching on a dedicated connection. - let _pause_guard = self.scheduler.pause(self).await?; - - // Start a new dedicated connection. - let mut connection = Imap::new_configured(self, channel::bounded(1).1).await?; - let mut session = connection.prepare(self).await?; - - // Fetch IMAP folders. - let folder = connection.folder.clone(); - connection - .fetch_move_delete(self, &mut session, &folder) + self.scheduler + .fetch_from_all_transports_at_once(self) .await?; - - // Update quota (to send warning if full) - but only check it once in a while. - // note: For now this only checks quota of primary transport, - // because background check only checks primary transport at the moment - if self - .quota_needs_update( - session.transport_id(), - DC_BACKGROUND_FETCH_QUOTA_CHECK_RATELIMIT, - ) - .await - && let Err(err) = self.update_recent_quota(&mut session, &folder).await - { - warn!(self, "Failed to update quota: {err:#}."); - } } info!( self, - "background_fetch done for {address} took {:?}.", + "background_fetch done, took {:?}.", time_elapsed(&time_start), ); diff --git a/src/imap.rs b/src/imap.rs index c211b8b3d6..a39e13356a 100644 --- a/src/imap.rs +++ b/src/imap.rs @@ -245,18 +245,6 @@ impl Imap { }) } - /// Creates new disconnected IMAP client using configured parameters. - pub async fn new_configured( - context: &Context, - idle_interrupt_receiver: Receiver<()>, - ) -> Result { - let (transport_id, param) = ConfiguredLoginParam::load(context) - .await? - .context("Not configured")?; - let imap = Self::new(context, transport_id, param, idle_interrupt_receiver).await?; - Ok(imap) - } - /// Returns transport ID of the IMAP client. pub fn transport_id(&self) -> u32 { self.transport_id diff --git a/src/scheduler.rs b/src/scheduler.rs index 30a27eb8de..6997a7336d 100644 --- a/src/scheduler.rs +++ b/src/scheduler.rs @@ -281,6 +281,40 @@ impl SchedulerState { scheduler.interrupt_recently_seen(contact_id, timestamp); } } + + /// Fetches from all transports at once, each on a dedicated connection. + /// + /// IO is paused while fetching so that the scheduler does not connect as well. + pub(crate) async fn fetch_from_all_transports_at_once(&self, context: &Context) -> Result<()> { + let _pause_guard = self.pause(context).await?; + + let mut futures = Vec::new(); + for (transport_id, param, _is_published) in ConfiguredLoginParam::load_all(context).await? { + futures.push(async move { + if let Err(err) = fetch_from_transport(context, transport_id, param).await { + warn!(context, "Transport {transport_id}: fetch failed: {err:#}."); + } + }); + } + futures::future::join_all(futures).await; + Ok(()) + } +} + +async fn fetch_from_transport( + context: &Context, + transport_id: u32, + param: ConfiguredLoginParam, +) -> Result<()> { + // A single fetch has nothing to interrupt. + let (_, idle_interrupt_receiver) = channel::bounded(1); + let mut connection = Imap::new(context, transport_id, param, idle_interrupt_receiver).await?; + let mut session = connection.prepare(context).await?; + + let folder = connection.folder.clone(); + connection + .fetch_move_delete(context, &mut session, &folder) + .await } #[derive(Debug, Default)]