mirror of
https://github.com/rustmailer/bichon.git
synced 2026-08-31 01:52:30 +00:00
fix(core): count deduplicated imports as duplicates, not successes
A BLAKE3 content-hash hit in extract_envelope_core returned the same
Ok(()) as a real import, so every import surface counted silently
skipped messages as successes and the documented duplicates field
stayed 0 forever.
Make the outcome explicit: extract_envelope_core now returns
ExtractOutcome::{Imported, Duplicate} and every import loop (batch
/import, upload EML, MBOX, PST) counts Duplicate into duplicates
instead of success. total = success + duplicates + failed holds on
every path; duplicates produce no failed_details entries and an
all-duplicates run reports Completed. The batch endpoint status check
treats duplicates as processed work so a duplicates-plus-failures run
keeps reporting Completed as before. The SMTP receiver and IMAP sync
ignore the outcome.
This commit is contained in:
@@ -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<ExtractOutcome> {
|
||||
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<ExtractOutcome> {
|
||||
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<ExtractOutcome> {
|
||||
//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(
|
||||
|
||||
@@ -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<FailedItemDetail> = 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<FailedItemDetail> = 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<FailedItemDetail>) {
|
||||
) -> (usize, usize, Vec<FailedItemDetail>) {
|
||||
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<FailedItemDetail> = 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,
|
||||
};
|
||||
|
||||
@@ -17,7 +17,7 @@
|
||||
// along with this program. If not, see <http://www.gnu.org/licenses/>.
|
||||
|
||||
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<F>(
|
||||
account_id: u64,
|
||||
total: usize,
|
||||
success_count: &mut usize,
|
||||
duplicate_count: &mut usize,
|
||||
failed_details: &mut Vec<super::FailedItemDetail>,
|
||||
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<F>(
|
||||
account_id,
|
||||
total,
|
||||
success_count,
|
||||
duplicate_count,
|
||||
failed_details,
|
||||
index,
|
||||
progress_cb,
|
||||
@@ -351,11 +353,12 @@ fn process_folder_with_progress_inner<F>(
|
||||
account_id: u64,
|
||||
total: usize,
|
||||
success_count: &mut usize,
|
||||
duplicate_count: &mut usize,
|
||||
failed_details: &mut Vec<super::FailedItemDetail>,
|
||||
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<F>(
|
||||
account_id,
|
||||
total,
|
||||
success_count,
|
||||
duplicate_count,
|
||||
failed_details,
|
||||
index,
|
||||
progress_cb,
|
||||
@@ -437,9 +441,12 @@ fn process_folder_with_progress_inner<F>(
|
||||
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<F>(
|
||||
|
||||
// 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<F>(
|
||||
account_id,
|
||||
total,
|
||||
success_count,
|
||||
duplicate_count,
|
||||
failed_details,
|
||||
index,
|
||||
progress_cb,
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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 {}: {:?}",
|
||||
|
||||
Reference in New Issue
Block a user