Merge pull request 'export: batch Email/import, upload in parallel, and show progress' (#10) from feat/export-progress-batching into main
This commit was merged in pull request #10.
This commit is contained in:
@@ -244,6 +244,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.
|
||||
|
||||
|
||||
@@ -97,6 +97,10 @@ impl Inner {
|
||||
#[derive(Debug, Clone, Copy)]
|
||||
enum Kind {
|
||||
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,
|
||||
}
|
||||
|
||||
@@ -220,6 +224,23 @@ impl HttpClient {
|
||||
.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(
|
||||
&self,
|
||||
upload_url: &str,
|
||||
@@ -298,6 +319,16 @@ impl HttpClient {
|
||||
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 => {
|
||||
attempt += 1;
|
||||
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) {
|
||||
Disposition::Fatal => return Err(err),
|
||||
Disposition::Retryable => {
|
||||
@@ -659,6 +691,90 @@ mod tests {
|
||||
use super::*;
|
||||
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]
|
||||
fn invalid_certificates_are_accepted_only_for_the_named_host() {
|
||||
let client = HttpClient::new(
|
||||
|
||||
@@ -131,6 +131,14 @@ impl Request {
|
||||
let value = client.post_json(api_url, &self.envelope()?)?;
|
||||
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)]
|
||||
|
||||
@@ -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());
|
||||
}
|
||||
}
|
||||
|
||||
+329
-57
@@ -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<Plan, Error> {
|
||||
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<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 (targets, target_keys) = target_emails(net).map_err(Error::from)?;
|
||||
|
||||
let mut local: Vec<(i64, EmailRow)> = {
|
||||
let mut stmt = ctx
|
||||
@@ -134,26 +97,332 @@ pub fn reconcile(
|
||||
|
||||
let migrated = maps.targets_of(ObjectType::Mailbox);
|
||||
let mut updates: Vec<(String, Value)> = Vec::new();
|
||||
let mut creates: Vec<usize> = 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(
|
||||
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<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.row,
|
||||
counts,
|
||||
logger,
|
||||
),
|
||||
);
|
||||
}
|
||||
}
|
||||
update_batch(net, ty, updates, counts, logger);
|
||||
|
||||
Ok(Plan::default())
|
||||
}
|
||||
|
||||
/// 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> {
|
||||
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<SingleImport, JmapErro
|
||||
.unwrap_or("")
|
||||
.to_owned();
|
||||
if error_type == "alreadyExists" {
|
||||
return Ok(SingleImport::Skipped);
|
||||
return SingleImport::Skipped;
|
||||
}
|
||||
return Ok(SingleImport::NotCreated {
|
||||
return SingleImport::NotCreated {
|
||||
error_type,
|
||||
detail: err.to_string(),
|
||||
});
|
||||
};
|
||||
}
|
||||
if mr
|
||||
.args
|
||||
.get("created")
|
||||
.and_then(Value::as_object)
|
||||
.is_some_and(|c| !c.is_empty())
|
||||
{
|
||||
return Ok(SingleImport::Created);
|
||||
let created = mr.args.get("created").and_then(Value::as_object);
|
||||
if created.is_some_and(|c| c.contains_key(cid) || (in_call == 1 && !c.is_empty())) {
|
||||
return SingleImport::Created;
|
||||
}
|
||||
Ok(SingleImport::NotCreated {
|
||||
SingleImport::NotCreated {
|
||||
error_type: String::new(),
|
||||
detail: format!("Email/import returned neither created nor notCreated for {cid}"),
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
|
||||
@@ -17,6 +17,7 @@ use crate::logging::Logger;
|
||||
use crate::sync::import_jmap::mapping::{
|
||||
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::{Context, TypeCounts};
|
||||
use crate::types::ObjectType;
|
||||
@@ -79,7 +80,13 @@ pub fn reconcile(
|
||||
let mut matched_uids: HashSet<String> = 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());
|
||||
|
||||
@@ -17,6 +17,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;
|
||||
|
||||
@@ -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);
|
||||
}
|
||||
}
|
||||
@@ -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);
|
||||
}
|
||||
@@ -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: <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]
|
||||
fn export_email_blob_not_found_reuploads_and_retries() {
|
||||
let mut server = mockito::Server::new();
|
||||
|
||||
Reference in New Issue
Block a user