use std::collections::HashSet;
use std::fs;
use std::io::Write;
use std::path::{Path, PathBuf};
use async_imap::imap_proto::types::BodyStructure;
use futures::TryStreamExt;
use crate::commands::keyring::email::imap_client::ImapSession;
const MANIFEST_FILE_NAME: &str = ".manifest";
const CHECKPOINT_FILE_NAME: &str = ".job-checkpoint";
pub(crate) struct ManifestEntry {
pub mailbox: String,
pub uid: u32,
pub size: u64,
pub attachments: u32,
}
fn count_attachments(structure: &BodyStructure) -> u32 {
match structure {
BodyStructure::Multipart { bodies, .. } => bodies.iter().map(count_attachments).sum(),
BodyStructure::Basic { common, .. }
| BodyStructure::Text { common, .. }
| BodyStructure::Message { common, .. } => {
let has_attachment_disposition = common
.disposition
.as_ref()
.is_some_and(|d| d.ty.eq_ignore_ascii_case("attachment"));
let has_content_type_name = common.ty.params.as_ref().is_some_and(|params| {
params
.iter()
.any(|(key, _)| key.eq_ignore_ascii_case("name"))
});
u32::from(has_attachment_disposition || has_content_type_name)
}
}
}
async fn fetch_manifest_batch(
session: &mut ImapSession,
mailbox_name: &str,
uids: &[u32],
) -> Result<Vec<ManifestEntry>, String> {
let uid_set = uids
.iter()
.map(|uid| uid.to_string())
.collect::<Vec<_>>()
.join(",");
let mut fetches = session
.uid_fetch(&uid_set, "(UID RFC822.SIZE BODYSTRUCTURE)")
.await
.map_err(|err| format!("failed to fetch message sizes in '{mailbox_name}': {err}"))?;
let mut entries = Vec::new();
while let Some(fetch) = fetches
.try_next()
.await
.map_err(|err| format!("failed to fetch message sizes in '{mailbox_name}': {err}"))?
{
let (Some(uid), Some(size)) = (fetch.uid, fetch.size) else {
continue;
};
let attachments = fetch.bodystructure().map_or(0, count_attachments);
entries.push(ManifestEntry {
mailbox: mailbox_name.to_string(),
uid,
size: u64::from(size),
attachments,
});
}
Ok(entries)
}
enum BisectOutcome {
Fetched(Vec<ManifestEntry>),
Poisoned(ManifestEntry),
Retry(Vec<u32>, Vec<u32>),
}
fn bisect_step(
mailbox_name: &str,
batch: &[u32],
result: Result<Vec<ManifestEntry>, String>,
) -> BisectOutcome {
match result {
Ok(entries) => BisectOutcome::Fetched(entries),
Err(_) if batch.len() == 1 => BisectOutcome::Poisoned(ManifestEntry {
mailbox: mailbox_name.to_string(),
uid: batch[0],
size: 0,
attachments: 0,
}),
Err(_) => {
let mid = batch.len() / 2;
BisectOutcome::Retry(batch[..mid].to_vec(), batch[mid..].to_vec())
}
}
}
pub(crate) async fn pull_manifest(
session: &mut ImapSession,
mailbox_name: &str,
pending: &[u32],
) -> Result<Vec<ManifestEntry>, String> {
if pending.is_empty() {
return Ok(Vec::new());
}
let mut work = vec![pending.to_vec()];
let mut entries = Vec::new();
while let Some(batch) = work.pop() {
let result = fetch_manifest_batch(session, mailbox_name, &batch).await;
match bisect_step(mailbox_name, &batch, result) {
BisectOutcome::Fetched(fetched) => entries.extend(fetched),
BisectOutcome::Poisoned(entry) => {
eprintln!(
"warning: skipping UID {} in '{mailbox_name}': BODYSTRUCTURE could not be \
parsed (likely a MESSAGE/GLOBAL, MESSAGE/DELIVERY-STATUS, or other \
non-RFC822 embedded-message part -- a known imap-proto limitation)",
entry.uid
);
entries.push(entry);
}
BisectOutcome::Retry(left, right) => {
work.push(right);
work.push(left);
}
}
}
Ok(entries)
}
pub(crate) fn save_manifest(staging_dir: &Path, entries: &[ManifestEntry]) -> Result<(), String> {
let path = staging_dir.join(MANIFEST_FILE_NAME);
let mut contents = String::new();
for entry in entries {
contents.push_str(&format!(
"{}\t{}\t{}\t{}\n",
entry.mailbox, entry.uid, entry.size, entry.attachments
));
}
fs::write(&path, contents).map_err(|err| format!("failed to write {}: {err}", path.display()))
}
pub(crate) struct CheckpointEntry {
pub mailbox: String,
pub uid: u32,
pub message_hash: String,
pub md_staged_relpath: String,
pub desired_md_name: String,
pub mailbox_tag: String,
pub attachments: Vec<(String, String)>,
}
pub(crate) fn append_checkpoint(staging_dir: &Path, entry: &CheckpointEntry) -> Result<(), String> {
let path = staging_dir.join(CHECKPOINT_FILE_NAME);
let mut file = fs::OpenOptions::new()
.create(true)
.append(true)
.open(&path)
.map_err(|err| format!("failed to open {}: {err}", path.display()))?;
writeln!(file, "{}", format_checkpoint_line(entry))
.map_err(|err| format!("failed to write {}: {err}", path.display()))
}
fn format_checkpoint_line(entry: &CheckpointEntry) -> String {
let attachments = entry
.attachments
.iter()
.map(|(hash, relpath)| format!("{hash}={relpath}"))
.collect::<Vec<_>>()
.join(";");
format!(
"{}\t{}\t{}\t{}\t{}\t{}\t{}",
entry.mailbox,
entry.uid,
entry.message_hash,
entry.md_staged_relpath,
entry.desired_md_name,
entry.mailbox_tag,
attachments
)
}
pub(crate) fn clear_checkpoint_for_mailbox(
staging_dir: &Path,
mailbox: &str,
) -> Result<(), String> {
let entries = load_checkpoint(staging_dir)?;
let contents = entries
.iter()
.filter(|entry| entry.mailbox != mailbox)
.map(|entry| format!("{}\n", format_checkpoint_line(entry)))
.collect::<String>();
let path = staging_dir.join(CHECKPOINT_FILE_NAME);
fs::write(&path, contents).map_err(|err| format!("failed to write {}: {err}", path.display()))
}
pub(crate) fn load_checkpoint(staging_dir: &Path) -> Result<Vec<CheckpointEntry>, String> {
let path = staging_dir.join(CHECKPOINT_FILE_NAME);
let contents = match fs::read_to_string(&path) {
Ok(contents) => contents,
Err(err) if err.kind() == std::io::ErrorKind::NotFound => return Ok(Vec::new()),
Err(err) => return Err(format!("failed to read {}: {err}", path.display())),
};
Ok(contents
.lines()
.filter_map(|line| {
let mut parts = line.split('\t');
let mailbox = parts.next()?.to_string();
let uid = parts.next()?.parse().ok()?;
let message_hash = parts.next()?.to_string();
let md_staged_relpath = parts.next()?.to_string();
let desired_md_name = parts.next()?.to_string();
let mailbox_tag = parts.next()?.to_string();
let attachments_field = parts.next()?;
let attachments = if attachments_field.is_empty() {
Vec::new()
} else {
attachments_field
.split(';')
.filter_map(|pair| pair.split_once('='))
.map(|(hash, relpath)| (hash.to_string(), relpath.to_string()))
.collect()
};
Some(CheckpointEntry {
mailbox,
uid,
message_hash,
md_staged_relpath,
desired_md_name,
mailbox_tag,
attachments,
})
})
.collect())
}
pub(crate) fn done_uids(entries: &[CheckpointEntry]) -> HashSet<(String, u32)> {
entries
.iter()
.map(|entry| (entry.mailbox.clone(), entry.uid))
.collect()
}
pub(crate) struct Batch {
pub mailbox: String,
pub mailbox_relpath: PathBuf,
pub uids: Vec<u32>,
}
const MIN_BATCHES_PER_WORKER: usize = 4;
pub(crate) fn split_into_batches(uids: &[u32], concurrency: usize) -> Vec<Vec<u32>> {
if uids.is_empty() {
return Vec::new();
}
let batch_size = (uids.len() / (concurrency.max(1) * MIN_BATCHES_PER_WORKER)).max(1);
uids.chunks(batch_size).map(<[u32]>::to_vec).collect()
}
#[cfg(test)]
mod tests {
use async_imap::imap_proto::types::{
BodyContentCommon, BodyContentSinglePart, ContentDisposition, ContentEncoding, ContentType,
};
use super::*;
fn entry(uid: u32) -> ManifestEntry {
ManifestEntry {
mailbox: "INBOX".to_string(),
uid,
size: 100,
attachments: 0,
}
}
#[test]
fn bisect_step_success_returns_fetched() {
let outcome = bisect_step("INBOX", &[1, 2, 3], Ok(vec![entry(1), entry(2), entry(3)]));
assert!(matches!(outcome, BisectOutcome::Fetched(entries) if entries.len() == 3));
}
#[test]
fn bisect_step_single_uid_failure_is_poisoned_with_zeroed_placeholder() {
let outcome = bisect_step("INBOX", &[42], Err("parse error".to_string()));
let BisectOutcome::Poisoned(placeholder) = outcome else {
panic!("expected Poisoned");
};
assert_eq!(placeholder.mailbox, "INBOX");
assert_eq!(placeholder.uid, 42);
assert_eq!(placeholder.size, 0);
assert_eq!(placeholder.attachments, 0);
}
#[test]
fn bisect_step_multi_uid_failure_splits_in_half() {
let outcome = bisect_step("INBOX", &[1, 2, 3, 4], Err("parse error".to_string()));
let BisectOutcome::Retry(left, right) = outcome else {
panic!("expected Retry");
};
assert_eq!(left, vec![1, 2]);
assert_eq!(right, vec![3, 4]);
}
#[test]
fn bisect_step_odd_length_failure_splits_with_extra_uid_on_the_right() {
let outcome = bisect_step("INBOX", &[1, 2, 3], Err("parse error".to_string()));
let BisectOutcome::Retry(left, right) = outcome else {
panic!("expected Retry");
};
assert_eq!(left, vec![1]);
assert_eq!(right, vec![2, 3]);
}
#[test]
fn bisect_step_chained_isolates_a_single_poisoned_uid_among_healthy_ones() {
let fake_fetch = |batch: &[u32]| -> Result<Vec<ManifestEntry>, String> {
if batch.contains(&3) {
Err("parse error".to_string())
} else {
Ok(batch.iter().map(|&uid| entry(uid)).collect())
}
};
let mut work = vec![vec![1u32, 2, 3, 4]];
let mut fetched = Vec::new();
let mut poisoned = Vec::new();
while let Some(batch) = work.pop() {
match bisect_step("INBOX", &batch, fake_fetch(&batch)) {
BisectOutcome::Fetched(entries) => fetched.extend(entries),
BisectOutcome::Poisoned(entry) => poisoned.push(entry),
BisectOutcome::Retry(left, right) => {
work.push(right);
work.push(left);
}
}
}
assert_eq!(poisoned.len(), 1);
assert_eq!(poisoned[0].uid, 3);
let mut fetched_uids: Vec<u32> = fetched.iter().map(|e| e.uid).collect();
fetched_uids.sort_unstable();
assert_eq!(fetched_uids, vec![1, 2, 4]);
}
#[test]
fn save_manifest_writes_tab_separated_four_column_format() {
let dir = tempfile::tempdir().unwrap();
let entries = vec![
ManifestEntry {
mailbox: "INBOX".to_string(),
uid: 5,
size: 1024,
attachments: 2,
},
ManifestEntry {
mailbox: "Sent Items".to_string(),
uid: 9,
size: 2048,
attachments: 0,
},
];
save_manifest(dir.path(), &entries).unwrap();
let contents = fs::read_to_string(dir.path().join(MANIFEST_FILE_NAME)).unwrap();
assert_eq!(contents, "INBOX\t5\t1024\t2\nSent Items\t9\t2048\t0\n");
}
#[test]
fn save_manifest_is_a_full_snapshot_not_an_append() {
let dir = tempfile::tempdir().unwrap();
save_manifest(
dir.path(),
&[ManifestEntry {
mailbox: "INBOX".to_string(),
uid: 1,
size: 10,
attachments: 0,
}],
)
.unwrap();
save_manifest(
dir.path(),
&[ManifestEntry {
mailbox: "INBOX".to_string(),
uid: 2,
size: 20,
attachments: 1,
}],
)
.unwrap();
let contents = fs::read_to_string(dir.path().join(MANIFEST_FILE_NAME)).unwrap();
assert_eq!(contents, "INBOX\t2\t20\t1\n");
}
#[test]
fn count_attachments_bare_part_with_no_disposition_is_zero() {
let structure = BodyStructure::Text {
common: BodyContentCommon {
ty: ContentType {
ty: "text".into(),
subtype: "plain".into(),
params: None,
},
disposition: None,
language: None,
location: None,
},
other: BodyContentSinglePart {
id: None,
md5: None,
description: None,
transfer_encoding: ContentEncoding::SevenBit,
octets: 100,
},
lines: 5,
extension: None,
};
assert_eq!(count_attachments(&structure), 0);
}
#[test]
fn count_attachments_content_type_name_param_counts_without_disposition() {
let structure = BodyStructure::Basic {
common: BodyContentCommon {
ty: ContentType {
ty: "application".into(),
subtype: "pdf".into(),
params: Some(vec![("name".into(), "invoice.pdf".into())]),
},
disposition: None,
language: None,
location: None,
},
other: BodyContentSinglePart {
id: None,
md5: None,
description: None,
transfer_encoding: ContentEncoding::SevenBit,
octets: 100,
},
extension: None,
};
assert_eq!(count_attachments(&structure), 1);
}
#[test]
fn count_attachments_multipart_sums_attachment_dispositions_recursively() {
fn basic_part(disposition: Option<&'static str>) -> BodyStructure<'static> {
BodyStructure::Basic {
common: BodyContentCommon {
ty: ContentType {
ty: "application".into(),
subtype: "octet-stream".into(),
params: None,
},
disposition: disposition.map(|ty| ContentDisposition {
ty: ty.into(),
params: None,
}),
language: None,
location: None,
},
other: BodyContentSinglePart {
id: None,
md5: None,
description: None,
transfer_encoding: ContentEncoding::SevenBit,
octets: 100,
},
extension: None,
}
}
let nested = BodyStructure::Multipart {
common: BodyContentCommon {
ty: ContentType {
ty: "multipart".into(),
subtype: "mixed".into(),
params: None,
},
disposition: None,
language: None,
location: None,
},
bodies: vec![
basic_part(None),
basic_part(Some("Attachment")),
basic_part(Some("attachment")),
],
extension: None,
};
let top = BodyStructure::Multipart {
common: BodyContentCommon {
ty: ContentType {
ty: "multipart".into(),
subtype: "mixed".into(),
params: None,
},
disposition: None,
language: None,
location: None,
},
bodies: vec![basic_part(Some("inline")), nested],
extension: None,
};
assert_eq!(count_attachments(&top), 2);
}
#[test]
fn load_missing_checkpoint_is_empty() {
let dir = tempfile::tempdir().unwrap();
assert!(load_checkpoint(dir.path()).unwrap().is_empty());
}
#[test]
fn checkpoint_round_trips_across_appends() {
let dir = tempfile::tempdir().unwrap();
append_checkpoint(
dir.path(),
&CheckpointEntry {
mailbox: "INBOX".to_string(),
uid: 5,
message_hash: "abc123".to_string(),
md_staged_relpath: "inbox/5.md".to_string(),
desired_md_name: "2024-01-26-hello.md".to_string(),
mailbox_tag: "mailbox/inbox".to_string(),
attachments: Vec::new(),
},
)
.unwrap();
append_checkpoint(
dir.path(),
&CheckpointEntry {
mailbox: "INBOX".to_string(),
uid: 9,
message_hash: "def456".to_string(),
md_staged_relpath: "inbox/9.md".to_string(),
desired_md_name: "2024-01-27-world.md".to_string(),
mailbox_tag: "mailbox/inbox".to_string(),
attachments: vec![
("hash1".to_string(), "inbox/9/attachments/a.pdf".to_string()),
("hash2".to_string(), "inbox/9/attachments/b.png".to_string()),
],
},
)
.unwrap();
let loaded = load_checkpoint(dir.path()).unwrap();
assert_eq!(loaded.len(), 2);
assert_eq!(loaded[0].uid, 5);
assert!(loaded[0].attachments.is_empty());
assert_eq!(loaded[1].uid, 9);
assert_eq!(loaded[1].message_hash, "def456");
assert_eq!(loaded[1].md_staged_relpath, "inbox/9.md");
assert_eq!(loaded[1].mailbox_tag, "mailbox/inbox");
assert_eq!(
loaded[1].attachments,
vec![
("hash1".to_string(), "inbox/9/attachments/a.pdf".to_string()),
("hash2".to_string(), "inbox/9/attachments/b.png".to_string()),
]
);
}
#[test]
fn clear_checkpoint_for_mailbox_drops_only_matching_entries() {
let dir = tempfile::tempdir().unwrap();
append_checkpoint(
dir.path(),
&CheckpointEntry {
mailbox: "INBOX".to_string(),
uid: 1,
message_hash: "hash1".to_string(),
md_staged_relpath: "inbox/1.md".to_string(),
desired_md_name: "1.md".to_string(),
mailbox_tag: "mailbox/inbox".to_string(),
attachments: Vec::new(),
},
)
.unwrap();
append_checkpoint(
dir.path(),
&CheckpointEntry {
mailbox: "Archive".to_string(),
uid: 2,
message_hash: "hash2".to_string(),
md_staged_relpath: "archive/2.md".to_string(),
desired_md_name: "2.md".to_string(),
mailbox_tag: "mailbox/archive".to_string(),
attachments: Vec::new(),
},
)
.unwrap();
clear_checkpoint_for_mailbox(dir.path(), "INBOX").unwrap();
let loaded = load_checkpoint(dir.path()).unwrap();
assert_eq!(loaded.len(), 1);
assert_eq!(loaded[0].mailbox, "Archive");
}
#[test]
fn load_checkpoint_skips_malformed_lines() {
let dir = tempfile::tempdir().unwrap();
fs::write(
dir.path().join(CHECKPOINT_FILE_NAME),
"INBOX\t5\thash\tinbox/5.md\t5.md\tmailbox/inbox\t\nnot-enough-fields\nINBOX\t6\thash2\tinbox/6.md\t6.md\tmailbox/inbox\t\n",
)
.unwrap();
let loaded = load_checkpoint(dir.path()).unwrap();
assert_eq!(loaded.len(), 2);
assert_eq!(loaded[0].uid, 5);
assert_eq!(loaded[1].uid, 6);
}
#[test]
fn done_uids_builds_composite_key_set() {
let entries = vec![
CheckpointEntry {
mailbox: "INBOX".to_string(),
uid: 5,
message_hash: "abc".to_string(),
md_staged_relpath: "inbox/5.md".to_string(),
desired_md_name: "2024-01-26-hello.md".to_string(),
mailbox_tag: "mailbox/inbox".to_string(),
attachments: Vec::new(),
},
CheckpointEntry {
mailbox: "Archive".to_string(),
uid: 5,
message_hash: "def".to_string(),
md_staged_relpath: "archive/5.md".to_string(),
desired_md_name: "2024-01-26-hello.md".to_string(),
mailbox_tag: "mailbox/archive".to_string(),
attachments: Vec::new(),
},
];
let done = done_uids(&entries);
assert!(done.contains(&("INBOX".to_string(), 5)));
assert!(done.contains(&("Archive".to_string(), 5)));
assert!(!done.contains(&("INBOX".to_string(), 6)));
}
#[test]
fn split_into_batches_empty_input_is_empty() {
assert!(split_into_batches(&[], 4).is_empty());
}
#[test]
fn split_into_batches_covers_every_uid_with_no_duplicates() {
let uids: Vec<u32> = (1..=1000).collect();
let batches = split_into_batches(&uids, 4);
let mut seen: Vec<u32> = batches.into_iter().flatten().collect();
seen.sort_unstable();
assert_eq!(seen, uids);
}
#[test]
fn split_into_batches_count_scales_with_concurrency() {
let uids: Vec<u32> = (1..=1000).collect();
let low_concurrency_batches = split_into_batches(&uids, 1).len();
let high_concurrency_batches = split_into_batches(&uids, 8).len();
assert!(low_concurrency_batches > 1);
assert!(high_concurrency_batches > low_concurrency_batches);
}
#[test]
fn split_into_batches_single_uid_is_one_batch() {
assert_eq!(split_into_batches(&[42], 8), vec![vec![42]]);
}
#[test]
fn split_into_batches_more_batches_than_concurrency_for_a_large_mailbox() {
let uids: Vec<u32> = (1..=5000).collect();
let concurrency = 4;
let batch_count = split_into_batches(&uids, concurrency).len();
assert!(batch_count > concurrency);
}
}