export: batch Email/import, upload in parallel, and show progress #10
@@ -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.
|
||||||
|
|
||||||
|
|||||||
@@ -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(
|
||||||
|
|||||||
@@ -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)]
|
||||||
|
|||||||
@@ -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());
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
+332
-60
@@ -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,28 +97,334 @@ 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),
|
||||||
net,
|
|
||||||
&mut uploader,
|
|
||||||
maps,
|
|
||||||
unit.local_id,
|
|
||||||
&unit.row,
|
|
||||||
counts,
|
|
||||||
logger,
|
|
||||||
),
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
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);
|
update_batch(net, ty, updates, counts, logger);
|
||||||
|
|
||||||
Ok(Plan::default())
|
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,
|
||||||
|
);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
/// 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
|
||||||
/// together. A source that files one message in several folders (IMAP and
|
/// together. A source that files one message in several folders (IMAP and
|
||||||
/// Maildir copies, Gmail labels) leaves one archive row per folder; on the
|
/// Maildir copies, Gmail labels) leaves one archive row per folder; on the
|
||||||
@@ -521,6 +790,13 @@ fn send_single_import(
|
|||||||
|
|
||||||
fn interpret_import(mr: &MethodCall, cid: &str) -> Result<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)]
|
||||||
|
|||||||
@@ -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());
|
||||||
|
|||||||
@@ -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;
|
||||||
|
|||||||
@@ -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);
|
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();
|
||||||
|
|||||||
Reference in New Issue
Block a user