diff --git a/crates/bridge/src/event_handler.rs b/crates/bridge/src/event_handler.rs index ff4bdcd..396fafe 100644 --- a/crates/bridge/src/event_handler.rs +++ b/crates/bridge/src/event_handler.rs @@ -349,6 +349,7 @@ async fn apply_mail_set_entry_create( details: details.clone(), rfc2822: rfc2822.clone(), uid: 0, + attachments_pending: false, }; // Persist the body straight away so the IMAP layer can serve FETCH // BODY[] without a placeholder, and the prefetch loop skips it. @@ -388,6 +389,7 @@ async fn apply_mail_set_entry_create( details: None, rfc2822: None, uid: 0, + attachments_pending: false, }; assign_uid_and_upsert(store, local_store, target_folder, &elem_id, &mut stored).await; }, diff --git a/crates/bridge/src/imap/session.rs b/crates/bridge/src/imap/session.rs index 9bd8c4b..de68876 100644 --- a/crates/bridge/src/imap/session.rs +++ b/crates/bridge/src/imap/session.rs @@ -1396,8 +1396,8 @@ mod tests { .set_folder( "inbox", vec![ - StoredMail { mail: m1, details: None, rfc2822: None, uid: 1 }, - StoredMail { mail: m2, details: None, rfc2822: None, uid: 2 }, + StoredMail { mail: m1, details: None, rfc2822: None, uid: 1, attachments_pending: false }, + StoredMail { mail: m2, details: None, rfc2822: None, uid: 2, attachments_pending: false }, ], ) .await; @@ -1435,6 +1435,7 @@ mod tests { details: None, rfc2822: None, uid: (i + 1) as u32, + attachments_pending: false, }) .collect(); store.set_folder("inbox", stored).await; @@ -1619,8 +1620,8 @@ mod tests { let rfc2 = crate::mail::mail_to_rfc2822(&m2, Some(&d2), &[]); store.set_folder_list(vec![inbox_folder()]).await; store.set_folder("inbox", vec![ - StoredMail { mail: m1, details: Some(d1), rfc2822: Some(rfc1), uid: 1 }, - StoredMail { mail: m2, details: Some(d2), rfc2822: Some(rfc2), uid: 2 }, + StoredMail { mail: m1, details: Some(d1), rfc2822: Some(rfc1), uid: 1, attachments_pending: false }, + StoredMail { mail: m2, details: Some(d2), rfc2822: Some(rfc2), uid: 2, attachments_pending: false }, ]).await; let mut session = ImapSession::new(store, backend, None); @@ -1945,8 +1946,8 @@ mod tests { .set_folder( "inbox", vec![ - StoredMail { mail: m1, details: None, rfc2822: None, uid: 1 }, - StoredMail { mail: m2, details: None, rfc2822: None, uid: 2 }, + StoredMail { mail: m1, details: None, rfc2822: None, uid: 1, attachments_pending: false }, + StoredMail { mail: m2, details: None, rfc2822: None, uid: 2, attachments_pending: false }, ], ) .await; @@ -1966,7 +1967,7 @@ mod tests { store .set_folder( "inbox", - vec![StoredMail { mail: m1, details: None, rfc2822: None, uid: 1 }], + vec![StoredMail { mail: m1, details: None, rfc2822: None, uid: 1, attachments_pending: false }], ) .await; let resp = session.check_new_mail().await; @@ -1989,8 +1990,8 @@ mod tests { .set_folder( "inbox", vec![ - StoredMail { mail: old1, details: None, rfc2822: None, uid: 1 }, - StoredMail { mail: new3, details: None, rfc2822: None, uid: 3 }, + StoredMail { mail: old1, details: None, rfc2822: None, uid: 1, attachments_pending: false }, + StoredMail { mail: new3, details: None, rfc2822: None, uid: 3, attachments_pending: false }, ], ) .await; @@ -2013,7 +2014,7 @@ mod tests { store .set_folder( "inbox", - vec![StoredMail { mail: m1, details: None, rfc2822: None, uid: 1 }], + vec![StoredMail { mail: m1, details: None, rfc2822: None, uid: 1, attachments_pending: false }], ) .await; let resp = session.check_new_mail().await; diff --git a/crates/bridge/src/sync.rs b/crates/bridge/src/sync.rs index f16df78..73f6883 100644 --- a/crates/bridge/src/sync.rs +++ b/crates/bridge/src/sync.rs @@ -21,6 +21,12 @@ pub struct StoredMail { pub rfc2822: Option, /// Stable IMAP UID within the folder, persisted across restarts. pub uid: u32, + /// True while the mail's body has been rendered without its attachment + /// parts because the server returned them with unresolved encryption + /// metadata (a fresh-mail race). The `prefetch_loop` keeps retrying + /// just the attachment load until it succeeds — at which point the + /// `.eml` is rewritten as a proper `multipart/mixed`. + pub attachments_pending: bool, } pub struct MailStore { @@ -285,6 +291,7 @@ impl MailStore { element_id: &str, details: MailDetails, rfc2822: String, + attachments_pending: bool, ) { let mut folders = self.folders.write().await; if let Some(folder) = folders.get_mut(folder_id) { @@ -298,6 +305,31 @@ impl MailStore { }) { m.details = Some(details); m.rfc2822 = Some(rfc2822); + m.attachments_pending = attachments_pending; + } + } + drop(folders); + self.bump_generation(); + } + + /// Rewrite just the cached RFC 2822 for a mail (and clear its pending + /// flag) after a successful attachment-only retry. The mail's details + /// are already in store from the original prefetch; we only need to + /// refresh the body so the next IMAP FETCH BODY[] returns the + /// multipart envelope. + async fn update_mail_rfc2822(&self, folder_id: &str, element_id: &str, rfc2822: String) { + let mut folders = self.folders.write().await; + if let Some(folder) = folders.get_mut(folder_id) { + if let Some(m) = folder.iter_mut().find(|m| { + m.mail + ._id + .as_ref() + .map(|id| id.element_id.to_string()) + .as_deref() + == Some(element_id) + }) { + m.rfc2822 = Some(rfc2822); + m.attachments_pending = false; } } drop(folders); @@ -395,6 +427,12 @@ async fn prefetch_loop( mut shutdown: watch::Receiver, ) { let mut store_watch = store.subscribe(); + // Per-mail throttle for attachment retries — keys are element ids, + // values are the instant we last attempted the retry. We hold this map + // in the prefetch_loop's stack so it survives across store generations + // but resets on bridge restart. The session-key-transient race typically + // resolves within a couple of minutes, so a 60s minimum gap is plenty. + let mut attachment_retry_history: HashMap = HashMap::new(); loop { if *shutdown.borrow() { return; @@ -403,19 +441,44 @@ async fn prefetch_loop( if *shutdown.borrow() { return; } - prefetch_details(store, local_store, backend, &folder).await; + prefetch_details( + store, + local_store, + backend, + &folder, + &mut attachment_retry_history, + ) + .await; } - // Wait for the store to change again (or shutdown). `changed()` on a - // freshly-borrowed receiver returns immediately if a bump happened - // during the prefetch pass above, so we never miss a generation. + // If any attachment retry is still in flight (throttle is set) we + // must wake up after the throttle window even if nothing else + // touches the store — otherwise the retry loop relies on another + // mail arriving to be re-triggered. + let next_tick = if attachment_retry_history.is_empty() { + None + } else { + Some(ATTACHMENT_RETRY_THROTTLE) + }; tokio::select! { _ = store_watch.changed() => {} _ = shutdown.changed() => return, + _ = async { + match next_tick { + Some(d) => tokio::time::sleep(d).await, + None => std::future::pending::<()>().await, + } + } => {} } } } +/// Minimum interval between two attempts to reload a mail's attachments +/// after a session-key-transient failure. Stops a tight retry loop while +/// still recovering within ~1 minute once the server propagates the +/// `_ownerEncSessionKey` / `_ownerGroup` fields. +const ATTACHMENT_RETRY_THROTTLE: Duration = Duration::from_secs(60); + async fn load_cached_folder( store: &MailStore, local_store: &LocalStore, @@ -465,6 +528,7 @@ async fn load_cached_folder( details: None, rfc2822, uid: meta.uid as u32, + attachments_pending: false, }); } @@ -527,6 +591,7 @@ pub(crate) async fn sync_folder( details: existing.details.clone(), rfc2822: existing.rfc2822.clone(), uid, + attachments_pending: existing.attachments_pending, }); } else { let rfc2822 = mail_to_rfc2822(mail, None, &[]); @@ -535,6 +600,7 @@ pub(crate) async fn sync_folder( details: None, rfc2822: Some(rfc2822), uid, + attachments_pending: false, }); } @@ -580,30 +646,64 @@ async fn prefetch_details( local_store: &LocalStore, backend: &dyn MailBackend, folder: &FolderInfo, + attachment_retry_history: &mut HashMap, ) { let mails = store.get_folder(&folder.id).await; + + // First pass — fresh mails still missing their `.eml` on disk. let api_needed: Vec = mails - .into_iter() + .iter() .filter(|m| m.details.is_none()) .filter_map(|m| { let eid = m.mail._id.as_ref()?.element_id.to_string(); if local_store.has_eml(&eid) { None } else { - Some(m.mail) + Some(m.mail.clone()) } }) .collect(); - if api_needed.is_empty() { + // Second pass — mails whose body is already cached but whose attachments + // failed to load on the first try because the server hadn't propagated + // their session-key metadata yet. Throttled per element id so we don't + // hammer the server every store mutation. + let now = std::time::Instant::now(); + let attachment_retry_needed: Vec<(StoredMail, MailDetails)> = mails + .iter() + .filter_map(|m| { + if !m.attachments_pending { + return None; + } + let details = m.details.clone()?; + let eid = m.mail._id.as_ref()?.element_id.to_string(); + if let Some(last) = attachment_retry_history.get(&eid) { + if now.duration_since(*last) < ATTACHMENT_RETRY_THROTTLE { + return None; + } + } + Some((m.clone(), details)) + }) + .collect(); + + if api_needed.is_empty() && attachment_retry_needed.is_empty() { return; } - debug!( - "Pre-fetching {} mail details for {}", - api_needed.len(), - folder.imap_path - ); + if !api_needed.is_empty() { + debug!( + "Pre-fetching {} mail details for {}", + api_needed.len(), + folder.imap_path + ); + } + if !attachment_retry_needed.is_empty() { + debug!( + "Retrying attachments for {} mail(s) in {}", + attachment_retry_needed.len(), + folder.imap_path + ); + } for mail in &api_needed { tokio::time::sleep(INTER_REQUEST_DELAY).await; @@ -611,22 +711,28 @@ async fn prefetch_details( let result = retry(|| backend.load_mail_details(mail)).await; match result { Ok(Some(details)) => { - // Best-effort attachment fetch: if any blob fails to load we - // still ship the body — the user can re-open the mail later - // and the prefetch sweep will retry. Empty `mail.attachments` - // short-circuits inside the backend so this is free for the - // overwhelmingly common "no attachments" case. - let attachments_owned = match backend.load_attachments(mail).await { - Ok(atts) => atts, - Err(e) => { - warn!( - "Failed to load attachments for {}: {e}", - mail._id - .as_ref() - .map(|id| id.element_id.to_string()) - .unwrap_or_default() - ); - Vec::new() + let mut attachments_pending = false; + let attachments_owned = if mail.attachments.is_empty() { + Vec::new() + } else { + match backend.load_attachments(mail).await { + Ok(atts) => atts, + Err(e) => { + attachments_pending = is_transient_attachment_error(&e); + warn!( + "Failed to load attachments for {}: {e} (will {})", + mail._id + .as_ref() + .map(|id| id.element_id.to_string()) + .unwrap_or_default(), + if attachments_pending { + "retry on next prefetch sweep" + } else { + "ship body-only" + } + ); + Vec::new() + } } }; let attachment_refs: Vec<(&TutanotaFile, &[u8])> = attachments_owned @@ -643,9 +749,18 @@ async fn prefetch_details( if let Err(e) = local_store.mark_has_details(&eid) { warn!("Failed to mark has_details {}: {}", eid, e); } + if attachments_pending { + attachment_retry_history.insert(eid.clone(), now); + } store - .update_mail_details(&folder.id, &eid, details, rfc2822) + .update_mail_details( + &folder.id, + &eid, + details, + rfc2822, + attachments_pending, + ) .await; } } @@ -670,6 +785,64 @@ async fn prefetch_details( } } } + + // Attachment-only retries. We already have the body & details cached + // in the store; only the multipart envelope needs rebuilding. + for (stored, details) in &attachment_retry_needed { + tokio::time::sleep(INTER_REQUEST_DELAY).await; + let eid = match stored.mail._id.as_ref() { + Some(id) => id.element_id.to_string(), + None => continue, + }; + attachment_retry_history.insert(eid.clone(), std::time::Instant::now()); + + match backend.load_attachments(&stored.mail).await { + Ok(attachments) => { + let refs: Vec<(&TutanotaFile, &[u8])> = attachments + .iter() + .map(|(f, d)| (f, d.as_slice())) + .collect(); + let rfc2822 = mail_to_rfc2822(&stored.mail, Some(details), &refs); + if let Err(e) = local_store.write_eml(&eid, &rfc2822) { + warn!("Failed to rewrite eml {}: {}", eid, e); + } + store + .update_mail_rfc2822(&folder.id, &eid, rfc2822) + .await; + attachment_retry_history.remove(&eid); + debug!("Attachment retry succeeded for {}", eid); + } + Err(e) if is_transient_attachment_error(&e) => { + debug!( + "Attachment retry for {} still transient: {e} (next attempt in {:?})", + eid, ATTACHMENT_RETRY_THROTTLE + ); + } + Err(e) => { + warn!( + "Attachment retry for {} failed permanently: {e} — giving up", + eid + ); + // Clear the pending flag so we don't keep retrying a + // permanent failure (e.g. server returned 404). + store + .update_mail_rfc2822(&folder.id, &eid, stored.rfc2822.clone().unwrap_or_default()) + .await; + attachment_retry_history.remove(&eid); + } + } + } +} + +/// True when an attachment-load failure is the well-known +/// "freshly-created File without propagated session-key" race that the +/// `TutaSession::load_file_with_retry` helper logs as +/// "still not decryptable after retries". Matching on the string keeps us +/// decoupled from the SDK error variant tree. +fn is_transient_attachment_error(err: &str) -> bool { + err.contains("missing owner key/group data") + || err.contains("Session key resolution failure") + || err.contains("still not decryptable") } async fn retry(mut f: F) -> Result @@ -791,6 +964,7 @@ mod tests { details: None, rfc2822: None, uid, + attachments_pending: false, } } diff --git a/crates/bridge/src/tuta.rs b/crates/bridge/src/tuta.rs index 50cf666..e3690a9 100644 --- a/crates/bridge/src/tuta.rs +++ b/crates/bridge/src/tuta.rs @@ -329,18 +329,22 @@ impl TutaSession { /// freshly-created attachments whose `_ownerEncSessionKey` / /// `_ownerGroup` haven't been published yet — the SDK's session-key /// resolution then fails with "instance missing owner key/group data". - /// 1s + 3s + 6s ≈ 10s of tolerance covers the propagation lag in - /// practice; anything still failing after that bubbles up so the - /// `prefetch_details` caller can decide to ship a body-only envelope. + /// 2s + 4s + 8s + 16s ≈ 30s of tolerance covers the propagation lag in + /// practice (self-sends can take >20s because the same mailbox + /// produces both the Sent entity and the Inbox copy); anything still + /// failing after that bubbles up so the `prefetch_details` caller can + /// decide to ship a body-only envelope and let the next prefetch + /// sweep (triggered by another store mutation) take another shot. async fn load_file_with_retry( &self, file_id: &IdTupleGenerated, ) -> Result { let crypto_client = self.crypto_client(); let backoffs = [ - std::time::Duration::from_secs(1), - std::time::Duration::from_secs(3), - std::time::Duration::from_secs(6), + std::time::Duration::from_secs(2), + std::time::Duration::from_secs(4), + std::time::Duration::from_secs(8), + std::time::Duration::from_secs(16), ]; let mut iter = backoffs.iter(); loop {