diff --git a/src/modules/cache/vendor/gmail/sync/flow.rs b/src/modules/cache/vendor/gmail/sync/flow.rs index 07c13d6..25868fc 100644 --- a/src/modules/cache/vendor/gmail/sync/flow.rs +++ b/src/modules/cache/vendor/gmail/sync/flow.rs @@ -36,7 +36,7 @@ pub async fn fetch_and_save_since_date( // Each page returns message IDs, and we still need to fetch message details individually. let mut page_token: Option = None; let mut page = 1; // Used only for tracking sync progress - let semaphore = Arc::new(Semaphore::new(10)); + let semaphore = Arc::new(Semaphore::new(5)); let mut history_ids = Vec::new(); loop { let resp = GmailClient::list_messages( @@ -163,7 +163,7 @@ pub async fn fetch_and_save_full_label( .await?; // Update page_token returned by Gmail API page_token = resp.next_page_token; - // Concurrently fetch message details for this page, with concurrency limited to 10 + // Concurrently fetch message details for this page, with concurrency limited to 5 if let Some(messages) = resp.messages { let mut batch_messages = Vec::with_capacity(ENVELOPE_BATCH_SIZE as usize); if initial { diff --git a/src/modules/cache/vendor/tests.rs b/src/modules/cache/vendor/tests.rs index 6335c12..900962d 100644 --- a/src/modules/cache/vendor/tests.rs +++ b/src/modules/cache/vendor/tests.rs @@ -20,7 +20,7 @@ use crate::{ vendor::gmail::{ model::{ history::HistoryList, - messages::{FullMessage, MessageList, MessageMeta, PartBody}, + messages::{MessageList, MessageMeta, PartBody}, }, sync::envelope::GmailEnvelope, }, @@ -28,7 +28,6 @@ use crate::{ common::{rustls::RustMailerTls, Addr}, context::Initialize, grpc::service::rustmailer_grpc::{GetOAuth2TokensRequest, OAuth2ServiceClient}, - message::content::FullMessageContent, }, rustmailer_version, }; diff --git a/src/modules/hook/http/mod.rs b/src/modules/hook/http/mod.rs index 8a9c83b..1c491d6 100644 --- a/src/modules/hook/http/mod.rs +++ b/src/modules/hook/http/mod.rs @@ -35,7 +35,7 @@ impl HttpClient { fn base_builder() -> reqwest::ClientBuilder { reqwest::ClientBuilder::new() .user_agent(rustmailer_version!()) - .timeout(Duration::from_secs(30)) + .timeout(Duration::from_secs(60)) .connect_timeout(Duration::from_secs(10)) } @@ -113,60 +113,87 @@ impl HttpClient { /// Wrapper around the Gmail API `GET` request to fetch data. pub async fn get(&self, url: &str, access_token: &str) -> RustMailerResult { - let res = self - .client - .get(url) - .header(AUTHORIZATION, format!("Bearer {}", access_token)) - .header(CONTENT_TYPE, "application/json") - .send() - .await - .map_err(|e| { - raise_error!( - format!("Request failed: {:#?}", e), - ErrorCode::InternalError - ) - })?; + let mut attempt = 0; + let max_attempts = 3; + let mut delay_ms = 500; - if res.status().is_success() { - let json: serde_json::Value = res.json().await.map_err(|e| { - raise_error!( - format!("Failed to parse response: {:#?}", e), - ErrorCode::InternalError - ) - })?; - Ok(json) - } else { - let status = res.status(); - let text = res.text().await.map_err(|e| { - raise_error!( - format!("Failed to read error response: {:#?}", e), - ErrorCode::InternalError - ) - })?; - if matches!(status, StatusCode::NOT_FOUND) || matches!(status, StatusCode::BAD_REQUEST) - { - error!( - status = ?status, - url = %url, - response = %text, - "Gmail API client error" - ); - return Err(raise_error!( - format!( - "Gmail API returned client error (status {}) for {}. Response: {}", - status, url, text - ), - ErrorCode::GmailApiInvalidHistoryId - )); + loop { + attempt += 1; + let res_result = self + .client + .get(url) + .header(AUTHORIZATION, format!("Bearer {}", access_token)) + .header(CONTENT_TYPE, "application/json") + .send() + .await; + + match res_result { + Ok(res) => { + if res.status().is_success() { + let json: serde_json::Value = res.json().await.map_err(|e| { + raise_error!( + format!("Failed to parse response: {:#?}", e), + ErrorCode::InternalError + ) + })?; + return Ok(json); + } else { + let status = res.status(); + let text = res.text().await.map_err(|e| { + raise_error!( + format!("Failed to read error response: {:#?}", e), + ErrorCode::InternalError + ) + })?; + + if matches!(status, StatusCode::NOT_FOUND | StatusCode::BAD_REQUEST) { + error!( + status = ?status, + url = %url, + response = %text, + "Gmail API client error" + ); + return Err(raise_error!( + format!( + "Gmail API returned client error (status {}) for {}. Response: {}", + status, url, text + ), + ErrorCode::GmailApiInvalidHistoryId + )); + } + + return Err(raise_error!( + format!( + "Gmail API call to {} failed with status {}: {}", + url, status, text + ), + ErrorCode::GmailApiCallFailed + )); + } + } + Err(e) => { + if e.is_timeout() || e.is_connect() { + if attempt < max_attempts { + tokio::time::sleep(std::time::Duration::from_millis(delay_ms)).await; + delay_ms *= 2; + continue; + } else { + return Err(raise_error!( + format!( + "Request to {} failed after {} attempts: {:#?}", + url, attempt, e + ), + ErrorCode::GmailApiCallFailed + )); + } + } else { + return Err(raise_error!( + format!("Request to {} failed: {:#?}", url, e), + ErrorCode::InternalError + )); + } + } } - // Return the error with status and response text for more context - Err(raise_error!( - format!( - "Gmail API call to {} failed with status {}: {}", - url, status, text - ), - ErrorCode::GmailApiCallFailed - )) } }