Fix sync issues
This commit is contained in:
+31
-1
@@ -136,13 +136,43 @@ fn try_treat_url_as_home_or_collection(
|
||||
if collections.is_empty() {
|
||||
return Ok(None);
|
||||
}
|
||||
let principal_url = resolve_principal_url(client, kind, url)?;
|
||||
Ok(Some(Discovery {
|
||||
principal_url: None,
|
||||
principal_url,
|
||||
home_set_url: url.to_owned(),
|
||||
collections,
|
||||
}))
|
||||
}
|
||||
|
||||
fn resolve_principal_url(
|
||||
client: &DavClient,
|
||||
kind: DavKind,
|
||||
url: &str,
|
||||
) -> Result<Option<String>, DiscoveryError> {
|
||||
if matches!(kind, DavKind::Webdav) {
|
||||
return Ok(None);
|
||||
}
|
||||
let body = xml::propfind_current_user_principal();
|
||||
let ms = match client.propfind_responses(url, 0, &body, url) {
|
||||
Ok(ms) => ms,
|
||||
Err(JmapError::HttpStatus { status, .. }) if (400..600).contains(&status) => {
|
||||
return Ok(None);
|
||||
}
|
||||
Err(JmapError::RetriesExhausted(_)) => return Ok(None),
|
||||
Err(e) => return Err(DiscoveryError::Transport(e)),
|
||||
};
|
||||
let final_url = ms.final_url.clone();
|
||||
let Some(principal_href) = ms
|
||||
.responses
|
||||
.iter()
|
||||
.find_map(|r| r.props.current_user_principal.as_deref())
|
||||
else {
|
||||
return Ok(None);
|
||||
};
|
||||
let resolved = absolute_url(&final_url, &normalise(&final_url, principal_href)?)?;
|
||||
Ok(Some(resolved))
|
||||
}
|
||||
|
||||
fn try_via_principal(
|
||||
client: &DavClient,
|
||||
kind: DavKind,
|
||||
|
||||
+2
-2
@@ -98,8 +98,8 @@ fn fail(err: &Error) -> i32 {
|
||||
fn report(summary: &Summary) {
|
||||
for (type_name, counts) in &summary.per_type {
|
||||
println!(
|
||||
"{type_name}: created={} fetched={} deleted={} skipped={} failed={}",
|
||||
counts.created, counts.fetched, counts.deleted, counts.skipped, counts.failed
|
||||
"{type_name}: created={} fetched={} updated={} deleted={} skipped={} failed={}",
|
||||
counts.created, counts.fetched, counts.updated, counts.deleted, counts.skipped, counts.failed
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
+13
-34
@@ -4,19 +4,16 @@
|
||||
* SPDX-License-Identifier: Apache-2.0 OR MIT
|
||||
*/
|
||||
|
||||
use std::collections::HashMap;
|
||||
use std::collections::{HashMap, HashSet};
|
||||
|
||||
use serde_json::{Value, json};
|
||||
|
||||
use super::common::{create_batch, jid, target_get_all};
|
||||
use super::{Maps, Net, Plan, Uploader};
|
||||
use crate::db;
|
||||
use crate::error::Error;
|
||||
use crate::jmap::blobxfer;
|
||||
use crate::jmap::request::Request;
|
||||
use crate::logging::Logger;
|
||||
use crate::sync::import_jmap::mapping::{SIEVE_SELECT, row_to_sieve_script};
|
||||
use crate::sync::keys::blake3_bytes;
|
||||
use crate::sync::{Context, TypeCounts};
|
||||
use crate::types::ObjectType;
|
||||
|
||||
@@ -30,21 +27,12 @@ pub fn reconcile(
|
||||
let ty = ObjectType::SieveScript;
|
||||
let targets = target_get_all(net, ty).map_err(Error::from)?;
|
||||
|
||||
let mut target_by_key: HashMap<[u8; 32], String> = HashMap::new();
|
||||
let mut target_by_name: HashMap<String, String> = HashMap::new();
|
||||
for t in &targets {
|
||||
let (Some(id), Some(blob)) = (jid(t), t.get("blobId").and_then(Value::as_str)) else {
|
||||
let (Some(id), Some(name)) = (jid(t), t.get("name").and_then(Value::as_str)) else {
|
||||
continue;
|
||||
};
|
||||
let bytes = blobxfer::download_bytes(
|
||||
&net.client,
|
||||
&net.session,
|
||||
&net.account,
|
||||
blob,
|
||||
"application/sieve",
|
||||
"script",
|
||||
)
|
||||
.map_err(Error::from)?;
|
||||
target_by_key.insert(blake3_bytes(&bytes), id);
|
||||
target_by_name.insert(name.to_owned(), id);
|
||||
}
|
||||
|
||||
let locals: Vec<(i64, Option<String>, bool, i64)> = {
|
||||
@@ -71,13 +59,10 @@ pub fn reconcile(
|
||||
let mut uploader = Uploader::new(net, &ctx.conn);
|
||||
|
||||
for (local, name, is_active, blob_local) in &locals {
|
||||
let bytes = db::blobs::blob_bytes(&ctx.conn, *blob_local)
|
||||
.map_err(|e| Error::Partial(e.to_string()))?
|
||||
.ok_or_else(|| Error::Partial("sieve blob missing".to_owned()))?;
|
||||
let key = blake3_bytes(&bytes);
|
||||
let target_id = if let Some(id) = target_by_key.get(&key) {
|
||||
let matched = name.as_ref().and_then(|n| target_by_name.get(n)).cloned();
|
||||
let target_id = if let Some(id) = matched {
|
||||
counts.skipped += 1;
|
||||
id.clone()
|
||||
id
|
||||
} else {
|
||||
let blob_id = uploader
|
||||
.upload_with(*blob_local, "application/sieve")
|
||||
@@ -92,7 +77,9 @@ pub fn reconcile(
|
||||
match outcome.created.first().and_then(|(_, v)| jid(v)) {
|
||||
Some(id) => {
|
||||
counts.created += 1;
|
||||
target_by_key.insert(key, id.clone());
|
||||
if let Some(n) = name {
|
||||
target_by_name.insert(n.clone(), id.clone());
|
||||
}
|
||||
id
|
||||
}
|
||||
None => {
|
||||
@@ -128,18 +115,10 @@ pub fn reconcile(
|
||||
}
|
||||
}
|
||||
|
||||
let local_keys: std::collections::HashSet<[u8; 32]> = locals
|
||||
let local_names: HashSet<String> = locals.iter().filter_map(|(_, n, _, _)| n.clone()).collect();
|
||||
let mut prune_candidates: Vec<String> = target_by_name
|
||||
.iter()
|
||||
.filter_map(|(_, _, _, b)| {
|
||||
db::blobs::blob_bytes(&ctx.conn, *b)
|
||||
.ok()
|
||||
.flatten()
|
||||
.map(|by| blake3_bytes(&by))
|
||||
})
|
||||
.collect();
|
||||
let mut prune_candidates: Vec<String> = target_by_key
|
||||
.iter()
|
||||
.filter(|(k, _)| !local_keys.contains(*k))
|
||||
.filter(|(name, _)| !local_names.contains(*name))
|
||||
.map(|(_, id)| id.clone())
|
||||
.collect();
|
||||
prune_candidates.sort();
|
||||
|
||||
@@ -826,6 +826,18 @@ fn reconcile_folder(
|
||||
tx.commit()?;
|
||||
}
|
||||
|
||||
if !diff.present.is_empty() {
|
||||
counts.updated += refresh_present_flags(
|
||||
conn,
|
||||
client,
|
||||
control_ctx,
|
||||
&folder.name,
|
||||
uidvalidity,
|
||||
&diff.present,
|
||||
opts,
|
||||
)?;
|
||||
}
|
||||
|
||||
let now = OffsetDateTime::now_utc()
|
||||
.format(&Rfc3339)
|
||||
.unwrap_or_else(|_| String::from("1970-01-01T00:00:00Z"));
|
||||
@@ -961,6 +973,93 @@ fn insert_single_message(
|
||||
Ok(())
|
||||
}
|
||||
|
||||
fn refresh_present_flags(
|
||||
conn: &mut Connection,
|
||||
client: &mut ImapClient,
|
||||
control_ctx: &ControlCtx,
|
||||
folder: &str,
|
||||
uidvalidity: u32,
|
||||
present: &[u32],
|
||||
opts: RunOpts,
|
||||
) -> Result<u64, Error> {
|
||||
let RunOpts {
|
||||
source_id,
|
||||
fetch_batch,
|
||||
include_deleted,
|
||||
logger: _,
|
||||
} = opts;
|
||||
let mut updated: u64 = 0;
|
||||
let tx = conn.transaction()?;
|
||||
for batch in chunks(present, fetch_batch.max(1)) {
|
||||
let set = command::format_uid_set(batch, true);
|
||||
let resp = control_run_collect(
|
||||
client,
|
||||
control_ctx,
|
||||
&command::uid_fetch(&set, &["UID", "FLAGS"]),
|
||||
)?;
|
||||
for u in &resp.untagged {
|
||||
let Some(attrs) = fetch::extract(u) else {
|
||||
continue;
|
||||
};
|
||||
let Some(uid) = attrs.uid else {
|
||||
continue;
|
||||
};
|
||||
let translation = translate_flags(&attrs.flags, include_deleted);
|
||||
if translation.has_deleted_flag && !include_deleted {
|
||||
continue;
|
||||
}
|
||||
let Some((local_id, existing)) =
|
||||
load_email_keywords(&tx, source_id, folder, uidvalidity, uid)?
|
||||
else {
|
||||
continue;
|
||||
};
|
||||
if keyword_set_differs(&existing, &translation.keywords) {
|
||||
let json = Value::Array(
|
||||
translation
|
||||
.keywords
|
||||
.iter()
|
||||
.cloned()
|
||||
.map(Value::String)
|
||||
.collect(),
|
||||
);
|
||||
tx.execute(
|
||||
"UPDATE emails SET keywords = ?1 WHERE id = ?2",
|
||||
params![json.to_string(), local_id],
|
||||
)?;
|
||||
updated += 1;
|
||||
}
|
||||
}
|
||||
}
|
||||
tx.commit()?;
|
||||
Ok(updated)
|
||||
}
|
||||
|
||||
fn load_email_keywords(
|
||||
tx: &rusqlite::Transaction<'_>,
|
||||
source_id: i64,
|
||||
folder: &str,
|
||||
uidvalidity: u32,
|
||||
uid: u32,
|
||||
) -> Result<Option<(i64, Vec<String>)>, Error> {
|
||||
let Some(local_id) = db::imap_ids::local_for_email(tx, source_id, folder, uidvalidity, uid)?
|
||||
else {
|
||||
return Ok(None);
|
||||
};
|
||||
let kw_json: String = tx.query_row(
|
||||
"SELECT keywords FROM emails WHERE id = ?1",
|
||||
params![local_id],
|
||||
|r| r.get(0),
|
||||
)?;
|
||||
let kws: Vec<String> = serde_json::from_str(&kw_json).unwrap_or_default();
|
||||
Ok(Some((local_id, kws)))
|
||||
}
|
||||
|
||||
fn keyword_set_differs(a: &[String], b: &[String]) -> bool {
|
||||
let sa: HashSet<&str> = a.iter().map(String::as_str).collect();
|
||||
let sb: HashSet<&str> = b.iter().map(String::as_str).collect();
|
||||
sa != sb
|
||||
}
|
||||
|
||||
fn dry_run_summary(
|
||||
conn: &Connection,
|
||||
source_id: Option<i64>,
|
||||
|
||||
@@ -75,6 +75,7 @@ pub struct ExportConfig {
|
||||
pub struct TypeCounts {
|
||||
pub created: u64,
|
||||
pub fetched: u64,
|
||||
pub updated: u64,
|
||||
pub deleted: u64,
|
||||
pub skipped: u64,
|
||||
pub failed: u64,
|
||||
|
||||
Reference in New Issue
Block a user