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)]