Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
135 changes: 94 additions & 41 deletions crates/flectar-mail-core/src/jmap/sync.rs
Original file line number Diff line number Diff line change
Expand Up @@ -540,22 +540,51 @@ async fn sync_emails(ctx: &SyncCtx, config: &AccountConfig, c: &ConnectedClient)
.core_capabilities()
.map(|v| v.max_objects_in_get().clamp(1, 1000))
.unwrap_or(256);
let mut emails = Vec::new();
for chunk in ids.chunks(max) {
let mut request = c.client.build();
request
.get_email()
.account_id(&c.account_id)
.ids(chunk.iter().cloned())
.properties(header_properties());
let mut response: EmailGetResponse =
request.send_get_email().await.map_err(client::map_error)?;
destroyed.extend(response.take_not_found());
emails.extend(response.take_list());
// Each Email/get page is committed on its own while the next one is in
// flight: a first sync fills the mailbox newest first instead of showing
// nothing until the whole history is downloaded. The email state is only
// checkpointed once every page is stored, so an interrupted run restarts.
let mut remote_ids = HashSet::with_capacity(ids.len());
let mut chunks = ids.chunks(max);
let mut fetched = match chunks.next() {
Some(chunk) => Some(fetch_headers(c, chunk).await?),
None => None,
};
while let Some((emails, not_found)) = fetched.take() {
destroyed.extend(not_found);
remote_ids.extend(
emails
.iter()
.filter_map(|email| email.id().map(str::to_owned)),
);
let persist = persist_emails(ctx, config, emails);
fetched = match chunks.next() {
Some(chunk) => {
let (next, persisted) = tokio::join!(fetch_headers(c, chunk), persist);
persisted?;
Some(next?)
}
None => {
persist.await?;
None
}
};
}
destroyed.sort();
destroyed.dedup();
persist_emails(ctx, config, emails, destroyed, checkpoint_state, full).await
finish_email_sync(ctx, config, remote_ids, destroyed, checkpoint_state, full).await
}

async fn fetch_headers(c: &ConnectedClient, ids: &[String]) -> Result<(Vec<Email>, Vec<String>)> {
let mut request = c.client.build();
request
.get_email()
.account_id(&c.account_id)
.ids(ids.iter().cloned())
.properties(header_properties());
let mut response: EmailGetResponse =
request.send_get_email().await.map_err(client::map_error)?;
Ok((response.take_list(), response.take_not_found()))
}

fn addresses(values: Option<&[jmap_client::email::EmailAddress]>) -> Vec<Address> {
Expand Down Expand Up @@ -585,20 +614,8 @@ fn first(values: Option<&[String]>) -> Option<String> {
values.and_then(|values| values.first()).cloned()
}

async fn persist_emails(
ctx: &SyncCtx,
config: &AccountConfig,
emails: Vec<Email>,
destroyed: Vec<String>,
state: String,
full: bool,
) -> Result<()> {
async fn persist_emails(ctx: &SyncCtx, config: &AccountConfig, emails: Vec<Email>) -> Result<()> {
let account_id = config.id;
let cutoff = config.settings.mail_history.cutoff_ms_at(now_ms());
let remote_ids = emails
.iter()
.filter_map(|email| email.id().map(str::to_owned))
.collect::<HashSet<_>>();
let (thread_ids, stale_paths) = ctx
.db
.write(move |conn| {
Expand All @@ -614,20 +631,6 @@ async fn persist_emails(
let mut changed_threads = HashSet::new();
let mut stale_paths = Vec::new();

for remote_id in destroyed {
if let Some(row) = repo::messages::by_jmap_id(&tx, account_id, &remote_id)?
&& !repo::actions::has_pending_remote_creation(&tx, row.id)?
{
if let Some(thread_id) = row.thread_id {
changed_threads.insert(thread_id);
}
if let Some(path) = row.raw_path {
stale_paths.push(path);
}
repo::messages::delete(&tx, row.id)?;
}
}

for email in emails {
let remote_id = email
.id()
Expand Down Expand Up @@ -891,6 +894,48 @@ async fn persist_emails(
changed_threads.insert(thread_id);
}

for thread_id in &changed_threads {
repo::threads::recompute(&tx, *thread_id)?;
}
tx.commit()?;
Ok((changed_threads.into_iter().collect::<Vec<_>>(), stale_paths))
})
.await?;
publish_changes(ctx, thread_ids, stale_paths).await;
Ok(())
}

/// Applies remote deletions and full-sync reconciliation before checkpointing.
/// Waiting for every header page keeps surviving threads and their snoozes alive.
async fn finish_email_sync(
ctx: &SyncCtx,
config: &AccountConfig,
remote_ids: HashSet<String>,
destroyed: Vec<String>,
state: String,
full: bool,
) -> Result<()> {
let account_id = config.id;
let cutoff = config.settings.mail_history.cutoff_ms_at(now_ms());
let (thread_ids, stale_paths) = ctx
.db
.write(move |conn| {
let tx = conn.transaction()?;
let mut changed_threads = HashSet::new();
let mut stale_paths = Vec::new();
for remote_id in destroyed {
if let Some(row) = repo::messages::by_jmap_id(&tx, account_id, &remote_id)?
&& !repo::actions::has_pending_remote_creation(&tx, row.id)?
{
if let Some(thread_id) = row.thread_id {
changed_threads.insert(thread_id);
}
if let Some(path) = row.raw_path {
stale_paths.push(path);
}
repo::messages::delete(&tx, row.id)?;
}
}
if full {
let mut stmt = tx.prepare(
"SELECT id,thread_id,raw_path,jmap_id,date FROM messages
Expand Down Expand Up @@ -937,13 +982,17 @@ async fn persist_emails(
Ok((changed_threads.into_iter().collect::<Vec<_>>(), stale_paths))
})
.await?;
publish_changes(ctx, thread_ids, stale_paths).await;
Ok(())
}

async fn publish_changes(ctx: &SyncCtx, thread_ids: Vec<i64>, stale_paths: Vec<String>) {
for path in stale_paths {
let _ = tokio::fs::remove_file(path).await;
}
if !thread_ids.is_empty() {
ctx.bus.emit(CoreEvent::MailUpdated { thread_ids });
}
Ok(())
}

fn persist_attachment_parts(
Expand Down Expand Up @@ -2260,6 +2309,10 @@ async fn fetch_attachment(
crate::mime::extract_attachment(&raw, &part_id).map(|(bytes, _)| bytes)
}

#[cfg(test)]
#[path = "sync_tests.rs"]
mod sync_tests;

#[cfg(test)]
mod tests {
use super::*;
Expand Down
Loading