pub mod dedup;
pub mod manifest;
pub mod sink;
pub mod transform;
pub mod wizard;
mod worker;
use std::collections::HashSet;
use std::fs;
use std::path::PathBuf;
use async_imap::types::NameAttribute;
use futures::TryStreamExt;
use indicatif::MultiProgress;
use crate::commands::keyring::bucket::store::BucketConfig;
use crate::commands::keyring::email::identity::Identity;
use crate::core::crypto::Aes256GcmSivEncryptor;
use crate::core::job::Job;
use manifest::Batch;
use worker::JobSummary;
pub(crate) struct IdentityContext {
pub identity: Identity,
pub secret: String,
pub staging_dir: PathBuf,
pub output_dir: PathBuf,
}
#[derive(Debug, Default)]
pub(crate) struct IdentityManifestSummary {
pub alias: String,
pub mailboxes: usize,
pub pending_messages: usize,
pub pending_attachments: usize,
pub pending_bytes: u64,
}
pub(crate) struct PendingMailbox {
pub mailbox: String,
pub mailbox_relpath: PathBuf,
pub uids: Vec<u32>,
}
#[tracing::instrument(skip(ctx, multi_progress), fields(identity = %ctx.identity.alias, email = %ctx.identity.email))]
pub(crate) async fn gather_pending(
ctx: &IdentityContext,
multi_progress: &MultiProgress,
) -> Result<(Vec<PendingMailbox>, IdentityManifestSummary), String> {
let _ = multi_progress.println(format!("Connecting to {}...", ctx.identity.alias));
let mut session = worker::connect_with_retry(ctx).await?;
let names: Vec<_> = session
.list(None, Some("*"))
.await
.map_err(|err| format!("failed to list mailboxes: {err}"))?
.try_collect()
.await
.map_err(|err| format!("failed to list mailboxes: {err}"))?;
let mailboxes: Vec<(String, Option<String>)> = names
.iter()
.filter(|name| !name.attributes().contains(&NameAttribute::NoSelect))
.map(|name| {
(
name.name().to_string(),
name.delimiter().map(str::to_string),
)
})
.collect();
let checkpoint_entries = manifest::load_checkpoint(&ctx.staging_dir)?;
let done = manifest::done_uids(&checkpoint_entries);
let mut pending_mailboxes = Vec::new();
let mut fresh_manifest = Vec::new();
let mut summary = IdentityManifestSummary {
alias: ctx.identity.alias.clone(),
..Default::default()
};
let bar = sink::new_progress_bar(
format!("{} manifest", ctx.identity.alias),
mailboxes.len() as u64,
multi_progress,
);
for (mailbox_name, delimiter) in &mailboxes {
let mailbox_relpath = sink::sanitize_mailbox_path(mailbox_name, delimiter.as_deref());
let mailbox_dir = ctx.staging_dir.join(&mailbox_relpath);
fs::create_dir_all(&mailbox_dir)
.map_err(|err| format!("failed to create {}: {err}", mailbox_dir.display()))?;
let mailbox_response = session
.examine(mailbox_name)
.await
.map_err(|err| format!("failed to open '{mailbox_name}' read-only: {err}"))?;
let current_validity = mailbox_response.uid_validity.unwrap_or(0);
if sink::is_stale(sink::read_uidvalidity(&mailbox_dir), current_validity) {
sink::clear_eml_files(&mailbox_dir)?;
manifest::clear_checkpoint_for_mailbox(&ctx.staging_dir, mailbox_name)?;
}
sink::write_uidvalidity(&mailbox_dir, current_validity)?;
let server_uids: HashSet<u32> = session
.uid_search("ALL")
.await
.map_err(|err| format!("failed to search '{mailbox_name}': {err}"))?;
let done_here: HashSet<u32> = done
.iter()
.filter(|(mailbox, _)| mailbox == mailbox_name)
.map(|(_, uid)| *uid)
.collect();
let pending_uids = sink::missing_uids(&server_uids, &done_here);
if pending_uids.is_empty() {
bar.inc(1);
continue;
}
let mailbox_manifest =
manifest::pull_manifest(&mut session, mailbox_name, &pending_uids).await?;
summary.mailboxes += 1;
summary.pending_messages += mailbox_manifest.len();
summary.pending_bytes += mailbox_manifest.iter().map(|entry| entry.size).sum::<u64>();
summary.pending_attachments += mailbox_manifest
.iter()
.map(|entry| entry.attachments as usize)
.sum::<usize>();
fresh_manifest.extend(mailbox_manifest);
pending_mailboxes.push(PendingMailbox {
mailbox: mailbox_name.clone(),
mailbox_relpath,
uids: pending_uids,
});
bar.inc(1);
}
bar.finish();
let _ = session.logout().await;
manifest::save_manifest(&ctx.staging_dir, &fresh_manifest)?;
Ok((pending_mailboxes, summary))
}
pub(crate) fn batches_from_pending(pending: &[PendingMailbox], concurrency: usize) -> Vec<Batch> {
pending
.iter()
.flat_map(|mailbox| {
manifest::split_into_batches(&mailbox.uids, concurrency)
.into_iter()
.map(|uids| Batch {
mailbox: mailbox.mailbox.clone(),
mailbox_relpath: mailbox.mailbox_relpath.clone(),
uids,
})
})
.collect()
}
pub(crate) const DEFAULT_MAX_CONNECTIONS_PER_IDENTITY: usize = 6;
pub(crate) struct EmailSyncJob {
pub contexts: Vec<IdentityContext>,
pub remote: Option<(BucketConfig, String)>,
pub encryptor: Option<Aes256GcmSivEncryptor>,
pub max_connections_per_identity: usize,
}
pub(crate) struct EmailSyncPlan {
pub pending_by_identity: Vec<Vec<PendingMailbox>>,
pub manifest_summaries: Vec<IdentityManifestSummary>,
}
impl Job for EmailSyncJob {
type Plan = EmailSyncPlan;
type Summary = JobSummary;
async fn gather(&self) -> Result<EmailSyncPlan, String> {
let multi_progress = MultiProgress::new();
let mut pending_by_identity = Vec::with_capacity(self.contexts.len());
let mut manifest_summaries = Vec::with_capacity(self.contexts.len());
for ctx in &self.contexts {
let (pending, summary) = gather_pending(ctx, &multi_progress).await?;
pending_by_identity.push(pending);
manifest_summaries.push(summary);
}
Ok(EmailSyncPlan {
pending_by_identity,
manifest_summaries,
})
}
async fn run(
self,
plan: EmailSyncPlan,
concurrency: usize,
upload_concurrency: usize,
) -> Result<JobSummary, String> {
let EmailSyncJob {
contexts,
remote,
encryptor,
max_connections_per_identity,
} = self;
let remote_ref = remote
.as_ref()
.map(|(bucket_config, secret)| (bucket_config, secret.as_str()));
worker::run_email_sync_job(
contexts,
plan.pending_by_identity,
concurrency,
upload_concurrency,
max_connections_per_identity,
remote_ref,
encryptor.as_ref(),
)
.await
}
}