From c2e909a09f7888c3cf0063cf9ca5b1e5a3a6a6e7 Mon Sep 17 00:00:00 2001 From: mdecimus Date: Tue, 1 Aug 2023 19:11:59 +0200 Subject: [PATCH] Automatic retry for import/export blob downloads (#14) --- crates/cli/src/modules/export.rs | 22 ++++- crates/cli/src/modules/import.rs | 130 +++++++++++++++++------------ crates/cli/src/modules/mod.rs | 2 + crates/imap/src/op/append.rs | 2 +- crates/jmap/src/api/config.rs | 2 + crates/jmap/src/api/http.rs | 2 +- crates/jmap/src/email/import.rs | 2 +- crates/jmap/src/email/set.rs | 2 +- crates/jmap/src/lib.rs | 3 + crates/jmap/src/services/ingest.rs | 2 +- crates/jmap/src/sieve/ingest.rs | 2 +- 11 files changed, 107 insertions(+), 64 deletions(-) diff --git a/crates/cli/src/modules/export.rs b/crates/cli/src/modules/export.rs index 78e91b5f..a59c9ede 100644 --- a/crates/cli/src/modules/export.rs +++ b/crates/cli/src/modules/export.rs @@ -38,6 +38,8 @@ use jmap_client::{ use serde::Serialize; use tokio::io::AsyncWriteExt; +use crate::modules::RETRY_ATTEMPTS; + use super::{cli::ExportCommands, name_to_id, UnwrapResult}; pub async fn cmd_export(mut client: Client, command: ExportCommands) { @@ -94,10 +96,22 @@ pub async fn cmd_export(mut client: Client, command: ExportCommands) { blob_path.push(&blob_id); futures.push(async move { - let bytes = client - .download(&blob_id) - .await - .unwrap_result("download blob"); + let mut retry_count = 0; + + let bytes = loop { + match client.download(&blob_id).await { + Ok(bytes) => break bytes, + Err(_) if retry_count < RETRY_ATTEMPTS => { + tokio::time::sleep(std::time::Duration::from_secs(1)).await; + retry_count += 1; + } + result => { + result.unwrap_result("download blob"); + return; + } + } + }; + tokio::fs::OpenOptions::new() .create(true) .write(true) diff --git a/crates/cli/src/modules/import.rs b/crates/cli/src/modules/import.rs index 393d98e7..5869d120 100644 --- a/crates/cli/src/modules/import.rs +++ b/crates/cli/src/modules/import.rs @@ -46,7 +46,7 @@ use mail_parser::mailbox::{ use serde::de::DeserializeOwned; use tokio::{fs::File, io::AsyncReadExt}; -use crate::modules::{name_to_id, UnwrapResult}; +use crate::modules::{name_to_id, UnwrapResult, RETRY_ATTEMPTS}; use super::{ cli::{ImportCommands, MailboxFormat}, @@ -331,43 +331,54 @@ pub async fn cmd_import(mut client: Client, command: ImportCommands) { pbs.1 += 1; } - if let Err(err) = client - .email_import( - message.contents, - [mailbox_id.as_ref()], - if !message.flags.is_empty() { - message - .flags - .into_iter() - .map(|f| match f { - maildir::Flag::Passed => "$passed", - maildir::Flag::Replied => "$answered", - maildir::Flag::Seen => "$seen", - maildir::Flag::Trashed => "$deleted", - maildir::Flag::Draft => "$draft", - maildir::Flag::Flagged => "$flagged", - }) - .into() - } else { - None - }, - if message.internal_date > 0 { - (message.internal_date as i64).into() - } else { - None - }, - ) - .await - { - failures.lock().unwrap().push(format!( - concat!( - "Failed to import message {} ", - "with identifier '{}': {}" - ), - message_num, message.identifier, err - )); - } else { - total_imported.fetch_add(1, Ordering::Relaxed); + let mut retry_count = 0; + loop { + match client + .email_import( + message.contents.clone(), + [mailbox_id.as_ref()], + if !message.flags.is_empty() { + message + .flags + .iter() + .map(|f| match f { + maildir::Flag::Passed => "$passed", + maildir::Flag::Replied => "$answered", + maildir::Flag::Seen => "$seen", + maildir::Flag::Trashed => "$deleted", + maildir::Flag::Draft => "$draft", + maildir::Flag::Flagged => "$flagged", + }) + .into() + } else { + None + }, + if message.internal_date > 0 { + (message.internal_date as i64).into() + } else { + None + }, + ) + .await + { + Ok(_) => { + total_imported.fetch_add(1, Ordering::Relaxed); + } + Err(_) if retry_count < RETRY_ATTEMPTS => { + retry_count += 1; + continue; + } + Err(err) => { + failures.lock().unwrap().push(format!( + concat!( + "Failed to import message {} ", + "with identifier '{}': {}" + ), + message_num, message.identifier, err + )); + } + } + break; } }); @@ -648,22 +659,33 @@ async fn import_emails( } } - if let Err(err) = client - .email_import( - contents, - mailboxes, - if !keywords.is_empty() { - Some(keywords) - } else { - None - }, - email.received_at(), - ) - .await - { - eprintln!("Failed to import emailId {id}: {err}"); - } else { - total_imported.fetch_add(1, Ordering::Relaxed); + let mut retry_count = 0; + loop { + match client + .email_import( + contents.clone(), + mailboxes.clone(), + if !keywords.is_empty() { + Some(keywords.clone()) + } else { + None + }, + email.received_at(), + ) + .await + { + Ok(_) => { + total_imported.fetch_add(1, Ordering::Relaxed); + } + Err(_) if retry_count < RETRY_ATTEMPTS => { + retry_count += 1; + continue; + } + Err(err) => { + eprintln!("Failed to import emailId {id}: {err}"); + } + } + break; } }); diff --git a/crates/cli/src/modules/mod.rs b/crates/cli/src/modules/mod.rs index d6693209..00dbc37d 100644 --- a/crates/cli/src/modules/mod.rs +++ b/crates/cli/src/modules/mod.rs @@ -38,6 +38,8 @@ pub mod import; pub mod queue; pub mod report; +const RETRY_ATTEMPTS: usize = 5; + pub trait UnwrapResult { fn unwrap_result(self, action: &str) -> T; } diff --git a/crates/imap/src/op/append.rs b/crates/imap/src/op/append.rs index 556b78ff..e23316d4 100644 --- a/crates/imap/src/op/append.rs +++ b/crates/imap/src/op/append.rs @@ -142,7 +142,7 @@ impl SessionData { keywords: message.flags.into_iter().map(Keyword::from).collect(), received_at: message.received_at.map(|d| d as u64), skip_duplicates: false, - encrypt: true, + encrypt: self.jmap.config.encrypt && self.jmap.config.encrypt_append, }) .await { diff --git a/crates/jmap/src/api/config.rs b/crates/jmap/src/api/config.rs index 8d60d7c4..ac25b386 100644 --- a/crates/jmap/src/api/config.rs +++ b/crates/jmap/src/api/config.rs @@ -138,6 +138,8 @@ impl crate::Config { principal_allow_lookups: settings .property("jmap.principal.allow-lookups")? .unwrap_or(true), + encrypt: settings.property_or_static("jmap.encryption.enable", "true")?, + encrypt_append: settings.property_or_static("jmap.encryption.append", "false")?, }; config.add_capabilites(settings); Ok(config) diff --git a/crates/jmap/src/api/http.rs b/crates/jmap/src/api/http.rs index 69bd2cec..219e3558 100644 --- a/crates/jmap/src/api/http.rs +++ b/crates/jmap/src/api/http.rs @@ -237,7 +237,7 @@ pub async fn parse_jmap_request( _ => (), } } - "crypto" => match *req.method() { + "crypto" if jmap.config.encrypt => match *req.method() { Method::GET => { return jmap.handle_crypto_update(&mut req).await; } diff --git a/crates/jmap/src/email/import.rs b/crates/jmap/src/email/import.rs index 277f3088..06ec02d6 100644 --- a/crates/jmap/src/email/import.rs +++ b/crates/jmap/src/email/import.rs @@ -141,7 +141,7 @@ impl JMAP { keywords: email.keywords, received_at: email.received_at.map(|r| r.into()), skip_duplicates: false, - encrypt: true, + encrypt: self.config.encrypt && self.config.encrypt_append, }) .await { diff --git a/crates/jmap/src/email/set.rs b/crates/jmap/src/email/set.rs index ad4eee11..ea5b3ce1 100644 --- a/crates/jmap/src/email/set.rs +++ b/crates/jmap/src/email/set.rs @@ -737,7 +737,7 @@ impl JMAP { keywords, received_at, skip_duplicates: false, - encrypt: false, + encrypt: self.config.encrypt && self.config.encrypt_append, }) .await { diff --git a/crates/jmap/src/lib.rs b/crates/jmap/src/lib.rs index b8c2989c..fa297178 100644 --- a/crates/jmap/src/lib.rs +++ b/crates/jmap/src/lib.rs @@ -150,6 +150,9 @@ pub struct Config { pub oauth_expiry_refresh_token_renew: u64, pub oauth_max_auth_attempts: u32, + pub encrypt: bool, + pub encrypt_append: bool, + pub principal_allow_lookups: bool, pub capabilities: BaseCapabilities, diff --git a/crates/jmap/src/services/ingest.rs b/crates/jmap/src/services/ingest.rs index bad64d89..e52b548c 100644 --- a/crates/jmap/src/services/ingest.rs +++ b/crates/jmap/src/services/ingest.rs @@ -104,7 +104,7 @@ impl JMAP { keywords: vec![], received_at: None, skip_duplicates: true, - encrypt: true, + encrypt: self.config.encrypt, }) .await } diff --git a/crates/jmap/src/sieve/ingest.rs b/crates/jmap/src/sieve/ingest.rs index eafd87a5..8c64b29a 100644 --- a/crates/jmap/src/sieve/ingest.rs +++ b/crates/jmap/src/sieve/ingest.rs @@ -450,7 +450,7 @@ impl JMAP { keywords: sieve_message.flags, received_at: None, skip_duplicates: true, - encrypt: true, + encrypt: self.config.encrypt, }) .await {