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.
This commit is contained in:
@@ -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();
|
||||||
|
|||||||
@@ -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!(
|
||||||
|
|||||||
@@ -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();
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
@@ -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();
|
||||||
|
|||||||
@@ -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);
|
||||||
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user