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, }