import: keep what an interrupted IMAP or Maildir import wrote #8

Merged
jcoffey-dev merged 1 commits from fix/imap-import-resilience into main 2026-09-30 19:33:18 +00:00
5 changed files with 467 additions and 27 deletions
+130 -9
View File
@@ -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
*/ */
@@ -45,6 +46,9 @@ const EMAIL_TYPE: &str = "email";
#[derive(Clone, Copy)] #[derive(Clone, Copy)]
pub(super) struct RunOpts { pub(super) struct RunOpts {
source_id: i64, source_id: i64,
/// The folder generation this folder's fetch jobs carry; see
/// `WorkerPool::cancel_before`.
generation: u64,
fetch_batch: usize, fetch_batch: usize,
include_deleted: bool, include_deleted: bool,
logger: Logger, logger: Logger,
@@ -398,6 +402,7 @@ fn run_into(
let opts = RunOpts { let opts = RunOpts {
source_id, source_id,
generation: 0,
fetch_batch: config.fetch_batch.max(1), fetch_batch: config.fetch_batch.max(1),
include_deleted: config.include_deleted, include_deleted: config.include_deleted,
logger, logger,
@@ -422,13 +427,17 @@ fn run_into(
if i > 0 { if i > 0 {
let _ = control_run_collect(&mut client, &control_ctx, "NOOP"); 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( match reconcile_folder(
&mut conn, &mut conn,
&mut client, &mut client,
&control_ctx, &control_ctx,
&pool, &pool,
folder, folder,
opts, RunOpts { generation, ..opts },
&mut email_counts, &mut email_counts,
) { ) {
Ok(()) => {} Ok(()) => {}
@@ -695,6 +704,7 @@ fn reconcile_folder(
) -> Result<(), Error> { ) -> Result<(), Error> {
let RunOpts { let RunOpts {
source_id, source_id,
generation,
fetch_batch, fetch_batch,
include_deleted: _, include_deleted: _,
logger, logger,
@@ -805,6 +815,7 @@ fn reconcile_folder(
let n_batches = batches.len(); let n_batches = batches.len();
for batch in &batches { for batch in &batches {
pool.submit(FetchJob { pool.submit(FetchJob {
generation,
folder: folder.name.clone(), folder: folder.name.clone(),
wire_name: folder.wire_name.clone(), wire_name: folder.wire_name.clone(),
uidvalidity, uidvalidity,
@@ -813,12 +824,16 @@ fn reconcile_folder(
} }
let keepalive_interval = std::time::Duration::from_secs(45); let keepalive_interval = std::time::Duration::from_secs(45);
let mut chunks_done: usize = 0; let mut chunks_done: usize = 0;
let tx = conn.transaction()?;
let target = FetchTarget { let target = FetchTarget {
folder: folder.name.as_str(), folder: folder.name.as_str(),
uidvalidity, uidvalidity,
mailbox_local, 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 { while chunks_done < n_batches {
let event = loop { let event = loop {
match pool.recv_timeout(keepalive_interval) { match pool.recv_timeout(keepalive_interval) {
@@ -831,15 +846,15 @@ fn reconcile_folder(
} }
} }
}; };
match event { match route_event(event, generation) {
FetchEvent::Item { attrs, .. } => { Routed::Stale => {}
insert_single_message(&tx, &target, &attrs, opts, counts)?; Routed::Item(attrs) => {
insert_recording_failure(&mut tx, &target, &attrs, opts, counts)?;
} }
FetchEvent::ChunkDone { Routed::ChunkDone {
folder: chunk_folder, folder: chunk_folder,
outcome,
uids_requested, uids_requested,
.. outcome,
} => { } => {
chunks_done += 1; chunks_done += 1;
if let Err(e) = outcome { if let Err(e) = outcome {
@@ -850,6 +865,8 @@ fn reconcile_folder(
); );
counts.failed += uids_requested.len() as u64; counts.failed += uids_requested.len() as u64;
} }
tx.commit()?;
tx = conn.transaction()?;
} }
} }
} }
@@ -920,6 +937,78 @@ fn delete_vanished_emails(
Ok(()) Ok(())
} }
pub(super) enum Routed {
/// From an earlier folder generation: dropped, never filed here.
Stale,
Item(fetch::FetchAttrs),
ChunkDone {
folder: String,
uids_requested: Vec<u32>,
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(super) struct FetchTarget<'a> {
pub folder: &'a str, pub folder: &'a str,
pub uidvalidity: u32, pub uidvalidity: u32,
@@ -927,7 +1016,7 @@ pub(super) struct FetchTarget<'a> {
} }
fn insert_single_message( fn insert_single_message(
tx: &rusqlite::Transaction<'_>, tx: &Connection,
target: &FetchTarget<'_>, target: &FetchTarget<'_>,
attrs: &fetch::FetchAttrs, attrs: &fetch::FetchAttrs,
opts: RunOpts, opts: RunOpts,
@@ -935,6 +1024,7 @@ fn insert_single_message(
) -> Result<(), Error> { ) -> Result<(), Error> {
let RunOpts { let RunOpts {
source_id, source_id,
generation: _,
fetch_batch: _, fetch_batch: _,
include_deleted, include_deleted,
logger, logger,
@@ -1012,6 +1102,7 @@ fn refresh_present_flags(
) -> Result<u64, Error> { ) -> Result<u64, Error> {
let RunOpts { let RunOpts {
source_id, source_id,
generation: _,
fetch_batch, fetch_batch,
include_deleted, include_deleted,
logger: _, logger: _,
@@ -1237,6 +1328,36 @@ fn dry_run_summary(
mod tests { mod tests {
use super::*; 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] #[test]
fn parse_endpoint_imaps_defaults_to_993() { fn parse_endpoint_imaps_defaults_to_993() {
let e = parse_endpoint("imaps://mail.example.com").unwrap(); let e = parse_endpoint("imaps://mail.example.com").unwrap();
+32 -13
View File
@@ -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
*/ */
@@ -58,20 +59,22 @@ pub fn imap_internaldate_to_rfc3339(s: &str) -> Result<String, Error> {
Ok(out) 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<u32, Error> { fn month_to_num(s: &str) -> Result<u32, Error> {
let m = match s { let m = match s.to_ascii_lowercase().as_str() {
"Jan" => 1, "jan" => 1,
"Feb" => 2, "feb" => 2,
"Mar" => 3, "mar" => 3,
"Apr" => 4, "apr" => 4,
"May" => 5, "may" => 5,
"Jun" => 6, "jun" => 6,
"Jul" => 7, "jul" => 7,
"Aug" => 8, "aug" => 8,
"Sep" => 9, "sep" => 9,
"Oct" => 10, "oct" => 10,
"Nov" => 11, "nov" => 11,
"Dec" => 12, "dec" => 12,
other => return Err(Error::Partial(format!("INTERNALDATE month {other:?}"))), other => return Err(Error::Partial(format!("INTERNALDATE month {other:?}"))),
}; };
Ok(m) Ok(m)
@@ -99,6 +102,22 @@ fn parse_zone(s: &str) -> Result<(char, u32, u32), Error> {
mod tests { mod tests {
use super::*; 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] #[test]
fn utc_zone_becomes_z() { fn utc_zone_becomes_z() {
assert_eq!( assert_eq!(
+120 -2
View File
@@ -1,11 +1,13 @@
/* /*
* 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
*/ */
use std::panic::{AssertUnwindSafe, catch_unwind}; use std::panic::{AssertUnwindSafe, catch_unwind};
use std::sync::Arc; use std::sync::Arc;
use std::sync::atomic::{AtomicU64, Ordering};
use std::thread; use std::thread;
use crossbeam_channel::{Receiver, Sender, unbounded}; use crossbeam_channel::{Receiver, Sender, unbounded};
@@ -23,7 +25,12 @@ use super::fetch::FetchAttrs;
pub const HARD_CAP: usize = 8; 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 struct FetchJob {
pub generation: u64,
pub folder: String, pub folder: String,
pub wire_name: String, pub wire_name: String,
pub uidvalidity: u32, pub uidvalidity: u32,
@@ -32,11 +39,13 @@ pub struct FetchJob {
pub enum FetchEvent { pub enum FetchEvent {
Item { Item {
generation: u64,
folder: String, folder: String,
uidvalidity: u32, uidvalidity: u32,
attrs: FetchAttrs, attrs: FetchAttrs,
}, },
ChunkDone { ChunkDone {
generation: u64,
folder: String, folder: String,
uidvalidity: u32, uidvalidity: u32,
uids_requested: Vec<u32>, uids_requested: Vec<u32>,
@@ -59,6 +68,7 @@ pub struct WorkerPool {
job_tx: Sender<FetchJob>, job_tx: Sender<FetchJob>,
event_rx: Receiver<FetchEvent>, event_rx: Receiver<FetchEvent>,
handles: Vec<thread::JoinHandle<()>>, handles: Vec<thread::JoinHandle<()>>,
cancel_below: Arc<AtomicU64>,
} }
impl WorkerPool { impl WorkerPool {
@@ -68,13 +78,15 @@ impl WorkerPool {
let (event_tx, event_rx) = unbounded::<FetchEvent>(); let (event_tx, event_rx) = unbounded::<FetchEvent>();
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));
for _ in 0..size { for _ in 0..size {
let args = args.clone(); let args = args.clone();
let job_rx = job_rx.clone(); let job_rx = job_rx.clone();
let event_tx = event_tx.clone(); let event_tx = event_tx.clone();
let cancel_below = cancel_below.clone();
let handle = thread::spawn(move || { let handle = thread::spawn(move || {
worker_loop(args, job_rx, event_tx); worker_loop(args, job_rx, event_tx, cancel_below);
}); });
handles.push(handle); handles.push(handle);
} }
@@ -83,9 +95,17 @@ impl WorkerPool {
job_tx, job_tx,
event_rx, event_rx,
handles, 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) { pub fn submit(&self, job: FetchJob) {
let _ = self.job_tx.send(job); let _ = self.job_tx.send(job);
} }
@@ -101,21 +121,43 @@ impl WorkerPool {
self.event_rx.recv_timeout(timeout) 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) { pub fn shutdown(self) {
self.cancel_before(u64::MAX);
drop(self.job_tx); drop(self.job_tx);
while self.event_rx.recv().is_ok() {}
for h in self.handles { for h in self.handles {
let _ = h.join(); let _ = h.join();
} }
} }
} }
fn worker_loop(args: Arc<WorkerArgs>, job_rx: Receiver<FetchJob>, event_tx: Sender<FetchEvent>) { fn worker_loop(
args: Arc<WorkerArgs>,
job_rx: Receiver<FetchJob>,
event_tx: Sender<FetchEvent>,
cancel_below: Arc<AtomicU64>,
) {
let mut client: Option<ImapClient> = None; let mut client: Option<ImapClient> = None;
let mut current_folder: Option<String> = None; let mut current_folder: Option<String> = None;
while let Ok(job) = job_rx.recv() { while let Ok(job) = job_rx.recv() {
let job_gen = job.generation;
let job_folder = job.folder.clone(); let job_folder = job.folder.clone();
let job_uv = job.uidvalidity; let job_uv = job.uidvalidity;
let job_uids = job.uids.clone(); 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 event_tx_for_job = event_tx.clone();
let outcome = match catch_unwind(AssertUnwindSafe(|| { let outcome = match catch_unwind(AssertUnwindSafe(|| {
run_job_with_retry( run_job_with_retry(
@@ -134,6 +176,7 @@ fn worker_loop(args: Arc<WorkerArgs>, job_rx: Receiver<FetchJob>, event_tx: Send
} }
}; };
let _ = event_tx.send(FetchEvent::ChunkDone { let _ = event_tx.send(FetchEvent::ChunkDone {
generation: job_gen,
folder: job_folder, folder: job_folder,
uidvalidity: job_uv, uidvalidity: job_uv,
uids_requested: job_uids, uids_requested: job_uids,
@@ -235,6 +278,7 @@ fn run_one_job(
*current_folder = Some(job.folder.clone()); *current_folder = Some(job.folder.clone());
} }
let set = command::format_uid_set(&job.uids, true); let set = command::format_uid_set(&job.uids, true);
let generation = job.generation;
let folder = job.folder.clone(); let folder = job.folder.clone();
let uv = job.uidvalidity; let uv = job.uidvalidity;
client.run_streamed( client.run_streamed(
@@ -247,6 +291,7 @@ fn run_one_job(
&& let Some(attrs) = super::fetch::extract(&u) && let Some(attrs) = super::fetch::extract(&u)
{ {
let _ = event_tx.send(FetchEvent::Item { let _ = event_tx.send(FetchEvent::Item {
generation,
folder: folder.clone(), folder: folder.clone(),
uidvalidity: uv, uidvalidity: uv,
attrs, attrs,
@@ -256,3 +301,76 @@ fn run_one_job(
)?; )?;
Ok(()) 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();
}
}
+55 -3
View File
@@ -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
*/ */
@@ -140,7 +141,7 @@ pub struct InsertContext<'a> {
} }
pub fn insert_new( pub fn insert_new(
tx: &rusqlite::Transaction<'_>, tx: &Connection,
ctx: InsertContext<'_>, ctx: InsertContext<'_>,
entry: &DiskEntry, entry: &DiskEntry,
) -> Result<Option<i64>, InsertError> { ) -> Result<Option<i64>, InsertError> {
@@ -268,6 +269,10 @@ pub fn delete_vanished(
const PROGRESS_TICK: u64 = 1000; 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( pub fn apply_folder(
conn: &mut Connection, conn: &mut Connection,
ctx: InsertContext<'_>, ctx: InsertContext<'_>,
@@ -277,15 +282,23 @@ pub fn apply_folder(
) -> Result<(), crate::error::Error> { ) -> Result<(), crate::error::Error> {
let stored_keywords = let stored_keywords =
load_present_keywords(conn, ctx.source_id, ctx.folder).unwrap_or_default(); 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 total_new = diff.new.len() as u64;
let mut inserted: u64 = 0; let mut inserted: u64 = 0;
for entry in &diff.new { 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(_)) => { Ok(Some(_)) => {
sp.commit()?;
counts.created += 1; counts.created += 1;
counts.fetched += 1; counts.fetched += 1;
inserted += 1; inserted += 1;
if inserted.is_multiple_of(COMMIT_EVERY) {
tx.commit()?;
tx = conn.transaction()?;
}
if inserted.is_multiple_of(PROGRESS_TICK) if inserted.is_multiple_of(PROGRESS_TICK)
&& logger.enabled(crate::logging::LEVEL_PROGRESS) && logger.enabled(crate::logging::LEVEL_PROGRESS)
{ {
@@ -296,9 +309,11 @@ pub fn apply_folder(
} }
} }
Ok(None) => { Ok(None) => {
sp.commit()?;
counts.skipped += 1; counts.skipped += 1;
} }
Err(e) => { Err(e) => {
drop(sp);
logger.warn(&format!( logger.warn(&format!(
"maildir {folder:?}/{name}: {e}", "maildir {folder:?}/{name}: {e}",
folder = ctx.folder, 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] #[test]
fn list_folder_returns_cur_and_new_skipping_tmp() { fn list_folder_returns_cur_and_new_skipping_tmp() {
let td = tempfile::tempdir().unwrap(); let td = tempfile::tempdir().unwrap();
+130
View File
@@ -2005,3 +2005,133 @@ fn assert_name_selected_as_listed(listed: &'static str, stored: &str, archive_na
); );
let _ = std::fs::remove_file(&archive); 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: <a@h>\r\nSubject: a\r\n\r\none";
const BODY_B: &[u8] = b"From: a@b\r\nMessage-ID: <b@h>\r\nSubject: b\r\n\r\ntwo";
const BODY_C: &[u8] = b"From: a@b\r\nMessage-ID: <c@h>\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<i64> = conn
.prepare("SELECT uid FROM sync_id_imap WHERE type_name = 'email' ORDER BY uid")
.unwrap()
.query_map([], |r| r.get(0))
.unwrap()
.collect::<Result<_, _>>()
.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);
}