import: bound the memory an IMAP or EWS import holds
Two things let a large mailbox fill memory. IMAP workers handed every
fetched message to the archive writer through an unbounded queue, so fast
workers could hold whole folders' worth of bodies while the single writer
caught up. And fetch batches were sized by count only: 20 items per EWS
GetItem, 8 connections at a time, is a few megabytes of ordinary mail and
several gigabytes of large attachments.
- The IMAP event queue is bounded at two events per worker, so a worker
waits for the writer instead of running ahead of it.
- Fetch batches are bounded by bytes as well as by count, with a shared
helper, `sync::batch::by_count_and_bytes`: 32 MiB by default, and a
single larger message goes alone.
- IMAP learns each new message's RFC822.SIZE on the control connection,
in 1000-UID metadata fetches, before the body fetch. A server that won't
say leaves the chunks sized by count. `--fetch-batch-mib` sets the cap.
- EWS asks FindItem for `item:Size` and packs GetItem batches by it,
fetched `--ews-connections` batches at a time. `--ews-getitem-batch-mib`
sets the cap. Items from an incremental SyncFolderItems run carry no
size and stay batched by count.
- docs/usage.md describes both options.
This commit is contained in:
+17
-2
@@ -83,7 +83,7 @@ inbuxa-migrate import imap \
|
|||||||
[--include <REGEX>...] [--exclude <REGEX>...] [--exclude-special <ROLE>...] \
|
[--include <REGEX>...] [--exclude <REGEX>...] [--exclude-special <ROLE>...] \
|
||||||
[--folder <NAME>...] [--subscribed-only] [--noautomap] \
|
[--folder <NAME>...] [--subscribed-only] [--noautomap] \
|
||||||
[--include-deleted] [--allow-cleartext] [--compress] \
|
[--include-deleted] [--allow-cleartext] [--compress] \
|
||||||
[--fetch-batch <N>] [--imap-connections <1..8>] \
|
[--fetch-batch <N>] [--fetch-batch-mib <MIB>] [--imap-connections <1..8>] \
|
||||||
<ARCHIVE>
|
<ARCHIVE>
|
||||||
```
|
```
|
||||||
|
|
||||||
@@ -91,6 +91,14 @@ Imports mail, and only mail, from any IMAP server. Folders are chosen with
|
|||||||
`--include` and `--exclude` patterns, or by exact name with `--folder`, but
|
`--include` and `--exclude` patterns, or by exact name with `--folder`, but
|
||||||
not both. `--exclude-special` drops folders by SPECIAL-USE role.
|
not both. `--exclude-special` drops folders by SPECIAL-USE role.
|
||||||
|
|
||||||
|
Messages are fetched in chunks of at most `--fetch-batch` messages and
|
||||||
|
`--fetch-batch-mib` MiB (32 by default), so a folder of large attachments is
|
||||||
|
fetched a little at a time, like any other. A single message larger than the
|
||||||
|
cap is fetched on its own. Each chunk is written to the archive as it
|
||||||
|
arrives: an interrupted import keeps what it fetched, and the next run picks
|
||||||
|
up from there. A message that can't be imported is reported with its folder
|
||||||
|
and UID, and the rest of the folder carries on.
|
||||||
|
|
||||||
### CalDAV
|
### CalDAV
|
||||||
|
|
||||||
```
|
```
|
||||||
@@ -170,7 +178,8 @@ inbuxa-migrate import exchange-ews \
|
|||||||
(--auth-basic <USER> [--auth-password <PASS>] \
|
(--auth-basic <USER> [--auth-password <PASS>] \
|
||||||
| --auth-bearer [TOKEN] [--ews-tenant <T> --ews-client-id <ID> \
|
| --auth-bearer [TOKEN] [--ews-tenant <T> --ews-client-id <ID> \
|
||||||
(--ews-device-code | --ews-client-secret <SECRET>)]) \
|
(--ews-device-code | --ews-client-secret <SECRET>)]) \
|
||||||
[--ews-connections <1..8>] [--ews-getitem-batch <N>] [--ews-attachment-batch <N>] \
|
[--ews-connections <1..8>] [--ews-getitem-batch <N>] [--ews-getitem-batch-mib <MIB>] \
|
||||||
|
[--ews-attachment-batch <N>] \
|
||||||
[--ews-no-syncfolderitems] \
|
[--ews-no-syncfolderitems] \
|
||||||
<ARCHIVE>
|
<ARCHIVE>
|
||||||
```
|
```
|
||||||
@@ -180,6 +189,12 @@ Imports a mailbox from an on-premises Exchange Server through EWS. Without
|
|||||||
Basic, with a bearer token acquired beforehand, with OAuth's interactive
|
Basic, with a bearer token acquired beforehand, with OAuth's interactive
|
||||||
device-code flow, or with app-only client credentials.
|
device-code flow, or with app-only client credentials.
|
||||||
|
|
||||||
|
Items are fetched in GetItem batches of at most `--ews-getitem-batch` items
|
||||||
|
and `--ews-getitem-batch-mib` MiB (32 by default), `--ews-connections` at a
|
||||||
|
time. The byte cap applies where Exchange reports each item's size, which it
|
||||||
|
does on a full listing; items found through an incremental sync are batched
|
||||||
|
by count.
|
||||||
|
|
||||||
For Exchange Online, use `exchange-graph` instead. Microsoft is retiring EWS in Exchange Online: from October 1, 2026 it is blocked unless a tenant administrator sets `EwsEnabled` to `True` and adds the client id to `EwsAllowedAppIDs`, and on April 1, 2027 it is switched off for every tenant. On-premises Exchange Server is not affected.
|
For Exchange Online, use `exchange-graph` instead. Microsoft is retiring EWS in Exchange Online: from October 1, 2026 it is blocked unless a tenant administrator sets `EwsEnabled` to `True` and adds the client id to `EwsAllowedAppIDs`, and on April 1, 2027 it is switched off for every tenant. On-premises Exchange Server is not affected.
|
||||||
|
|
||||||
### Microsoft Exchange (Graph)
|
### Microsoft Exchange (Graph)
|
||||||
|
|||||||
+21
@@ -277,6 +277,14 @@ struct ImapImportArgs {
|
|||||||
)]
|
)]
|
||||||
fetch_batch: usize,
|
fetch_batch: usize,
|
||||||
|
|
||||||
|
#[arg(
|
||||||
|
long,
|
||||||
|
value_name = "MIB",
|
||||||
|
default_value_t = 32,
|
||||||
|
help = "Most message bytes per body FETCH chunk, in MiB (one larger message goes alone)"
|
||||||
|
)]
|
||||||
|
fetch_batch_mib: u64,
|
||||||
|
|
||||||
#[arg(
|
#[arg(
|
||||||
long,
|
long,
|
||||||
value_name = "N",
|
value_name = "N",
|
||||||
@@ -804,6 +812,7 @@ fn resolve_imap_import(args: ImapImportArgs) -> Result<Action, Error> {
|
|||||||
automap: !args.noautomap,
|
automap: !args.noautomap,
|
||||||
include_deleted: args.include_deleted,
|
include_deleted: args.include_deleted,
|
||||||
fetch_batch: args.fetch_batch,
|
fetch_batch: args.fetch_batch,
|
||||||
|
fetch_batch_bytes: args.fetch_batch_mib.max(1).saturating_mul(1024 * 1024),
|
||||||
imap_connections,
|
imap_connections,
|
||||||
allow_source_change: args.allow_source_change,
|
allow_source_change: args.allow_source_change,
|
||||||
},
|
},
|
||||||
@@ -973,6 +982,14 @@ pub struct ExchangeEwsImportArgs {
|
|||||||
)]
|
)]
|
||||||
ews_getitem_batch: usize,
|
ews_getitem_batch: usize,
|
||||||
|
|
||||||
|
#[arg(
|
||||||
|
long,
|
||||||
|
value_name = "MIB",
|
||||||
|
default_value_t = 32,
|
||||||
|
help = "Most item bytes per GetItem batch, in MiB (one larger item goes alone)"
|
||||||
|
)]
|
||||||
|
ews_getitem_batch_mib: u64,
|
||||||
|
|
||||||
#[arg(
|
#[arg(
|
||||||
long,
|
long,
|
||||||
value_name = "N",
|
value_name = "N",
|
||||||
@@ -1022,6 +1039,10 @@ fn resolve_exchange_ews_import(args: ExchangeEwsImportArgs) -> Result<Action, Er
|
|||||||
auth,
|
auth,
|
||||||
ews_connections,
|
ews_connections,
|
||||||
getitem_batch: args.ews_getitem_batch.max(1),
|
getitem_batch: args.ews_getitem_batch.max(1),
|
||||||
|
getitem_batch_bytes: args
|
||||||
|
.ews_getitem_batch_mib
|
||||||
|
.max(1)
|
||||||
|
.saturating_mul(1024 * 1024),
|
||||||
attachment_batch: args.ews_attachment_batch.max(1),
|
attachment_batch: args.ews_attachment_batch.max(1),
|
||||||
use_syncfolderitems: !args.ews_no_syncfolderitems,
|
use_syncfolderitems: !args.ews_no_syncfolderitems,
|
||||||
allow_source_change: args.allow_source_change,
|
allow_source_change: args.allow_source_change,
|
||||||
|
|||||||
@@ -503,6 +503,8 @@ pub struct FindItemResponse {
|
|||||||
pub struct ItemEntry {
|
pub struct ItemEntry {
|
||||||
pub element: String,
|
pub element: String,
|
||||||
pub id: ItemId,
|
pub id: ItemId,
|
||||||
|
/// `item:Size` in bytes, when the server returned it.
|
||||||
|
pub size: Option<u64>,
|
||||||
}
|
}
|
||||||
|
|
||||||
pub fn parse_find_item_response(body: &[u8]) -> Result<FindItemResponse, EwsError> {
|
pub fn parse_find_item_response(body: &[u8]) -> Result<FindItemResponse, EwsError> {
|
||||||
@@ -511,6 +513,7 @@ pub fn parse_find_item_response(body: &[u8]) -> Result<FindItemResponse, EwsErro
|
|||||||
let mut buf = Vec::new();
|
let mut buf = Vec::new();
|
||||||
let mut out = FindItemResponse::default();
|
let mut out = FindItemResponse::default();
|
||||||
let mut in_root = false;
|
let mut in_root = false;
|
||||||
|
let mut reading_size = false;
|
||||||
loop {
|
loop {
|
||||||
buf.clear();
|
buf.clear();
|
||||||
let (ns, ev) = xml.read_resolved_event_into(&mut buf)?;
|
let (ns, ev) = xml.read_resolved_event_into(&mut buf)?;
|
||||||
@@ -531,7 +534,13 @@ pub fn parse_find_item_response(body: &[u8]) -> Result<FindItemResponse, EwsErro
|
|||||||
out.items.push(ItemEntry {
|
out.items.push(ItemEntry {
|
||||||
element: local.clone(),
|
element: local.clone(),
|
||||||
id: ItemId::default(),
|
id: ItemId::default(),
|
||||||
|
size: None,
|
||||||
});
|
});
|
||||||
|
} else if local.eq_ignore_ascii_case("Size")
|
||||||
|
&& matches!(ev, Event::Start(_))
|
||||||
|
&& !out.items.is_empty()
|
||||||
|
{
|
||||||
|
reading_size = true;
|
||||||
} else if local.eq_ignore_ascii_case("ItemId")
|
} else if local.eq_ignore_ascii_case("ItemId")
|
||||||
&& let Some(last) = out.items.last_mut()
|
&& let Some(last) = out.items.last_mut()
|
||||||
{
|
{
|
||||||
@@ -539,10 +548,18 @@ pub fn parse_find_item_response(body: &[u8]) -> Result<FindItemResponse, EwsErro
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
Event::Text(ref t) if reading_size => {
|
||||||
|
if let Some(last) = out.items.last_mut() {
|
||||||
|
let text: &str = t;
|
||||||
|
last.size = text.trim().parse().ok();
|
||||||
|
}
|
||||||
|
}
|
||||||
Event::End(e) => {
|
Event::End(e) => {
|
||||||
let local = e.local_name().as_ref().to_owned();
|
let local = e.local_name().as_ref().to_owned();
|
||||||
if local.eq_ignore_ascii_case("RootFolder") {
|
if local.eq_ignore_ascii_case("RootFolder") {
|
||||||
in_root = false;
|
in_root = false;
|
||||||
|
} else if local.eq_ignore_ascii_case("Size") {
|
||||||
|
reading_size = false;
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
Event::Eof => break,
|
Event::Eof => break,
|
||||||
|
|||||||
+16
-1
@@ -1,5 +1,6 @@
|
|||||||
/*
|
/*
|
||||||
* SPDX-FileCopyrightText: 2020 Stalwart Labs LLC <[email protected]>
|
* SPDX-FileCopyrightText: 2020 Stalwart Labs LLC <[email protected]>
|
||||||
|
* SPDX-FileCopyrightText: 2026 John Coffey <[email protected]>
|
||||||
*
|
*
|
||||||
* SPDX-License-Identifier: Apache-2.0 OR MIT
|
* SPDX-License-Identifier: Apache-2.0 OR MIT
|
||||||
*/
|
*/
|
||||||
@@ -83,7 +84,12 @@ pub fn find_item_body(
|
|||||||
let mut out = String::with_capacity(512);
|
let mut out = String::with_capacity(512);
|
||||||
out.push_str("<m:FindItem Traversal=\"");
|
out.push_str("<m:FindItem Traversal=\"");
|
||||||
out.push_str(traversal.as_str());
|
out.push_str(traversal.as_str());
|
||||||
out.push_str("\"><m:ItemShape><t:BaseShape>IdOnly</t:BaseShape></m:ItemShape>");
|
// item:Size lets GetItem batches be split by bytes as well as by count.
|
||||||
|
out.push_str(
|
||||||
|
"\"><m:ItemShape><t:BaseShape>IdOnly</t:BaseShape>\
|
||||||
|
<t:AdditionalProperties><t:FieldURI FieldURI=\"item:Size\"/></t:AdditionalProperties>\
|
||||||
|
</m:ItemShape>",
|
||||||
|
);
|
||||||
out.push_str("<m:IndexedPageItemView MaxEntriesReturned=\"");
|
out.push_str("<m:IndexedPageItemView MaxEntriesReturned=\"");
|
||||||
out.push_str(&page_size.to_string());
|
out.push_str(&page_size.to_string());
|
||||||
out.push_str("\" Offset=\"");
|
out.push_str("\" Offset=\"");
|
||||||
@@ -290,6 +296,15 @@ mod tests {
|
|||||||
assert!(body.contains("<t:DistinguishedFolderId Id=\"archiveroot\"/>"));
|
assert!(body.contains("<t:DistinguishedFolderId Id=\"archiveroot\"/>"));
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn find_item_asks_for_item_size() {
|
||||||
|
let folder = FolderId::new("FID", "FCK");
|
||||||
|
let body = find_item_body(FolderRef::Concrete(&folder), Traversal::Shallow, 0, 50);
|
||||||
|
assert!(body.contains(
|
||||||
|
"<t:BaseShape>IdOnly</t:BaseShape><t:AdditionalProperties><t:FieldURI FieldURI=\"item:Size\"/></t:AdditionalProperties></m:ItemShape>"
|
||||||
|
));
|
||||||
|
}
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
fn find_item_paginates_with_offset_and_page_size() {
|
fn find_item_paginates_with_offset_and_page_size() {
|
||||||
let folder = FolderId::new("FID", "FCK");
|
let folder = FolderId::new("FID", "FCK");
|
||||||
|
|||||||
@@ -0,0 +1,92 @@
|
|||||||
|
/*
|
||||||
|
* SPDX-FileCopyrightText: 2026 John Coffey <[email protected]>
|
||||||
|
*
|
||||||
|
* SPDX-License-Identifier: Apache-2.0 OR MIT
|
||||||
|
*/
|
||||||
|
|
||||||
|
//! Fetch batches bounded by bytes as well as by count. A batch of a few
|
||||||
|
//! hundred messages is small for ordinary mail and many gigabytes for a
|
||||||
|
//! mailbox of large attachments; capping the bytes too keeps what a batch
|
||||||
|
//! holds in memory about the same whatever the mail is like.
|
||||||
|
|
||||||
|
/// The default byte cap for one fetch batch.
|
||||||
|
pub const DEFAULT_BATCH_BYTES: u64 = 32 * 1024 * 1024;
|
||||||
|
|
||||||
|
/// Splits `items` into contiguous batches of at most `max_count` items and
|
||||||
|
/// at most `max_bytes` bytes, as reported by `size`. An item whose size is
|
||||||
|
/// unknown counts as 0 bytes, so without sizes this is batching by count. An
|
||||||
|
/// item larger than `max_bytes` goes in a batch of its own: every batch has
|
||||||
|
/// at least one item.
|
||||||
|
pub fn by_count_and_bytes<T>(
|
||||||
|
items: &[T],
|
||||||
|
size: impl Fn(&T) -> u64,
|
||||||
|
max_count: usize,
|
||||||
|
max_bytes: u64,
|
||||||
|
) -> Vec<&[T]> {
|
||||||
|
let max_count = max_count.max(1);
|
||||||
|
let max_bytes = max_bytes.max(1);
|
||||||
|
let mut out = Vec::new();
|
||||||
|
let mut start = 0usize;
|
||||||
|
let mut bytes = 0u64;
|
||||||
|
for (i, item) in items.iter().enumerate() {
|
||||||
|
let s = size(item);
|
||||||
|
let count = i - start;
|
||||||
|
if count > 0 && (count >= max_count || bytes.saturating_add(s) > max_bytes) {
|
||||||
|
out.push(&items[start..i]);
|
||||||
|
start = i;
|
||||||
|
bytes = 0;
|
||||||
|
}
|
||||||
|
bytes = bytes.saturating_add(s);
|
||||||
|
}
|
||||||
|
if start < items.len() {
|
||||||
|
out.push(&items[start..]);
|
||||||
|
}
|
||||||
|
out
|
||||||
|
}
|
||||||
|
|
||||||
|
#[cfg(test)]
|
||||||
|
mod tests {
|
||||||
|
use super::*;
|
||||||
|
|
||||||
|
fn batches(sizes: &[u64], max_count: usize, max_bytes: u64) -> Vec<Vec<u64>> {
|
||||||
|
by_count_and_bytes(sizes, |s| *s, max_count, max_bytes)
|
||||||
|
.into_iter()
|
||||||
|
.map(|b| b.to_vec())
|
||||||
|
.collect()
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn respects_the_byte_cap() {
|
||||||
|
assert_eq!(
|
||||||
|
batches(&[10, 10, 10, 10, 10], 100, 25),
|
||||||
|
vec![vec![10, 10], vec![10, 10], vec![10]]
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn respects_the_count_cap() {
|
||||||
|
assert_eq!(
|
||||||
|
batches(&[1, 1, 1, 1, 1], 2, 1000),
|
||||||
|
vec![vec![1, 1], vec![1, 1], vec![1]]
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn an_item_over_the_cap_goes_alone() {
|
||||||
|
assert_eq!(
|
||||||
|
batches(&[5, 500, 5, 5], 100, 20),
|
||||||
|
vec![vec![5], vec![500], vec![5, 5]]
|
||||||
|
);
|
||||||
|
assert_eq!(batches(&[500], 100, 20), vec![vec![500]]);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn unknown_sizes_batch_by_count() {
|
||||||
|
assert_eq!(batches(&[0, 0, 0], 2, 1), vec![vec![0, 0], vec![0]]);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn nothing_in_nothing_out() {
|
||||||
|
assert!(batches(&[], 10, 10).is_empty());
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -1,5 +1,6 @@
|
|||||||
/*
|
/*
|
||||||
* SPDX-FileCopyrightText: 2020 Stalwart Labs LLC <[email protected]>
|
* SPDX-FileCopyrightText: 2020 Stalwart Labs LLC <[email protected]>
|
||||||
|
* SPDX-FileCopyrightText: 2026 John Coffey <[email protected]>
|
||||||
*
|
*
|
||||||
* SPDX-License-Identifier: Apache-2.0 OR MIT
|
* SPDX-License-Identifier: Apache-2.0 OR MIT
|
||||||
*/
|
*/
|
||||||
@@ -104,47 +105,53 @@ fn reconcile_one(
|
|||||||
to_fetch.push(id.clone());
|
to_fetch.push(id.clone());
|
||||||
}
|
}
|
||||||
if !to_fetch.is_empty() {
|
if !to_fetch.is_empty() {
|
||||||
let failed_items = for_each_fetched_item(ctx, ItemShape::CalendarItem, &to_fetch, |msg| {
|
let failed_items = for_each_fetched_item(
|
||||||
if !msg.success {
|
ctx,
|
||||||
|
ItemShape::CalendarItem,
|
||||||
|
&to_fetch,
|
||||||
|
&outcome.sizes,
|
||||||
|
|msg| {
|
||||||
|
if !msg.success {
|
||||||
|
if matches!(
|
||||||
|
msg.response_code,
|
||||||
|
crate::exchange_ews::types::ResponseCode::ItemNotFound
|
||||||
|
) {
|
||||||
|
counts.skipped += 1;
|
||||||
|
} else {
|
||||||
|
counts.failed += 1;
|
||||||
|
ctx.logger
|
||||||
|
.warn(&format!("GetItem (calendar) error: {}", msg.response_code));
|
||||||
|
}
|
||||||
|
return Ok(());
|
||||||
|
}
|
||||||
|
let parsed = parse_calendar_item(&msg.inner_xml).map_err(Error::from)?;
|
||||||
|
if parsed.id.id.is_empty() {
|
||||||
|
counts.failed += 1;
|
||||||
|
return Ok(());
|
||||||
|
}
|
||||||
if matches!(
|
if matches!(
|
||||||
msg.response_code,
|
parsed.calendar_item_type,
|
||||||
crate::exchange_ews::types::ResponseCode::ItemNotFound
|
Some(CalendarItemType::Occurrence) | Some(CalendarItemType::Exception)
|
||||||
) {
|
) {
|
||||||
counts.skipped += 1;
|
counts.skipped += 1;
|
||||||
} else {
|
return Ok(());
|
||||||
counts.failed += 1;
|
|
||||||
ctx.logger
|
|
||||||
.warn(&format!("GetItem (calendar) error: {}", msg.response_code));
|
|
||||||
}
|
}
|
||||||
return Ok(());
|
let existing = plan
|
||||||
}
|
.present_changed
|
||||||
let parsed = parse_calendar_item(&msg.inner_xml).map_err(Error::from)?;
|
.iter()
|
||||||
if parsed.id.id.is_empty() {
|
.find(|(id, _)| id.id == parsed.id.id)
|
||||||
counts.failed += 1;
|
.map(|(_, local)| *local);
|
||||||
return Ok(());
|
apply_event(
|
||||||
}
|
conn,
|
||||||
if matches!(
|
ctx,
|
||||||
parsed.calendar_item_type,
|
&parsed,
|
||||||
Some(CalendarItemType::Occurrence) | Some(CalendarItemType::Exception)
|
local_folder_id,
|
||||||
) {
|
&folder.id,
|
||||||
counts.skipped += 1;
|
existing,
|
||||||
return Ok(());
|
counts,
|
||||||
}
|
)
|
||||||
let existing = plan
|
},
|
||||||
.present_changed
|
)?;
|
||||||
.iter()
|
|
||||||
.find(|(id, _)| id.id == parsed.id.id)
|
|
||||||
.map(|(_, local)| *local);
|
|
||||||
apply_event(
|
|
||||||
conn,
|
|
||||||
ctx,
|
|
||||||
&parsed,
|
|
||||||
local_folder_id,
|
|
||||||
&folder.id,
|
|
||||||
existing,
|
|
||||||
counts,
|
|
||||||
)
|
|
||||||
})?;
|
|
||||||
counts.failed += failed_items;
|
counts.failed += failed_items;
|
||||||
}
|
}
|
||||||
delete_vanished(
|
delete_vanished(
|
||||||
|
|||||||
@@ -1,5 +1,6 @@
|
|||||||
/*
|
/*
|
||||||
* SPDX-FileCopyrightText: 2020 Stalwart Labs LLC <[email protected]>
|
* SPDX-FileCopyrightText: 2020 Stalwart Labs LLC <[email protected]>
|
||||||
|
* SPDX-FileCopyrightText: 2026 John Coffey <[email protected]>
|
||||||
*
|
*
|
||||||
* SPDX-License-Identifier: Apache-2.0 OR MIT
|
* SPDX-License-Identifier: Apache-2.0 OR MIT
|
||||||
*/
|
*/
|
||||||
@@ -96,40 +97,41 @@ fn reconcile_one(
|
|||||||
to_fetch.push(id.clone());
|
to_fetch.push(id.clone());
|
||||||
}
|
}
|
||||||
if !to_fetch.is_empty() {
|
if !to_fetch.is_empty() {
|
||||||
let failed_items = for_each_fetched_item(ctx, ItemShape::Contact, &to_fetch, |msg| {
|
let failed_items =
|
||||||
if !msg.success {
|
for_each_fetched_item(ctx, ItemShape::Contact, &to_fetch, &outcome.sizes, |msg| {
|
||||||
if matches!(
|
if !msg.success {
|
||||||
msg.response_code,
|
if matches!(
|
||||||
crate::exchange_ews::types::ResponseCode::ItemNotFound
|
msg.response_code,
|
||||||
) {
|
crate::exchange_ews::types::ResponseCode::ItemNotFound
|
||||||
counts.skipped += 1;
|
) {
|
||||||
} else {
|
counts.skipped += 1;
|
||||||
counts.failed += 1;
|
} else {
|
||||||
ctx.logger
|
counts.failed += 1;
|
||||||
.warn(&format!("GetItem (contact) error: {}", msg.response_code));
|
ctx.logger
|
||||||
|
.warn(&format!("GetItem (contact) error: {}", msg.response_code));
|
||||||
|
}
|
||||||
|
return Ok(());
|
||||||
}
|
}
|
||||||
return Ok(());
|
let parsed = parse_contact_item(&msg.inner_xml).map_err(Error::from)?;
|
||||||
}
|
if parsed.id.id.is_empty() {
|
||||||
let parsed = parse_contact_item(&msg.inner_xml).map_err(Error::from)?;
|
counts.failed += 1;
|
||||||
if parsed.id.id.is_empty() {
|
return Ok(());
|
||||||
counts.failed += 1;
|
}
|
||||||
return Ok(());
|
let existing = plan
|
||||||
}
|
.present_changed
|
||||||
let existing = plan
|
.iter()
|
||||||
.present_changed
|
.find(|(id, _)| id.id == parsed.id.id)
|
||||||
.iter()
|
.map(|(_, local)| *local);
|
||||||
.find(|(id, _)| id.id == parsed.id.id)
|
apply_contact(
|
||||||
.map(|(_, local)| *local);
|
conn,
|
||||||
apply_contact(
|
ctx,
|
||||||
conn,
|
&parsed,
|
||||||
ctx,
|
local_folder_id,
|
||||||
&parsed,
|
&folder.id,
|
||||||
local_folder_id,
|
existing,
|
||||||
&folder.id,
|
counts,
|
||||||
existing,
|
)
|
||||||
counts,
|
})?;
|
||||||
)
|
|
||||||
})?;
|
|
||||||
counts.failed += failed_items;
|
counts.failed += failed_items;
|
||||||
}
|
}
|
||||||
delete_vanished(
|
delete_vanished(
|
||||||
|
|||||||
@@ -37,6 +37,8 @@ pub struct EwsImportConfig {
|
|||||||
pub auth: EwsAuth,
|
pub auth: EwsAuth,
|
||||||
pub ews_connections: usize,
|
pub ews_connections: usize,
|
||||||
pub getitem_batch: usize,
|
pub getitem_batch: usize,
|
||||||
|
/// Byte cap for one GetItem batch; see `sync::batch`.
|
||||||
|
pub getitem_batch_bytes: u64,
|
||||||
pub attachment_batch: usize,
|
pub attachment_batch: usize,
|
||||||
pub use_syncfolderitems: bool,
|
pub use_syncfolderitems: bool,
|
||||||
pub allow_source_change: bool,
|
pub allow_source_change: bool,
|
||||||
@@ -147,6 +149,7 @@ pub fn run(common: CommonConfig, config: EwsImportConfig) -> Result<Summary, Err
|
|||||||
url: &session_url,
|
url: &session_url,
|
||||||
source_id,
|
source_id,
|
||||||
batch_size: config.getitem_batch.max(1),
|
batch_size: config.getitem_batch.max(1),
|
||||||
|
batch_bytes: config.getitem_batch_bytes.max(1),
|
||||||
attachment_batch: config.attachment_batch.max(1),
|
attachment_batch: config.attachment_batch.max(1),
|
||||||
connections: config.ews_connections.clamp(1, 8),
|
connections: config.ews_connections.clamp(1, 8),
|
||||||
use_syncfolderitems: config.use_syncfolderitems,
|
use_syncfolderitems: config.use_syncfolderitems,
|
||||||
|
|||||||
@@ -1,5 +1,6 @@
|
|||||||
/*
|
/*
|
||||||
* SPDX-FileCopyrightText: 2020 Stalwart Labs LLC <[email protected]>
|
* SPDX-FileCopyrightText: 2020 Stalwart Labs LLC <[email protected]>
|
||||||
|
* SPDX-FileCopyrightText: 2026 John Coffey <[email protected]>
|
||||||
*
|
*
|
||||||
* SPDX-License-Identifier: Apache-2.0 OR MIT
|
* SPDX-License-Identifier: Apache-2.0 OR MIT
|
||||||
*/
|
*/
|
||||||
@@ -23,6 +24,8 @@ pub struct ItemRunCtx<'a> {
|
|||||||
pub url: &'a str,
|
pub url: &'a str,
|
||||||
pub source_id: i64,
|
pub source_id: i64,
|
||||||
pub batch_size: usize,
|
pub batch_size: usize,
|
||||||
|
/// Byte cap for one GetItem batch, by `item:Size`; see `sync::batch`.
|
||||||
|
pub batch_bytes: u64,
|
||||||
pub attachment_batch: usize,
|
pub attachment_batch: usize,
|
||||||
pub connections: usize,
|
pub connections: usize,
|
||||||
pub use_syncfolderitems: bool,
|
pub use_syncfolderitems: bool,
|
||||||
@@ -47,6 +50,9 @@ pub struct EnumeratedItem {
|
|||||||
pub struct EnumerationOutcome {
|
pub struct EnumerationOutcome {
|
||||||
pub items: Vec<EnumeratedItem>,
|
pub items: Vec<EnumeratedItem>,
|
||||||
pub mode: EnumerationMode,
|
pub mode: EnumerationMode,
|
||||||
|
/// `item:Size` by item id, where the listing gave it. FindItem does;
|
||||||
|
/// SyncFolderItems changes do not, and are batched by count.
|
||||||
|
pub sizes: HashMap<String, u64>,
|
||||||
}
|
}
|
||||||
|
|
||||||
#[derive(Debug, Clone)]
|
#[derive(Debug, Clone)]
|
||||||
@@ -70,9 +76,10 @@ pub fn enumerate_folder(
|
|||||||
{
|
{
|
||||||
return Ok(outcome);
|
return Ok(outcome);
|
||||||
}
|
}
|
||||||
enumerate_via_find_item(ctx, folder).map(|items| EnumerationOutcome {
|
enumerate_via_find_item(ctx, folder).map(|(items, sizes)| EnumerationOutcome {
|
||||||
items,
|
items,
|
||||||
mode: EnumerationMode::Full,
|
mode: EnumerationMode::Full,
|
||||||
|
sizes,
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -151,14 +158,16 @@ fn try_sync_folder_items(
|
|||||||
deletions,
|
deletions,
|
||||||
new_sync_state: sync_state,
|
new_sync_state: sync_state,
|
||||||
},
|
},
|
||||||
|
sizes: HashMap::new(),
|
||||||
}))
|
}))
|
||||||
}
|
}
|
||||||
|
|
||||||
fn enumerate_via_find_item(
|
fn enumerate_via_find_item(
|
||||||
ctx: &ItemRunCtx<'_>,
|
ctx: &ItemRunCtx<'_>,
|
||||||
folder: &FolderId,
|
folder: &FolderId,
|
||||||
) -> Result<Vec<EnumeratedItem>, EwsError> {
|
) -> Result<(Vec<EnumeratedItem>, HashMap<String, u64>), EwsError> {
|
||||||
let mut items: Vec<EnumeratedItem> = Vec::new();
|
let mut items: Vec<EnumeratedItem> = Vec::new();
|
||||||
|
let mut sizes: HashMap<String, u64> = HashMap::new();
|
||||||
let mut offset: u32 = 0;
|
let mut offset: u32 = 0;
|
||||||
let page_size: u32 = 500;
|
let page_size: u32 = 500;
|
||||||
loop {
|
loop {
|
||||||
@@ -172,6 +181,9 @@ fn enumerate_via_find_item(
|
|||||||
let parsed = parse_find_item_response(&resp.body)?;
|
let parsed = parse_find_item_response(&resp.body)?;
|
||||||
let returned = parsed.items.len() as u32;
|
let returned = parsed.items.len() as u32;
|
||||||
for entry in parsed.items {
|
for entry in parsed.items {
|
||||||
|
if let Some(size) = entry.size {
|
||||||
|
sizes.insert(entry.id.id.clone(), size);
|
||||||
|
}
|
||||||
items.push(EnumeratedItem {
|
items.push(EnumeratedItem {
|
||||||
element: entry.element,
|
element: entry.element,
|
||||||
id: entry.id,
|
id: entry.id,
|
||||||
@@ -188,7 +200,7 @@ fn enumerate_via_find_item(
|
|||||||
break;
|
break;
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
Ok(items)
|
Ok((items, sizes))
|
||||||
}
|
}
|
||||||
|
|
||||||
#[derive(Debug, Clone)]
|
#[derive(Debug, Clone)]
|
||||||
@@ -285,13 +297,22 @@ pub fn get_items(
|
|||||||
shape: ItemShape,
|
shape: ItemShape,
|
||||||
ids: &[ItemId],
|
ids: &[ItemId],
|
||||||
) -> Result<GetItemBatchOutcome, EwsError> {
|
) -> Result<GetItemBatchOutcome, EwsError> {
|
||||||
let batch = ctx.batch_size.max(1);
|
let batches: Vec<&[ItemId]> = ids.chunks(ctx.batch_size.max(1)).collect();
|
||||||
|
get_item_batches(ctx, shape, &batches)
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Runs one GetItem per batch, over up to `connections` connections.
|
||||||
|
pub fn get_item_batches(
|
||||||
|
ctx: &ItemRunCtx<'_>,
|
||||||
|
shape: ItemShape,
|
||||||
|
batches: &[&[ItemId]],
|
||||||
|
) -> Result<GetItemBatchOutcome, EwsError> {
|
||||||
let workers = ctx.connections.clamp(1, 8);
|
let workers = ctx.connections.clamp(1, 8);
|
||||||
let version = ctx.client.server_version();
|
let version = ctx.client.server_version();
|
||||||
let mut failed_items: u64 = 0;
|
let mut failed_items: u64 = 0;
|
||||||
if workers <= 1 || ids.len() <= batch {
|
if workers <= 1 || batches.len() <= 1 {
|
||||||
let mut all = Vec::new();
|
let mut all = Vec::new();
|
||||||
for chunk in ids.chunks(batch) {
|
for chunk in batches.iter().copied() {
|
||||||
let body = get_item_body(shape, chunk, version);
|
let body = get_item_body(shape, chunk, version);
|
||||||
match ctx.client.call(ctx.url, "GetItem", &body) {
|
match ctx.client.call(ctx.url, "GetItem", &body) {
|
||||||
Ok(resp) => match parse_response_messages(&resp.body, "GetItemResponseMessage") {
|
Ok(resp) => match parse_response_messages(&resp.body, "GetItemResponseMessage") {
|
||||||
@@ -337,7 +358,7 @@ pub fn get_items(
|
|||||||
(n, result)
|
(n, result)
|
||||||
});
|
});
|
||||||
let mut submitted = 0usize;
|
let mut submitted = 0usize;
|
||||||
for chunk in ids.chunks(batch) {
|
for chunk in batches {
|
||||||
pool.submit(chunk.to_vec());
|
pool.submit(chunk.to_vec());
|
||||||
submitted += 1;
|
submitted += 1;
|
||||||
}
|
}
|
||||||
@@ -368,21 +389,31 @@ pub fn get_items(
|
|||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Fetches `ids` in GetItem batches of at most `batch_size` items and, where
|
||||||
|
/// `sizes` knows them, at most `batch_bytes` bytes, and hands each message
|
||||||
|
/// on as it is parsed. Batches are fetched `connections` at a time, so what
|
||||||
|
/// is held at once is about `batch_bytes` per connection, however large the
|
||||||
|
/// mail is.
|
||||||
pub fn for_each_fetched_item<F>(
|
pub fn for_each_fetched_item<F>(
|
||||||
ctx: &ItemRunCtx<'_>,
|
ctx: &ItemRunCtx<'_>,
|
||||||
shape: ItemShape,
|
shape: ItemShape,
|
||||||
ids: &[ItemId],
|
ids: &[ItemId],
|
||||||
|
sizes: &HashMap<String, u64>,
|
||||||
mut on_message: F,
|
mut on_message: F,
|
||||||
) -> Result<u64, Error>
|
) -> Result<u64, Error>
|
||||||
where
|
where
|
||||||
F: FnMut(crate::exchange_ews::parse::ResponseMessage) -> Result<(), Error>,
|
F: FnMut(crate::exchange_ews::parse::ResponseMessage) -> Result<(), Error>,
|
||||||
{
|
{
|
||||||
let batch = ctx.batch_size.max(1);
|
|
||||||
let workers = ctx.connections.clamp(1, 8);
|
let workers = ctx.connections.clamp(1, 8);
|
||||||
let window = batch.saturating_mul(workers).max(batch);
|
let batches = crate::sync::batch::by_count_and_bytes(
|
||||||
|
ids,
|
||||||
|
|id| sizes.get(&id.id).copied().unwrap_or(0),
|
||||||
|
ctx.batch_size,
|
||||||
|
ctx.batch_bytes,
|
||||||
|
);
|
||||||
let mut failed_items = 0u64;
|
let mut failed_items = 0u64;
|
||||||
for win in ids.chunks(window) {
|
for win in batches.chunks(workers) {
|
||||||
let outcome = get_items(ctx, shape, win).map_err(Error::from)?;
|
let outcome = get_item_batches(ctx, shape, win).map_err(Error::from)?;
|
||||||
failed_items = failed_items.saturating_add(outcome.failed_items);
|
failed_items = failed_items.saturating_add(outcome.failed_items);
|
||||||
for msg in outcome.messages {
|
for msg in outcome.messages {
|
||||||
on_message(msg)?;
|
on_message(msg)?;
|
||||||
@@ -444,6 +475,7 @@ mod tests {
|
|||||||
element: "Message".to_owned(),
|
element: "Message".to_owned(),
|
||||||
id: ItemId::new("A", "ck-2"),
|
id: ItemId::new("A", "ck-2"),
|
||||||
}],
|
}],
|
||||||
|
sizes: HashMap::new(),
|
||||||
mode: EnumerationMode::Delta {
|
mode: EnumerationMode::Delta {
|
||||||
deletions: vec!["Z".to_owned()],
|
deletions: vec!["Z".to_owned()],
|
||||||
new_sync_state: "STATE2".to_owned(),
|
new_sync_state: "STATE2".to_owned(),
|
||||||
|
|||||||
@@ -1,5 +1,6 @@
|
|||||||
/*
|
/*
|
||||||
* SPDX-FileCopyrightText: 2020 Stalwart Labs LLC <[email protected]>
|
* SPDX-FileCopyrightText: 2020 Stalwart Labs LLC <[email protected]>
|
||||||
|
* SPDX-FileCopyrightText: 2026 John Coffey <[email protected]>
|
||||||
*
|
*
|
||||||
* SPDX-License-Identifier: Apache-2.0 OR MIT
|
* SPDX-License-Identifier: Apache-2.0 OR MIT
|
||||||
*/
|
*/
|
||||||
@@ -95,42 +96,43 @@ fn reconcile_one_folder(
|
|||||||
to_fetch.push(id.clone());
|
to_fetch.push(id.clone());
|
||||||
}
|
}
|
||||||
if !to_fetch.is_empty() {
|
if !to_fetch.is_empty() {
|
||||||
let failed_items = for_each_fetched_item(ctx, ItemShape::Message, &to_fetch, |msg| {
|
let failed_items =
|
||||||
if !msg.success {
|
for_each_fetched_item(ctx, ItemShape::Message, &to_fetch, &outcome.sizes, |msg| {
|
||||||
if matches!(
|
if !msg.success {
|
||||||
msg.response_code,
|
if matches!(
|
||||||
crate::exchange_ews::types::ResponseCode::ItemNotFound
|
msg.response_code,
|
||||||
) {
|
crate::exchange_ews::types::ResponseCode::ItemNotFound
|
||||||
counts.skipped += 1;
|
) {
|
||||||
} else {
|
counts.skipped += 1;
|
||||||
counts.failed += 1;
|
} else {
|
||||||
ctx.logger.warn(&format!(
|
counts.failed += 1;
|
||||||
"GetItem (message) error: {} {}",
|
ctx.logger.warn(&format!(
|
||||||
msg.response_code, msg.message_text
|
"GetItem (message) error: {} {}",
|
||||||
));
|
msg.response_code, msg.message_text
|
||||||
|
));
|
||||||
|
}
|
||||||
|
return Ok(());
|
||||||
}
|
}
|
||||||
return Ok(());
|
let parsed = parse_message_item(&msg.inner_xml).map_err(Error::from)?;
|
||||||
}
|
if parsed.id.id.is_empty() {
|
||||||
let parsed = parse_message_item(&msg.inner_xml).map_err(Error::from)?;
|
counts.failed += 1;
|
||||||
if parsed.id.id.is_empty() {
|
return Ok(());
|
||||||
counts.failed += 1;
|
}
|
||||||
return Ok(());
|
let existing = plan
|
||||||
}
|
.present_changed
|
||||||
let existing = plan
|
.iter()
|
||||||
.present_changed
|
.find(|(id, _)| id.id == parsed.id.id)
|
||||||
.iter()
|
.map(|(_, local)| *local);
|
||||||
.find(|(id, _)| id.id == parsed.id.id)
|
apply_message(
|
||||||
.map(|(_, local)| *local);
|
conn,
|
||||||
apply_message(
|
ctx,
|
||||||
conn,
|
&parsed,
|
||||||
ctx,
|
local_folder_id,
|
||||||
&parsed,
|
&folder.id,
|
||||||
local_folder_id,
|
existing,
|
||||||
&folder.id,
|
counts,
|
||||||
existing,
|
)
|
||||||
counts,
|
})?;
|
||||||
)
|
|
||||||
})?;
|
|
||||||
counts.failed += failed_items;
|
counts.failed += failed_items;
|
||||||
}
|
}
|
||||||
delete_vanished(
|
delete_vanished(
|
||||||
|
|||||||
@@ -50,6 +50,7 @@ pub(super) struct RunOpts {
|
|||||||
/// `WorkerPool::cancel_before`.
|
/// `WorkerPool::cancel_before`.
|
||||||
generation: u64,
|
generation: u64,
|
||||||
fetch_batch: usize,
|
fetch_batch: usize,
|
||||||
|
fetch_batch_bytes: u64,
|
||||||
include_deleted: bool,
|
include_deleted: bool,
|
||||||
logger: Logger,
|
logger: Logger,
|
||||||
}
|
}
|
||||||
@@ -199,6 +200,8 @@ pub struct ImapImportConfig {
|
|||||||
pub automap: bool,
|
pub automap: bool,
|
||||||
pub include_deleted: bool,
|
pub include_deleted: bool,
|
||||||
pub fetch_batch: usize,
|
pub fetch_batch: usize,
|
||||||
|
/// Byte cap for one fetch chunk, by RFC822.SIZE; see `sync::batch`.
|
||||||
|
pub fetch_batch_bytes: u64,
|
||||||
pub imap_connections: usize,
|
pub imap_connections: usize,
|
||||||
pub allow_source_change: bool,
|
pub allow_source_change: bool,
|
||||||
}
|
}
|
||||||
@@ -404,6 +407,7 @@ fn run_into(
|
|||||||
source_id,
|
source_id,
|
||||||
generation: 0,
|
generation: 0,
|
||||||
fetch_batch: config.fetch_batch.max(1),
|
fetch_batch: config.fetch_batch.max(1),
|
||||||
|
fetch_batch_bytes: config.fetch_batch_bytes.max(1),
|
||||||
include_deleted: config.include_deleted,
|
include_deleted: config.include_deleted,
|
||||||
logger,
|
logger,
|
||||||
};
|
};
|
||||||
@@ -706,6 +710,7 @@ fn reconcile_folder(
|
|||||||
source_id,
|
source_id,
|
||||||
generation,
|
generation,
|
||||||
fetch_batch,
|
fetch_batch,
|
||||||
|
fetch_batch_bytes: _,
|
||||||
include_deleted: _,
|
include_deleted: _,
|
||||||
logger,
|
logger,
|
||||||
} = opts;
|
} = opts;
|
||||||
@@ -811,7 +816,13 @@ fn reconcile_folder(
|
|||||||
folder.name
|
folder.name
|
||||||
))
|
))
|
||||||
})?;
|
})?;
|
||||||
let batches: Vec<&[u32]> = chunks(&diff.new, fetch_batch);
|
let sizes = fetch_sizes(client, control_ctx, &folder.name, &diff.new, logger);
|
||||||
|
let batches: Vec<&[u32]> = crate::sync::batch::by_count_and_bytes(
|
||||||
|
&diff.new,
|
||||||
|
|uid| sizes.get(uid).copied().unwrap_or(0),
|
||||||
|
fetch_batch,
|
||||||
|
opts.fetch_batch_bytes,
|
||||||
|
);
|
||||||
let n_batches = batches.len();
|
let n_batches = batches.len();
|
||||||
for batch in &batches {
|
for batch in &batches {
|
||||||
pool.submit(FetchJob {
|
pool.submit(FetchJob {
|
||||||
@@ -892,6 +903,47 @@ fn reconcile_folder(
|
|||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// RFC822.SIZE for each of `uids`, fetched on the control connection in
|
||||||
|
/// large metadata-only chunks, so the body fetch can be split by bytes.
|
||||||
|
/// Best effort: a server that won't answer just leaves the chunks sized by
|
||||||
|
/// count.
|
||||||
|
fn fetch_sizes(
|
||||||
|
client: &mut ImapClient,
|
||||||
|
ctx: &ControlCtx,
|
||||||
|
folder: &str,
|
||||||
|
uids: &[u32],
|
||||||
|
logger: Logger,
|
||||||
|
) -> HashMap<u32, u64> {
|
||||||
|
const SIZE_CHUNK: usize = 1000;
|
||||||
|
let mut sizes = HashMap::with_capacity(uids.len());
|
||||||
|
for chunk in chunks(uids, SIZE_CHUNK) {
|
||||||
|
let set = command::format_uid_set(chunk, true);
|
||||||
|
let cmd = command::uid_fetch(&set, &["UID", "RFC822.SIZE"]);
|
||||||
|
match call_with_retry(client, ctx, |c| c.run_collect(&cmd)) {
|
||||||
|
Ok(resp) => {
|
||||||
|
for u in &resp.untagged {
|
||||||
|
if let Some(attrs) = fetch::extract(u)
|
||||||
|
&& let (Some(uid), Some(size)) = (attrs.uid, attrs.size)
|
||||||
|
{
|
||||||
|
sizes.insert(uid, size);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
Err(e) => {
|
||||||
|
log_at(
|
||||||
|
logger,
|
||||||
|
LEVEL_DEFAULT,
|
||||||
|
&format!(
|
||||||
|
"folder {folder:?}: message sizes unavailable ({e}); fetching in chunks by count only"
|
||||||
|
),
|
||||||
|
);
|
||||||
|
break;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
sizes
|
||||||
|
}
|
||||||
|
|
||||||
fn wipe_folder_emails(
|
fn wipe_folder_emails(
|
||||||
conn: &mut Connection,
|
conn: &mut Connection,
|
||||||
source_id: i64,
|
source_id: i64,
|
||||||
@@ -1026,6 +1078,7 @@ fn insert_single_message(
|
|||||||
source_id,
|
source_id,
|
||||||
generation: _,
|
generation: _,
|
||||||
fetch_batch: _,
|
fetch_batch: _,
|
||||||
|
fetch_batch_bytes: _,
|
||||||
include_deleted,
|
include_deleted,
|
||||||
logger,
|
logger,
|
||||||
} = opts;
|
} = opts;
|
||||||
@@ -1104,6 +1157,7 @@ fn refresh_present_flags(
|
|||||||
source_id,
|
source_id,
|
||||||
generation: _,
|
generation: _,
|
||||||
fetch_batch,
|
fetch_batch,
|
||||||
|
fetch_batch_bytes: _,
|
||||||
include_deleted,
|
include_deleted,
|
||||||
logger: _,
|
logger: _,
|
||||||
} = opts;
|
} = opts;
|
||||||
|
|||||||
@@ -10,7 +10,7 @@ use std::sync::Arc;
|
|||||||
use std::sync::atomic::{AtomicU64, Ordering};
|
use std::sync::atomic::{AtomicU64, Ordering};
|
||||||
use std::thread;
|
use std::thread;
|
||||||
|
|
||||||
use crossbeam_channel::{Receiver, Sender, unbounded};
|
use crossbeam_channel::{Receiver, Sender, bounded, unbounded};
|
||||||
|
|
||||||
use crate::imap::client::{ConnectMode, ImapClient};
|
use crate::imap::client::{ConnectMode, ImapClient};
|
||||||
use crate::imap::command;
|
use crate::imap::command;
|
||||||
@@ -75,7 +75,11 @@ impl WorkerPool {
|
|||||||
pub fn start(args: WorkerArgs, pool_size: usize) -> Result<WorkerPool, ImapError> {
|
pub fn start(args: WorkerArgs, pool_size: usize) -> Result<WorkerPool, ImapError> {
|
||||||
let size = pool_size.clamp(1, HARD_CAP);
|
let size = pool_size.clamp(1, HARD_CAP);
|
||||||
let (job_tx, job_rx) = unbounded::<FetchJob>();
|
let (job_tx, job_rx) = unbounded::<FetchJob>();
|
||||||
let (event_tx, event_rx) = unbounded::<FetchEvent>();
|
// Jobs are only lists of UIDs, but each event carries a whole message:
|
||||||
|
// bounding the events stops fast workers from running ahead of the
|
||||||
|
// single archive writer, so memory holds at most a couple of messages
|
||||||
|
// per worker rather than whole folders.
|
||||||
|
let (event_tx, event_rx) = bounded::<FetchEvent>(size * 2);
|
||||||
let mut handles = Vec::with_capacity(size);
|
let mut handles = Vec::with_capacity(size);
|
||||||
let args = Arc::new(args);
|
let args = Arc::new(args);
|
||||||
let cancel_below = Arc::new(AtomicU64::new(0));
|
let cancel_below = Arc::new(AtomicU64::new(0));
|
||||||
|
|||||||
@@ -5,6 +5,7 @@
|
|||||||
* SPDX-License-Identifier: Apache-2.0 OR MIT
|
* SPDX-License-Identifier: Apache-2.0 OR MIT
|
||||||
*/
|
*/
|
||||||
|
|
||||||
|
pub mod batch;
|
||||||
pub mod emailmeta;
|
pub mod emailmeta;
|
||||||
pub mod export;
|
pub mod export;
|
||||||
pub mod import_dav;
|
pub mod import_dav;
|
||||||
|
|||||||
@@ -38,6 +38,7 @@ fn imap_config(account: &Account, imap: &Endpoint) -> ImapImportConfig {
|
|||||||
automap: true,
|
automap: true,
|
||||||
include_deleted: false,
|
include_deleted: false,
|
||||||
fetch_batch: 64,
|
fetch_batch: 64,
|
||||||
|
fetch_batch_bytes: inbuxa_migrate::sync::batch::DEFAULT_BATCH_BYTES,
|
||||||
imap_connections: 2,
|
imap_connections: 2,
|
||||||
allow_source_change: false,
|
allow_source_change: false,
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -43,6 +43,7 @@ fn imap_config(account: &Account, imap: &integration::Endpoint) -> ImapImportCon
|
|||||||
automap: true,
|
automap: true,
|
||||||
include_deleted: false,
|
include_deleted: false,
|
||||||
fetch_batch: 64,
|
fetch_batch: 64,
|
||||||
|
fetch_batch_bytes: inbuxa_migrate::sync::batch::DEFAULT_BATCH_BYTES,
|
||||||
imap_connections: 2,
|
imap_connections: 2,
|
||||||
allow_source_change: false,
|
allow_source_change: false,
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -864,6 +864,7 @@ fn for_each_fetched_item_streams_every_id_across_windows() {
|
|||||||
url: &url,
|
url: &url,
|
||||||
source_id: 1,
|
source_id: 1,
|
||||||
batch_size: 1,
|
batch_size: 1,
|
||||||
|
batch_bytes: inbuxa_migrate::sync::batch::DEFAULT_BATCH_BYTES,
|
||||||
attachment_batch: 1,
|
attachment_batch: 1,
|
||||||
connections: 2,
|
connections: 2,
|
||||||
use_syncfolderitems: false,
|
use_syncfolderitems: false,
|
||||||
@@ -873,18 +874,86 @@ fn for_each_fetched_item_streams_every_id_across_windows() {
|
|||||||
let ids: Vec<ItemId> = (0..5).map(|i| ItemId::new(format!("I{i}"), "K")).collect();
|
let ids: Vec<ItemId> = (0..5).map(|i| ItemId::new(format!("I{i}"), "K")).collect();
|
||||||
|
|
||||||
let mut delivered = 0usize;
|
let mut delivered = 0usize;
|
||||||
let failed = for_each_fetched_item(&ctx, ItemShape::Message, &ids, |msg| {
|
let failed =
|
||||||
assert!(msg.success);
|
for_each_fetched_item(&ctx, ItemShape::Message, &ids, &Default::default(), |msg| {
|
||||||
delivered += 1;
|
assert!(msg.success);
|
||||||
Ok(())
|
delivered += 1;
|
||||||
})
|
Ok(())
|
||||||
.expect("streaming fetch should succeed");
|
})
|
||||||
|
.expect("streaming fetch should succeed");
|
||||||
|
|
||||||
assert_eq!(delivered, 5, "every id must be delivered exactly once");
|
assert_eq!(delivered, 5, "every id must be delivered exactly once");
|
||||||
assert_eq!(failed, 0);
|
assert_eq!(failed, 0);
|
||||||
_m.assert();
|
_m.assert();
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn getitem_batches_are_split_by_bytes_when_sizes_are_known() {
|
||||||
|
use inbuxa_migrate::logging::Logger;
|
||||||
|
use inbuxa_migrate::sync::import_exchange_ews::items::{ItemRunCtx, for_each_fetched_item};
|
||||||
|
use std::collections::HashMap;
|
||||||
|
|
||||||
|
let mut server = mockito::Server::new();
|
||||||
|
let url = format!("{}/EWS/Exchange.asmx", server.url());
|
||||||
|
let one_message = envelope(&format!(
|
||||||
|
"<m:GetItemResponse{NS}><m:ResponseMessages><m:GetItemResponseMessage ResponseClass=\"Success\">\
|
||||||
|
<m:ResponseCode>NoError</m:ResponseCode><m:Items><t:Message><t:ItemId Id=\"X\" ChangeKey=\"K\"/></t:Message></m:Items>\
|
||||||
|
</m:GetItemResponseMessage></m:ResponseMessages></m:GetItemResponse>"
|
||||||
|
));
|
||||||
|
// Ten items fit one batch by count, but at 20 bytes each and a 30-byte
|
||||||
|
// cap every item goes alone: four GetItem calls, not one.
|
||||||
|
let m = server
|
||||||
|
.mock("POST", "/EWS/Exchange.asmx")
|
||||||
|
.with_status(200)
|
||||||
|
.with_header("content-type", TXT_XML)
|
||||||
|
.with_body(&one_message)
|
||||||
|
.expect(4)
|
||||||
|
.create();
|
||||||
|
let c = client(0);
|
||||||
|
let ctx = ItemRunCtx {
|
||||||
|
client: &c,
|
||||||
|
url: &url,
|
||||||
|
source_id: 1,
|
||||||
|
batch_size: 10,
|
||||||
|
batch_bytes: 30,
|
||||||
|
attachment_batch: 1,
|
||||||
|
connections: 2,
|
||||||
|
use_syncfolderitems: false,
|
||||||
|
sync_batch: 512,
|
||||||
|
logger: Logger::new(0),
|
||||||
|
};
|
||||||
|
let ids: Vec<ItemId> = (0..4).map(|i| ItemId::new(format!("I{i}"), "K")).collect();
|
||||||
|
let sizes: HashMap<String, u64> = ids.iter().map(|id| (id.id.clone(), 20)).collect();
|
||||||
|
let mut delivered = 0usize;
|
||||||
|
for_each_fetched_item(&ctx, ItemShape::Message, &ids, &sizes, |_| {
|
||||||
|
delivered += 1;
|
||||||
|
Ok(())
|
||||||
|
})
|
||||||
|
.expect("fetch");
|
||||||
|
assert_eq!(delivered, 4);
|
||||||
|
m.assert();
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn find_item_reports_each_items_size() {
|
||||||
|
use inbuxa_migrate::exchange_ews::parse::parse_find_item_response;
|
||||||
|
let body = envelope(&format!(
|
||||||
|
"<m:FindItemResponse{NS}><m:ResponseMessages><m:FindItemResponseMessage ResponseClass=\"Success\">\
|
||||||
|
<m:ResponseCode>NoError</m:ResponseCode>\
|
||||||
|
<m:RootFolder TotalItemsInView=\"2\" IncludesLastItemInRange=\"true\"><t:Items>\
|
||||||
|
<t:Message><t:ItemId Id=\"A\" ChangeKey=\"K\"/><t:Size>1234</t:Size></t:Message>\
|
||||||
|
<t:Message><t:ItemId Id=\"B\" ChangeKey=\"K\"/></t:Message>\
|
||||||
|
</t:Items></m:RootFolder></m:FindItemResponseMessage></m:ResponseMessages></m:FindItemResponse>"
|
||||||
|
));
|
||||||
|
let r = parse_find_item_response(body.as_bytes()).unwrap();
|
||||||
|
assert_eq!(r.items.len(), 2);
|
||||||
|
assert_eq!(
|
||||||
|
(r.items[0].id.id.as_str(), r.items[0].size),
|
||||||
|
("A", Some(1234))
|
||||||
|
);
|
||||||
|
assert_eq!((r.items[1].id.id.as_str(), r.items[1].size), ("B", None));
|
||||||
|
}
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
fn warning_response_class_is_treated_as_success_in_mock() {
|
fn warning_response_class_is_treated_as_success_in_mock() {
|
||||||
let body = envelope(&format!(
|
let body = envelope(&format!(
|
||||||
|
|||||||
@@ -169,6 +169,7 @@ fn run_import(
|
|||||||
automap: true,
|
automap: true,
|
||||||
include_deleted: false,
|
include_deleted: false,
|
||||||
fetch_batch: 256,
|
fetch_batch: 256,
|
||||||
|
fetch_batch_bytes: inbuxa_migrate::sync::batch::DEFAULT_BATCH_BYTES,
|
||||||
imap_connections: 1,
|
imap_connections: 1,
|
||||||
allow_source_change: false,
|
allow_source_change: false,
|
||||||
};
|
};
|
||||||
@@ -2135,3 +2136,88 @@ fn a_failed_chunk_keeps_the_ones_before_it_and_a_rerun_fetches_only_what_is_miss
|
|||||||
db::init::apply_schema(&conn).unwrap();
|
db::init::apply_schema(&conn).unwrap();
|
||||||
assert_eq!(count(&conn, "emails"), 2);
|
assert_eq!(count(&conn, "emails"), 2);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn body_fetches_are_split_by_message_size() {
|
||||||
|
// Three 20-byte messages under a 30-byte cap: the control connection
|
||||||
|
// learns the sizes first, and each body is then fetched on its own.
|
||||||
|
let control: Script = Box::new(|conn: &mut MockConn| -> std::io::Result<()> {
|
||||||
|
auth_preamble(conn, "IMAP4rev2 LITERAL+ AUTH=PLAIN")?;
|
||||||
|
let (tag, _) = conn.read_command()?;
|
||||||
|
conn.write_line("* LIST () \"/\" \"INBOX\"")?;
|
||||||
|
conn.write_line(&format!("{tag} OK"))?;
|
||||||
|
let (tag, _) = conn.read_command()?;
|
||||||
|
conn.write_line(&format!("{tag} OK"))?;
|
||||||
|
let (tag, _) = conn.read_command()?;
|
||||||
|
write_select(conn, &tag, 555, 4, 3)?;
|
||||||
|
let (tag, _) = conn.read_command()?;
|
||||||
|
conn.write_line("* SEARCH 1 2 3")?;
|
||||||
|
conn.write_line(&format!("{tag} OK"))?;
|
||||||
|
let (tag, cmd) = conn.read_command()?;
|
||||||
|
assert_eq!(cmd, "UID FETCH 1:3 (UID RFC822.SIZE)");
|
||||||
|
for uid in 1..=3 {
|
||||||
|
conn.write_line(&format!("* {uid} FETCH (UID {uid} RFC822.SIZE 20)"))?;
|
||||||
|
}
|
||||||
|
conn.write_line(&format!("{tag} OK"))?;
|
||||||
|
drain_until_close(conn);
|
||||||
|
Ok(())
|
||||||
|
});
|
||||||
|
let fetches = std::sync::Arc::new(Mutex::new(Vec::<String>::new()));
|
||||||
|
let worker: Script = {
|
||||||
|
let fetches = fetches.clone();
|
||||||
|
Box::new(move |conn: &mut MockConn| -> std::io::Result<()> {
|
||||||
|
auth_preamble(conn, "IMAP4rev2 LITERAL+ AUTH=PLAIN")?;
|
||||||
|
let (tag, _) = conn.read_command()?;
|
||||||
|
write_select(conn, &tag, 555, 4, 3)?;
|
||||||
|
let bodies: [&[u8]; 3] = [BODY_A, BODY_B, BODY_C];
|
||||||
|
for _ in 0..3 {
|
||||||
|
let (tag, cmd) = conn.read_command()?;
|
||||||
|
let uid: u32 = cmd
|
||||||
|
.strip_prefix("UID FETCH ")
|
||||||
|
.and_then(|r| r.split(' ').next())
|
||||||
|
.and_then(|n| n.parse().ok())
|
||||||
|
.unwrap_or_else(|| panic!("expected a single-UID fetch, got {cmd}"));
|
||||||
|
fetches.lock().unwrap().push(cmd.clone());
|
||||||
|
write_fetch_message(conn, uid, uid, bodies[(uid - 1) as usize])?;
|
||||||
|
conn.write_line(&format!("{tag} OK"))?;
|
||||||
|
}
|
||||||
|
drain_until_close(conn);
|
||||||
|
Ok(())
|
||||||
|
})
|
||||||
|
};
|
||||||
|
let server = MockImap::start_scripts(vec![control, worker]);
|
||||||
|
let archive = tempfile("byte-chunks");
|
||||||
|
let summary =
|
||||||
|
run_import(&server, "alice", archive, |c| c.fetch_batch_bytes = 30).expect("import");
|
||||||
|
assert_eq!(email_counts(&summary).created, 3, "summary={summary:?}");
|
||||||
|
assert_eq!(
|
||||||
|
fetches.lock().unwrap().len(),
|
||||||
|
3,
|
||||||
|
"one body fetch per message"
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn a_chunk_larger_than_the_event_queue_completes() {
|
||||||
|
// One worker holds at most two events in flight; a ten-message chunk
|
||||||
|
// has to wait on the archive writer, not deadlock on it.
|
||||||
|
const UIDS: &[u32] = &[1, 2, 3, 4, 5, 6, 7, 8, 9, 10];
|
||||||
|
let worker: Script = Box::new(|conn: &mut MockConn| -> std::io::Result<()> {
|
||||||
|
auth_preamble(conn, "IMAP4rev2 LITERAL+ AUTH=PLAIN")?;
|
||||||
|
let (tag, _) = conn.read_command()?;
|
||||||
|
write_select(conn, &tag, 999, 11, 10)?;
|
||||||
|
let (tag, cmd) = conn.read_command()?;
|
||||||
|
assert!(cmd.starts_with("UID FETCH 1:10 "), "got {cmd}");
|
||||||
|
for uid in UIDS {
|
||||||
|
let body = format!("From: a@b\r\nMessage-ID: <{uid}@h>\r\n\r\nbody {uid}");
|
||||||
|
write_fetch_message(conn, *uid, *uid, body.as_bytes())?;
|
||||||
|
}
|
||||||
|
conn.write_line(&format!("{tag} OK"))?;
|
||||||
|
drain_until_close(conn);
|
||||||
|
Ok(())
|
||||||
|
});
|
||||||
|
let server = MockImap::start_scripts(vec![control_script_one_folder(999, 11, UIDS), worker]);
|
||||||
|
let archive = tempfile("backpressure");
|
||||||
|
let summary = run_import(&server, "alice", archive, |_| {}).expect("import");
|
||||||
|
assert_eq!(email_counts(&summary).created, 10, "summary={summary:?}");
|
||||||
|
}
|
||||||
|
|||||||
@@ -71,6 +71,7 @@ fn imap_basic_config(localpart: &str) -> ImapImportConfig {
|
|||||||
automap: true,
|
automap: true,
|
||||||
include_deleted: false,
|
include_deleted: false,
|
||||||
fetch_batch: 256,
|
fetch_batch: 256,
|
||||||
|
fetch_batch_bytes: inbuxa_migrate::sync::batch::DEFAULT_BATCH_BYTES,
|
||||||
imap_connections: 4,
|
imap_connections: 4,
|
||||||
allow_source_change: false,
|
allow_source_change: false,
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user