From 2f33cd76a159527b9790a4f2b289cc6a72a5b24f Mon Sep 17 00:00:00 2001 From: John Coffey Date: Wed, 30 Sep 2026 12:24:37 -0700 Subject: [PATCH] import: keep what an interrupted IMAP or Maildir import wrote IMAP import wrote a whole folder in one transaction and stopped the folder at the first message it could not import. A crash near the end of a large INBOX kept nothing, and one bad INTERNALDATE lost the rest of the folder. Worse, when a folder stopped early, fetches still in flight for it could be filed into the next folder's mailbox. - The transaction is committed after every fetch chunk. Each message is written in its own savepoint, so what is committed is always whole, and a rerun fetches only the UIDs still missing. - A message that cannot be imported is rolled back on its own, logged with its folder and UID, and counted as failed; the folder carries on, and the message stays out of the UID map so the next run tries it again. Archive and I/O errors still stop the run. - Fetch jobs and events carry a folder generation. Moving to a new folder cancels queued work for older ones, and any event from an older generation is dropped, never filed. Shutdown drains in-flight events before joining the workers, so it cannot hang on a blocked worker. - INTERNALDATE month names are matched in any case. - Maildir import gets the same per-message savepoint, and commits every 500 new messages instead of once per folder. --- src/sync/import_imap/coordinator.rs | 139 +++++++++++++++++++++++++-- src/sync/import_imap/internaldate.rs | 45 ++++++--- src/sync/import_imap/pool.rs | 122 ++++++++++++++++++++++- src/sync/import_maildir/messages.rs | 58 ++++++++++- tests/mock_imap.rs | 130 +++++++++++++++++++++++++ 5 files changed, 467 insertions(+), 27 deletions(-) diff --git a/src/sync/import_imap/coordinator.rs b/src/sync/import_imap/coordinator.rs index b5e4c33..2b2d8a2 100644 --- a/src/sync/import_imap/coordinator.rs +++ b/src/sync/import_imap/coordinator.rs @@ -1,5 +1,6 @@ /* * SPDX-FileCopyrightText: 2020 Stalwart Labs LLC + * SPDX-FileCopyrightText: 2026 John Coffey * * SPDX-License-Identifier: Apache-2.0 OR MIT */ @@ -45,6 +46,9 @@ const EMAIL_TYPE: &str = "email"; #[derive(Clone, Copy)] pub(super) struct RunOpts { source_id: i64, + /// The folder generation this folder's fetch jobs carry; see + /// `WorkerPool::cancel_before`. + generation: u64, fetch_batch: usize, include_deleted: bool, logger: Logger, @@ -398,6 +402,7 @@ fn run_into( let opts = RunOpts { source_id, + generation: 0, fetch_batch: config.fetch_batch.max(1), include_deleted: config.include_deleted, logger, @@ -422,13 +427,17 @@ fn run_into( if i > 0 { let _ = control_run_collect(&mut client, &control_ctx, "NOOP"); } + // Each folder gets a new generation; anything still queued or in + // flight for an earlier folder is skipped or dropped from here on. + let generation = i as u64 + 1; + pool.cancel_before(generation); match reconcile_folder( &mut conn, &mut client, &control_ctx, &pool, folder, - opts, + RunOpts { generation, ..opts }, &mut email_counts, ) { Ok(()) => {} @@ -695,6 +704,7 @@ fn reconcile_folder( ) -> Result<(), Error> { let RunOpts { source_id, + generation, fetch_batch, include_deleted: _, logger, @@ -805,6 +815,7 @@ fn reconcile_folder( let n_batches = batches.len(); for batch in &batches { pool.submit(FetchJob { + generation, folder: folder.name.clone(), wire_name: folder.wire_name.clone(), uidvalidity, @@ -813,12 +824,16 @@ fn reconcile_folder( } let keepalive_interval = std::time::Duration::from_secs(45); let mut chunks_done: usize = 0; - let tx = conn.transaction()?; let target = FetchTarget { folder: folder.name.as_str(), uidvalidity, mailbox_local, }; + // Committed after every chunk, not once per folder: each message is + // written whole (see `insert_recording_failure`), so what is + // committed is always consistent, a crash keeps it, and the next run + // fetches only the UIDs that are still missing. + let mut tx = conn.transaction()?; while chunks_done < n_batches { let event = loop { match pool.recv_timeout(keepalive_interval) { @@ -831,15 +846,15 @@ fn reconcile_folder( } } }; - match event { - FetchEvent::Item { attrs, .. } => { - insert_single_message(&tx, &target, &attrs, opts, counts)?; + match route_event(event, generation) { + Routed::Stale => {} + Routed::Item(attrs) => { + insert_recording_failure(&mut tx, &target, &attrs, opts, counts)?; } - FetchEvent::ChunkDone { + Routed::ChunkDone { folder: chunk_folder, - outcome, uids_requested, - .. + outcome, } => { chunks_done += 1; if let Err(e) = outcome { @@ -850,6 +865,8 @@ fn reconcile_folder( ); counts.failed += uids_requested.len() as u64; } + tx.commit()?; + tx = conn.transaction()?; } } } @@ -920,6 +937,78 @@ fn delete_vanished_emails( Ok(()) } +pub(super) enum Routed { + /// From an earlier folder generation: dropped, never filed here. + Stale, + Item(fetch::FetchAttrs), + ChunkDone { + folder: String, + uids_requested: Vec, + outcome: Result<(), ImapError>, + }, +} + +pub(super) fn route_event(event: FetchEvent, generation: u64) -> Routed { + match event { + FetchEvent::Item { + generation: g, + attrs, + .. + } if g == generation => Routed::Item(attrs), + FetchEvent::ChunkDone { + generation: g, + folder, + uids_requested, + outcome, + .. + } if g == generation => Routed::ChunkDone { + folder, + uids_requested, + outcome, + }, + _ => Routed::Stale, + } +} + +/// Writes one message inside its own savepoint. A message that cannot be +/// imported -- an INTERNALDATE that will not parse, say -- is rolled back +/// on its own, logged with its folder and UID, counted as failed, and the +/// folder carries on; it stays out of the UID map, so the next run tries +/// it again. Archive and I/O errors still stop the run. +fn insert_recording_failure( + tx: &mut rusqlite::Transaction<'_>, + target: &FetchTarget<'_>, + attrs: &fetch::FetchAttrs, + opts: RunOpts, + counts: &mut TypeCounts, +) -> Result<(), Error> { + let sp = tx.savepoint()?; + match insert_single_message(&sp, target, attrs, opts, counts) { + Ok(()) => { + sp.commit()?; + Ok(()) + } + Err(e) if e.aborts_run() => Err(e), + Err(e) => { + drop(sp); + log_at( + opts.logger, + LEVEL_DEFAULT, + &format!( + "folder {:?} uid {}: not imported: {e}", + target.folder, + attrs + .uid + .map(|u| u.to_string()) + .unwrap_or_else(|| "?".to_owned()) + ), + ); + counts.failed += 1; + Ok(()) + } + } +} + pub(super) struct FetchTarget<'a> { pub folder: &'a str, pub uidvalidity: u32, @@ -927,7 +1016,7 @@ pub(super) struct FetchTarget<'a> { } fn insert_single_message( - tx: &rusqlite::Transaction<'_>, + tx: &Connection, target: &FetchTarget<'_>, attrs: &fetch::FetchAttrs, opts: RunOpts, @@ -935,6 +1024,7 @@ fn insert_single_message( ) -> Result<(), Error> { let RunOpts { source_id, + generation: _, fetch_batch: _, include_deleted, logger, @@ -1012,6 +1102,7 @@ fn refresh_present_flags( ) -> Result { let RunOpts { source_id, + generation: _, fetch_batch, include_deleted, logger: _, @@ -1237,6 +1328,36 @@ fn dry_run_summary( mod tests { use super::*; + fn item(generation: u64, uid: u32) -> FetchEvent { + FetchEvent::Item { + generation, + folder: "F".to_owned(), + uidvalidity: 1, + attrs: fetch::FetchAttrs { + uid: Some(uid), + ..Default::default() + }, + } + } + + #[test] + fn events_from_another_generation_are_dropped() { + assert!(matches!(route_event(item(1, 5), 2), Routed::Stale)); + assert!(matches!(route_event(item(3, 5), 2), Routed::Stale)); + match route_event(item(2, 5), 2) { + Routed::Item(attrs) => assert_eq!(attrs.uid, Some(5)), + _ => panic!("current-generation item was not routed"), + } + let stale_done = FetchEvent::ChunkDone { + generation: 1, + folder: "Old".to_owned(), + uidvalidity: 1, + uids_requested: vec![5], + outcome: Ok(()), + }; + assert!(matches!(route_event(stale_done, 2), Routed::Stale)); + } + #[test] fn parse_endpoint_imaps_defaults_to_993() { let e = parse_endpoint("imaps://mail.example.com").unwrap(); diff --git a/src/sync/import_imap/internaldate.rs b/src/sync/import_imap/internaldate.rs index 039a7dd..6ac8093 100644 --- a/src/sync/import_imap/internaldate.rs +++ b/src/sync/import_imap/internaldate.rs @@ -1,5 +1,6 @@ /* * SPDX-FileCopyrightText: 2020 Stalwart Labs LLC + * SPDX-FileCopyrightText: 2026 John Coffey * * SPDX-License-Identifier: Apache-2.0 OR MIT */ @@ -58,20 +59,22 @@ pub fn imap_internaldate_to_rfc3339(s: &str) -> Result { Ok(out) } +// RFC 3501 spells the month "Jan", but servers are not all that careful, +// and a date is not worth losing a message over: match any case. fn month_to_num(s: &str) -> Result { - let m = match s { - "Jan" => 1, - "Feb" => 2, - "Mar" => 3, - "Apr" => 4, - "May" => 5, - "Jun" => 6, - "Jul" => 7, - "Aug" => 8, - "Sep" => 9, - "Oct" => 10, - "Nov" => 11, - "Dec" => 12, + let m = match s.to_ascii_lowercase().as_str() { + "jan" => 1, + "feb" => 2, + "mar" => 3, + "apr" => 4, + "may" => 5, + "jun" => 6, + "jul" => 7, + "aug" => 8, + "sep" => 9, + "oct" => 10, + "nov" => 11, + "dec" => 12, other => return Err(Error::Partial(format!("INTERNALDATE month {other:?}"))), }; Ok(m) @@ -99,6 +102,22 @@ fn parse_zone(s: &str) -> Result<(char, u32, u32), Error> { mod tests { use super::*; + #[test] + fn month_matches_any_case() { + for d in [ + "12-May-2025 10:00:00 +0000", + "12-may-2025 10:00:00 +0000", + "12-MAY-2025 10:00:00 +0000", + ] { + assert_eq!( + imap_internaldate_to_rfc3339(d).unwrap(), + "2025-05-12T10:00:00Z", + "{d}" + ); + } + assert!(imap_internaldate_to_rfc3339("12-Mai-2025 10:00:00 +0000").is_err()); + } + #[test] fn utc_zone_becomes_z() { assert_eq!( diff --git a/src/sync/import_imap/pool.rs b/src/sync/import_imap/pool.rs index ff9900f..ecfddfd 100644 --- a/src/sync/import_imap/pool.rs +++ b/src/sync/import_imap/pool.rs @@ -1,11 +1,13 @@ /* * SPDX-FileCopyrightText: 2020 Stalwart Labs LLC + * SPDX-FileCopyrightText: 2026 John Coffey * * SPDX-License-Identifier: Apache-2.0 OR MIT */ use std::panic::{AssertUnwindSafe, catch_unwind}; use std::sync::Arc; +use std::sync::atomic::{AtomicU64, Ordering}; use std::thread; use crossbeam_channel::{Receiver, Sender, unbounded}; @@ -23,7 +25,12 @@ use super::fetch::FetchAttrs; pub const HARD_CAP: usize = 8; +/// A job or event from an older folder generation than the one the +/// coordinator is working on belongs to a folder it has already given up +/// on. Workers skip such jobs without fetching, and the coordinator drops +/// such events, so nothing from one folder can be filed into the next. pub struct FetchJob { + pub generation: u64, pub folder: String, pub wire_name: String, pub uidvalidity: u32, @@ -32,11 +39,13 @@ pub struct FetchJob { pub enum FetchEvent { Item { + generation: u64, folder: String, uidvalidity: u32, attrs: FetchAttrs, }, ChunkDone { + generation: u64, folder: String, uidvalidity: u32, uids_requested: Vec, @@ -59,6 +68,7 @@ pub struct WorkerPool { job_tx: Sender, event_rx: Receiver, handles: Vec>, + cancel_below: Arc, } impl WorkerPool { @@ -68,13 +78,15 @@ impl WorkerPool { let (event_tx, event_rx) = unbounded::(); let mut handles = Vec::with_capacity(size); let args = Arc::new(args); + let cancel_below = Arc::new(AtomicU64::new(0)); for _ in 0..size { let args = args.clone(); let job_rx = job_rx.clone(); let event_tx = event_tx.clone(); + let cancel_below = cancel_below.clone(); let handle = thread::spawn(move || { - worker_loop(args, job_rx, event_tx); + worker_loop(args, job_rx, event_tx, cancel_below); }); handles.push(handle); } @@ -83,9 +95,17 @@ impl WorkerPool { job_tx, event_rx, handles, + cancel_below, }) } + /// Jobs of any generation below `generation` are skipped from now on: + /// called when the coordinator moves to a new folder, so work still + /// queued for one it abandoned is not fetched. + pub fn cancel_before(&self, generation: u64) { + self.cancel_below.fetch_max(generation, Ordering::SeqCst); + } + pub fn submit(&self, job: FetchJob) { let _ = self.job_tx.send(job); } @@ -101,21 +121,43 @@ impl WorkerPool { self.event_rx.recv_timeout(timeout) } + /// Stops the workers. Events still in flight are drained and dropped + /// first: a worker blocked handing over an event the coordinator will + /// never read (after a folder was abandoned) would otherwise never + /// finish, and joining it would hang. pub fn shutdown(self) { + self.cancel_before(u64::MAX); drop(self.job_tx); + while self.event_rx.recv().is_ok() {} for h in self.handles { let _ = h.join(); } } } -fn worker_loop(args: Arc, job_rx: Receiver, event_tx: Sender) { +fn worker_loop( + args: Arc, + job_rx: Receiver, + event_tx: Sender, + cancel_below: Arc, +) { let mut client: Option = None; let mut current_folder: Option = None; while let Ok(job) = job_rx.recv() { + let job_gen = job.generation; let job_folder = job.folder.clone(); let job_uv = job.uidvalidity; let job_uids = job.uids.clone(); + if job_gen < cancel_below.load(Ordering::SeqCst) { + let _ = event_tx.send(FetchEvent::ChunkDone { + generation: job_gen, + folder: job_folder, + uidvalidity: job_uv, + uids_requested: job_uids, + outcome: Err(ImapError::Protocol("cancelled: folder abandoned".into())), + }); + continue; + } let event_tx_for_job = event_tx.clone(); let outcome = match catch_unwind(AssertUnwindSafe(|| { run_job_with_retry( @@ -134,6 +176,7 @@ fn worker_loop(args: Arc, job_rx: Receiver, event_tx: Send } }; let _ = event_tx.send(FetchEvent::ChunkDone { + generation: job_gen, folder: job_folder, uidvalidity: job_uv, uids_requested: job_uids, @@ -235,6 +278,7 @@ fn run_one_job( *current_folder = Some(job.folder.clone()); } let set = command::format_uid_set(&job.uids, true); + let generation = job.generation; let folder = job.folder.clone(); let uv = job.uidvalidity; client.run_streamed( @@ -247,6 +291,7 @@ fn run_one_job( && let Some(attrs) = super::fetch::extract(&u) { let _ = event_tx.send(FetchEvent::Item { + generation, folder: folder.clone(), uidvalidity: uv, attrs, @@ -256,3 +301,76 @@ fn run_one_job( )?; Ok(()) } + +#[cfg(test)] +mod tests { + use super::*; + use std::time::Duration; + + fn unreachable_args() -> WorkerArgs { + WorkerArgs { + connector: Arc::new(Connector::new(false).expect("connector")), + endpoint: Arc::new(Endpoint { + host: "127.0.0.1".to_owned(), + port: 1, + implicit_tls: false, + }), + mode: ConnectMode::Plain, + auth: ImapAuth::Basic { + user: "u".to_owned(), + password: "p".to_owned(), + }, + compress: false, + policy: RetryPolicy::new(0), + backoff: BackoffState::new(), + logger: Logger::from_flags(false, 0), + } + } + + #[test] + fn a_job_from_an_abandoned_generation_is_skipped_without_fetching() { + // Port 1 refuses connections: a job that were actually run would come + // back as a connection error, not as a cancellation. + let pool = WorkerPool::start(unreachable_args(), 1).expect("pool"); + pool.cancel_before(2); + pool.submit(FetchJob { + generation: 1, + folder: "Old".to_owned(), + wire_name: "Old".to_owned(), + uidvalidity: 7, + uids: vec![1, 2, 3], + }); + match pool.recv_timeout(Duration::from_secs(5)).expect("event") { + FetchEvent::ChunkDone { + generation, + folder, + uids_requested, + outcome, + .. + } => { + assert_eq!(generation, 1); + assert_eq!(folder, "Old"); + assert_eq!(uids_requested, vec![1, 2, 3]); + let err = outcome.expect_err("cancelled"); + assert!(err.to_string().contains("cancelled"), "{err}"); + } + FetchEvent::Item { .. } => panic!("a cancelled job fetched something"), + } + pool.shutdown(); + } + + #[test] + fn shutdown_returns_with_work_still_queued() { + let pool = WorkerPool::start(unreachable_args(), 2).expect("pool"); + for g in 0..20u64 { + pool.submit(FetchJob { + generation: g, + folder: format!("F{g}"), + wire_name: format!("F{g}"), + uidvalidity: 1, + uids: vec![1], + }); + } + pool.shutdown(); + } +} diff --git a/src/sync/import_maildir/messages.rs b/src/sync/import_maildir/messages.rs index c223b0a..e09c6e4 100644 --- a/src/sync/import_maildir/messages.rs +++ b/src/sync/import_maildir/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 */ @@ -140,7 +141,7 @@ pub struct InsertContext<'a> { } pub fn insert_new( - tx: &rusqlite::Transaction<'_>, + tx: &Connection, ctx: InsertContext<'_>, entry: &DiskEntry, ) -> Result, InsertError> { @@ -268,6 +269,10 @@ pub fn delete_vanished( const PROGRESS_TICK: u64 = 1000; +/// New messages are committed in groups of this many rather than once per +/// folder, so an interrupted import of a large folder keeps what it wrote. +const COMMIT_EVERY: u64 = 500; + pub fn apply_folder( conn: &mut Connection, ctx: InsertContext<'_>, @@ -277,15 +282,23 @@ pub fn apply_folder( ) -> Result<(), crate::error::Error> { let stored_keywords = load_present_keywords(conn, ctx.source_id, ctx.folder).unwrap_or_default(); - let tx = conn.transaction()?; + let mut tx = conn.transaction()?; let total_new = diff.new.len() as u64; let mut inserted: u64 = 0; for entry in &diff.new { - match insert_new(&tx, ctx, entry) { + // Each message in its own savepoint: one that fails part-way leaves + // nothing behind, not an email row without its id mapping. + let sp = tx.savepoint()?; + match insert_new(&sp, ctx, entry) { Ok(Some(_)) => { + sp.commit()?; counts.created += 1; counts.fetched += 1; inserted += 1; + if inserted.is_multiple_of(COMMIT_EVERY) { + tx.commit()?; + tx = conn.transaction()?; + } if inserted.is_multiple_of(PROGRESS_TICK) && logger.enabled(crate::logging::LEVEL_PROGRESS) { @@ -296,9 +309,11 @@ pub fn apply_folder( } } Ok(None) => { + sp.commit()?; counts.skipped += 1; } Err(e) => { + drop(sp); logger.warn(&format!( "maildir {folder:?}/{name}: {e}", folder = ctx.folder, @@ -423,6 +438,43 @@ mod tests { } } + #[test] + fn a_message_that_cannot_be_read_is_recorded_and_the_rest_are_imported() { + let td = tempfile::tempdir().unwrap(); + ensure_folder_skel(td.path()); + write_maildir_message(td.path(), "cur", "1.M0.host:2,S", b"Subject: a\r\n\r\na"); + let gone = write_maildir_message(td.path(), "cur", "2.M0.host:2,S", b"Subject: b\r\n\r\nb"); + write_maildir_message(td.path(), "cur", "3.M0.host:2,S", b"Subject: c\r\n\r\nc"); + let listing = list_folder(td.path()).unwrap(); + // Gone between the listing and the read, as a file being moved is. + fs::remove_file(gone).unwrap(); + let (mut c, sid) = fresh_archive(); + let d = diff(listing.entries, &HashMap::new()); + let ctx = InsertContext { + source_id: sid, + folder: "INBOX", + mailbox_local: 1, + include_deleted: false, + }; + let mut counts = TypeCounts::default(); + apply_folder( + &mut c, + ctx, + d, + &mut counts, + crate::logging::Logger::from_flags(false, 0), + ) + .unwrap(); + assert_eq!((counts.created, counts.failed), (2, 1)); + let emails: i64 = c + .query_row("SELECT COUNT(*) FROM emails", [], |r| r.get(0)) + .unwrap(); + let mapped: i64 = c + .query_row("SELECT COUNT(*) FROM sync_id_maildir", [], |r| r.get(0)) + .unwrap(); + assert_eq!((emails, mapped), (2, 2), "no email row without its mapping"); + } + #[test] fn list_folder_returns_cur_and_new_skipping_tmp() { let td = tempfile::tempdir().unwrap(); diff --git a/tests/mock_imap.rs b/tests/mock_imap.rs index b3ab835..118fb03 100644 --- a/tests/mock_imap.rs +++ b/tests/mock_imap.rs @@ -2005,3 +2005,133 @@ fn assert_name_selected_as_listed(listed: &'static str, stored: &str, archive_na ); let _ = std::fs::remove_file(&archive); } + +fn write_fetch_message_dated( + conn: &mut MockConn, + seq: u32, + uid: u32, + internaldate: &str, + body: &[u8], +) -> std::io::Result<()> { + let header = format!( + "* {seq} FETCH (UID {uid} FLAGS (\\Seen) INTERNALDATE \"{internaldate}\" RFC822.SIZE {} BODY[] {{{}}}\r\n", + body.len(), + body.len() + ); + conn.write_raw(header.as_bytes())?; + conn.write_raw(body)?; + conn.write_raw(b")\r\n") +} + +const BODY_A: &[u8] = b"From: a@b\r\nMessage-ID: \r\nSubject: a\r\n\r\none"; +const BODY_B: &[u8] = b"From: a@b\r\nMessage-ID: \r\nSubject: b\r\n\r\ntwo"; +const BODY_C: &[u8] = b"From: a@b\r\nMessage-ID: \r\nSubject: c\r\n\r\nthree"; + +fn email_counts(summary: &inbuxa_migrate::sync::Summary) -> inbuxa_migrate::sync::TypeCounts { + summary + .per_type + .iter() + .find(|(k, _)| *k == "email") + .map(|(_, c)| c.clone()) + .expect("email counts") +} + +#[test] +fn a_message_that_will_not_import_is_recorded_and_the_folder_carries_on() { + 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, 12345, 4, 3)?; + let (tag, cmd) = conn.read_command()?; + assert!(cmd.starts_with("UID FETCH"), "got {cmd}"); + write_fetch_message_dated(conn, 1, 1, "12-May-2025 10:00:00 +0000", BODY_A)?; + // No such month: this one cannot be imported. + write_fetch_message_dated(conn, 2, 2, "12-Mai-2025 10:00:00 +0000", BODY_B)?; + // Lower case is only untidy, and is imported. + write_fetch_message_dated(conn, 3, 3, "12-may-2025 10:00:00 +0000", BODY_C)?; + conn.write_line(&format!("{tag} OK"))?; + drain_until_close(conn); + Ok(()) + }); + let server = MockImap::start_scripts(vec![ + control_script_one_folder(12345, 4, &[1, 2, 3]), + worker, + ]); + let archive = tempfile("bad-message"); + let summary = run_import(&server, "alice", archive.clone(), |_| {}).expect("import"); + let email = email_counts(&summary); + assert_eq!(email.created, 2, "summary={summary:?}"); + assert_eq!(email.failed, 1, "summary={summary:?}"); + let conn = Connection::open(&archive).unwrap(); + db::init::apply_schema(&conn).unwrap(); + assert_eq!(count(&conn, "emails"), 2); + let uids: Vec = conn + .prepare("SELECT uid FROM sync_id_imap WHERE type_name = 'email' ORDER BY uid") + .unwrap() + .query_map([], |r| r.get(0)) + .unwrap() + .collect::>() + .unwrap(); + assert_eq!( + uids, + vec![1, 3], + "the failed message must stay out of the UID map" + ); +} + +#[test] +fn a_failed_chunk_keeps_the_ones_before_it_and_a_rerun_fetches_only_what_is_missing() { + // The archive is opened with an exclusive lock, so the commit after each + // chunk cannot be watched from outside while a run is going; what can be + // checked is its effect. First run, one UID per chunk: chunk 1 arrives, + // chunk 2 fails the way a dying server would. + let archive = tempfile("chunk-commit"); + let worker_1: 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, 777, 3, 2)?; + let (tag, cmd) = conn.read_command()?; + assert!(cmd.starts_with("UID FETCH 1 "), "got {cmd}"); + write_fetch_message(conn, 1, 1, BODY_A)?; + conn.write_line(&format!("{tag} OK"))?; + let (tag, cmd) = conn.read_command()?; + assert!(cmd.starts_with("UID FETCH 2 "), "got {cmd}"); + conn.write_line(&format!("{tag} NO [SERVERBUG] gone"))?; + drain_until_close(conn); + Ok(()) + }); + // Second run: only UID 2 is missing, so only UID 2 may be fetched. + let worker_2: 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, 777, 3, 2)?; + let (tag, cmd) = conn.read_command()?; + assert!( + cmd.starts_with("UID FETCH 2 "), + "rerun refetched more than UID 2: {cmd}" + ); + write_fetch_message(conn, 2, 2, BODY_B)?; + conn.write_line(&format!("{tag} OK"))?; + drain_until_close(conn); + Ok(()) + }); + let server = MockImap::start_scripts(vec![ + control_script_one_folder(777, 3, &[1, 2]), + worker_1, + control_script_one_folder(777, 3, &[1, 2]), + worker_2, + ]); + + let first = + run_import(&server, "alice", archive.clone(), |c| c.fetch_batch = 1).expect("first run"); + let e1 = email_counts(&first); + assert_eq!((e1.created, e1.failed), (1, 1), "first={first:?}"); + + let second = + run_import(&server, "alice", archive.clone(), |c| c.fetch_batch = 1).expect("second run"); + let e2 = email_counts(&second); + assert_eq!((e2.created, e2.failed), (1, 0), "second={second:?}"); + let conn = Connection::open(&archive).unwrap(); + db::init::apply_schema(&conn).unwrap(); + assert_eq!(count(&conn, "emails"), 2); +} -- 2.54.0