export: bring matched items up to date on every run
Export matched each item against the target and then skipped it, so a second run -- the usual final pass of a cutover -- never carried anything that had changed at the source since the first: read and flagged state, moves between folders, edited contacts, events and Sieve scripts. It reported them as skipped and exited 0, while the usage guide said matched items were updated. Matched items are now updated, with one batched /set per type: - Email: keywords are set to the archive's, added and removed, compared case-insensitively. Memberships of folders this run migrated are added and removed to match; folders that exist only on the target are left alone, and a message is never left in no folder. Properties the server did not report are not touched. - Contacts and events: when both copies carry `updated`, the archive's is written only if it is newer, compared as instants so an offset or a fraction of a second is not taken for a change; otherwise each property the archive writes is compared, and those that differ are sent whole. - Sieve scripts: the target's copy is downloaded and compared byte for byte with what export would write -- after renaming Stalwart's vendor names for an inbuxa target -- and replaced when it differs, so a renamed script is not re-uploaded on every run. Updated items are counted as `updated`; unchanged ones stay `skipped`. The usage guide now describes this.
This commit is contained in:
+59
-1
@@ -5,7 +5,7 @@
|
||||
* SPDX-License-Identifier: Apache-2.0 OR MIT
|
||||
*/
|
||||
|
||||
use std::collections::HashMap;
|
||||
use std::collections::{HashMap, HashSet};
|
||||
use std::io::{IsTerminal, Write};
|
||||
|
||||
use rusqlite::Connection;
|
||||
@@ -49,6 +49,14 @@ impl Maps {
|
||||
fn insert(&mut self, ty: ObjectType, local: i64, target: JmapId) {
|
||||
self.m.entry(ty).or_default().insert(local, target);
|
||||
}
|
||||
|
||||
/// Every target id this run mapped for `ty`: the objects it migrated.
|
||||
fn targets_of(&self, ty: ObjectType) -> HashSet<String> {
|
||||
self.m
|
||||
.get(&ty)
|
||||
.map(|m| m.values().map(|id| id.0.clone()).collect())
|
||||
.unwrap_or_default()
|
||||
}
|
||||
}
|
||||
|
||||
impl TargetResolver for Maps {
|
||||
@@ -533,6 +541,56 @@ mod common {
|
||||
)
|
||||
}
|
||||
|
||||
/// Sends `updates` (target id, patch) as batched `/set` calls and counts
|
||||
/// the result into `counts`. A dry run counts them as updated and sends
|
||||
/// nothing.
|
||||
pub fn update_batch(
|
||||
net: &Net,
|
||||
ty: ObjectType,
|
||||
updates: Vec<(String, Value)>,
|
||||
counts: &mut TypeCounts,
|
||||
logger: &Logger,
|
||||
) {
|
||||
if updates.is_empty() {
|
||||
return;
|
||||
}
|
||||
if net.dry_run {
|
||||
counts.updated += updates.len() as u64;
|
||||
return;
|
||||
}
|
||||
let total = updates.len() as u64;
|
||||
let mut map = Map::new();
|
||||
for (id, patch) in updates {
|
||||
map.insert(id, patch);
|
||||
}
|
||||
match set_call(
|
||||
&net.client,
|
||||
&net.api,
|
||||
&net.account,
|
||||
ty.jmap_name(),
|
||||
SetRequest {
|
||||
update: Some(Value::Object(map)),
|
||||
..Default::default()
|
||||
},
|
||||
&net.limits,
|
||||
) {
|
||||
Ok(outcome) => {
|
||||
counts.updated += outcome.updated.len() as u64;
|
||||
for (id, err) in &outcome.not_updated {
|
||||
logger.warn(&format!("{}/set {id} not updated: {err}", ty.jmap_name()));
|
||||
counts.failed += 1;
|
||||
}
|
||||
}
|
||||
Err(e) => {
|
||||
logger.warn(&format!(
|
||||
"{}/set: updating {total} object(s) failed: {e}",
|
||||
ty.jmap_name()
|
||||
));
|
||||
counts.failed += total;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
fn blob_not_found(outcome: &crate::jmap::request::SetOutcome, cid: &str) -> bool {
|
||||
outcome.not_created.iter().any(|(c, err)| {
|
||||
c == cid && err.get("type").and_then(Value::as_str) == Some("blobNotFound")
|
||||
|
||||
+141
-67
@@ -9,12 +9,12 @@ use std::collections::{HashMap, HashSet};
|
||||
|
||||
use serde_json::{Map, Value, json};
|
||||
|
||||
use super::common::{jid, target_query_get};
|
||||
use super::common::{jid, target_query_get, update_batch};
|
||||
use super::{Maps, Net, Plan, Uploader};
|
||||
use crate::error::Error;
|
||||
use crate::jmap::error::JmapError;
|
||||
use crate::jmap::request::{
|
||||
MethodCall, Request, SetRequest, check_method_error, get_objects, retry_method_call, set_call,
|
||||
MethodCall, Request, check_method_error, get_objects, retry_method_call,
|
||||
};
|
||||
use crate::jmap::retry::MethodCallKind;
|
||||
use crate::jmap::wire::JmapId;
|
||||
@@ -62,8 +62,12 @@ pub fn reconcile(
|
||||
) -> Result<Plan, Error> {
|
||||
let ty = ObjectType::Email;
|
||||
|
||||
let target_min = target_query_get(net, ty, Some(&["messageId", "size", "mailboxIds"]))
|
||||
.map_err(Error::from)?;
|
||||
let target_min = target_query_get(
|
||||
net,
|
||||
ty,
|
||||
Some(&["messageId", "size", "mailboxIds", "keywords"]),
|
||||
)
|
||||
.map_err(Error::from)?;
|
||||
let mut indices: Vec<EmailIndex> = target_min.iter().map(server_index).collect();
|
||||
|
||||
let fallback_ids: Vec<JmapId> = target_min
|
||||
@@ -128,11 +132,12 @@ pub fn reconcile(
|
||||
.collect();
|
||||
let pairs = pair_with_targets(&local_keys, &sizes, &target_keys, &targets);
|
||||
|
||||
let mut membership_updates: Vec<(String, Value)> = Vec::new();
|
||||
let migrated = maps.targets_of(ObjectType::Mailbox);
|
||||
let mut updates: Vec<(String, Value)> = Vec::new();
|
||||
for (i, unit) in units.iter().enumerate() {
|
||||
match pairs[i] {
|
||||
Some(t) => match missing_memberships(&unit.row, &targets[t], maps) {
|
||||
Some(patch) => membership_updates.push((targets[t].id.clone(), patch)),
|
||||
Some(t) => match email_patch(&unit.row, &targets[t], maps, &migrated) {
|
||||
Some(patch) => updates.push((targets[t].id.clone(), patch)),
|
||||
None => counts.skipped += 1,
|
||||
},
|
||||
None => export_one(
|
||||
@@ -146,7 +151,7 @@ pub fn reconcile(
|
||||
),
|
||||
}
|
||||
}
|
||||
send_membership_updates(net, membership_updates, counts, logger);
|
||||
update_batch(net, ty, updates, counts, logger);
|
||||
|
||||
Ok(Plan::default())
|
||||
}
|
||||
@@ -197,6 +202,7 @@ struct TargetEmail {
|
||||
id: String,
|
||||
size: Option<u64>,
|
||||
mailboxes: Option<HashSet<String>>,
|
||||
keywords: Option<HashSet<String>>,
|
||||
}
|
||||
|
||||
impl TargetEmail {
|
||||
@@ -208,6 +214,10 @@ impl TargetEmail {
|
||||
.get("mailboxIds")
|
||||
.and_then(Value::as_object)
|
||||
.map(|m| m.keys().cloned().collect()),
|
||||
keywords: v
|
||||
.get("keywords")
|
||||
.and_then(Value::as_object)
|
||||
.map(|m| m.keys().map(|k| k.to_lowercase()).collect()),
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -254,65 +264,58 @@ fn pair_with_targets(
|
||||
out
|
||||
}
|
||||
|
||||
/// The `Email/set` patch adding the folders a matched email is missing on the
|
||||
/// target, or `None` when it is already in all of them. Folders that exist
|
||||
/// only on the target are left alone.
|
||||
fn missing_memberships(row: &EmailRow, target: &TargetEmail, maps: &Maps) -> Option<Value> {
|
||||
let have = target.mailboxes.as_ref()?;
|
||||
/// The `Email/set` patch that brings a matched email on the target in line
|
||||
/// with the archive, or `None` when it already is. The source is taken as
|
||||
/// the truth for what it covers: keywords are added and removed to match, and
|
||||
/// so are memberships of folders this run migrated. Folders that exist only
|
||||
/// on the target are left alone, an email is never left in no folder, and
|
||||
/// whatever the server did not report is not touched.
|
||||
fn email_patch(
|
||||
row: &EmailRow,
|
||||
target: &TargetEmail,
|
||||
maps: &Maps,
|
||||
migrated: &HashSet<String>,
|
||||
) -> Option<Value> {
|
||||
let mut patch = Map::new();
|
||||
for ml in &row.mailbox_locals {
|
||||
if let Some(t) = maps.target(ObjectType::Mailbox, *ml)
|
||||
&& !have.contains(&t.0)
|
||||
{
|
||||
patch.insert(format!("mailboxIds/{}", t.0), Value::Bool(true));
|
||||
if let Some(have) = &target.mailboxes {
|
||||
let want: HashSet<String> = row
|
||||
.mailbox_locals
|
||||
.iter()
|
||||
.filter_map(|ml| maps.target(ObjectType::Mailbox, *ml).map(|t| t.0))
|
||||
.collect();
|
||||
let add: Vec<&String> = want.iter().filter(|t| !have.contains(*t)).collect();
|
||||
let remove: Vec<&String> = have
|
||||
.iter()
|
||||
.filter(|t| migrated.contains(*t) && !want.contains(*t))
|
||||
.collect();
|
||||
let left = have.len() - remove.len() + add.len();
|
||||
for t in add {
|
||||
patch.insert(
|
||||
format!("mailboxIds/{}", pointer_escape(t)),
|
||||
Value::Bool(true),
|
||||
);
|
||||
}
|
||||
if left > 0 {
|
||||
for t in remove {
|
||||
patch.insert(format!("mailboxIds/{}", pointer_escape(t)), Value::Null);
|
||||
}
|
||||
}
|
||||
}
|
||||
if let Some(have) = &target.keywords {
|
||||
let want: HashSet<String> = row.keywords.iter().map(|k| k.to_lowercase()).collect();
|
||||
for k in want.difference(have) {
|
||||
patch.insert(format!("keywords/{}", pointer_escape(k)), Value::Bool(true));
|
||||
}
|
||||
for k in have.difference(&want) {
|
||||
patch.insert(format!("keywords/{}", pointer_escape(k)), Value::Null);
|
||||
}
|
||||
}
|
||||
(!patch.is_empty()).then_some(Value::Object(patch))
|
||||
}
|
||||
|
||||
fn send_membership_updates(
|
||||
net: &Net,
|
||||
updates: Vec<(String, Value)>,
|
||||
counts: &mut TypeCounts,
|
||||
logger: &Logger,
|
||||
) {
|
||||
if updates.is_empty() {
|
||||
return;
|
||||
}
|
||||
if net.dry_run {
|
||||
counts.updated += updates.len() as u64;
|
||||
return;
|
||||
}
|
||||
let total = updates.len() as u64;
|
||||
let mut map = Map::new();
|
||||
for (id, patch) in updates {
|
||||
map.insert(id, patch);
|
||||
}
|
||||
match set_call(
|
||||
&net.client,
|
||||
&net.api,
|
||||
&net.account,
|
||||
ObjectType::Email.jmap_name(),
|
||||
SetRequest {
|
||||
update: Some(Value::Object(map)),
|
||||
..Default::default()
|
||||
},
|
||||
&net.limits,
|
||||
) {
|
||||
Ok(outcome) => {
|
||||
counts.updated += outcome.updated.len() as u64;
|
||||
for (id, err) in &outcome.not_updated {
|
||||
logger.warn(&format!("Email/set {id}: folders not added: {err}"));
|
||||
counts.failed += 1;
|
||||
}
|
||||
}
|
||||
Err(e) => {
|
||||
logger.warn(&format!(
|
||||
"Email/set: adding folders to {total} email(s) failed: {e}"
|
||||
));
|
||||
counts.failed += total;
|
||||
}
|
||||
}
|
||||
/// Escapes one JSON Pointer segment (RFC 6901), as JMAP patch paths use.
|
||||
fn pointer_escape(segment: &str) -> String {
|
||||
segment.replace('~', "~0").replace('/', "~1")
|
||||
}
|
||||
|
||||
fn build_mailbox_ids(row: &EmailRow, maps: &Maps) -> Option<Map<String, Value>> {
|
||||
@@ -570,9 +573,14 @@ mod tests {
|
||||
id: id.to_owned(),
|
||||
size,
|
||||
mailboxes: mailboxes.map(|m| m.iter().map(|s| (*s).to_owned()).collect()),
|
||||
keywords: None,
|
||||
}
|
||||
}
|
||||
|
||||
fn set(items: &[&str]) -> HashSet<String> {
|
||||
items.iter().map(|s| (*s).to_owned()).collect()
|
||||
}
|
||||
|
||||
fn mid(m: &str) -> EmailKey {
|
||||
EmailKey::MessageId(m.to_owned())
|
||||
}
|
||||
@@ -628,19 +636,85 @@ mod tests {
|
||||
assert_eq!(pairs, vec![Some(0), None]);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn missing_memberships_adds_only_the_absent_migrated_folders() {
|
||||
fn two_folder_maps() -> Maps {
|
||||
let mut maps = Maps::default();
|
||||
maps.insert(ObjectType::Mailbox, 1, JmapId("T1".into()));
|
||||
maps.insert(ObjectType::Mailbox, 2, JmapId("T2".into()));
|
||||
maps
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn patch_adds_the_absent_migrated_folders() {
|
||||
let maps = two_folder_maps();
|
||||
let r = row(10, &[1, 2, 3], &[]);
|
||||
let patch = missing_memberships(&r, &target("E", None, Some(&["T1", "Own"])), &maps)
|
||||
.expect("T2 is missing");
|
||||
let patch = email_patch(
|
||||
&r,
|
||||
&target("E", None, Some(&["T1", "Own"])),
|
||||
&maps,
|
||||
&set(&["T1", "T2"]),
|
||||
)
|
||||
.expect("T2 is missing");
|
||||
assert_eq!(patch, json!({ "mailboxIds/T2": true }));
|
||||
assert!(missing_memberships(&r, &target("E", None, Some(&["T1", "T2"])), &maps).is_none());
|
||||
assert!(
|
||||
missing_memberships(&r, &target("E", None, None), &maps).is_none(),
|
||||
email_patch(
|
||||
&r,
|
||||
&target("E", None, Some(&["T1", "T2"])),
|
||||
&maps,
|
||||
&set(&["T1", "T2"])
|
||||
)
|
||||
.is_none()
|
||||
);
|
||||
assert!(
|
||||
email_patch(&r, &target("E", None, None), &maps, &set(&["T1", "T2"])).is_none(),
|
||||
"unknown membership is left alone"
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn patch_moves_between_migrated_folders_but_keeps_target_only_ones() {
|
||||
let maps = two_folder_maps();
|
||||
let r = row(10, &[2], &[]);
|
||||
let patch = email_patch(
|
||||
&r,
|
||||
&target("E", None, Some(&["T1", "Own"])),
|
||||
&maps,
|
||||
&set(&["T1", "T2"]),
|
||||
)
|
||||
.unwrap();
|
||||
assert_eq!(
|
||||
patch,
|
||||
json!({ "mailboxIds/T2": true, "mailboxIds/T1": null })
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn patch_never_leaves_an_email_in_no_folder() {
|
||||
let maps = two_folder_maps();
|
||||
let r = row(10, &[9], &[]);
|
||||
assert!(
|
||||
email_patch(
|
||||
&r,
|
||||
&target("E", None, Some(&["T1"])),
|
||||
&maps,
|
||||
&set(&["T1", "T2"])
|
||||
)
|
||||
.is_none(),
|
||||
"the only folder is not removed when nothing replaces it"
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn patch_syncs_keywords_both_ways_case_insensitively() {
|
||||
let maps = two_folder_maps();
|
||||
let r = row(10, &[1], &["$Seen", "work/urgent"]);
|
||||
let mut t = target("E", None, Some(&["T1"]));
|
||||
t.keywords = Some(set(&["$seen", "$flagged"]));
|
||||
let patch = email_patch(&r, &t, &maps, &set(&["T1", "T2"])).unwrap();
|
||||
assert_eq!(
|
||||
patch,
|
||||
json!({ "keywords/work~1urgent": true, "keywords/$flagged": null })
|
||||
);
|
||||
t.keywords = Some(set(&["$seen", "work/urgent"]));
|
||||
assert!(email_patch(&r, &t, &maps, &set(&["T1", "T2"])).is_none());
|
||||
}
|
||||
}
|
||||
|
||||
+84
-23
@@ -9,10 +9,11 @@ use std::collections::{HashMap, HashSet};
|
||||
|
||||
use serde_json::{Value, json};
|
||||
|
||||
use super::common::{create_batch, jid, retry_if_blob_missing, target_get_all};
|
||||
use super::common::{create_batch, jid, retry_if_blob_missing, target_get_all, update_batch};
|
||||
use super::sieve_names;
|
||||
use super::{Maps, Net, Plan, Uploader};
|
||||
use crate::error::Error;
|
||||
use crate::jmap::blobxfer;
|
||||
use crate::jmap::request::{Request, check_method_error};
|
||||
use crate::logging::{LEVEL_DEFAULT, Logger};
|
||||
use crate::sync::import_jmap::mapping::BlobBytes;
|
||||
@@ -31,10 +32,14 @@ pub fn reconcile(
|
||||
let targets = target_get_all(net, ty).map_err(Error::from)?;
|
||||
|
||||
let mut target_by_name: HashMap<String, String> = HashMap::new();
|
||||
let mut target_blob: HashMap<String, String> = HashMap::new();
|
||||
for t in &targets {
|
||||
let (Some(id), Some(name)) = (jid(t), t.get("name").and_then(Value::as_str)) else {
|
||||
continue;
|
||||
};
|
||||
if let Some(blob) = t.get("blobId").and_then(Value::as_str) {
|
||||
target_blob.insert(id.clone(), blob.to_owned());
|
||||
}
|
||||
target_by_name.insert(name.to_owned(), id);
|
||||
}
|
||||
|
||||
@@ -66,19 +71,55 @@ pub fn reconcile(
|
||||
.find(|(_, _, a, _)| *a)
|
||||
.map(|(_, n, _, _)| n.clone().unwrap_or_default());
|
||||
|
||||
let mut updates: Vec<(String, Value)> = Vec::new();
|
||||
|
||||
for (local, name, is_active, blob_local) in &locals {
|
||||
let matched = name.as_ref().and_then(|n| target_by_name.get(n)).cloned();
|
||||
let label = name.as_deref().unwrap_or("(unnamed)");
|
||||
let rewritten = if rename_vendor {
|
||||
renamed_script(&uploader, *blob_local)?
|
||||
} else {
|
||||
None
|
||||
};
|
||||
let target_id = if let Some(id) = matched {
|
||||
counts.skipped += 1;
|
||||
// Compare what would be written -- the renamed bytes where the
|
||||
// script needed renaming -- so an unchanged script stays unchanged.
|
||||
let ours = match &rewritten {
|
||||
Some((bytes, _)) => bytes.clone(),
|
||||
None => uploader.bytes(*blob_local).map_err(Error::from)?,
|
||||
};
|
||||
match content_differs(net, &ours, target_blob.get(&id)) {
|
||||
Ok(false) => counts.skipped += 1,
|
||||
Ok(true) => {
|
||||
let blob = match &rewritten {
|
||||
Some((bytes, renamed)) => {
|
||||
log_renames(label, renamed, logger);
|
||||
uploader.upload_bytes_as(*blob_local, "application/sieve", bytes)
|
||||
}
|
||||
None => uploader.upload_with(*blob_local, "application/sieve"),
|
||||
};
|
||||
match blob {
|
||||
Ok(b) => updates.push((id.clone(), json!({ "blobId": b.0 }))),
|
||||
Err(e) => {
|
||||
logger.warn(&format!(
|
||||
"SieveScript {label}: upload for update failed: {e}"
|
||||
));
|
||||
counts.failed += 1;
|
||||
}
|
||||
}
|
||||
}
|
||||
Err(e) => {
|
||||
logger.warn(&format!("SieveScript {label}: not compared: {e}"));
|
||||
counts.skipped += 1;
|
||||
}
|
||||
}
|
||||
id
|
||||
} else {
|
||||
let cid = format!("c{local}");
|
||||
let label = name.as_deref().unwrap_or("(unnamed)");
|
||||
let rewritten = if rename_vendor {
|
||||
renamed_script(&uploader, *blob_local, label, logger)?
|
||||
} else {
|
||||
None
|
||||
};
|
||||
if let Some((_, renamed)) = &rewritten {
|
||||
log_renames(label, renamed, logger);
|
||||
}
|
||||
let rewritten = rewritten.as_ref().map(|(bytes, _)| bytes);
|
||||
let build = |up: &mut Uploader<'_>| -> Result<Value, Error> {
|
||||
let blob_id = match &rewritten {
|
||||
Some(bytes) => up.upload_bytes_as(*blob_local, "application/sieve", bytes),
|
||||
@@ -120,6 +161,8 @@ pub fn reconcile(
|
||||
}
|
||||
}
|
||||
|
||||
update_batch(net, ty, updates, counts, logger);
|
||||
|
||||
if active_target.is_none() && locals.iter().all(|(_, _, a, _)| !*a) {
|
||||
deactivate = true;
|
||||
}
|
||||
@@ -190,22 +233,40 @@ fn target_sieve_extensions(net: &Net) -> Vec<String> {
|
||||
.unwrap_or_default()
|
||||
}
|
||||
|
||||
/// A script's bytes after renaming, and the names that were renamed.
|
||||
type Renamed = (Vec<u8>, Vec<String>);
|
||||
|
||||
/// The script's bytes with Stalwart's vendor names renamed for an inbuxa
|
||||
/// target, or `None` when it needs no change. Each rename is logged.
|
||||
fn renamed_script(
|
||||
uploader: &Uploader<'_>,
|
||||
blob_local: i64,
|
||||
label: &str,
|
||||
logger: &Logger,
|
||||
) -> Result<Option<Vec<u8>>, Error> {
|
||||
/// target, and the names renamed, or `None` when it needs no change.
|
||||
fn renamed_script(uploader: &Uploader<'_>, blob_local: i64) -> Result<Option<Renamed>, Error> {
|
||||
let bytes = uploader.bytes(blob_local).map_err(Error::from)?;
|
||||
Ok(sieve_names::rewrite(&bytes).map(|(out, renamed)| {
|
||||
for old in &renamed {
|
||||
let new = old.replacen("vnd.stalwart.", "vnd.inbuxa.", 1);
|
||||
if logger.enabled(LEVEL_DEFAULT) {
|
||||
eprintln!("export: SieveScript {label}: renamed {old} to {new}");
|
||||
}
|
||||
Ok(sieve_names::rewrite(&bytes))
|
||||
}
|
||||
|
||||
/// Prints each rename made to a script about to be written.
|
||||
fn log_renames(label: &str, renamed: &[String], logger: &Logger) {
|
||||
for old in renamed {
|
||||
let new = old.replacen("vnd.stalwart.", "vnd.inbuxa.", 1);
|
||||
if logger.enabled(LEVEL_DEFAULT) {
|
||||
eprintln!("export: SieveScript {label}: renamed {old} to {new}");
|
||||
}
|
||||
out
|
||||
}))
|
||||
}
|
||||
}
|
||||
|
||||
/// Whether the target's copy of a script differs from `ours`. A target that
|
||||
/// reports no blob is taken as different, so ours is written.
|
||||
fn content_differs(net: &Net, ours: &[u8], target_blob: Option<&String>) -> Result<bool, Error> {
|
||||
let Some(target_blob) = target_blob else {
|
||||
return Ok(true);
|
||||
};
|
||||
let theirs = blobxfer::download_bytes(
|
||||
&net.client,
|
||||
&net.session,
|
||||
&net.account,
|
||||
target_blob,
|
||||
"application/sieve",
|
||||
"script.sieve",
|
||||
)
|
||||
.map_err(Error::from)?;
|
||||
Ok(ours != theirs.as_slice())
|
||||
}
|
||||
|
||||
@@ -1,15 +1,16 @@
|
||||
/*
|
||||
* SPDX-FileCopyrightText: 2020 Stalwart Labs LLC <[email protected]>
|
||||
* SPDX-FileCopyrightText: 2026 John Coffey <[email protected]>
|
||||
*
|
||||
* SPDX-License-Identifier: Apache-2.0 OR MIT
|
||||
*/
|
||||
|
||||
use std::collections::HashSet;
|
||||
use std::collections::{HashMap, HashSet};
|
||||
use std::fmt::Write as _;
|
||||
|
||||
use serde_json::Value;
|
||||
use serde_json::{Map, Value};
|
||||
|
||||
use super::common::{create_batch, jid, target_query_get};
|
||||
use super::common::{create_batch, jid, target_query_get, update_batch};
|
||||
use super::{Maps, Net, Plan, Uploader};
|
||||
use crate::error::Error;
|
||||
use crate::logging::Logger;
|
||||
@@ -42,10 +43,10 @@ pub fn reconcile(
|
||||
logger: &Logger,
|
||||
) -> Result<Plan, Error> {
|
||||
let targets = target_query_get(net, ty, None).map_err(Error::from)?;
|
||||
let mut by_uid: std::collections::HashMap<String, String> = std::collections::HashMap::new();
|
||||
let mut by_uid: HashMap<String, (String, &Value)> = HashMap::new();
|
||||
for t in &targets {
|
||||
if let (Some(uid), Some(id)) = (target_uid(t), jid(t)) {
|
||||
by_uid.entry(uid).or_insert(id);
|
||||
by_uid.entry(uid).or_insert((id, t));
|
||||
}
|
||||
}
|
||||
|
||||
@@ -76,12 +77,23 @@ pub fn reconcile(
|
||||
};
|
||||
|
||||
let mut matched_uids: HashSet<String> = HashSet::new();
|
||||
let mut updates: Vec<(String, Value)> = Vec::new();
|
||||
let blobs = Uploader::new(net, &ctx.conn);
|
||||
for (local, uid) in &rows {
|
||||
if let Some(tid) = by_uid.get(uid) {
|
||||
if let Some((tid, existing)) = by_uid.get(uid) {
|
||||
maps.insert(ty, *local, crate::jmap::wire::JmapId(tid.clone()));
|
||||
matched_uids.insert(uid.clone());
|
||||
counts.skipped += 1;
|
||||
match build_wire(ctx, ty, *local, maps, &blobs) {
|
||||
Ok(wire) => match changed_properties(&wire, existing) {
|
||||
Some(patch) => updates.push((tid.clone(), patch)),
|
||||
None => counts.skipped += 1,
|
||||
},
|
||||
Err(e) if e.aborts_run() => return Err(e),
|
||||
Err(e) => {
|
||||
logger.warn(&format!("{} not compared: {e}", describe(ty, *local, uid)));
|
||||
counts.skipped += 1;
|
||||
}
|
||||
}
|
||||
continue;
|
||||
}
|
||||
let cid = format!("c{local}");
|
||||
@@ -123,6 +135,8 @@ pub fn reconcile(
|
||||
}
|
||||
}
|
||||
|
||||
update_batch(net, ty, updates, counts, logger);
|
||||
|
||||
let objs: Vec<TargetObj> = targets
|
||||
.iter()
|
||||
.filter_map(|t| {
|
||||
@@ -172,3 +186,80 @@ fn build_wire(
|
||||
calendar_event_to_wire(&cal, dr != 0, ud != 0, &data, maps, blobs).map_err(Error::from)
|
||||
}
|
||||
}
|
||||
|
||||
/// An `updated` value as a point in time, so that offsets and fractional
|
||||
/// seconds compare as the same instant. `None` if it does not parse.
|
||||
fn parse_updated(s: &str) -> Option<time::OffsetDateTime> {
|
||||
time::OffsetDateTime::parse(s, &time::format_description::well_known::Rfc3339).ok()
|
||||
}
|
||||
|
||||
/// The update that makes `target` match the archive's `wire` object, or
|
||||
/// `None` when nothing changed. When both carry `updated`, it decides: the
|
||||
/// archive's copy wins only if it is newer. Otherwise each property the
|
||||
/// archive writes is compared, and those that differ are sent whole.
|
||||
fn changed_properties(wire: &Value, target: &Value) -> Option<Value> {
|
||||
let wire = wire.as_object()?;
|
||||
let stamp = |v: Option<&Value>| v.and_then(Value::as_str).and_then(parse_updated);
|
||||
if let (Some(ours), Some(theirs)) = (stamp(wire.get("updated")), stamp(target.get("updated")))
|
||||
&& ours <= theirs
|
||||
{
|
||||
return None;
|
||||
}
|
||||
let mut patch = Map::new();
|
||||
for (k, v) in wire {
|
||||
if k == "uid" || k == "id" {
|
||||
continue;
|
||||
}
|
||||
if target.get(k) != Some(v) {
|
||||
patch.insert(k.clone(), v.clone());
|
||||
}
|
||||
}
|
||||
(!patch.is_empty()).then_some(Value::Object(patch))
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
use serde_json::json;
|
||||
|
||||
#[test]
|
||||
fn unchanged_object_needs_no_update() {
|
||||
let wire = json!({"uid": "u", "name": {"full": "Ann"}, "addressBookIds": {"A": true}});
|
||||
let target = json!({"id": "T", "uid": "u", "name": {"full": "Ann"},
|
||||
"addressBookIds": {"A": true}, "extra": 1});
|
||||
assert_eq!(changed_properties(&wire, &target), None);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn changed_properties_are_sent_whole() {
|
||||
let wire = json!({"uid": "u", "name": {"full": "Ann B"}, "addressBookIds": {"A": true}});
|
||||
let target = json!({"id": "T", "uid": "u", "name": {"full": "Ann"},
|
||||
"addressBookIds": {"A": true}});
|
||||
assert_eq!(
|
||||
changed_properties(&wire, &target),
|
||||
Some(json!({"name": {"full": "Ann B"}}))
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn updated_compares_instants_not_strings() {
|
||||
let target = json!({"uid": "u", "title": "old", "updated": "2026-01-02T00:00:00Z"});
|
||||
let same_instant = json!({"uid": "u", "title": "new",
|
||||
"updated": "2026-01-02T01:00:00.000+01:00"});
|
||||
assert_eq!(changed_properties(&same_instant, &target), None);
|
||||
let later = json!({"uid": "u", "title": "new", "updated": "2026-01-02T00:00:00.5Z"});
|
||||
assert!(changed_properties(&later, &target).is_some());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn updated_decides_when_both_sides_carry_it() {
|
||||
let target = json!({"uid": "u", "title": "old", "updated": "2026-01-02T00:00:00Z"});
|
||||
let older = json!({"uid": "u", "title": "new", "updated": "2026-01-01T00:00:00Z"});
|
||||
assert_eq!(changed_properties(&older, &target), None, "target is newer");
|
||||
let newer = json!({"uid": "u", "title": "new", "updated": "2026-01-03T00:00:00Z"});
|
||||
assert_eq!(
|
||||
changed_properties(&newer, &target),
|
||||
Some(json!({"title": "new", "updated": "2026-01-03T00:00:00Z"}))
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user