diff --git a/docs/usage.md b/docs/usage.md index b6982c3..3cf96ff 100644 --- a/docs/usage.md +++ b/docs/usage.md @@ -196,6 +196,12 @@ By default it only adds and updates -- items that match are updated, the rest are created, and anything already on the target that the archive does not cover is left alone. +Email is matched by Message-ID, or without one by sender, subject, date and +recipients; where several messages share one, size decides. A message the +source kept in several folders -- IMAP and Maildir copies, Gmail labels -- is +written once, in all of them, and a later run adds any folder it is still +missing on the target. + `--prune` also deletes what is on the target and not in the archive. It asks first; `--yes` answers for it, for scripts. Export speaks JMAP only. diff --git a/src/sync/export/email.rs b/src/sync/export/email.rs index 1d8c7ab..9b23e59 100644 --- a/src/sync/export/email.rs +++ b/src/sync/export/email.rs @@ -1,5 +1,6 @@ /* * SPDX-FileCopyrightText: 2020 Stalwart Labs LLC + * SPDX-FileCopyrightText: 2026 John Coffey * * SPDX-License-Identifier: Apache-2.0 OR MIT */ @@ -13,7 +14,7 @@ use super::{Maps, Net, Plan, Uploader}; use crate::error::Error; use crate::jmap::error::JmapError; use crate::jmap::request::{ - MethodCall, Request, check_method_error, get_objects, retry_method_call, + MethodCall, Request, SetRequest, check_method_error, get_objects, retry_method_call, set_call, }; use crate::jmap::retry::MethodCallKind; use crate::jmap::wire::JmapId; @@ -61,7 +62,8 @@ pub fn reconcile( ) -> Result { let ty = ObjectType::Email; - let target_min = target_query_get(net, ty, Some(&["messageId"])).map_err(Error::from)?; + let target_min = target_query_get(net, ty, Some(&["messageId", "size", "mailboxIds"])) + .map_err(Error::from)?; let mut indices: Vec = target_min.iter().map(server_index).collect(); let fallback_ids: Vec = target_min @@ -92,9 +94,10 @@ pub fn reconcile( } } } - let target_keys: HashSet = email_keys(&indices).into_iter().collect(); + let targets: Vec = target_min.iter().map(TargetEmail::from_value).collect(); + let target_keys = email_keys(&indices); - let local: Vec<(i64, EmailRow)> = { + let mut local: Vec<(i64, EmailRow)> = { let mut stmt = ctx .conn .prepare(EMAIL_SELECT) @@ -109,26 +112,209 @@ pub fn reconcile( .map(|(id, r)| Ok((id, r.map_err(Error::from)?))) .collect::>()? }; + local.sort_by_key(|(id, _)| *id); + let units = fold_by_blob(local); - let local_indices: Vec = local + let local_indices: Vec = units .iter() - .map(|(_, r)| index_from_json(&r.message_match)) + .map(|u| index_from_json(&u.row.message_match)) .collect(); let local_keys = email_keys(&local_indices); let mut uploader = Uploader::new(net, &ctx.conn); - for (i, key) in local_keys.iter().enumerate() { - if target_keys.contains(key) { - counts.skipped += 1; - continue; + let sizes: Vec> = units + .iter() + .map(|u| uploader.blob_len(u.row.blob_local_id)) + .collect(); + let pairs = pair_with_targets(&local_keys, &sizes, &target_keys, &targets); + + let mut membership_updates: Vec<(String, Value)> = Vec::new(); + for (i, unit) in units.iter().enumerate() { + match pairs[i] { + Some(t) => match missing_memberships(&unit.row, &targets[t], maps) { + Some(patch) => membership_updates.push((targets[t].id.clone(), patch)), + None => counts.skipped += 1, + }, + None => export_one( + net, + &mut uploader, + maps, + unit.local_id, + &unit.row, + counts, + logger, + ), } - let (local_id, row) = &local[i]; - export_one(net, &mut uploader, maps, *local_id, row, counts, logger); } + send_membership_updates(net, membership_updates, counts, logger); Ok(Plan::default()) } +/// One message to write: the archive rows that hold the same bytes, folded +/// together. A source that files one message in several folders (IMAP and +/// Maildir copies, Gmail labels) leaves one archive row per folder; on the +/// target it is one email in all of them. Keywords are the union of the +/// copies', so a message read or flagged in any folder stays so. +struct Unit { + local_id: i64, + row: EmailRow, +} + +fn fold_by_blob(local: Vec<(i64, EmailRow)>) -> Vec { + let mut order: Vec = Vec::new(); + let mut by_blob: HashMap = HashMap::new(); + for (local_id, row) in local { + match by_blob.get_mut(&row.blob_local_id) { + Some(unit) => { + for m in row.mailbox_locals { + if !unit.row.mailbox_locals.contains(&m) { + unit.row.mailbox_locals.push(m); + } + } + for k in row.keywords { + if !unit.row.keywords.contains(&k) { + unit.row.keywords.push(k); + } + } + } + None => { + order.push(row.blob_local_id); + by_blob.insert(row.blob_local_id, Unit { local_id, row }); + } + } + } + order + .into_iter() + .filter_map(|b| by_blob.remove(&b)) + .collect() +} + +/// What export needs to know about an email already on the target. `size` +/// and `mailboxes` are `None` when the server did not return them, and then +/// nothing is inferred from their absence. +struct TargetEmail { + id: String, + size: Option, + mailboxes: Option>, +} + +impl TargetEmail { + fn from_value(v: &Value) -> TargetEmail { + TargetEmail { + id: jid(v).unwrap_or_default(), + size: v.get("size").and_then(Value::as_u64), + mailboxes: v + .get("mailboxIds") + .and_then(Value::as_object) + .map(|m| m.keys().cloned().collect()), + } + } +} + +/// Pairs each local message with at most one target email. Messages are +/// matched by key (Message-ID, or the fallback digest); where several share a +/// key -- genuinely different messages with the same Message-ID, or copies +/// already on the target -- equal size decides first, then order, so two +/// different messages are never folded onto one target email. +fn pair_with_targets( + local_keys: &[EmailKey], + local_sizes: &[Option], + target_keys: &[EmailKey], + targets: &[TargetEmail], +) -> Vec> { + let mut by_key: HashMap<&EmailKey, Vec> = HashMap::new(); + for (t, key) in target_keys.iter().enumerate() { + by_key.entry(key).or_default().push(t); + } + let mut groups: HashMap<&EmailKey, Vec> = HashMap::new(); + for (i, key) in local_keys.iter().enumerate() { + groups.entry(key).or_default().push(i); + } + let mut out = vec![None; local_keys.len()]; + for (key, members) in groups { + let Some(candidates) = by_key.get_mut(key) else { + continue; + }; + for &i in &members { + let Some(size) = local_sizes[i] else { continue }; + if let Some(pos) = candidates + .iter() + .position(|&t| targets[t].size == Some(size)) + { + out[i] = Some(candidates.remove(pos)); + } + } + for &i in &members { + if out[i].is_none() && !candidates.is_empty() { + out[i] = Some(candidates.remove(0)); + } + } + } + out +} + +/// The `Email/set` patch adding the folders a matched email is missing on the +/// target, or `None` when it is already in all of them. Folders that exist +/// only on the target are left alone. +fn missing_memberships(row: &EmailRow, target: &TargetEmail, maps: &Maps) -> Option { + let have = target.mailboxes.as_ref()?; + let mut patch = Map::new(); + for ml in &row.mailbox_locals { + if let Some(t) = maps.target(ObjectType::Mailbox, *ml) + && !have.contains(&t.0) + { + patch.insert(format!("mailboxIds/{}", t.0), Value::Bool(true)); + } + } + (!patch.is_empty()).then_some(Value::Object(patch)) +} + +fn send_membership_updates( + net: &Net, + updates: Vec<(String, Value)>, + counts: &mut TypeCounts, + logger: &Logger, +) { + if updates.is_empty() { + return; + } + if net.dry_run { + counts.updated += updates.len() as u64; + return; + } + let total = updates.len() as u64; + let mut map = Map::new(); + for (id, patch) in updates { + map.insert(id, patch); + } + match set_call( + &net.client, + &net.api, + &net.account, + ObjectType::Email.jmap_name(), + SetRequest { + update: Some(Value::Object(map)), + ..Default::default() + }, + &net.limits, + ) { + Ok(outcome) => { + counts.updated += outcome.updated.len() as u64; + for (id, err) in &outcome.not_updated { + logger.warn(&format!("Email/set {id}: folders not added: {err}")); + counts.failed += 1; + } + } + Err(e) => { + logger.warn(&format!( + "Email/set: adding folders to {total} email(s) failed: {e}" + )); + counts.failed += total; + } + } +} + fn build_mailbox_ids(row: &EmailRow, maps: &Maps) -> Option> { let mut mids = Map::new(); for ml in &row.mailbox_locals { @@ -364,3 +550,97 @@ fn interpret_import(mr: &MethodCall, cid: &str) -> Result EmailRow { + EmailRow { + blob_local_id: blob, + received_at: "2020-01-01T00:00:00Z".to_owned(), + mailbox_locals: mailboxes.to_vec(), + keywords: keywords.iter().map(|k| (*k).to_owned()).collect(), + message_match: "{}".to_owned(), + } + } + + fn target(id: &str, size: Option, mailboxes: Option<&[&str]>) -> TargetEmail { + TargetEmail { + id: id.to_owned(), + size, + mailboxes: mailboxes.map(|m| m.iter().map(|s| (*s).to_owned()).collect()), + } + } + + fn mid(m: &str) -> EmailKey { + EmailKey::MessageId(m.to_owned()) + } + + #[test] + fn copies_of_one_message_fold_into_one_unit_in_every_folder() { + let units = fold_by_blob(vec![ + (1, row(10, &[1], &["$seen"])), + (2, row(10, &[2], &["$flagged", "$seen"])), + (3, row(11, &[1], &[])), + ]); + assert_eq!(units.len(), 2); + assert_eq!(units[0].local_id, 1); + assert_eq!(units[0].row.mailbox_locals, vec![1, 2]); + assert_eq!(units[0].row.keywords, vec!["$seen", "$flagged"]); + assert_eq!(units[1].row.blob_local_id, 11); + } + + #[test] + fn different_messages_sharing_a_message_id_pair_by_size() { + let pairs = pair_with_targets( + &[mid("m@h"), mid("m@h")], + &[Some(100), Some(200)], + &[mid("m@h")], + &[target("T", Some(200), None)], + ); + assert_eq!( + pairs, + vec![None, Some(0)], + "only the 200-byte one is on the target" + ); + } + + #[test] + fn size_pairing_ignores_target_order() { + let pairs = pair_with_targets( + &[mid("m@h"), mid("m@h")], + &[Some(100), Some(200)], + &[mid("m@h"), mid("m@h")], + &[target("A", Some(200), None), target("B", Some(100), None)], + ); + assert_eq!(pairs, vec![Some(1), Some(0)]); + } + + #[test] + fn a_lone_message_pairs_by_key_when_the_target_gives_no_size() { + let pairs = pair_with_targets( + &[mid("m@h"), mid("n@h")], + &[Some(100), Some(50)], + &[mid("m@h")], + &[target("T", None, None)], + ); + assert_eq!(pairs, vec![Some(0), None]); + } + + #[test] + fn missing_memberships_adds_only_the_absent_migrated_folders() { + let mut maps = Maps::default(); + maps.insert(ObjectType::Mailbox, 1, JmapId("T1".into())); + maps.insert(ObjectType::Mailbox, 2, JmapId("T2".into())); + let r = row(10, &[1, 2, 3], &[]); + let patch = missing_memberships(&r, &target("E", None, Some(&["T1", "Own"])), &maps) + .expect("T2 is missing"); + assert_eq!(patch, json!({ "mailboxIds/T2": true })); + assert!(missing_memberships(&r, &target("E", None, Some(&["T1", "T2"])), &maps).is_none()); + assert!( + missing_memberships(&r, &target("E", None, None), &maps).is_none(), + "unknown membership is left alone" + ); + } +} diff --git a/tests/mock_sync.rs b/tests/mock_sync.rs index 2bd9546..d96063f 100644 --- a/tests/mock_sync.rs +++ b/tests/mock_sync.rs @@ -645,16 +645,14 @@ fn export_email_blob_not_found_reuploads_and_retries() { [], ) .unwrap(); - let raw = b"From: a@x\r\nSubject: dup\r\nMessage-ID: \r\n\r\nbody"; + let raw = b"From: a@x\r\nSubject: stale\r\nMessage-ID: \r\n\r\nbody"; let blob = db::blobs::intern_blob(&conn, raw).unwrap(); - for _ in 0..2 { - conn.execute( - "INSERT INTO emails (blob_id,received_at,mailbox_ids,keywords) - VALUES (?1,'2020-01-01T00:00:00Z','[1]','[]')", - rusqlite::params![blob], - ) - .unwrap(); - } + conn.execute( + "INSERT INTO emails (blob_id,received_at,mailbox_ids,keywords) + VALUES (?1,'2020-01-01T00:00:00Z','[1]','[]')", + rusqlite::params![blob], + ) + .unwrap(); } let _root = server.mock("GET", "/").with_status(404).create(); @@ -716,43 +714,30 @@ fn export_email_blob_not_found_reuploads_and_retries() { .expect(1) .create(); - let imp_e1 = server + let imp_stale = server .mock("POST", api) .match_body(Matcher::AllOf(vec![ Matcher::Regex("Email/import".into()), Matcher::Regex("e1".into()), - ])) - .with_body( - json!({"methodResponses":[["Email/import",{"accountId":"w", - "created":{"e1":{"id":"x1","blobId":"UP1","threadId":"t","size":10}}},"i"]]}) - .to_string(), - ) - .expect(1) - .create(); - let imp_e2_stale = server - .mock("POST", api) - .match_body(Matcher::AllOf(vec![ - Matcher::Regex("Email/import".into()), - Matcher::Regex("e2".into()), Matcher::Regex("UP1".into()), ])) .with_body( json!({"methodResponses":[["Email/import",{"accountId":"w", - "notCreated":{"e2":{"type":"blobNotFound"}}},"i"]]}) + "notCreated":{"e1":{"type":"blobNotFound"}}},"i"]]}) .to_string(), ) .expect(1) .create(); - let imp_e2_fresh = server + let imp_fresh = server .mock("POST", api) .match_body(Matcher::AllOf(vec![ Matcher::Regex("Email/import".into()), - Matcher::Regex("e2".into()), + Matcher::Regex("e1".into()), Matcher::Regex("UP2".into()), ])) .with_body( json!({"methodResponses":[["Email/import",{"accountId":"w", - "created":{"e2":{"id":"x2","blobId":"UP2","threadId":"t","size":10}}},"i"]]}) + "created":{"e1":{"id":"x1","blobId":"UP2","threadId":"t","size":10}}},"i"]]}) .to_string(), ) .expect(1) @@ -789,16 +774,18 @@ fn export_email_blob_not_found_reuploads_and_retries() { .find(|(t, _)| *t == "Email") .map(|(_, c)| c.clone()) .expect("email counts"); - assert_eq!(email.created, 2, "both emails end up created"); + assert_eq!( + email.created, 1, + "the stale blob is re-uploaded and the email created" + ); assert_eq!(email.failed, 0, "blobNotFound self-heals, not a failure"); assert_eq!(email.skipped, 0); assert!(!summary.any_failed()); up1.assert(); up2.assert(); - imp_e1.assert(); - imp_e2_stale.assert(); - imp_e2_fresh.assert(); + imp_stale.assert(); + imp_fresh.assert(); let _ = std::fs::remove_file(&archive); } @@ -4459,3 +4446,258 @@ fn export_email_fatal_method_error_is_not_retried() { let _ = std::fs::remove_file(&archive); } + +/// Archive with an Inbox (1) and an Archive folder (2), and a target whose +/// Mailbox/get returns both as T1 and T2, so no folder is created. +fn two_folder_archive_and_target( + server: &mut mockito::Server, + archive: &Path, +) -> Vec { + let conn = db::init::open(archive).unwrap(); + conn.execute( + "INSERT INTO mailboxes (id,name,parent_id,role) VALUES + (1,'Inbox',NULL,'inbox'), (2,'Archive',NULL,'archive')", + [], + ) + .unwrap(); + let base = server.url(); + let api = "/jmap/api"; + vec![ + server.mock("GET", "/").with_status(404).create(), + server + .mock("GET", "/.well-known/jmap") + .with_body(session_body_full(&base)) + .expect_at_least(1) + .create(), + anchor_terminator(server, api, "Mailbox"), + anchor_terminator(server, api, "Email"), + server + .mock("POST", api) + .match_body(Matcher::Regex("Mailbox/query".into())) + .with_body( + json!({"methodResponses":[["Mailbox/query", + {"accountId":"w","ids":["T1","T2"]},"q"]]}) + .to_string(), + ) + .expect_at_least(1) + .create(), + server + .mock("POST", api) + .match_body(Matcher::Regex("Mailbox/get".into())) + .with_body( + json!({"methodResponses":[["Mailbox/get",{"accountId":"w","list":[ + {"id":"T1","name":"Inbox","role":"inbox","parentId":null,"myRights":{"mayDelete":true}}, + {"id":"T2","name":"Archive","role":"archive","parentId":null,"myRights":{"mayDelete":true}} + ],"notFound":[]},"g"]]}) + .to_string(), + ) + .expect_at_least(1) + .create(), + ] +} + +fn insert_email_copy(archive: &Path, raw: &[u8], mailbox: i64) { + let conn = db::init::open(archive).unwrap(); + let blob = db::blobs::intern_blob(&conn, raw).unwrap(); + let mm = inbuxa_migrate::sync::keys::index_to_json( + &inbuxa_migrate::sync::emailmeta::email_index_from_blob(raw), + ); + conn.execute( + "INSERT INTO emails (blob_id,received_at,mailbox_ids,keywords,message_match) + VALUES (?1,'2020-01-01T00:00:00Z',?2,'[]',?3)", + rusqlite::params![blob, format!("[{mailbox}]"), mm], + ) + .unwrap(); +} + +const ONE_MESSAGE: &[u8] = b"From: a@x\r\nSubject: both\r\nMessage-ID: \r\n\r\nbody"; + +#[test] +fn export_message_in_two_folders_lands_once_in_both() { + let mut server = mockito::Server::new(); + let base = server.url(); + let api = "/jmap/api"; + let archive = tmp(); + let _setup = two_folder_archive_and_target(&mut server, &archive); + insert_email_copy(&archive, ONE_MESSAGE, 1); + insert_email_copy(&archive, ONE_MESSAGE, 2); + + let _eq = server + .mock("POST", api) + .match_body(Matcher::Regex("Email/query".into())) + .with_body( + json!({"methodResponses":[["Email/query",{"accountId":"w","ids":[]},"q"]]}).to_string(), + ) + .expect(1) + .create(); + let upload = server + .mock("POST", Matcher::Regex("/jmap/upload/".into())) + .with_body(json!({"blobId":"BUP"}).to_string()) + .expect(1) + .create(); + let import = server + .mock("POST", api) + .match_body(Matcher::AllOf(vec![ + Matcher::Regex("Email/import".into()), + Matcher::Regex("\"T1\":true".into()), + Matcher::Regex("\"T2\":true".into()), + ])) + .with_body( + json!({"methodResponses":[["Email/import",{"accountId":"w", + "created":{"e1":{"id":"Y1","blobId":"BUP","threadId":"t","size":10}}},"i"]]}) + .to_string(), + ) + .expect(1) + .create(); + + let summary = sync::export::run( + common(&archive), + export_cfg_objects(&base, vec![ObjectType::Mailbox, ObjectType::Email]), + ) + .expect("export"); + let email = email_counts(&summary); + assert_eq!(email.created, 1, "one email, not one per folder"); + assert_eq!(email.failed, 0); + upload.assert(); + import.assert(); + let _ = std::fs::remove_file(&archive); +} + +#[test] +fn export_resumed_adds_the_folders_a_matched_message_is_missing() { + let mut server = mockito::Server::new(); + let base = server.url(); + let api = "/jmap/api"; + let archive = tmp(); + let _setup = two_folder_archive_and_target(&mut server, &archive); + insert_email_copy(&archive, ONE_MESSAGE, 1); + insert_email_copy(&archive, ONE_MESSAGE, 2); + + let _eq = server + .mock("POST", api) + .match_body(Matcher::Regex("Email/query".into())) + .with_body( + json!({"methodResponses":[["Email/query",{"accountId":"w","ids":["X1"]},"q"]]}) + .to_string(), + ) + .expect(1) + .create(); + let _eg = server + .mock("POST", api) + .match_body(Matcher::Regex("Email/get".into())) + .with_body( + json!({"methodResponses":[["Email/get",{"accountId":"w","list":[ + {"id":"X1","messageId":["both@h"],"size":ONE_MESSAGE.len(),"mailboxIds":{"T1":true}} + ],"notFound":[]},"g"]]}) + .to_string(), + ) + .expect(1) + .create(); + let no_upload = server + .mock("POST", Matcher::Regex("/jmap/upload/".into())) + .expect(0) + .create(); + let no_import = server + .mock("POST", api) + .match_body(Matcher::Regex("Email/import".into())) + .expect(0) + .create(); + let set = server + .mock("POST", api) + .match_body(Matcher::AllOf(vec![ + Matcher::Regex("Email/set".into()), + Matcher::Regex("\"mailboxIds/T2\":true".into()), + ])) + .with_body( + json!({"methodResponses":[["Email/set",{"accountId":"w","updated":{"X1":null}},"s"]]}) + .to_string(), + ) + .expect(1) + .create(); + + let summary = sync::export::run( + common(&archive), + export_cfg_objects(&base, vec![ObjectType::Mailbox, ObjectType::Email]), + ) + .expect("export"); + let email = email_counts(&summary); + assert_eq!(email.updated, 1, "the Archive membership is added"); + assert_eq!(email.created, 0); + assert_eq!(email.failed, 0); + set.assert(); + no_upload.assert(); + no_import.assert(); + let _ = std::fs::remove_file(&archive); +} + +#[test] +fn export_different_messages_sharing_a_message_id_are_not_merged() { + let mut server = mockito::Server::new(); + let base = server.url(); + let api = "/jmap/api"; + let archive = tmp(); + let _setup = two_folder_archive_and_target(&mut server, &archive); + let first: &[u8] = b"From: a@x\r\nSubject: one\r\nMessage-ID: \r\n\r\nfirst body"; + let second: &[u8] = + b"From: a@x\r\nSubject: one\r\nMessage-ID: \r\n\r\na different, longer second body"; + insert_email_copy(&archive, first, 1); + insert_email_copy(&archive, second, 1); + + let _eq = server + .mock("POST", api) + .match_body(Matcher::Regex("Email/query".into())) + .with_body( + json!({"methodResponses":[["Email/query",{"accountId":"w","ids":["X1"]},"q"]]}) + .to_string(), + ) + .expect(1) + .create(); + let _eg = server + .mock("POST", api) + .match_body(Matcher::Regex("Email/get".into())) + .with_body( + json!({"methodResponses":[["Email/get",{"accountId":"w","list":[ + {"id":"X1","messageId":["same@h"],"size":second.len(),"mailboxIds":{"T1":true}} + ],"notFound":[]},"g"]]}) + .to_string(), + ) + .expect(1) + .create(); + let upload = server + .mock("POST", Matcher::Regex("/jmap/upload/".into())) + .with_body(json!({"blobId":"BUP"}).to_string()) + .expect(1) + .create(); + let import = server + .mock("POST", api) + .match_body(Matcher::Regex("Email/import".into())) + .with_body( + json!({"methodResponses":[["Email/import",{"accountId":"w", + "created":{"e1":{"id":"Y1","blobId":"BUP","threadId":"t","size":10}}},"i"]]}) + .to_string(), + ) + .expect(1) + .create(); + let no_set = server + .mock("POST", api) + .match_body(Matcher::Regex("Email/set".into())) + .expect(0) + .create(); + + let summary = sync::export::run( + common(&archive), + export_cfg_objects(&base, vec![ObjectType::Mailbox, ObjectType::Email]), + ) + .expect("export"); + let email = email_counts(&summary); + assert_eq!( + email.created, 1, + "the first message is not on the target yet" + ); + assert_eq!(email.skipped, 1, "the second is, matched by size"); + assert_eq!(email.failed, 0); + upload.assert(); + import.assert(); + no_set.assert(); + let _ = std::fs::remove_file(&archive); +}