export: batch Email/import, upload in parallel, and show progress
ci / test (pull_request) Skipped
github/ci (branch) GitHub Actions
ci / github (pull_request) Successful in 3m11s
ci / announce (pull_request) Skipped

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.
This commit is contained in:
2026-09-30 12:39:57 -07:00
parent 4ce772d6e5
commit 3f832e17b2
8 changed files with 1084 additions and 204 deletions
+192
View File
@@ -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<Result<JmapId, JmapError>> {
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<i64> = 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<i64, JmapError> = 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<J, S, R>(
jobs: &[J],
workers: usize,
init: impl Fn() -> S + Sync,
f: impl Fn(&mut S, &J) -> R + Sync,
) -> Vec<R>
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<Vec<Option<R>>> = 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<Vec<u8>, 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<Summary, Error>
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<Summary, Error>
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<Summary, Error>
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<usize> = (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::<Vec<_>>());
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<i32> = run_bounded(&[], 4, || (), |_, j: &i32| *j);
assert!(out.is_empty());
}
}