From 3f832e17b21079e1c69662fcdf8c223e2a1c7130 Mon Sep 17 00:00:00 2001 From: John Coffey Date: Wed, 30 Sep 2026 12:39:57 -0700 Subject: [PATCH] export: batch Email/import, upload in parallel, and show progress Export wrote one message per request and said nothing while it did, so a large mailbox took hours of silence. Messages now go in batches of up to the target's maxObjectsInSet (at most 50), their blobs uploaded several at once, up to maxConcurrentUpload and no more than --threads; each upload thread reads the archive through its own read-only connection. Every message is still counted on its own: one the target rejects fails alone, a request too large is split, and a method error on the whole call is retried a message at a time. A batch is sent once. If it ends without a clear answer -- a dropped connection, a gateway timeout, a partial failure -- the target is read again, the messages that arrived count as created, and only the rest are imported again, so none is doubled. A progress line (count, rate, time left) is printed every few seconds for mail, contacts and events, and each type ends with a line of what was created, updated, left unchanged and failed. --- docs/usage.md | 14 ++ src/sync/export.rs | 192 ++++++++++++++++++ src/sync/export/email.rs | 392 +++++++++++++++++++++++++++++++------ src/sync/export/uidtype.rs | 7 + src/sync/mod.rs | 1 + src/sync/progress.rs | 196 +++++++++++++++++++ tests/mock_export_batch.rs | 342 ++++++++++++++++++++++++++++++++ tests/mock_sync.rs | 144 -------------- 8 files changed, 1084 insertions(+), 204 deletions(-) create mode 100644 src/sync/progress.rs create mode 100644 tests/mock_export_batch.rs diff --git a/docs/usage.md b/docs/usage.md index 4802889..5ce94b1 100644 --- a/docs/usage.md +++ b/docs/usage.md @@ -229,6 +229,20 @@ what changed at the source in between. - **Sieve scripts** are matched by name, and the target's is replaced when its content differs from the archive's. +Messages go in batches, as many to an `Email/import` as the target's +`maxObjectsInSet` allows, up to 50, with their blobs uploaded several at a +time -- the target's `maxConcurrentUpload`, and no more than `--threads`. +Each message is still counted on its own: one the target rejects fails +alone, and the rest of its batch lands. A batch that ends without a clear +answer -- a dropped connection, a gateway timeout -- is never sent again as +it was. The target is read first, the messages that arrived are counted as +created, and only the rest are imported again, so none is ever doubled. + +While it runs, export prints a line every few seconds for mail, contacts and +events: how many of how many, how fast, and about how long is left. Each +type ends with a line of what was created, updated, left unchanged and +failed. + `--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.rs b/src/sync/export.rs index 143a2ec..4e29017 100644 --- a/src/sync/export.rs +++ b/src/sync/export.rs @@ -7,6 +7,9 @@ use std::collections::{HashMap, HashSet}; use std::io::{IsTerminal, Write}; +use std::path::PathBuf; +use std::sync::Mutex; +use std::sync::atomic::{AtomicUsize, Ordering}; use rusqlite::Connection; use serde_json::{Map, Value, json}; @@ -134,6 +137,77 @@ impl<'a> Uploader<'a> { Ok(id) } + /// Uploads several stored blobs at once, on up to `Net::upload_workers` + /// threads, and returns each one's result in the order given. Each thread + /// reads its blobs through its own read-only connection to the archive, so + /// at most one blob per thread is held in memory. With one worker, or in a + /// dry run, it is `upload_with` in a loop. + fn upload_many( + &mut self, + local_ids: &[i64], + content_type: &str, + ) -> Vec> { + if self.net.dry_run || self.net.upload_workers <= 1 { + return local_ids + .iter() + .map(|id| self.upload_with(*id, content_type)) + .collect(); + } + self.touched.extend_from_slice(local_ids); + let mut todo: Vec = Vec::new(); + for id in local_ids { + if !self.cache.contains_key(id) && !todo.contains(id) { + todo.push(*id); + } + } + let net = self.net; + let results = run_bounded( + &todo, + net.upload_workers, + || { + rusqlite::Connection::open_with_flags( + &net.archive, + rusqlite::OpenFlags::SQLITE_OPEN_READ_ONLY, + ) + }, + |conn, local_id| { + let conn = conn + .as_ref() + .map_err(|e| JmapError::malformed(format!("archive not readable: {e}")))?; + let bytes = db::blobs::blob_bytes(conn, *local_id)?.ok_or_else(|| { + JmapError::malformed(format!("blob local id {local_id} missing")) + })?; + blobxfer::upload_bytes( + &net.client, + &net.session, + &net.account, + content_type, + &bytes, + ) + }, + ); + let mut failed: HashMap = HashMap::new(); + for (local_id, result) in todo.into_iter().zip(results) { + match result { + Ok(id) => { + self.cache.insert(local_id, id); + } + Err(e) => { + failed.insert(local_id, e); + } + } + } + local_ids + .iter() + .map(|id| match self.cache.get(id) { + Some(blob) => Ok(blob.clone()), + None => Err(failed.get(id).map(clone_error).unwrap_or_else(|| { + JmapError::malformed(format!("blob local id {id} not uploaded")) + })), + }) + .collect() + } + fn invalidate(&mut self, local_id: i64) { self.cache.remove(&local_id); } @@ -147,6 +221,58 @@ impl<'a> Uploader<'a> { } } +/// A copy of an upload error for each archive row that shares the blob. The +/// errors that carry meaning for the caller -- the size limits -- keep their +/// kind; the rest keep their message. +fn clone_error(e: &JmapError) -> JmapError { + match e { + JmapError::RequestTooLarge => JmapError::RequestTooLarge, + JmapError::SingleObjectTooLarge(m) => JmapError::SingleObjectTooLarge(m.clone()), + other => JmapError::Transport(other.to_string()), + } +} + +/// Runs `f` over `jobs` on at most `workers` threads and returns the results +/// in job order. Each thread builds its own state once with `init`, such as a +/// connection of its own to the archive. +fn run_bounded( + jobs: &[J], + workers: usize, + init: impl Fn() -> S + Sync, + f: impl Fn(&mut S, &J) -> R + Sync, +) -> Vec +where + J: Sync, + R: Send, +{ + let workers = workers.clamp(1, jobs.len().max(1)); + if workers == 1 { + let mut state = init(); + return jobs.iter().map(|j| f(&mut state, j)).collect(); + } + let next = AtomicUsize::new(0); + let slots: Mutex>> = Mutex::new((0..jobs.len()).map(|_| None).collect()); + std::thread::scope(|scope| { + for _ in 0..workers { + scope.spawn(|| { + let mut state = init(); + loop { + let i = next.fetch_add(1, Ordering::SeqCst); + let Some(job) = jobs.get(i) else { break }; + let r = f(&mut state, job); + slots.lock().expect("result slots")[i] = Some(r); + } + }); + } + }); + slots + .into_inner() + .expect("result slots") + .into_iter() + .map(|r| r.expect("every job ran")) + .collect() +} + impl BlobBytes for Uploader<'_> { fn bytes(&self, local_id: i64) -> Result, JmapError> { db::blobs::blob_bytes(self.conn, local_id)? @@ -162,6 +288,11 @@ struct Net { limits: Limits, session: Session, dry_run: bool, + /// The archive's path, for the upload threads' own connections. + archive: PathBuf, + /// Blobs uploaded at once: the server's `maxConcurrentUpload`, and no + /// more than `--threads`. + upload_workers: usize, } fn has_rows(conn: &Connection, ty: ObjectType) -> bool { @@ -185,6 +316,10 @@ pub fn run(common: CommonConfig, config: ExportConfig) -> Result limits: connected.limits, session: connected.session.clone(), dry_run: ctx.dry_run(), + archive: ctx.common.archive.clone(), + upload_workers: (connected.limits.max_concurrent_upload as usize) + .min(ctx.common.threads) + .max(1), }; let work = work_list(&ctx.conn, &config, &connected, &logger); @@ -198,6 +333,7 @@ pub fn run(common: CommonConfig, config: ExportConfig) -> Result if logger.enabled(LEVEL_DEFAULT) { eprintln!("export: {} ...", ty.jmap_name()); } + let started = std::time::Instant::now(); let mut counts = TypeCounts::default(); let res = reconcile_type( &ctx, @@ -217,6 +353,16 @@ pub fn run(common: CommonConfig, config: ExportConfig) -> Result Plan::default() } }; + if logger.enabled(LEVEL_DEFAULT) && !ctx.dry_run() { + eprintln!( + "{}", + crate::sync::progress::done_line( + &format!("export: {}", ty.jmap_name()), + &counts, + started.elapsed() + ) + ); + } plans.insert(*ty, plan); counts_per_type.insert(*ty, counts); } @@ -635,3 +781,49 @@ mod common { outcome } } + +#[cfg(test)] +mod pool_tests { + use super::run_bounded; + use std::sync::atomic::{AtomicUsize, Ordering}; + use std::time::Duration; + + #[test] + fn never_more_workers_at_once_than_the_cap_and_results_keep_job_order() { + let jobs: Vec = (0..24).collect(); + let running = AtomicUsize::new(0); + let most = AtomicUsize::new(0); + let inits = AtomicUsize::new(0); + let out = run_bounded( + &jobs, + 3, + || inits.fetch_add(1, Ordering::SeqCst), + |_, j| { + let now = running.fetch_add(1, Ordering::SeqCst) + 1; + most.fetch_max(now, Ordering::SeqCst); + std::thread::sleep(Duration::from_millis(5)); + running.fetch_sub(1, Ordering::SeqCst); + j * 2 + }, + ); + assert_eq!(out, jobs.iter().map(|j| j * 2).collect::>()); + let most = most.load(Ordering::SeqCst); + assert!(most <= 3, "{most} ran at once"); + assert!(most > 1, "work ran in parallel"); + assert_eq!( + inits.load(Ordering::SeqCst), + 3, + "state built once per worker" + ); + } + + #[test] + fn one_worker_or_one_job_runs_in_place() { + let out = run_bounded(&[1, 2, 3], 1, || (), |_, j| j + 1); + assert_eq!(out, vec![2, 3, 4]); + let out = run_bounded(&[7], 8, || (), |_, j| j + 1); + assert_eq!(out, vec![8]); + let out: Vec = run_bounded(&[], 4, || (), |_, j: &i32| *j); + assert!(out.is_empty()); + } +} diff --git a/src/sync/export/email.rs b/src/sync/export/email.rs index b58cb4a..8200fa8 100644 --- a/src/sync/export/email.rs +++ b/src/sync/export/email.rs @@ -21,6 +21,7 @@ use crate::jmap::wire::JmapId; use crate::logging::Logger; use crate::sync::import_jmap::mapping::{EMAIL_SELECT, EmailRow, TargetResolver, row_to_email}; use crate::sync::keys::{EmailIndex, EmailKey, email_index, email_keys, index_from_json}; +use crate::sync::progress::Progress; use crate::sync::{Context, TypeCounts}; use crate::types::ObjectType; @@ -61,45 +62,7 @@ pub fn reconcile( logger: &Logger, ) -> Result { let ty = ObjectType::Email; - - let target_min = target_query_get( - net, - ty, - Some(&["messageId", "size", "mailboxIds", "keywords"]), - ) - .map_err(Error::from)?; - let mut indices: Vec = target_min.iter().map(server_index).collect(); - - let fallback_ids: Vec = target_min - .iter() - .zip(indices.iter()) - .filter(|(_, i)| i.mids.is_empty()) - .filter_map(|(v, _)| jid(v).map(JmapId)) - .collect(); - if !fallback_ids.is_empty() { - let got = get_objects::( - &net.client, - &net.api, - &net.account, - ty.jmap_name(), - &fallback_ids, - Some(&["messageId", "from", "subject", "sentAt", "to"]), - &net.limits, - ) - .map_err(Error::from)?; - let by_id: HashMap = got - .list - .iter() - .filter_map(|v| jid(v).map(|i| (i, v))) - .collect(); - for (v, slot) in target_min.iter().zip(indices.iter_mut()) { - if let Some(full) = jid(v).and_then(|i| by_id.get(&i)) { - *slot = server_index(full); - } - } - } - let targets: Vec = target_min.iter().map(TargetEmail::from_value).collect(); - let target_keys = email_keys(&indices); + let (targets, target_keys) = target_emails(net).map_err(Error::from)?; let mut local: Vec<(i64, EmailRow)> = { let mut stmt = ctx @@ -134,28 +97,334 @@ pub fn reconcile( let migrated = maps.targets_of(ObjectType::Mailbox); let mut updates: Vec<(String, Value)> = Vec::new(); + let mut creates: Vec = Vec::new(); for (i, unit) in units.iter().enumerate() { match pairs[i] { Some(t) => match email_patch(&unit.row, &targets[t], maps, &migrated) { Some(patch) => 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, - ), + None => creates.push(i), } } + let mut progress = Progress::new("export: Email", creates.len() as u64, logger); + import_units( + net, + &mut uploader, + maps, + &units, + &creates, + counts, + logger, + &mut progress, + ); update_batch(net, ty, updates, counts, logger); Ok(Plan::default()) } +/// The emails already on the target, and the key each one matches by. +fn target_emails(net: &Net) -> Result<(Vec, Vec), JmapError> { + let ty = ObjectType::Email; + let target_min = target_query_get( + net, + ty, + Some(&["messageId", "size", "mailboxIds", "keywords"]), + )?; + let mut indices: Vec = target_min.iter().map(server_index).collect(); + + let fallback_ids: Vec = target_min + .iter() + .zip(indices.iter()) + .filter(|(_, i)| i.mids.is_empty()) + .filter_map(|(v, _)| jid(v).map(JmapId)) + .collect(); + if !fallback_ids.is_empty() { + let got = get_objects::( + &net.client, + &net.api, + &net.account, + ty.jmap_name(), + &fallback_ids, + Some(&["messageId", "from", "subject", "sentAt", "to"]), + &net.limits, + )?; + let by_id: HashMap = got + .list + .iter() + .filter_map(|v| jid(v).map(|i| (i, v))) + .collect(); + for (v, slot) in target_min.iter().zip(indices.iter_mut()) { + if let Some(full) = jid(v).and_then(|i| by_id.get(&i)) { + *slot = server_index(full); + } + } + } + let targets: Vec = target_min.iter().map(TargetEmail::from_value).collect(); + let keys = email_keys(&indices); + Ok((targets, keys)) +} + +/// The most emails one `Email/import` carries: the server's +/// `maxObjectsInSet`, but no more than this, so that a request that fails +/// without a clear answer leaves few messages to check. +const IMPORT_BATCH_CAP: usize = 50; + +/// One message ready to import: its creation id, its place in `units`, and +/// the `Email/import` entry. +struct Pending { + cid: String, + unit: usize, + item: Value, +} + +/// Writes the messages the target does not have yet. They go in batches of +/// up to `maxObjectsInSet` (capped by `IMPORT_BATCH_CAP`), their blobs +/// uploaded at once up to `maxConcurrentUpload`. Each message is still +/// counted on its own: one rejected in a batch fails alone. A batch that +/// fails without a clear answer is never resent blindly -- the target is +/// checked first, and only what did not arrive is imported again. +#[allow(clippy::too_many_arguments)] +fn import_units( + net: &Net, + uploader: &mut Uploader, + maps: &Maps, + units: &[Unit], + creates: &[usize], + counts: &mut TypeCounts, + logger: &Logger, + progress: &mut Progress, +) { + let batch = (net.limits.max_objects_in_set as usize).clamp(1, IMPORT_BATCH_CAP); + let mut unclear: Vec = Vec::new(); + for chunk in creates.chunks(batch) { + let mut ready: Vec<(usize, Map)> = Vec::new(); + for &i in chunk { + let row = &units[i].row; + match build_mailbox_ids(row, maps) { + Some(mids) => ready.push((i, mids)), + None => { + logger.warn(&format!( + "Email/import e{} ({}) skipped: mailbox not on target", + units[i].local_id, + blob_hint(uploader, row) + )); + counts.failed += 1; + } + } + } + let blobs: Vec = ready + .iter() + .map(|(i, _)| units[*i].row.blob_local_id) + .collect(); + let uploaded = uploader.upload_many(&blobs, "message/rfc822"); + let mut pending: Vec = Vec::new(); + for ((i, mids), result) in ready.into_iter().zip(uploaded) { + let row = &units[i].row; + let cid = format!("e{}", units[i].local_id); + match result { + Ok(blob) => pending.push(Pending { + cid, + unit: i, + item: import_item(blob.0, mids, build_keywords(row), &row.received_at), + }), + Err(e) => { + logger.warn(&format!( + "Email/import {cid} ({}) blob upload failed: {e}{}", + blob_hint(uploader, row), + size_note(&e) + )); + counts.failed += 1; + } + } + } + if net.dry_run { + counts.created += pending.len() as u64; + } else { + send_batch( + net, + uploader, + maps, + units, + pending, + counts, + logger, + &mut unclear, + ); + } + progress.add(chunk.len() as u64); + } + if !unclear.is_empty() { + settle_unclear(net, uploader, maps, units, unclear, counts, logger); + } +} + +/// Sends one `Email/import` for `batch` and counts each message's outcome. +/// A request too large for the server is split in two; a method error that +/// rejects the whole call is retried one message at a time, so the one at +/// fault fails alone. Anything that leaves it unclear whether the server +/// applied the call goes to `unclear`. +#[allow(clippy::too_many_arguments)] +fn send_batch( + net: &Net, + uploader: &mut Uploader, + maps: &Maps, + units: &[Unit], + batch: Vec, + counts: &mut TypeCounts, + logger: &Logger, + unclear: &mut Vec, +) { + if batch.is_empty() { + return; + } + let mut emails = Map::new(); + for p in &batch { + emails.insert(p.cid.clone(), p.item.clone()); + } + let mut req = Request::new(); + req.call( + "Email/import", + json!({ "accountId": net.account, "emails": Value::Object(emails) }), + "i", + ); + let cids: Vec<&str> = batch.iter().map(|p| p.cid.as_str()).collect(); + let sent = req.fits(&net.limits).and_then(|()| { + retry_method_call( + &net.client, + MethodCallKind::SingleObjectWrite, + logger, + || { + let resp = req.send_once(&net.client, &net.api)?; + let mr = resp.first()?; + check_method_error(mr)?; + Ok(cids + .iter() + .map(|cid| interpret_import_for(mr, cid, cids.len())) + .collect::>()) + }, + ) + }); + match sent { + Ok(outcomes) => { + for (p, outcome) in batch.into_iter().zip(outcomes) { + let row = &units[p.unit].row; + match outcome { + SingleImport::Created => counts.created += 1, + SingleImport::Skipped => counts.skipped += 1, + SingleImport::NotCreated { error_type, .. } if error_type == "blobNotFound" => { + retry_after_reupload(net, uploader, maps, &p.cid, row, counts, logger); + } + SingleImport::NotCreated { detail, .. } => { + logger.warn(&format!( + "Email/import {} ({}) failed: {detail}", + p.cid, + blob_hint(uploader, row) + )); + counts.failed += 1; + } + } + } + } + Err(JmapError::RequestTooLarge | JmapError::SingleObjectTooLarge(_)) if batch.len() > 1 => { + let mut batch = batch; + let second = batch.split_off(batch.len() / 2); + send_batch(net, uploader, maps, units, batch, counts, logger, unclear); + send_batch(net, uploader, maps, units, second, counts, logger, unclear); + } + Err(e) if applied_unknown(&e) => { + logger.warn(&format!( + "Email/import of {} message(s) ended without a clear answer ({e}); the target is checked before any is sent again", + batch.len() + )); + unclear.extend(batch); + } + Err(JmapError::Method { .. }) if batch.len() > 1 => { + for p in batch { + send_batch(net, uploader, maps, units, vec![p], counts, logger, unclear); + } + } + Err(e) => { + for p in batch { + logger.warn(&format!( + "Email/import {} ({}) send failed: {e}{}", + p.cid, + blob_hint(uploader, &units[p.unit].row), + size_note(&e) + )); + counts.failed += 1; + } + } + } +} + +/// Whether the server may have applied a call that failed with `e`: the +/// connection broke after the request was sent, the answer was unreadable, or +/// the server said it applied part of it. +fn applied_unknown(e: &JmapError) -> bool { + match e { + JmapError::Transport(_) + | JmapError::RetriesExhausted(_) + | JmapError::Malformed(_) + | JmapError::HttpStatus { .. } => true, + JmapError::Method { error_type, .. } => error_type == "serverPartialFail", + _ => false, + } +} + +/// Settles messages whose import ended without a clear answer: reads the +/// target again, counts those that arrived as created, and imports the rest +/// one at a time. If the target cannot be read, they are counted as failed -- +/// the next export matches whatever did arrive, so none is ever doubled. +fn settle_unclear( + net: &Net, + uploader: &mut Uploader, + maps: &Maps, + units: &[Unit], + unclear: Vec, + counts: &mut TypeCounts, + logger: &Logger, +) { + let (targets, target_keys) = match target_emails(net) { + Ok(t) => t, + Err(e) => { + logger.warn(&format!( + "could not read the target to settle {} message(s) ({e}); they count as failed, and the next export matches whatever arrived", + unclear.len() + )); + counts.failed += unclear.len() as u64; + return; + } + }; + let keys: Vec = email_keys( + &unclear + .iter() + .map(|p| index_from_json(&units[p.unit].row.message_match)) + .collect::>(), + ); + let sizes: Vec> = unclear + .iter() + .map(|p| uploader.blob_len(units[p.unit].row.blob_local_id)) + .collect(); + let pairs = pair_with_targets(&keys, &sizes, &target_keys, &targets); + for (p, found) in unclear.into_iter().zip(pairs) { + if found.is_some() { + counts.created += 1; + continue; + } + let unit = &units[p.unit]; + export_one( + net, + uploader, + maps, + unit.local_id, + &unit.row, + counts, + logger, + ); + } +} + /// 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 @@ -521,6 +790,13 @@ fn send_single_import( fn interpret_import(mr: &MethodCall, cid: &str) -> Result { check_method_error(mr)?; + Ok(interpret_import_for(mr, cid, 1)) +} + +/// One message's outcome in an `Email/import` answer. With a single message +/// in the call, any `created` entry is taken as its own, as servers may key it +/// differently. +fn interpret_import_for(mr: &MethodCall, cid: &str, in_call: usize) -> SingleImport { if let Some(err) = mr .args .get("notCreated") @@ -533,25 +809,21 @@ fn interpret_import(mr: &MethodCall, cid: &str) -> Result = HashSet::new(); let mut updates: Vec<(String, Value)> = Vec::new(); let blobs = Uploader::new(net, &ctx.conn); + let mut progress = Progress::new( + format!("export: {}", ty.jmap_name()), + rows.len() as u64, + logger, + ); for (local, uid) in &rows { + progress.add(1); if let Some((tid, existing)) = by_uid.get(uid) { maps.insert(ty, *local, crate::jmap::wire::JmapId(tid.clone())); matched_uids.insert(uid.clone()); diff --git a/src/sync/mod.rs b/src/sync/mod.rs index a9e7178..5764de4 100644 --- a/src/sync/mod.rs +++ b/src/sync/mod.rs @@ -16,6 +16,7 @@ pub mod import_maildir; pub mod import_managesieve; pub mod import_takeout; pub mod keys; +pub mod progress; pub mod prune; use std::path::PathBuf; diff --git a/src/sync/progress.rs b/src/sync/progress.rs new file mode 100644 index 0000000..cf6c0af --- /dev/null +++ b/src/sync/progress.rs @@ -0,0 +1,196 @@ +/* + * SPDX-FileCopyrightText: 2026 John Coffey + * + * SPDX-License-Identifier: Apache-2.0 OR MIT + */ + +//! A progress line for long runs: how many of how many, how fast, and about +//! how long is left. Printed to stderr at the default log level, at most once +//! per interval, so a large export shows it is moving without flooding the +//! terminal. + +use std::time::{Duration, Instant}; + +use crate::logging::{LEVEL_DEFAULT, Logger}; + +/// How often a progress line is printed while work continues. +pub const PROGRESS_INTERVAL: Duration = Duration::from_secs(5); + +pub struct Progress { + label: String, + total: u64, + done: u64, + started: Instant, + last: Instant, + interval: Duration, + enabled: bool, +} + +impl Progress { + /// `label` names the work, e.g. "export: Email". Nothing is printed when + /// the logger is quiet or there is nothing to do. + pub fn new(label: impl Into, total: u64, logger: &Logger) -> Progress { + let now = Instant::now(); + Progress { + label: label.into(), + total, + done: 0, + started: now, + last: now, + interval: PROGRESS_INTERVAL, + enabled: logger.enabled(LEVEL_DEFAULT) && total > 0, + } + } + + /// Records `n` more items done, and prints a line once the interval has + /// passed since the last one. + pub fn add(&mut self, n: u64) { + self.done = (self.done + n).min(self.total); + if !self.enabled { + return; + } + let now = Instant::now(); + if now.duration_since(self.last) >= self.interval { + self.last = now; + eprintln!("{}", self.line(now.duration_since(self.started))); + } + } + + fn line(&self, elapsed: Duration) -> String { + progress_line(&self.label, self.done, self.total, elapsed) + } +} + +/// `export: Email 1,200/5,000 (24%), 40/s, about 1m35s left`. The rate and +/// the time left are left out until there is enough to estimate them from. +pub fn progress_line(label: &str, done: u64, total: u64, elapsed: Duration) -> String { + let pct = (done * 100).checked_div(total).unwrap_or(100); + let mut out = format!("{label} {}/{} ({pct}%)", thousands(done), thousands(total)); + let secs = elapsed.as_secs_f64(); + if done > 0 && secs >= 1.0 { + let rate = done as f64 / secs; + out.push_str(&format!(", {}/s", format_rate(rate))); + let left = (total - done) as f64 / rate; + if total > done && left.is_finite() { + out.push_str(&format!( + ", about {} left", + format_duration(Duration::from_secs_f64(left)) + )); + } + } + out +} + +/// A duration as `45s`, `1m35s` or `2h03m`. +pub fn format_duration(d: Duration) -> String { + let total = d.as_secs(); + if total < 60 { + format!("{total}s") + } else if total < 3600 { + format!("{}m{:02}s", total / 60, total % 60) + } else { + format!("{}h{:02}m", total / 3600, (total % 3600) / 60) + } +} + +/// The line printed when a type is finished: `export: Email done: 120 +/// created, 3 updated, 5,000 unchanged, 0 failed (2m03s)`. +pub fn done_line(label: &str, counts: &crate::sync::TypeCounts, elapsed: Duration) -> String { + format!( + "{label} done: {} created, {} updated, {} unchanged, {} failed ({})", + thousands(counts.created), + thousands(counts.updated), + thousands(counts.skipped), + thousands(counts.failed), + format_duration(elapsed) + ) +} + +fn format_rate(rate: f64) -> String { + if rate >= 10.0 { + format!("{:.0}", rate) + } else { + format!("{:.1}", rate) + } +} + +pub fn thousands(n: u64) -> String { + let digits = n.to_string(); + let mut out = String::with_capacity(digits.len() + digits.len() / 3); + for (i, c) in digits.chars().enumerate() { + if i > 0 && (digits.len() - i).is_multiple_of(3) { + out.push(','); + } + out.push(c); + } + out +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn a_line_has_count_rate_and_time_left() { + let line = progress_line("export: Email", 1200, 5000, Duration::from_secs(30)); + assert!( + line.starts_with("export: Email 1,200/5,000 (24%)"), + "{line}" + ); + assert!(line.contains("40/s"), "{line}"); + assert!(line.contains("about 1m35s left"), "{line}"); + } + + #[test] + fn no_estimate_before_there_is_something_to_estimate_from() { + assert_eq!( + progress_line("export: Email", 0, 10, Duration::from_secs(5)), + "export: Email 0/10 (0%)" + ); + assert_eq!( + progress_line("export: Email", 3, 10, Duration::from_millis(200)), + "export: Email 3/10 (30%)" + ); + } + + #[test] + fn a_finished_run_shows_no_time_left() { + let line = progress_line("export: Email", 10, 10, Duration::from_secs(4)); + assert!(!line.contains("left"), "{line}"); + assert!(line.contains("(100%)"), "{line}"); + } + + #[test] + fn durations_and_counts_read_naturally() { + assert_eq!(format_duration(Duration::from_secs(45)), "45s"); + assert_eq!(format_duration(Duration::from_secs(95)), "1m35s"); + assert_eq!(format_duration(Duration::from_secs(7380)), "2h03m"); + assert_eq!(thousands(0), "0"); + assert_eq!(thousands(999), "999"); + assert_eq!(thousands(1_000), "1,000"); + assert_eq!(thousands(1_234_567), "1,234,567"); + } + + #[test] + fn the_done_line_names_every_count() { + let counts = crate::sync::TypeCounts { + created: 120, + updated: 3, + skipped: 5000, + failed: 1, + ..Default::default() + }; + assert_eq!( + done_line("export: Email", &counts, Duration::from_secs(123)), + "export: Email done: 120 created, 3 updated, 5,000 unchanged, 1 failed (2m03s)" + ); + } + + #[test] + fn a_quiet_logger_or_empty_job_prints_nothing() { + let p = Progress::new("x", 10, &Logger::from_flags(true, 0)); + assert!(!p.enabled); + let p = Progress::new("x", 0, &Logger::from_flags(false, 0)); + assert!(!p.enabled); + } +} diff --git a/tests/mock_export_batch.rs b/tests/mock_export_batch.rs new file mode 100644 index 0000000..b3baefc --- /dev/null +++ b/tests/mock_export_batch.rs @@ -0,0 +1,342 @@ +/* + * SPDX-FileCopyrightText: 2026 John Coffey + * + * SPDX-License-Identifier: Apache-2.0 OR MIT + */ + +//! Batched `Email/import` on export: batches honor the server's limits, one +//! rejected message fails alone, and a batch that ends without a clear answer +//! is settled against the target instead of being sent again. + +use std::path::{Path, PathBuf}; +use std::sync::Arc; +use std::sync::atomic::{AtomicUsize, Ordering}; + +use inbuxa_migrate::db; +use inbuxa_migrate::jmap::account::AccountSelector; +use inbuxa_migrate::jmap::http::Auth; +use inbuxa_migrate::logging::Logger; +use inbuxa_migrate::sync::{self, CommonConfig, ConnectConfig, ExportConfig, TypeCounts}; +use mockito::Matcher; +use serde_json::{Value, json}; + +const API: &str = "/jmap/api"; + +fn tmp() -> PathBuf { + static SEQ: AtomicUsize = AtomicUsize::new(0); + let n = SEQ.fetch_add(1, Ordering::Relaxed); + let mut p = std::env::temp_dir(); + p.push(format!( + "inbuxa-migrate-exportbatch-{}-{:?}-{n}.sqlite", + std::process::id(), + std::thread::current().id(), + )); + let _ = std::fs::remove_file(&p); + p +} + +fn session(base: &str, max_objects_in_set: u64, max_concurrent_upload: u64) -> String { + json!({ + "apiUrl": format!("{base}{API}"), + "uploadUrl": format!("{base}/jmap/upload/{{accountId}}/"), + "downloadUrl": format!("{base}/jmap/dl/{{accountId}}/{{blobId}}/{{type}}/{{name}}"), + "capabilities": { "urn:ietf:params:jmap:core": { + "maxObjectsInGet": 500, "maxObjectsInSet": max_objects_in_set, + "maxCallsInRequest": 16, "maxConcurrentRequests": 4, + "maxConcurrentUpload": max_concurrent_upload, + "maxSizeRequest": 10000000, "maxSizeUpload": 50000000 + } }, + "accounts": { "w": { "name": "alice", + "accountCapabilities": { "urn:ietf:params:jmap:mail": {} } } } + }) + .to_string() +} + +/// An archive with one Inbox and `n` distinct messages, `` .. ``. +fn seed(n: usize) -> PathBuf { + let archive = tmp(); + let conn = db::init::open(&archive).unwrap(); + conn.execute( + "INSERT INTO mailboxes (id,name,parent_id,role) VALUES (1,'Inbox',NULL,'inbox')", + [], + ) + .unwrap(); + for i in 1..=n { + let raw = format!("From: a@x\r\nSubject: m{i}\r\nMessage-ID: \r\n\r\nbody {i}"); + let blob = db::blobs::intern_blob(&conn, raw.as_bytes()).unwrap(); + let mm = inbuxa_migrate::sync::keys::index_to_json( + &inbuxa_migrate::sync::emailmeta::email_index_from_blob(raw.as_bytes()), + ); + conn.execute( + "INSERT INTO emails (blob_id,received_at,mailbox_ids,keywords,message_match) + VALUES (?1,'2020-01-01T00:00:00Z','[1]','[]',?2)", + rusqlite::params![blob, mm], + ) + .unwrap(); + } + archive +} + +/// The session, an Inbox already on the target, and uploads. The target's +/// email list is left to each test. +fn mock_target(server: &mut mockito::ServerGuard, session_body: String) -> Vec { + vec![ + server.mock("GET", "/").with_status(404).create(), + server + .mock("GET", "/.well-known/jmap") + .with_body(session_body) + .create(), + server + .mock("POST", API) + .match_body(Matcher::Regex("Mailbox/query".into())) + .with_body( + json!({"methodResponses":[["Mailbox/query", + {"accountId":"w","ids":["t1"]},"q"]]}) + .to_string(), + ) + .create(), + server + .mock("POST", API) + .match_body(Matcher::AllOf(vec![ + Matcher::Regex("Mailbox/query".into()), + Matcher::Regex("anchor".into()), + ])) + .with_body( + json!({"methodResponses":[["Mailbox/query", + {"accountId":"w","ids":[]},"q"]]}) + .to_string(), + ) + .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}}],"notFound":[]},"g"]]}) + .to_string(), + ) + .create(), + server + .mock("POST", Matcher::Regex("/jmap/upload/".into())) + .with_body(json!({"blobId":"UP"}).to_string()) + .create(), + ] +} + +fn empty_email_query(server: &mut mockito::ServerGuard) -> mockito::Mock { + server + .mock("POST", API) + .match_body(Matcher::Regex("Email/query".into())) + .with_body( + json!({"methodResponses":[["Email/query",{"accountId":"w","ids":[]},"q"]]}).to_string(), + ) + .create() +} + +/// The creation ids of the emails in an `Email/import` request body. +fn import_cids(body: &[u8]) -> Vec { + let v: Value = serde_json::from_slice(body).unwrap_or(Value::Null); + v["methodCalls"][0][1]["emails"] + .as_object() + .map(|m| m.keys().cloned().collect()) + .unwrap_or_default() +} + +/// An `Email/import` answer: each id in `cids` created, except those in +/// `rejected`, which come back as `invalidEmail`. +fn import_answer(cids: &[String], rejected: &[&str]) -> Vec { + let mut created = serde_json::Map::new(); + let mut not_created = serde_json::Map::new(); + for cid in cids { + if rejected.contains(&cid.as_str()) { + not_created.insert(cid.clone(), json!({"type":"invalidEmail"})); + } else { + created.insert( + cid.clone(), + json!({"id": format!("T{cid}"), "blobId":"b","threadId":"t","size":10}), + ); + } + } + json!({"methodResponses":[["Email/import", + {"accountId":"w","created":created,"notCreated":not_created},"i"]]}) + .to_string() + .into_bytes() +} + +fn run_export(archive: &Path, base: &str, threads: usize) -> TypeCounts { + let summary = sync::export::run( + CommonConfig { + archive: archive.to_path_buf(), + threads, + dry_run: false, + max_retries: 1, + allow_invalid_certs: false, + logger: Logger::from_flags(true, 0), + }, + ExportConfig { + connect: ConnectConfig { + url: base.to_owned(), + auth: Auth::Basic { + user: "u".into(), + password: "p".into(), + }, + account: AccountSelector::Id("w".into()), + }, + objects: None, + prune: false, + yes: true, + }, + ) + .expect("export run"); + summary + .per_type + .iter() + .find(|(t, _)| *t == "Email") + .map(|(_, c)| c.clone()) + .expect("email counts") +} + +#[test] +fn imports_are_batched_up_to_max_objects_in_set() { + let mut server = mockito::Server::new(); + let base = server.url(); + let archive = seed(5); + let _base_mocks = mock_target(&mut server, session(&base, 2, 4)); + let _eq = empty_email_query(&mut server); + + let sizes = Arc::new(std::sync::Mutex::new(Vec::new())); + let seen = sizes.clone(); + let imports = server + .mock("POST", API) + .match_body(Matcher::Regex("Email/import".into())) + .with_body_from_request(move |req| { + let cids = import_cids(req.body().unwrap()); + seen.lock().unwrap().push(cids.len()); + import_answer(&cids, &[]) + }) + .expect(3) + .create(); + + let email = run_export(&archive, &base, 4); + assert_eq!(email.created, 5); + assert_eq!(email.failed, 0); + imports.assert(); + let mut sizes = sizes.lock().unwrap().clone(); + sizes.sort_unstable(); + assert_eq!( + sizes, + vec![1, 2, 2], + "no call carries more than maxObjectsInSet" + ); + let _ = std::fs::remove_file(&archive); +} + +#[test] +fn a_rejected_message_in_a_batch_fails_alone() { + let mut server = mockito::Server::new(); + let base = server.url(); + let archive = seed(3); + let _base_mocks = mock_target(&mut server, session(&base, 50, 4)); + let _eq = empty_email_query(&mut server); + + let rejected = Arc::new(std::sync::Mutex::new(String::new())); + let pick = rejected.clone(); + let imports = server + .mock("POST", API) + .match_body(Matcher::Regex("Email/import".into())) + .with_body_from_request(move |req| { + let cids = import_cids(req.body().unwrap()); + let mut sorted = cids.clone(); + sorted.sort(); + let middle = sorted[1].clone(); + *pick.lock().unwrap() = middle.clone(); + import_answer(&cids, &[middle.as_str()]) + }) + .expect(1) + .create(); + + let email = run_export(&archive, &base, 1); + assert_eq!(email.created, 2, "the other two land"); + assert_eq!(email.failed, 1, "only {} fails", rejected.lock().unwrap()); + imports.assert(); + let _ = std::fs::remove_file(&archive); +} + +#[test] +fn an_unclear_batch_is_settled_against_the_target_not_resent() { + let mut server = mockito::Server::new(); + let base = server.url(); + let archive = seed(2); + let _base_mocks = mock_target(&mut server, session(&base, 50, 4)); + + // First look: the target is empty. After the unclear batch: m-1 arrived. + let queries = Arc::new(AtomicUsize::new(0)); + let q = queries.clone(); + let _eq = server + .mock("POST", API) + .match_body(Matcher::Regex("Email/query".into())) + .with_body_from_request(move |req| { + // A page after the first (it carries an anchor) is empty. + let paged = String::from_utf8_lossy(req.body().unwrap()).contains("anchor"); + let ids = if paged || q.fetch_add(1, Ordering::SeqCst) == 0 { + json!([]) + } else { + json!(["arrived"]) + }; + json!({"methodResponses":[["Email/query",{"accountId":"w","ids":ids},"q"]]}) + .to_string() + .into_bytes() + }) + .create(); + let _eg = server + .mock("POST", API) + .match_body(Matcher::Regex("Email/get".into())) + .with_body( + json!({"methodResponses":[["Email/get",{"accountId":"w","list":[ + {"id":"arrived","messageId":["m-1@h"],"mailboxIds":{"t1":true},"keywords":{}} + ],"notFound":[]},"g"]]}) + .to_string(), + ) + .create(); + + // The batch carries both messages and gets a gateway timeout, so it may + // or may not have been applied. mockito answers with the first matching + // mock still owed hits, so this answers the first import only; every + // later import is created by the next mock, which records what it sees. + let gateway = server + .mock("POST", API) + .match_body(Matcher::Regex("Email/import".into())) + .with_status(504) + .expect(1) + .create(); + let later = Arc::new(std::sync::Mutex::new(Vec::new())); + let later_seen = later.clone(); + let _created = server + .mock("POST", API) + .match_body(Matcher::Regex("Email/import".into())) + .with_body_from_request(move |req| { + let cids = import_cids(req.body().unwrap()); + later_seen.lock().unwrap().push(cids.clone()); + import_answer(&cids, &[]) + }) + .create(); + + let email = run_export(&archive, &base, 1); + gateway.assert(); + assert!( + queries.load(Ordering::SeqCst) >= 2, + "the target is read again before anything is resent" + ); + assert_eq!( + *later.lock().unwrap(), + vec![vec!["e2".to_owned()]], + "only the message that did not arrive is sent again, on its own" + ); + assert_eq!( + email.created, 2, + "one found on the target, one imported again" + ); + assert_eq!(email.failed, 0); + let _ = std::fs::remove_file(&archive); +} diff --git a/tests/mock_sync.rs b/tests/mock_sync.rs index e8b2dbc..829fdf7 100644 --- a/tests/mock_sync.rs +++ b/tests/mock_sync.rs @@ -487,150 +487,6 @@ fn export_mailbox_already_exists_maps_existing_id() { let _ = std::fs::remove_file(&archive); } -#[test] -fn email_export_sends_one_email_per_import_call() { - let mut server = mockito::Server::new(); - let base = server.url(); - let api = "/jmap/api"; - - let archive = tmp(); - { - let conn = db::init::open(&archive).unwrap(); - conn.execute( - "INSERT INTO mailboxes (id,name,parent_id,role) VALUES (1,'Inbox',NULL,'inbox')", - [], - ) - .unwrap(); - for n in 1..=2 { - let raw = - format!("From: a@x\r\nSubject: m{n}\r\nMessage-ID: \r\n\r\nbody {n}",); - let blob = db::blobs::intern_blob(&conn, raw.as_bytes()).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(); - let _wk = server - .mock("GET", "/.well-known/jmap") - .with_body(session_body(&base)) - .create(); - - let _mq = server - .mock("POST", api) - .match_body(Matcher::Regex("Mailbox/query".into())) - .with_body( - json!({"methodResponses":[["Mailbox/query", - {"accountId":"w","ids":["t1"]},"q"]]}) - .to_string(), - ) - .expect_at_least(1) - .create(); - let _mq_empty = server - .mock("POST", api) - .match_body(Matcher::AllOf(vec![ - Matcher::Regex("Mailbox/query".into()), - Matcher::Regex("anchor".into()), - ])) - .with_body( - json!({"methodResponses":[["Mailbox/query", - {"accountId":"w","ids":[]},"q"]]}) - .to_string(), - ) - .create(); - let _mg = 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}}],"notFound":[]},"g"]]}) - .to_string(), - ) - .expect_at_least(1) - .create(); - 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_at_least(1) - .create(); - - let _ups = server - .mock("POST", Matcher::Regex("/jmap/upload/".into())) - .with_body(json!({"blobId":"UP1"}).to_string()) - .expect(2) - .create(); - - let single_only = server - .mock("POST", api) - .match_body(Matcher::AllOf(vec![ - Matcher::Regex("Email/import".into()), - Matcher::Regex("e1".into()), - Matcher::Regex("e2".into()), - ])) - .expect(0) - .create(); - - let imports = server - .mock("POST", api) - .match_body(Matcher::Regex("Email/import".into())) - .with_body( - json!({"methodResponses":[["Email/import", - {"accountId":"w","created":{"e":{"id":"x","blobId":"b","threadId":"t","size":10}}},"i"]]}) - .to_string(), - ) - .expect(2) - .create(); - - let summary = sync::export::run( - CommonConfig { - archive: archive.clone(), - threads: 1, - dry_run: false, - max_retries: 0, - allow_invalid_certs: false, - logger: Logger::from_flags(true, 0), - }, - ExportConfig { - connect: ConnectConfig { - url: base.clone(), - auth: Auth::Basic { - user: "u".into(), - password: "p".into(), - }, - account: AccountSelector::Id("w".into()), - }, - objects: None, - prune: false, - yes: true, - }, - ) - .expect("export run"); - - let email = summary - .per_type - .iter() - .find(|(t, _)| *t == "Email") - .map(|(_, c)| c.clone()) - .expect("email counts"); - assert_eq!(email.created, 2, "both emails imported in per-item rounds"); - assert_eq!(email.failed, 0, "no per-unit failure"); - assert!(!summary.any_failed(), "no whole-run failure"); - - single_only.assert(); - imports.assert(); - let _ = std::fs::remove_file(&archive); -} - #[test] fn export_email_blob_not_found_reuploads_and_retries() { let mut server = mockito::Server::new();