diff --git a/crates/core/src/envelope/extractor.rs b/crates/core/src/envelope/extractor.rs index 88195a7..0b2e821 100644 --- a/crates/core/src/envelope/extractor.rs +++ b/crates/core/src/envelope/extractor.rs @@ -42,6 +42,15 @@ use tantivy::schema::Facet; use tracing::error; use uuid::Uuid; +/// The outcome of extracting an envelope: either the message was imported, +/// or it was skipped because its content hash was already archived. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +#[must_use] +pub enum ExtractOutcome { + Imported, + Duplicate, +} + pub async fn extract_envelope_and_store_it( fetch: Fetch, account_id: u64, @@ -64,14 +73,16 @@ pub async fn extract_envelope_and_store_it( } }; let size = fetch.size.unwrap_or(body.len() as u32); - extract_envelope_core(body, uid, size, internal_date, account_id, mailbox_id).await + extract_envelope_core(body, uid, size, internal_date, account_id, mailbox_id) + .await + .map(|_| ()) } pub async fn extract_envelope_from_eml( body: &[u8], account_id: u64, mailbox_id: u64, -) -> BichonResult<()> { +) -> BichonResult { extract_envelope_core(body, 0, body.len() as u32, 0, account_id, mailbox_id).await } @@ -79,7 +90,7 @@ pub async fn extract_envelope_from_smtp( body: &[u8], account_id: u64, mailbox_id: u64, -) -> BichonResult<()> { +) -> BichonResult { extract_envelope_core( body, 0, @@ -98,13 +109,13 @@ async fn extract_envelope_core( internal_date: i64, account_id: u64, mailbox_id: u64, -) -> BichonResult<()> { +) -> BichonResult { //The content hash of the original raw EML let email_content_hash = compute_content_hash(body); if DEDUP_CACHE.contains(account_id, mailbox_id, &email_content_hash) { tracing::debug!("Duplicate email detected"); //println!("Duplicate email detected"); - return Ok(()); + return Ok(ExtractOutcome::Duplicate); } let message: Message<'_> = MessageParser::new().parse(body).ok_or_else(|| { raise_error!( @@ -136,7 +147,7 @@ async fn extract_envelope_core( subject = subject.as_deref().unwrap_or("?"), "Email filtered out by archive rules" ); - return Ok(()); + return Ok(ExtractOutcome::Imported); } } } @@ -317,7 +328,7 @@ async fn extract_envelope_core( for doc in attachment_docs { ATTACHMENT_MANAGER.queue(doc).await; } - Ok(()) + Ok(ExtractOutcome::Imported) } pub fn extract_envelope_from_nested_message( diff --git a/crates/core/src/import/mod.rs b/crates/core/src/import/mod.rs index bd8db2a..6fc7b5c 100644 --- a/crates/core/src/import/mod.rs +++ b/crates/core/src/import/mod.rs @@ -34,7 +34,7 @@ use crate::{ { account::migration::{AccountModel, AccountType}, cache::imap::mailbox::{Attribute, AttributeEnum, MailBox}, - envelope::extractor::extract_envelope_from_eml, + envelope::extractor::{extract_envelope_from_eml, ExtractOutcome}, error::{BichonResult, code::ErrorCode}, settings::dir::DATA_DIR_MANAGER, utils::create_hash, @@ -132,6 +132,7 @@ impl ImportEmls { let account_id = account.id; let mut success_count = 0; + let mut duplicate_count = 0; let mut failed_details: Vec = Vec::new(); // Store failure details let total = request.emls.len(); @@ -169,9 +170,12 @@ impl ImportEmls { } match extract_envelope_from_eml(&decoded, account_id, mailbox_id).await { - Ok(_) => { + Ok(ExtractOutcome::Imported) => { success_count += 1; }, + Ok(ExtractOutcome::Duplicate) => { + duplicate_count += 1; + }, Err(e) => { let error_msg = format!( "Failed to extract envelope from EML at index {}: {:?}", @@ -194,7 +198,7 @@ impl ImportEmls { Ok(BatchEmlResult { total, success: success_count, - duplicates: 0, + duplicates: duplicate_count, failed: failed_count, failed_details, // Return the list of failure details }) @@ -566,7 +570,8 @@ fn process_eml_file( failed_details: vec![], }); - let (success_count, failed_details) = process_single_eml(&file_bytes, 0, account_id, mailbox_id); + let (success_count, duplicate_count, failed_details) = + process_single_eml(&file_bytes, 0, account_id, mailbox_id); // Clean up let _ = std::fs::remove_file(file_path); @@ -577,7 +582,7 @@ fn process_eml_file( format: "eml".to_string(), total, success: success_count, - duplicates: 0, + duplicates: duplicate_count, failed: failed_details.len(), failed_details, }; @@ -619,6 +624,7 @@ fn process_mbox_file( }); let mut success_count = 0usize; + let mut duplicate_count = 0usize; let mut failed_details: Vec = Vec::new(); for (index, entry) in mbox.iter().enumerate() { @@ -639,9 +645,12 @@ fn process_mbox_file( } match futures::executor::block_on(extract_envelope_from_eml(eml_bytes, account_id, mailbox_id)) { - Ok(_) => { + Ok(ExtractOutcome::Imported) => { success_count += 1; } + Ok(ExtractOutcome::Duplicate) => { + duplicate_count += 1; + } Err(e) => { failed_details.push(FailedItemDetail { index, @@ -658,7 +667,7 @@ fn process_mbox_file( format: "mbox".to_string(), total, success: success_count, - duplicates: 0, + duplicates: duplicate_count, failed: failed_details.len(), failed_details: failed_details.clone(), }); @@ -675,7 +684,7 @@ fn process_mbox_file( format: "mbox".to_string(), total, success: success_count, - duplicates: 0, + duplicates: duplicate_count, failed: failed_details.len(), failed_details, }; @@ -683,16 +692,16 @@ fn process_mbox_file( update_progress(import_id, final_progress); } -/// Process a single EML byte slice and return (success_count, failed_details). +/// Process a single EML byte slice and return (success_count, duplicate_count, failed_details). fn process_single_eml( eml_bytes: &[u8], index: usize, account_id: u64, mailbox_id: u64, -) -> (usize, Vec) { +) -> (usize, usize, Vec) { if eml_bytes.len() > MAX_SINGLE_EML_BYTES { let size_mb = eml_bytes.len() as f64 / 1024.0 / 1024.0; - return (0, vec![FailedItemDetail { + return (0, 0, vec![FailedItemDetail { index, error_message: format!( "Email is {:.1} MB (limit {} MB). Skipping.", @@ -703,8 +712,9 @@ fn process_single_eml( } match futures::executor::block_on(extract_envelope_from_eml(eml_bytes, account_id, mailbox_id)) { - Ok(_) => (1, vec![]), - Err(e) => (0, vec![FailedItemDetail { + Ok(ExtractOutcome::Imported) => (1, 0, vec![]), + Ok(ExtractOutcome::Duplicate) => (0, 1, vec![]), + Err(e) => (0, 0, vec![FailedItemDetail { index, error_message: format!("{:?}", e), }]), @@ -744,6 +754,7 @@ fn process_pst_upload( // Pass 2: process messages with progress updates let mut success_count: usize = 0; + let mut duplicate_count: usize = 0; let mut failed_details: Vec = Vec::new(); let mut index: usize = 0; @@ -783,16 +794,17 @@ fn process_pst_upload( account_id, total, // pass pre-counted total for accurate progress &mut success_count, + &mut duplicate_count, &mut failed_details, &mut index, - &|processed, actual_failed| { + &|success, duplicates, actual_failed| { update_progress(&import_id, ImportProgress { import_id: import_id.clone(), status: ImportStatus::Processing, format: format_str.clone(), total, - success: processed - actual_failed, - duplicates: 0, + success, + duplicates, failed: actual_failed, failed_details: vec![], }); @@ -808,7 +820,7 @@ fn process_pst_upload( format: "pst".to_string(), total, success: success_count, - duplicates: 0, + duplicates: duplicate_count, failed: failed_details.len(), failed_details, }; diff --git a/crates/core/src/import/pst/mod.rs b/crates/core/src/import/pst/mod.rs index 2fe7992..386e461 100644 --- a/crates/core/src/import/pst/mod.rs +++ b/crates/core/src/import/pst/mod.rs @@ -17,7 +17,7 @@ // along with this program. If not, see . use crate::base64_encode_url_safe; -use crate::envelope::extractor::extract_envelope_from_eml; +use crate::envelope::extractor::{extract_envelope_from_eml, ExtractOutcome}; use chrono::{DateTime, TimeZone, Utc}; use mail_send::mail_builder::headers::text::Text; use mail_send::mail_builder::MessageBuilder; @@ -327,11 +327,12 @@ pub fn process_folder_with_progress( account_id: u64, total: usize, success_count: &mut usize, + duplicate_count: &mut usize, failed_details: &mut Vec, index: &mut usize, progress_cb: &F, ) where - F: Fn(usize, usize), // (processed, failed) + F: Fn(usize, usize, usize), // (success, duplicates, failed) { process_folder_with_progress_inner( folder, @@ -339,6 +340,7 @@ pub fn process_folder_with_progress( account_id, total, success_count, + duplicate_count, failed_details, index, progress_cb, @@ -351,11 +353,12 @@ fn process_folder_with_progress_inner( account_id: u64, total: usize, success_count: &mut usize, + duplicate_count: &mut usize, failed_details: &mut Vec, index: &mut usize, progress_cb: &F, ) where - F: Fn(usize, usize), + F: Fn(usize, usize, usize), { let folder_name = folder .properties() @@ -386,6 +389,7 @@ fn process_folder_with_progress_inner( account_id, total, success_count, + duplicate_count, failed_details, index, progress_cb, @@ -437,9 +441,12 @@ fn process_folder_with_progress_inner( match futures::executor::block_on( extract_envelope_from_eml(&decoded, account_id, mailbox_id) ) { - Ok(_) => { + Ok(ExtractOutcome::Imported) => { *success_count += 1; } + Ok(ExtractOutcome::Duplicate) => { + *duplicate_count += 1; + } Err(e) => { failed_details.push(super::FailedItemDetail { index: *index, @@ -459,7 +466,7 @@ fn process_folder_with_progress_inner( // Report progress every 50 messages if batch_size % 50 == 0 { - progress_cb(*success_count + failed_details.len(), failed_details.len()); + progress_cb(*success_count, *duplicate_count, failed_details.len()); } } } @@ -475,6 +482,7 @@ fn process_folder_with_progress_inner( account_id, total, success_count, + duplicate_count, failed_details, index, progress_cb, diff --git a/crates/server/src/rest/api/import.rs b/crates/server/src/rest/api/import.rs index 100e415..8a07377 100644 --- a/crates/server/src/rest/api/import.rs +++ b/crates/server/src/rest/api/import.rs @@ -88,7 +88,7 @@ impl ImportApi { ), status: if result.failed == 0 { ImportStatus::Completed - } else if result.success == 0 { + } else if result.success == 0 && result.duplicates == 0 { ImportStatus::Failed } else { ImportStatus::Completed diff --git a/crates/smtp/src/server.rs b/crates/smtp/src/server.rs index ad38ff4..7b132f3 100644 --- a/crates/smtp/src/server.rs +++ b/crates/smtp/src/server.rs @@ -666,6 +666,7 @@ async fn parse_email(data: &[u8], session: &Session) -> BichonResult<()> { extract_envelope_from_smtp(data, rcpt.id, mailbox_id) .await + .map(|_| ()) .map_err(|e| { tracing::error!( "SMTP: Envelope extraction failed for {}: {:?}",