export: write a message once, in every folder it was in #4

Merged
jcoffey-dev merged 1 commits from fix/one-message-many-folders into main 2026-09-30 18:33:10 +00:00
3 changed files with 571 additions and 43 deletions
Showing only changes of commit 5ae0625ee1 - Show all commits
+6
View File
@@ -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.
+292 -12
View File
@@ -1,5 +1,6 @@
/*
* SPDX-FileCopyrightText: 2020 Stalwart Labs LLC <[email protected]>
* SPDX-FileCopyrightText: 2026 John Coffey <[email protected]>
*
* 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<Plan, Error> {
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<EmailIndex> = target_min.iter().map(server_index).collect();
let fallback_ids: Vec<JmapId> = target_min
@@ -92,9 +94,10 @@ pub fn reconcile(
}
}
}
let target_keys: HashSet<EmailKey> = email_keys(&indices).into_iter().collect();
let targets: Vec<TargetEmail> = 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::<Result<_, Error>>()?
};
local.sort_by_key(|(id, _)| *id);
let units = fold_by_blob(local);
let local_indices: Vec<EmailIndex> = local
let local_indices: Vec<EmailIndex> = 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<Option<u64>> = 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<Unit> {
let mut order: Vec<i64> = Vec::new();
let mut by_blob: HashMap<i64, Unit> = 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<u64>,
mailboxes: Option<HashSet<String>>,
}
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<u64>],
target_keys: &[EmailKey],
targets: &[TargetEmail],
) -> Vec<Option<usize>> {
let mut by_key: HashMap<&EmailKey, Vec<usize>> = HashMap::new();
for (t, key) in target_keys.iter().enumerate() {
by_key.entry(key).or_default().push(t);
}
let mut groups: HashMap<&EmailKey, Vec<usize>> = 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<Value> {
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<Map<String, Value>> {
let mut mids = Map::new();
for ml in &row.mailbox_locals {
@@ -364,3 +550,97 @@ fn interpret_import(mr: &MethodCall, cid: &str) -> Result<SingleImport, JmapErro
detail: format!("Email/import returned neither created nor notCreated for {cid}"),
})
}
#[cfg(test)]
mod tests {
use super::*;
fn row(blob: i64, mailboxes: &[i64], keywords: &[&str]) -> 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<u64>, 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"
);
}
}
+267 -25
View File
@@ -645,9 +645,8 @@ fn export_email_blob_not_found_reuploads_and_retries() {
[],
)
.unwrap();
let raw = b"From: a@x\r\nSubject: dup\r\nMessage-ID: <dup-1@h>\r\n\r\nbody";
let raw = b"From: a@x\r\nSubject: stale\r\nMessage-ID: <stale-1@h>\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]','[]')",
@@ -655,7 +654,6 @@ fn export_email_blob_not_found_reuploads_and_retries() {
)
.unwrap();
}
}
let _root = server.mock("GET", "/").with_status(404).create();
let _wk = server
@@ -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<mockito::Mock> {
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: <both@h>\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: <same@h>\r\n\r\nfirst body";
let second: &[u8] =
b"From: a@x\r\nSubject: one\r\nMessage-ID: <same@h>\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);
}