From 3cae9464f01e9decb7afa882d6ec490235d2e810 Mon Sep 17 00:00:00 2001 From: John Coffey Date: Wed, 30 Sep 2026 12:32:29 -0700 Subject: [PATCH] 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. --- docs/usage.md | 19 ++++- src/cli.rs | 21 +++++ src/exchange_ews/parse.rs | 17 ++++ src/exchange_ews/xml.rs | 17 +++- src/sync/batch.rs | 92 +++++++++++++++++++++ src/sync/import_exchange_ews/calendar.rs | 81 +++++++++--------- src/sync/import_exchange_ews/contacts.rs | 68 +++++++-------- src/sync/import_exchange_ews/coordinator.rs | 3 + src/sync/import_exchange_ews/items.rs | 54 +++++++++--- src/sync/import_exchange_ews/messages.rs | 72 ++++++++-------- src/sync/import_imap/coordinator.rs | 56 ++++++++++++- src/sync/import_imap/pool.rs | 8 +- src/sync/mod.rs | 1 + tests/integration_cyrus.rs | 1 + tests/integration_dovecot.rs | 1 + tests/mock_exchange_ews.rs | 81 ++++++++++++++++-- tests/mock_imap.rs | 86 +++++++++++++++++++ tests/sync_imap.rs | 1 + 18 files changed, 551 insertions(+), 128 deletions(-) create mode 100644 src/sync/batch.rs diff --git a/docs/usage.md b/docs/usage.md index 4802889..6d95e49 100644 --- a/docs/usage.md +++ b/docs/usage.md @@ -83,7 +83,7 @@ inbuxa-migrate import imap \ [--include ...] [--exclude ...] [--exclude-special ...] \ [--folder ...] [--subscribed-only] [--noautomap] \ [--include-deleted] [--allow-cleartext] [--compress] \ - [--fetch-batch ] [--imap-connections <1..8>] \ + [--fetch-batch ] [--fetch-batch-mib ] [--imap-connections <1..8>] \ ``` @@ -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 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 ``` @@ -170,7 +178,8 @@ inbuxa-migrate import exchange-ews \ (--auth-basic [--auth-password ] \ | --auth-bearer [TOKEN] [--ews-tenant --ews-client-id \ (--ews-device-code | --ews-client-secret )]) \ - [--ews-connections <1..8>] [--ews-getitem-batch ] [--ews-attachment-batch ] \ + [--ews-connections <1..8>] [--ews-getitem-batch ] [--ews-getitem-batch-mib ] \ + [--ews-attachment-batch ] \ [--ews-no-syncfolderitems] \ ``` @@ -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 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. ### Microsoft Exchange (Graph) diff --git a/src/cli.rs b/src/cli.rs index f585785..64e2cbf 100644 --- a/src/cli.rs +++ b/src/cli.rs @@ -277,6 +277,14 @@ struct ImapImportArgs { )] 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( long, value_name = "N", @@ -804,6 +812,7 @@ fn resolve_imap_import(args: ImapImportArgs) -> Result { automap: !args.noautomap, include_deleted: args.include_deleted, fetch_batch: args.fetch_batch, + fetch_batch_bytes: args.fetch_batch_mib.max(1).saturating_mul(1024 * 1024), imap_connections, allow_source_change: args.allow_source_change, }, @@ -973,6 +982,14 @@ pub struct ExchangeEwsImportArgs { )] 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( long, value_name = "N", @@ -1022,6 +1039,10 @@ fn resolve_exchange_ews_import(args: ExchangeEwsImportArgs) -> Result, } pub fn parse_find_item_response(body: &[u8]) -> Result { @@ -511,6 +513,7 @@ pub fn parse_find_item_response(body: &[u8]) -> Result Result Result { + if let Some(last) = out.items.last_mut() { + let text: &str = t; + last.size = text.trim().parse().ok(); + } + } Event::End(e) => { let local = e.local_name().as_ref().to_owned(); if local.eq_ignore_ascii_case("RootFolder") { in_root = false; + } else if local.eq_ignore_ascii_case("Size") { + reading_size = false; } } Event::Eof => break, diff --git a/src/exchange_ews/xml.rs b/src/exchange_ews/xml.rs index a19e931..ef56e25 100644 --- a/src/exchange_ews/xml.rs +++ b/src/exchange_ews/xml.rs @@ -1,5 +1,6 @@ /* * SPDX-FileCopyrightText: 2020 Stalwart Labs LLC + * SPDX-FileCopyrightText: 2026 John Coffey * * SPDX-License-Identifier: Apache-2.0 OR MIT */ @@ -83,7 +84,12 @@ pub fn find_item_body( let mut out = String::with_capacity(512); out.push_str("IdOnly"); + // item:Size lets GetItem batches be split by bytes as well as by count. + out.push_str( + "\">IdOnly\ + \ + ", + ); out.push_str("")); } + #[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( + "IdOnly" + )); + } + #[test] fn find_item_paginates_with_offset_and_page_size() { let folder = FolderId::new("FID", "FCK"); diff --git a/src/sync/batch.rs b/src/sync/batch.rs new file mode 100644 index 0000000..4a972a2 --- /dev/null +++ b/src/sync/batch.rs @@ -0,0 +1,92 @@ +/* + * SPDX-FileCopyrightText: 2026 John Coffey + * + * 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( + 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> { + 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()); + } +} diff --git a/src/sync/import_exchange_ews/calendar.rs b/src/sync/import_exchange_ews/calendar.rs index bfc7b40..381476e 100644 --- a/src/sync/import_exchange_ews/calendar.rs +++ b/src/sync/import_exchange_ews/calendar.rs @@ -1,5 +1,6 @@ /* * SPDX-FileCopyrightText: 2020 Stalwart Labs LLC + * SPDX-FileCopyrightText: 2026 John Coffey * * SPDX-License-Identifier: Apache-2.0 OR MIT */ @@ -104,47 +105,53 @@ fn reconcile_one( to_fetch.push(id.clone()); } if !to_fetch.is_empty() { - let failed_items = for_each_fetched_item(ctx, ItemShape::CalendarItem, &to_fetch, |msg| { - if !msg.success { + let failed_items = for_each_fetched_item( + 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!( - msg.response_code, - crate::exchange_ews::types::ResponseCode::ItemNotFound + parsed.calendar_item_type, + Some(CalendarItemType::Occurrence) | Some(CalendarItemType::Exception) ) { counts.skipped += 1; - } else { - counts.failed += 1; - ctx.logger - .warn(&format!("GetItem (calendar) error: {}", msg.response_code)); + return Ok(()); } - 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!( - parsed.calendar_item_type, - Some(CalendarItemType::Occurrence) | Some(CalendarItemType::Exception) - ) { - counts.skipped += 1; - return Ok(()); - } - 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, - ) - })?; + 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; } delete_vanished( diff --git a/src/sync/import_exchange_ews/contacts.rs b/src/sync/import_exchange_ews/contacts.rs index 07ebce5..e8b42b1 100644 --- a/src/sync/import_exchange_ews/contacts.rs +++ b/src/sync/import_exchange_ews/contacts.rs @@ -1,5 +1,6 @@ /* * SPDX-FileCopyrightText: 2020 Stalwart Labs LLC + * SPDX-FileCopyrightText: 2026 John Coffey * * SPDX-License-Identifier: Apache-2.0 OR MIT */ @@ -96,40 +97,41 @@ fn reconcile_one( to_fetch.push(id.clone()); } if !to_fetch.is_empty() { - let failed_items = for_each_fetched_item(ctx, ItemShape::Contact, &to_fetch, |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 (contact) error: {}", msg.response_code)); + let failed_items = + for_each_fetched_item(ctx, ItemShape::Contact, &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 (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() { - counts.failed += 1; - return Ok(()); - } - let existing = plan - .present_changed - .iter() - .find(|(id, _)| id.id == parsed.id.id) - .map(|(_, local)| *local); - apply_contact( - conn, - ctx, - &parsed, - local_folder_id, - &folder.id, - existing, - counts, - ) - })?; + let parsed = parse_contact_item(&msg.inner_xml).map_err(Error::from)?; + if parsed.id.id.is_empty() { + counts.failed += 1; + return Ok(()); + } + let existing = plan + .present_changed + .iter() + .find(|(id, _)| id.id == parsed.id.id) + .map(|(_, local)| *local); + apply_contact( + conn, + ctx, + &parsed, + local_folder_id, + &folder.id, + existing, + counts, + ) + })?; counts.failed += failed_items; } delete_vanished( diff --git a/src/sync/import_exchange_ews/coordinator.rs b/src/sync/import_exchange_ews/coordinator.rs index 053e77f..d40b6d4 100644 --- a/src/sync/import_exchange_ews/coordinator.rs +++ b/src/sync/import_exchange_ews/coordinator.rs @@ -37,6 +37,8 @@ pub struct EwsImportConfig { pub auth: EwsAuth, pub ews_connections: usize, pub getitem_batch: usize, + /// Byte cap for one GetItem batch; see `sync::batch`. + pub getitem_batch_bytes: u64, pub attachment_batch: usize, pub use_syncfolderitems: bool, pub allow_source_change: bool, @@ -147,6 +149,7 @@ pub fn run(common: CommonConfig, config: EwsImportConfig) -> Result + * SPDX-FileCopyrightText: 2026 John Coffey * * SPDX-License-Identifier: Apache-2.0 OR MIT */ @@ -23,6 +24,8 @@ pub struct ItemRunCtx<'a> { pub url: &'a str, pub source_id: i64, 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 connections: usize, pub use_syncfolderitems: bool, @@ -47,6 +50,9 @@ pub struct EnumeratedItem { pub struct EnumerationOutcome { pub items: Vec, 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, } #[derive(Debug, Clone)] @@ -70,9 +76,10 @@ pub fn enumerate_folder( { return Ok(outcome); } - enumerate_via_find_item(ctx, folder).map(|items| EnumerationOutcome { + enumerate_via_find_item(ctx, folder).map(|(items, sizes)| EnumerationOutcome { items, mode: EnumerationMode::Full, + sizes, }) } @@ -151,14 +158,16 @@ fn try_sync_folder_items( deletions, new_sync_state: sync_state, }, + sizes: HashMap::new(), })) } fn enumerate_via_find_item( ctx: &ItemRunCtx<'_>, folder: &FolderId, -) -> Result, EwsError> { +) -> Result<(Vec, HashMap), EwsError> { let mut items: Vec = Vec::new(); + let mut sizes: HashMap = HashMap::new(); let mut offset: u32 = 0; let page_size: u32 = 500; loop { @@ -172,6 +181,9 @@ fn enumerate_via_find_item( let parsed = parse_find_item_response(&resp.body)?; let returned = parsed.items.len() as u32; for entry in parsed.items { + if let Some(size) = entry.size { + sizes.insert(entry.id.id.clone(), size); + } items.push(EnumeratedItem { element: entry.element, id: entry.id, @@ -188,7 +200,7 @@ fn enumerate_via_find_item( break; } } - Ok(items) + Ok((items, sizes)) } #[derive(Debug, Clone)] @@ -285,13 +297,22 @@ pub fn get_items( shape: ItemShape, ids: &[ItemId], ) -> Result { - 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 { let workers = ctx.connections.clamp(1, 8); let version = ctx.client.server_version(); let mut failed_items: u64 = 0; - if workers <= 1 || ids.len() <= batch { + if workers <= 1 || batches.len() <= 1 { 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); match ctx.client.call(ctx.url, "GetItem", &body) { Ok(resp) => match parse_response_messages(&resp.body, "GetItemResponseMessage") { @@ -337,7 +358,7 @@ pub fn get_items( (n, result) }); let mut submitted = 0usize; - for chunk in ids.chunks(batch) { + for chunk in batches { pool.submit(chunk.to_vec()); 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( ctx: &ItemRunCtx<'_>, shape: ItemShape, ids: &[ItemId], + sizes: &HashMap, mut on_message: F, ) -> Result where F: FnMut(crate::exchange_ews::parse::ResponseMessage) -> Result<(), Error>, { - let batch = ctx.batch_size.max(1); 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; - for win in ids.chunks(window) { - let outcome = get_items(ctx, shape, win).map_err(Error::from)?; + for win in batches.chunks(workers) { + let outcome = get_item_batches(ctx, shape, win).map_err(Error::from)?; failed_items = failed_items.saturating_add(outcome.failed_items); for msg in outcome.messages { on_message(msg)?; @@ -444,6 +475,7 @@ mod tests { element: "Message".to_owned(), id: ItemId::new("A", "ck-2"), }], + sizes: HashMap::new(), mode: EnumerationMode::Delta { deletions: vec!["Z".to_owned()], new_sync_state: "STATE2".to_owned(), diff --git a/src/sync/import_exchange_ews/messages.rs b/src/sync/import_exchange_ews/messages.rs index 8d2acc5..c99ea48 100644 --- a/src/sync/import_exchange_ews/messages.rs +++ b/src/sync/import_exchange_ews/messages.rs @@ -1,5 +1,6 @@ /* * SPDX-FileCopyrightText: 2020 Stalwart Labs LLC + * SPDX-FileCopyrightText: 2026 John Coffey * * SPDX-License-Identifier: Apache-2.0 OR MIT */ @@ -95,42 +96,43 @@ fn reconcile_one_folder( to_fetch.push(id.clone()); } if !to_fetch.is_empty() { - let failed_items = for_each_fetched_item(ctx, ItemShape::Message, &to_fetch, |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 (message) error: {} {}", - msg.response_code, msg.message_text - )); + let failed_items = + for_each_fetched_item(ctx, ItemShape::Message, &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 (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() { - counts.failed += 1; - return Ok(()); - } - let existing = plan - .present_changed - .iter() - .find(|(id, _)| id.id == parsed.id.id) - .map(|(_, local)| *local); - apply_message( - conn, - ctx, - &parsed, - local_folder_id, - &folder.id, - existing, - counts, - ) - })?; + let parsed = parse_message_item(&msg.inner_xml).map_err(Error::from)?; + if parsed.id.id.is_empty() { + counts.failed += 1; + return Ok(()); + } + let existing = plan + .present_changed + .iter() + .find(|(id, _)| id.id == parsed.id.id) + .map(|(_, local)| *local); + apply_message( + conn, + ctx, + &parsed, + local_folder_id, + &folder.id, + existing, + counts, + ) + })?; counts.failed += failed_items; } delete_vanished( diff --git a/src/sync/import_imap/coordinator.rs b/src/sync/import_imap/coordinator.rs index 2b2d8a2..46ea844 100644 --- a/src/sync/import_imap/coordinator.rs +++ b/src/sync/import_imap/coordinator.rs @@ -50,6 +50,7 @@ pub(super) struct RunOpts { /// `WorkerPool::cancel_before`. generation: u64, fetch_batch: usize, + fetch_batch_bytes: u64, include_deleted: bool, logger: Logger, } @@ -199,6 +200,8 @@ pub struct ImapImportConfig { pub automap: bool, pub include_deleted: bool, 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 allow_source_change: bool, } @@ -404,6 +407,7 @@ fn run_into( source_id, generation: 0, fetch_batch: config.fetch_batch.max(1), + fetch_batch_bytes: config.fetch_batch_bytes.max(1), include_deleted: config.include_deleted, logger, }; @@ -706,6 +710,7 @@ fn reconcile_folder( source_id, generation, fetch_batch, + fetch_batch_bytes: _, include_deleted: _, logger, } = opts; @@ -811,7 +816,13 @@ fn reconcile_folder( 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(); for batch in &batches { pool.submit(FetchJob { @@ -892,6 +903,47 @@ fn reconcile_folder( 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 { + 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( conn: &mut Connection, source_id: i64, @@ -1026,6 +1078,7 @@ fn insert_single_message( source_id, generation: _, fetch_batch: _, + fetch_batch_bytes: _, include_deleted, logger, } = opts; @@ -1104,6 +1157,7 @@ fn refresh_present_flags( source_id, generation: _, fetch_batch, + fetch_batch_bytes: _, include_deleted, logger: _, } = opts; diff --git a/src/sync/import_imap/pool.rs b/src/sync/import_imap/pool.rs index ecfddfd..a99a760 100644 --- a/src/sync/import_imap/pool.rs +++ b/src/sync/import_imap/pool.rs @@ -10,7 +10,7 @@ use std::sync::Arc; use std::sync::atomic::{AtomicU64, Ordering}; 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::command; @@ -75,7 +75,11 @@ impl WorkerPool { pub fn start(args: WorkerArgs, pool_size: usize) -> Result { let size = pool_size.clamp(1, HARD_CAP); let (job_tx, job_rx) = unbounded::(); - let (event_tx, event_rx) = unbounded::(); + // 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::(size * 2); let mut handles = Vec::with_capacity(size); let args = Arc::new(args); let cancel_below = Arc::new(AtomicU64::new(0)); diff --git a/src/sync/mod.rs b/src/sync/mod.rs index a9e7178..45249d6 100644 --- a/src/sync/mod.rs +++ b/src/sync/mod.rs @@ -5,6 +5,7 @@ * SPDX-License-Identifier: Apache-2.0 OR MIT */ +pub mod batch; pub mod emailmeta; pub mod export; pub mod import_dav; diff --git a/tests/integration_cyrus.rs b/tests/integration_cyrus.rs index 2f78a79..58bab26 100644 --- a/tests/integration_cyrus.rs +++ b/tests/integration_cyrus.rs @@ -38,6 +38,7 @@ fn imap_config(account: &Account, imap: &Endpoint) -> ImapImportConfig { automap: true, include_deleted: false, fetch_batch: 64, + fetch_batch_bytes: inbuxa_migrate::sync::batch::DEFAULT_BATCH_BYTES, imap_connections: 2, allow_source_change: false, } diff --git a/tests/integration_dovecot.rs b/tests/integration_dovecot.rs index e6dcb13..b20fb2f 100644 --- a/tests/integration_dovecot.rs +++ b/tests/integration_dovecot.rs @@ -43,6 +43,7 @@ fn imap_config(account: &Account, imap: &integration::Endpoint) -> ImapImportCon automap: true, include_deleted: false, fetch_batch: 64, + fetch_batch_bytes: inbuxa_migrate::sync::batch::DEFAULT_BATCH_BYTES, imap_connections: 2, allow_source_change: false, } diff --git a/tests/mock_exchange_ews.rs b/tests/mock_exchange_ews.rs index c070fcc..a22fe2d 100644 --- a/tests/mock_exchange_ews.rs +++ b/tests/mock_exchange_ews.rs @@ -864,6 +864,7 @@ fn for_each_fetched_item_streams_every_id_across_windows() { url: &url, source_id: 1, batch_size: 1, + batch_bytes: inbuxa_migrate::sync::batch::DEFAULT_BATCH_BYTES, attachment_batch: 1, connections: 2, use_syncfolderitems: false, @@ -873,18 +874,86 @@ fn for_each_fetched_item_streams_every_id_across_windows() { let ids: Vec = (0..5).map(|i| ItemId::new(format!("I{i}"), "K")).collect(); let mut delivered = 0usize; - let failed = for_each_fetched_item(&ctx, ItemShape::Message, &ids, |msg| { - assert!(msg.success); - delivered += 1; - Ok(()) - }) - .expect("streaming fetch should succeed"); + let failed = + for_each_fetched_item(&ctx, ItemShape::Message, &ids, &Default::default(), |msg| { + assert!(msg.success); + delivered += 1; + Ok(()) + }) + .expect("streaming fetch should succeed"); assert_eq!(delivered, 5, "every id must be delivered exactly once"); assert_eq!(failed, 0); _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!( + "\ + NoError\ + " + )); + // 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 = (0..4).map(|i| ItemId::new(format!("I{i}"), "K")).collect(); + let sizes: HashMap = 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!( + "\ + NoError\ + \ + 1234\ + \ + " + )); + 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] fn warning_response_class_is_treated_as_success_in_mock() { let body = envelope(&format!( diff --git a/tests/mock_imap.rs b/tests/mock_imap.rs index 118fb03..b1714e7 100644 --- a/tests/mock_imap.rs +++ b/tests/mock_imap.rs @@ -169,6 +169,7 @@ fn run_import( automap: true, include_deleted: false, fetch_batch: 256, + fetch_batch_bytes: inbuxa_migrate::sync::batch::DEFAULT_BATCH_BYTES, imap_connections: 1, 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(); 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::::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:?}"); +} diff --git a/tests/sync_imap.rs b/tests/sync_imap.rs index c1b93fd..665e83b 100644 --- a/tests/sync_imap.rs +++ b/tests/sync_imap.rs @@ -71,6 +71,7 @@ fn imap_basic_config(localpart: &str) -> ImapImportConfig { automap: true, include_deleted: false, fetch_batch: 256, + fetch_batch_bytes: inbuxa_migrate::sync::batch::DEFAULT_BATCH_BYTES, imap_connections: 4, allow_source_change: false, } -- 2.54.0