use std::fs;
use std::path::{Path, PathBuf};
use indicatif::MultiProgress;
use crate::core::data::{self, ContentIndex, Dedup};
use super::manifest::CheckpointEntry;
use super::sink;
pub(crate) struct EmailDedup(pub(crate) ContentIndex);
impl Dedup for EmailDedup {
fn check(&self, hash: &str) -> Option<&str> {
self.0.check(hash)
}
fn commit(&mut self, hash: &str, relative_path: &str) -> Result<(), String> {
self.0.commit(hash, relative_path)
}
}
#[derive(Debug, Default, PartialEq, Eq)]
pub(crate) struct DedupSummary {
pub merged_messages: usize,
pub deduped_attachments: usize,
}
pub(crate) fn run_dedup_pass(
identity_dir: &Path,
staging_dir: &Path,
entries: &mut [CheckpointEntry],
message_index: &mut EmailDedup,
attachment_index: &mut EmailDedup,
multi_progress: &MultiProgress,
) -> Result<DedupSummary, String> {
entries.sort_by(|a, b| (&a.mailbox, a.uid).cmp(&(&b.mailbox, b.uid)));
let mut summary = DedupSummary::default();
let bar = sink::new_progress_bar(
"dedup".to_string(),
(entries.len() * 2) as u64,
multi_progress,
);
for entry in entries.iter() {
bar.inc(1);
if !staging_dir.join(&entry.md_staged_relpath).exists() {
continue;
}
match message_index.check(&entry.message_hash) {
Some(canonical_relpath) => {
let canonical_path = identity_dir.join(canonical_relpath);
match data::amend_frontmatter_for_duplicate(
&canonical_path,
&entry.mailbox_tag,
entry.uid,
) {
Ok(_) => {
tracing::debug!(
mailbox = %entry.mailbox,
uid = entry.uid,
action = "merged",
"message merged into existing canonical file"
);
summary.merged_messages += 1;
remove_staged_files(staging_dir, entry);
}
Err(err) => {
let _ = multi_progress.println(format!(
"Warning: canonical file for duplicate {} is missing or malformed: {err}, treating as canonical instead",
entry.md_staged_relpath
));
place_canonical_message(identity_dir, staging_dir, entry, message_index)?;
}
}
}
None => {
place_canonical_message(identity_dir, staging_dir, entry, message_index)?;
}
}
}
for entry in entries.iter() {
bar.inc(1);
let Some(message_canonical_relpath) = message_index.check(&entry.message_hash) else {
continue;
};
let md_path = identity_dir.join(message_canonical_relpath);
for (hash, staged_relpath) in &entry.attachments {
let staged_path = staged_attachment_path(staging_dir, entry, staged_relpath);
if !staged_path.exists() {
continue;
}
match attachment_index.check(hash) {
Some(canonical_relpath) => {
tracing::debug!(
mailbox = %entry.mailbox,
uid = entry.uid,
file = staged_relpath,
action = "attachment_reused",
"attachment content already present, reusing canonical copy"
);
summary.deduped_attachments += 1;
let _ = fs::remove_file(&staged_path);
data::rewrite_attachment_reference(
&md_path,
staged_relpath,
canonical_relpath,
)?;
}
None => {
let file_name = Path::new(staged_relpath)
.file_name()
.map(|name| name.to_string_lossy().into_owned())
.unwrap_or_else(|| staged_relpath.clone());
let attachments_dir = identity_dir.join("attachments");
fs::create_dir_all(&attachments_dir).map_err(|err| {
format!("failed to create {}: {err}", attachments_dir.display())
})?;
let final_path = data::unique_path(&attachments_dir.join(&file_name));
fs::rename(&staged_path, &final_path).map_err(|err| {
format!(
"failed to move {} to {}: {err}",
staged_path.display(),
final_path.display()
)
})?;
let final_relpath = format!(
"attachments/{}",
final_path.file_name().unwrap().to_string_lossy()
);
attachment_index.commit(hash, &final_relpath)?;
if final_relpath != *staged_relpath {
data::rewrite_attachment_reference(
&md_path,
staged_relpath,
&final_relpath,
)?;
}
}
}
}
}
bar.finish();
Ok(summary)
}
fn place_canonical_message(
identity_dir: &Path,
staging_dir: &Path,
entry: &CheckpointEntry,
message_index: &mut EmailDedup,
) -> Result<(), String> {
fs::create_dir_all(identity_dir)
.map_err(|err| format!("failed to create {}: {err}", identity_dir.display()))?;
let final_path = data::unique_path(&identity_dir.join(&entry.desired_md_name));
let staged_path = staging_dir.join(&entry.md_staged_relpath);
fs::rename(&staged_path, &final_path).map_err(|err| {
format!(
"failed to move {} to {}: {err}",
staged_path.display(),
final_path.display()
)
})?;
let final_relpath = final_path
.file_name()
.unwrap()
.to_string_lossy()
.into_owned();
message_index.commit(&entry.message_hash, &final_relpath)?;
Ok(())
}
fn remove_staged_files(staging_dir: &Path, entry: &CheckpointEntry) {
let _ = fs::remove_file(staging_dir.join(&entry.md_staged_relpath));
for (_, relpath) in &entry.attachments {
let _ = fs::remove_file(staged_attachment_path(staging_dir, entry, relpath));
}
}
fn staged_attachment_path(staging_dir: &Path, entry: &CheckpointEntry, relpath: &str) -> PathBuf {
let md_parent = Path::new(&entry.md_staged_relpath)
.parent()
.unwrap_or_else(|| Path::new(""));
let file_name = Path::new(relpath)
.file_name()
.map(|name| name.to_string_lossy().into_owned())
.unwrap_or_else(|| relpath.to_string());
staging_dir
.join(md_parent)
.join(entry.uid.to_string())
.join("attachments")
.join(file_name)
}
#[cfg(test)]
mod tests {
use super::*;
use crate::commands::job::email_sync::{manifest, transform};
fn entry(
mailbox: &str,
uid: u32,
message_hash: &str,
desired_md_name: &str,
attachments: Vec<(&str, &str)>,
) -> CheckpointEntry {
CheckpointEntry {
mailbox: mailbox.to_string(),
uid,
message_hash: message_hash.to_string(),
md_staged_relpath: format!("transformed/{mailbox}/{uid}.md"),
desired_md_name: desired_md_name.to_string(),
mailbox_tag: format!("mailbox/{}", mailbox.to_lowercase()),
attachments: attachments
.into_iter()
.map(|(hash, relpath)| (hash.to_string(), relpath.to_string()))
.collect(),
}
}
fn stage_message(staging_dir: &Path, entry: &CheckpointEntry, body: &str) {
let path = staging_dir.join(&entry.md_staged_relpath);
fs::create_dir_all(path.parent().unwrap()).unwrap();
fs::write(path, body).unwrap();
}
fn stage_attachment(staging_dir: &Path, relpath: &str, contents: &[u8]) {
let path = staging_dir.join(relpath);
fs::create_dir_all(path.parent().unwrap()).unwrap();
fs::write(path, contents).unwrap();
}
const FIXTURE_BODY: &str = "---\nfrom: \"a\"\ntags:\n - mailbox/inbox\n---\nbody";
fn indexes(staging: &Path) -> (EmailDedup, EmailDedup) {
(
EmailDedup(ContentIndex::load(staging, transform::MESSAGE_HASHES_FILE).unwrap()),
EmailDedup(ContentIndex::load(staging, transform::ATTACHMENT_HASHES_FILE).unwrap()),
)
}
#[test]
fn run_dedup_pass_places_a_lone_canonical_message() {
let staging = tempfile::tempdir().unwrap();
let identity_dir = tempfile::tempdir().unwrap();
let (mut message_index, mut attachment_index) = indexes(staging.path());
let e = entry("INBOX", 1, "hash-a", "2024-01-26-hello.md", vec![]);
stage_message(staging.path(), &e, FIXTURE_BODY);
let mut entries = vec![e];
let summary = run_dedup_pass(
identity_dir.path(),
staging.path(),
&mut entries,
&mut message_index,
&mut attachment_index,
&MultiProgress::new(),
)
.unwrap();
assert_eq!(summary.merged_messages, 0);
assert!(identity_dir.path().join("2024-01-26-hello.md").exists());
assert!(!staging.path().join("transformed/INBOX/1.md").exists());
assert_eq!(message_index.check("hash-a"), Some("2024-01-26-hello.md"));
}
#[test]
fn run_dedup_pass_merges_two_messages_with_the_same_hash() {
let staging = tempfile::tempdir().unwrap();
let identity_dir = tempfile::tempdir().unwrap();
let (mut message_index, mut attachment_index) = indexes(staging.path());
let first = entry("Archive", 2, "same-hash", "2024-01-26-hello.md", vec![]);
let second = entry("INBOX", 1, "same-hash", "2024-01-26-hello.md", vec![]);
stage_message(staging.path(), &first, FIXTURE_BODY);
stage_message(staging.path(), &second, FIXTURE_BODY);
let mut entries = vec![second, first];
let summary = run_dedup_pass(
identity_dir.path(),
staging.path(),
&mut entries,
&mut message_index,
&mut attachment_index,
&MultiProgress::new(),
)
.unwrap();
assert_eq!(summary.merged_messages, 1);
let md_files: Vec<_> = fs::read_dir(identity_dir.path())
.unwrap()
.filter_map(|e| e.ok())
.filter(|e| e.path().extension().and_then(|x| x.to_str()) == Some("md"))
.collect();
assert_eq!(md_files.len(), 1);
let contents = fs::read_to_string(md_files[0].path()).unwrap();
assert!(contents.contains("also-in:"));
assert!(contents.contains("mailbox/inbox#1"));
}
#[test]
fn run_dedup_pass_dedupes_cross_message_attachment_and_rewrites_reference() {
let staging = tempfile::tempdir().unwrap();
let identity_dir = tempfile::tempdir().unwrap();
let (mut message_index, mut attachment_index) = indexes(staging.path());
let first = entry(
"Archive",
1,
"hash-1",
"2024-01-26-first.md",
vec![("attach-hash", "attachments/a.pdf")],
);
let second = entry(
"INBOX",
2,
"hash-2",
"2024-01-27-second.md",
vec![("attach-hash", "attachments/a.pdf")],
);
let body_with_attachment = "---\nfrom: \"a\"\ntags:\n - mailbox/inbox\nattachments:\n - attachments/a.pdf\n---\nbody";
stage_message(staging.path(), &first, body_with_attachment);
stage_message(staging.path(), &second, body_with_attachment);
stage_attachment(
staging.path(),
"transformed/Archive/1/attachments/a.pdf",
b"content-1",
);
stage_attachment(
staging.path(),
"transformed/INBOX/2/attachments/a.pdf",
b"content-1",
);
let mut entries = vec![first, second];
let summary = run_dedup_pass(
identity_dir.path(),
staging.path(),
&mut entries,
&mut message_index,
&mut attachment_index,
&MultiProgress::new(),
)
.unwrap();
assert_eq!(summary.deduped_attachments, 1);
let attachments_dir = identity_dir.path().join("attachments");
assert_eq!(fs::read_dir(&attachments_dir).unwrap().count(), 1);
let second_md = identity_dir.path().join("2024-01-27-second.md");
let contents = fs::read_to_string(&second_md).unwrap();
assert!(contents.contains("attachments:\n - attachments/a.pdf"));
}
#[test]
fn run_dedup_pass_places_attachment_for_a_message_already_canonicalized_by_an_earlier_run() {
let staging = tempfile::tempdir().unwrap();
let identity_dir = tempfile::tempdir().unwrap();
let (mut message_index, mut attachment_index) = indexes(staging.path());
let e = entry(
"INBOX",
1,
"hash-a",
"2024-01-26-hello.md",
vec![("attach-hash", "attachments/a.pdf")],
);
fs::write(
identity_dir.path().join("2024-01-26-hello.md"),
FIXTURE_BODY,
)
.unwrap();
message_index
.commit("hash-a", "2024-01-26-hello.md")
.unwrap();
stage_attachment(
staging.path(),
"transformed/INBOX/1/attachments/a.pdf",
b"content",
);
let mut entries = vec![e];
let summary = run_dedup_pass(
identity_dir.path(),
staging.path(),
&mut entries,
&mut message_index,
&mut attachment_index,
&MultiProgress::new(),
)
.unwrap();
assert_eq!(summary.merged_messages, 0);
let attachments_dir = identity_dir.path().join("attachments");
assert_eq!(
fs::read_dir(&attachments_dir)
.unwrap_or_else(|err| panic!("{} should exist: {err}", attachments_dir.display()))
.count(),
1,
"the attachment should be placed even though its message was canonicalized before this call"
);
}
#[test]
fn run_dedup_pass_is_idempotent_when_rerun_on_the_same_entries() {
let staging = tempfile::tempdir().unwrap();
let identity_dir = tempfile::tempdir().unwrap();
let (mut message_index, mut attachment_index) = indexes(staging.path());
let e = entry("INBOX", 1, "hash-a", "2024-01-26-hello.md", vec![]);
stage_message(staging.path(), &e, FIXTURE_BODY);
let mut entries = vec![e];
run_dedup_pass(
identity_dir.path(),
staging.path(),
&mut entries,
&mut message_index,
&mut attachment_index,
&MultiProgress::new(),
)
.unwrap();
let after_first_run =
fs::read_to_string(identity_dir.path().join("2024-01-26-hello.md")).unwrap();
let summary = run_dedup_pass(
identity_dir.path(),
staging.path(),
&mut entries,
&mut message_index,
&mut attachment_index,
&MultiProgress::new(),
)
.unwrap();
assert_eq!(summary, DedupSummary::default());
assert_eq!(
fs::read_to_string(identity_dir.path().join("2024-01-26-hello.md")).unwrap(),
after_first_run,
"re-running the pass must not append a spurious self-referential also-in entry"
);
}
#[test]
fn run_dedup_pass_places_attachment_from_real_transform_and_checkpoint_round_trip() {
use crate::commands::keyring::email::identity::Identity;
use crate::commands::keyring::email::provider::Provider;
use crate::core::data::Transform;
use transform::EmailTransform;
let staging = tempfile::tempdir().unwrap();
let identity_dir = tempfile::tempdir().unwrap();
let input = tempfile::tempdir().unwrap();
let inbox = input.path().join("inbox");
fs::create_dir_all(&inbox).unwrap();
let eml = "From: Jane Doe <jane.doe@example.com>\r\n\
To: first.last@example.com\r\n\
Subject: Shipping\r\n\
Date: Fri, 26 Jan 2024 09:15:00 +0000\r\n\
MIME-Version: 1.0\r\n\
Content-Type: multipart/mixed; boundary=\"BOUNDARY\"\r\n\
\r\n\
--BOUNDARY\r\n\
Content-Type: text/plain; charset=utf-8\r\n\
\r\n\
Hello there!\r\n\
--BOUNDARY\r\n\
Content-Type: application/pdf\r\n\
Content-Disposition: attachment; filename=\"a.pdf\"\r\n\
Content-Transfer-Encoding: base64\r\n\
\r\n\
JVBERi0xLjQK\r\n\
--BOUNDARY--\r\n";
fs::write(inbox.join("1.eml"), eml).unwrap();
let transformer = EmailTransform {
identity: Identity {
alias: "first-last".to_string(),
email: "first.last@example.com".to_string(),
provider: Provider::Gmail,
host: "imap.gmail.com".to_string(),
port: 993,
max_imap_connections: None,
},
input_root: input.path().to_path_buf(),
staging_root: staging.path().to_path_buf(),
};
let outcome = transformer.transform(inbox.join("1.eml")).unwrap().unwrap();
assert_eq!(
outcome.attachments.len(),
1,
"fixture message should have exactly one attachment"
);
let checkpoint_entry = CheckpointEntry {
mailbox: "inbox".to_string(),
uid: 1,
message_hash: outcome.message_hash,
md_staged_relpath: outcome.md_staged_relpath,
desired_md_name: outcome.desired_md_name,
mailbox_tag: outcome.mailbox_tag,
attachments: outcome
.attachments
.into_iter()
.map(|attachment| (attachment.hash, attachment.staged_relpath))
.collect(),
};
manifest::append_checkpoint(staging.path(), &checkpoint_entry).unwrap();
let mut entries = manifest::load_checkpoint(staging.path()).unwrap();
let (mut message_index, mut attachment_index) = indexes(staging.path());
run_dedup_pass(
identity_dir.path(),
staging.path(),
&mut entries,
&mut message_index,
&mut attachment_index,
&MultiProgress::new(),
)
.unwrap();
let attachments_dir = identity_dir.path().join("attachments");
let placed = fs::read_dir(&attachments_dir)
.unwrap_or_else(|err| panic!("{} should exist: {err}", attachments_dir.display()))
.count();
assert_eq!(
placed, 1,
"the real staged attachment should be placed under identity_dir/attachments/"
);
}
}