export: batch Email/import, upload in parallel, and show progress #10

Merged
jcoffey-dev merged 2 commits from feat/export-progress-batching into main 2026-09-30 19:47:47 +00:00
10 changed files with 1208 additions and 204 deletions
+14
View File
@@ -229,6 +229,20 @@ what changed at the source in between.
- **Sieve scripts** are matched by name, and the target's is replaced when - **Sieve scripts** are matched by name, and the target's is replaced when
its content differs from the archive's. 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 `--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. first; `--yes` answers for it, for scripts. Export speaks JMAP only.
+116
View File
@@ -97,6 +97,10 @@ impl Inner {
#[derive(Debug, Clone, Copy)] #[derive(Debug, Clone, Copy)]
enum Kind { enum Kind {
Api, Api,
/// An API call that must not be sent twice: a write the server may
/// already have applied when the connection failed. A transport error is
/// returned to the caller, which checks the target instead of resending.
ApiOnce,
Upload, Upload,
} }
@@ -220,6 +224,23 @@ impl HttpClient {
.map_err(|e| JmapError::Malformed(format!("response is not valid json: {e}"))) .map_err(|e| JmapError::Malformed(format!("response is not valid json: {e}")))
} }
/// As `post_json`, but a transport failure is not retried: the request
/// may have reached the server, and sending it again could apply it twice.
/// Throttling and `503` answers, which mean it was not processed, are still
/// retried.
pub fn post_json_once(&self, url: &str, body: &Value) -> Result<Value, JmapError> {
let payload = serde_json::to_vec(body)?;
let raw = self.execute(
Kind::ApiOnce,
"POST",
url,
Some(&payload),
Some("application/json"),
)?;
serde_json::from_slice(&raw)
.map_err(|e| JmapError::Malformed(format!("response is not valid json: {e}")))
}
pub fn upload( pub fn upload(
&self, &self,
upload_url: &str, upload_url: &str,
@@ -298,6 +319,16 @@ impl HttpClient {
body: truncate(&body), body: truncate(&body),
}); });
} }
// A gateway error may come back after the server
// behind it applied the write: not safe to resend.
StatusOutcome::Retryable
if matches!(kind, Kind::ApiOnce) && matches!(status, 502 | 504) =>
{
return Err(JmapError::HttpStatus {
status,
body: truncate(&body),
});
}
StatusOutcome::Retryable => { StatusOutcome::Retryable => {
attempt += 1; attempt += 1;
self.inner.retries_total.fetch_add(1, Ordering::Relaxed); self.inner.retries_total.fetch_add(1, Ordering::Relaxed);
@@ -344,6 +375,7 @@ impl HttpClient {
} }
} }
} }
Attempt::Transport(err) if matches!(kind, Kind::ApiOnce) => return Err(err),
Attempt::Transport(err) => match transport_disposition(&err) { Attempt::Transport(err) => match transport_disposition(&err) {
Disposition::Fatal => return Err(err), Disposition::Fatal => return Err(err),
Disposition::Retryable => { Disposition::Retryable => {
@@ -659,6 +691,90 @@ mod tests {
use super::*; use super::*;
use crate::net::CertOverride; use crate::net::CertOverride;
/// A server that reads each request and hangs up without answering: a
/// transport failure after the request has been sent. Returns its URL and
/// the count of requests it has seen.
fn hang_up_server() -> (String, Arc<AtomicU64>) {
use std::io::Read;
let listener = std::net::TcpListener::bind("127.0.0.1:0").unwrap();
let url = format!("http://{}/jmap/api", listener.local_addr().unwrap());
let seen = Arc::new(AtomicU64::new(0));
let counter = seen.clone();
std::thread::spawn(move || {
for stream in listener.incoming() {
let Ok(mut stream) = stream else { break };
let mut buf = [0u8; 4096];
let _ = stream.read(&mut buf);
counter.fetch_add(1, Ordering::SeqCst);
}
});
(url, seen)
}
fn quick_retries(max_retries: u32) -> RetryPolicy {
RetryPolicy {
max_retries,
base: Duration::from_millis(1),
cap: Duration::from_millis(2),
}
}
#[test]
fn a_write_sent_once_is_not_resent_after_a_transport_failure() {
let (url, seen) = hang_up_server();
let client = HttpClient::new(
Auth::Bearer { token: "t".into() },
quick_retries(2),
CertOverride::none(),
);
let err = client
.post_json_once(&url, &serde_json::json!({}))
.unwrap_err();
assert!(matches!(err, JmapError::Transport(_)), "{err}");
assert_eq!(seen.load(Ordering::SeqCst), 1, "sent exactly once");
let (url, seen) = hang_up_server();
let _ = client.post_json(&url, &serde_json::json!({}));
assert_eq!(
seen.load(Ordering::SeqCst),
3,
"an ordinary call retries twice"
);
}
#[test]
fn a_write_sent_once_is_not_resent_after_a_gateway_timeout_but_is_after_503() {
let mut server = mockito::Server::new();
let url = format!("{}/jmap/api", server.url());
let client = HttpClient::new(
Auth::Bearer { token: "t".into() },
quick_retries(2),
CertOverride::none(),
);
let gateway = server
.mock("POST", "/jmap/api")
.with_status(504)
.expect(1)
.create();
let err = client
.post_json_once(&url, &serde_json::json!({}))
.unwrap_err();
assert!(
matches!(err, JmapError::HttpStatus { status: 504, .. }),
"{err}"
);
gateway.assert();
gateway.remove();
let busy = server
.mock("POST", "/jmap/api")
.with_status(503)
.expect(3)
.create();
let _ = client.post_json_once(&url, &serde_json::json!({}));
busy.assert();
}
#[test] #[test]
fn invalid_certificates_are_accepted_only_for_the_named_host() { fn invalid_certificates_are_accepted_only_for_the_named_host() {
let client = HttpClient::new( let client = HttpClient::new(
+8
View File
@@ -131,6 +131,14 @@ impl Request {
let value = client.post_json(api_url, &self.envelope()?)?; let value = client.post_json(api_url, &self.envelope()?)?;
Response::parse(value) Response::parse(value)
} }
/// As `send`, for a write that must not be applied twice: a transport
/// failure comes back as an error instead of being resent (see
/// `HttpClient::post_json_once`).
pub fn send_once(&self, client: &HttpClient, api_url: &str) -> Result<Response, JmapError> {
let value = client.post_json_once(api_url, &self.envelope()?)?;
Response::parse(value)
}
} }
#[derive(Debug)] #[derive(Debug)]
+192
View File
@@ -7,6 +7,9 @@
use std::collections::{HashMap, HashSet}; use std::collections::{HashMap, HashSet};
use std::io::{IsTerminal, Write}; use std::io::{IsTerminal, Write};
use std::path::PathBuf;
use std::sync::Mutex;
use std::sync::atomic::{AtomicUsize, Ordering};
use rusqlite::Connection; use rusqlite::Connection;
use serde_json::{Map, Value, json}; use serde_json::{Map, Value, json};
@@ -134,6 +137,77 @@ impl<'a> Uploader<'a> {
Ok(id) 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) { fn invalidate(&mut self, local_id: i64) {
self.cache.remove(&local_id); 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<'_> { impl BlobBytes for Uploader<'_> {
fn bytes(&self, local_id: i64) -> Result<Vec<u8>, JmapError> { fn bytes(&self, local_id: i64) -> Result<Vec<u8>, JmapError> {
db::blobs::blob_bytes(self.conn, local_id)? db::blobs::blob_bytes(self.conn, local_id)?
@@ -162,6 +288,11 @@ struct Net {
limits: Limits, limits: Limits,
session: Session, session: Session,
dry_run: bool, 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 { 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, limits: connected.limits,
session: connected.session.clone(), session: connected.session.clone(),
dry_run: ctx.dry_run(), 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); 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) { if logger.enabled(LEVEL_DEFAULT) {
eprintln!("export: {} ...", ty.jmap_name()); eprintln!("export: {} ...", ty.jmap_name());
} }
let started = std::time::Instant::now();
let mut counts = TypeCounts::default(); let mut counts = TypeCounts::default();
let res = reconcile_type( let res = reconcile_type(
&ctx, &ctx,
@@ -217,6 +353,16 @@ pub fn run(common: CommonConfig, config: ExportConfig) -> Result<Summary, Error>
Plan::default() 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); plans.insert(*ty, plan);
counts_per_type.insert(*ty, counts); counts_per_type.insert(*ty, counts);
} }
@@ -635,3 +781,49 @@ mod common {
outcome 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());
}
}
+329 -57
View File
@@ -21,6 +21,7 @@ use crate::jmap::wire::JmapId;
use crate::logging::Logger; use crate::logging::Logger;
use crate::sync::import_jmap::mapping::{EMAIL_SELECT, EmailRow, TargetResolver, row_to_email}; 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::keys::{EmailIndex, EmailKey, email_index, email_keys, index_from_json};
use crate::sync::progress::Progress;
use crate::sync::{Context, TypeCounts}; use crate::sync::{Context, TypeCounts};
use crate::types::ObjectType; use crate::types::ObjectType;
@@ -61,45 +62,7 @@ pub fn reconcile(
logger: &Logger, logger: &Logger,
) -> Result<Plan, Error> { ) -> Result<Plan, Error> {
let ty = ObjectType::Email; let ty = ObjectType::Email;
let (targets, target_keys) = target_emails(net).map_err(Error::from)?;
let target_min = target_query_get(
net,
ty,
Some(&["messageId", "size", "mailboxIds", "keywords"]),
)
.map_err(Error::from)?;
let mut indices: Vec<EmailIndex> = target_min.iter().map(server_index).collect();
let fallback_ids: Vec<JmapId> = 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::<Value>(
&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<String, &Value> = 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<TargetEmail> = target_min.iter().map(TargetEmail::from_value).collect();
let target_keys = email_keys(&indices);
let mut local: Vec<(i64, EmailRow)> = { let mut local: Vec<(i64, EmailRow)> = {
let mut stmt = ctx let mut stmt = ctx
@@ -134,26 +97,332 @@ pub fn reconcile(
let migrated = maps.targets_of(ObjectType::Mailbox); let migrated = maps.targets_of(ObjectType::Mailbox);
let mut updates: Vec<(String, Value)> = Vec::new(); let mut updates: Vec<(String, Value)> = Vec::new();
let mut creates: Vec<usize> = Vec::new();
for (i, unit) in units.iter().enumerate() { for (i, unit) in units.iter().enumerate() {
match pairs[i] { match pairs[i] {
Some(t) => match email_patch(&unit.row, &targets[t], maps, &migrated) { Some(t) => match email_patch(&unit.row, &targets[t], maps, &migrated) {
Some(patch) => updates.push((targets[t].id.clone(), patch)), Some(patch) => updates.push((targets[t].id.clone(), patch)),
None => counts.skipped += 1, None => counts.skipped += 1,
}, },
None => export_one( None => creates.push(i),
}
}
let mut progress = Progress::new("export: Email", creates.len() as u64, logger);
import_units(
net, net,
&mut uploader, &mut uploader,
maps, 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<TargetEmail>, Vec<EmailKey>), JmapError> {
let ty = ObjectType::Email;
let target_min = target_query_get(
net,
ty,
Some(&["messageId", "size", "mailboxIds", "keywords"]),
)?;
let mut indices: Vec<EmailIndex> = target_min.iter().map(server_index).collect();
let fallback_ids: Vec<JmapId> = 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::<Value>(
&net.client,
&net.api,
&net.account,
ty.jmap_name(),
&fallback_ids,
Some(&["messageId", "from", "subject", "sentAt", "to"]),
&net.limits,
)?;
let by_id: HashMap<String, &Value> = 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<TargetEmail> = 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<Pending> = Vec::new();
for chunk in creates.chunks(batch) {
let mut ready: Vec<(usize, Map<String, Value>)> = 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<i64> = ready
.iter()
.map(|(i, _)| units[*i].row.blob_local_id)
.collect();
let uploaded = uploader.upload_many(&blobs, "message/rfc822");
let mut pending: Vec<Pending> = 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<Pending>,
counts: &mut TypeCounts,
logger: &Logger,
unclear: &mut Vec<Pending>,
) {
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::<Vec<_>>())
},
)
});
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<Pending>,
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<EmailKey> = email_keys(
&unclear
.iter()
.map(|p| index_from_json(&units[p.unit].row.message_match))
.collect::<Vec<_>>(),
);
let sizes: Vec<Option<u64>> = 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.local_id,
&unit.row, &unit.row,
counts, counts,
logger, logger,
), );
} }
}
update_batch(net, ty, updates, counts, logger);
Ok(Plan::default())
} }
/// One message to write: the archive rows that hold the same bytes, folded /// One message to write: the archive rows that hold the same bytes, folded
@@ -521,6 +790,13 @@ fn send_single_import(
fn interpret_import(mr: &MethodCall, cid: &str) -> Result<SingleImport, JmapError> { fn interpret_import(mr: &MethodCall, cid: &str) -> Result<SingleImport, JmapError> {
check_method_error(mr)?; 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 if let Some(err) = mr
.args .args
.get("notCreated") .get("notCreated")
@@ -533,25 +809,21 @@ fn interpret_import(mr: &MethodCall, cid: &str) -> Result<SingleImport, JmapErro
.unwrap_or("") .unwrap_or("")
.to_owned(); .to_owned();
if error_type == "alreadyExists" { if error_type == "alreadyExists" {
return Ok(SingleImport::Skipped); return SingleImport::Skipped;
} }
return Ok(SingleImport::NotCreated { return SingleImport::NotCreated {
error_type, error_type,
detail: err.to_string(), detail: err.to_string(),
}); };
} }
if mr let created = mr.args.get("created").and_then(Value::as_object);
.args if created.is_some_and(|c| c.contains_key(cid) || (in_call == 1 && !c.is_empty())) {
.get("created") return SingleImport::Created;
.and_then(Value::as_object)
.is_some_and(|c| !c.is_empty())
{
return Ok(SingleImport::Created);
} }
Ok(SingleImport::NotCreated { SingleImport::NotCreated {
error_type: String::new(), error_type: String::new(),
detail: format!("Email/import returned neither created nor notCreated for {cid}"), detail: format!("Email/import returned neither created nor notCreated for {cid}"),
}) }
} }
#[cfg(test)] #[cfg(test)]
+7
View File
@@ -17,6 +17,7 @@ use crate::logging::Logger;
use crate::sync::import_jmap::mapping::{ use crate::sync::import_jmap::mapping::{
CALENDAR_EVENT_SELECT, CONTACT_CARD_SELECT, calendar_event_to_wire, contact_card_to_wire, CALENDAR_EVENT_SELECT, CONTACT_CARD_SELECT, calendar_event_to_wire, contact_card_to_wire,
}; };
use crate::sync::progress::Progress;
use crate::sync::prune::{TargetObj, candidates}; use crate::sync::prune::{TargetObj, candidates};
use crate::sync::{Context, TypeCounts}; use crate::sync::{Context, TypeCounts};
use crate::types::ObjectType; use crate::types::ObjectType;
@@ -79,7 +80,13 @@ pub fn reconcile(
let mut matched_uids: HashSet<String> = HashSet::new(); let mut matched_uids: HashSet<String> = HashSet::new();
let mut updates: Vec<(String, Value)> = Vec::new(); let mut updates: Vec<(String, Value)> = Vec::new();
let blobs = Uploader::new(net, &ctx.conn); 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 { for (local, uid) in &rows {
progress.add(1);
if let Some((tid, existing)) = by_uid.get(uid) { if let Some((tid, existing)) = by_uid.get(uid) {
maps.insert(ty, *local, crate::jmap::wire::JmapId(tid.clone())); maps.insert(ty, *local, crate::jmap::wire::JmapId(tid.clone()));
matched_uids.insert(uid.clone()); matched_uids.insert(uid.clone());
+1
View File
@@ -16,6 +16,7 @@ pub mod import_maildir;
pub mod import_managesieve; pub mod import_managesieve;
pub mod import_takeout; pub mod import_takeout;
pub mod keys; pub mod keys;
pub mod progress;
pub mod prune; pub mod prune;
use std::path::PathBuf; use std::path::PathBuf;
+196
View File
@@ -0,0 +1,196 @@
/*
* SPDX-FileCopyrightText: 2026 John Coffey <[email protected]>
*
* 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<String>, 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);
}
}
+342
View File
@@ -0,0 +1,342 @@
/*
* SPDX-FileCopyrightText: 2026 John Coffey <[email protected]>
*
* 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, `<m-1@h>` .. `<m-n@h>`.
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: <m-{i}@h>\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<mockito::Mock> {
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<String> {
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<u8> {
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);
}
-144
View File
@@ -487,150 +487,6 @@ fn export_mailbox_already_exists_maps_existing_id() {
let _ = std::fs::remove_file(&archive); 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: <m-{n}@h>\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] #[test]
fn export_email_blob_not_found_reuploads_and_retries() { fn export_email_blob_not_found_reuploads_and_retries() {
let mut server = mockito::Server::new(); let mut server = mockito::Server::new();