diff --git a/README.md b/README.md index ffa555e..8cb8350 100644 --- a/README.md +++ b/README.md @@ -9,13 +9,17 @@

- crates.io + build   - build + release   - docs.rs + npm   - license + downloads +   + platforms +   + license


@@ -288,13 +292,44 @@ Read-only dump of a local archive. This command never opens a network connection ## Testing +The default suite is hermetic (unit tests plus `mockito`-scripted JMAP/DAV/EWS/Graph behaviours) and needs no network or Docker: + ```sh cargo build cargo clippy --all-targets cargo test ``` -Live integration tests against a Stalwart server and container-based smoke tests against third-party servers (Dovecot, Cyrus, Radicale, Baikal, Apache `mod_dav`) are gated behind `--ignored` and require Docker plus the seeder fixtures. +### Live and integration tests (Docker required) + +Live integration tests against a Stalwart server and container-based tests against third-party servers (Dovecot, Cyrus, Radicale, Baikal, Apache `mod_dav`) are gated behind `--ignored`. **They require a running Docker daemon**: each test binary boots its own throwaway container via `testcontainers` (images are pulled automatically on first run), so Docker must be installed and `docker info` must succeed before invoking them. + +Run them per binary, and always with `--test-threads=1`: + +```sh +cargo test --test sync_jmap -- --ignored --test-threads=1 # live JMAP import/export/convergence/prune +cargo test --test sync_imap -- --ignored --test-threads=1 +cargo test --test sync_managesieve -- --ignored --test-threads=1 +cargo test --test sync_maildir -- --ignored --test-threads=1 +cargo test --test sync_caldav -- --ignored --test-threads=1 +cargo test --test sync_carddav -- --ignored --test-threads=1 +cargo test --test sync_webdav -- --ignored --test-threads=1 +cargo test --test live_stalwart -- --ignored --test-threads=1 +cargo test --test seed_smoke -- --ignored --test-threads=1 +cargo test --test seed_only -- --ignored --test-threads=1 + +# Third-party-server tests (one container each): +cargo test --test integration_radicale -- --ignored --test-threads=1 +cargo test --test integration_baikal -- --ignored --test-threads=1 +cargo test --test integration_webdav -- --ignored --test-threads=1 +cargo test --test integration_dovecot -- --ignored --test-threads=1 +cargo test --test integration_cyrus -- --ignored --test-threads=1 + +# Slow tests +cargo test --test mock_jmap -- --ignored +``` + +`--test-threads=1` is mandatory, not just advisory: within a binary every test shares a single per-binary container, and each test provisions then tears down the same disposable `vandelay.org` domain (and opens the archive with SQLite `EXCLUSIVE` locking). Separate binaries are isolated (each boots its own container on dynamic host ports), so plain `cargo test --test ` invocations are safe to run one after another. ## License diff --git a/src/db/mod.rs b/src/db/mod.rs index 1d3fdc7..dd1293d 100644 --- a/src/db/mod.rs +++ b/src/db/mod.rs @@ -15,4 +15,5 @@ pub mod init; pub mod maildir_ids; pub mod managesieve_ids; pub mod sources; +pub mod sync_state_jmap; pub mod takeout_ids; diff --git a/src/db/schema.sql b/src/db/schema.sql index 3ca3a57..ca0c041 100644 --- a/src/db/schema.sql +++ b/src/db/schema.sql @@ -28,6 +28,13 @@ CREATE TABLE IF NOT EXISTS sync_id_jmap ( CREATE INDEX IF NOT EXISTS sync_id_jmap_local_idx ON sync_id_jmap (type_name, local_id); +CREATE TABLE IF NOT EXISTS sync_state_jmap ( + source_id INTEGER NOT NULL REFERENCES sources(id) ON DELETE CASCADE, + type_name TEXT NOT NULL, + state TEXT NOT NULL, + PRIMARY KEY (source_id, type_name) +); + CREATE TABLE IF NOT EXISTS sync_id_imap ( source_id INTEGER NOT NULL REFERENCES sources(id) ON DELETE CASCADE, type_name TEXT NOT NULL CHECK (type_name IN ('mailbox','email')), diff --git a/src/db/sync_state_jmap.rs b/src/db/sync_state_jmap.rs new file mode 100644 index 0000000..7cd62db --- /dev/null +++ b/src/db/sync_state_jmap.rs @@ -0,0 +1,87 @@ +/* + * SPDX-FileCopyrightText: 2020 Stalwart Labs LLC + * + * SPDX-License-Identifier: Apache-2.0 OR MIT + */ + +use rusqlite::{Connection, OptionalExtension, params}; + +use crate::types::ObjectType; + +pub fn get( + conn: &Connection, + source_id: i64, + ty: ObjectType, +) -> Result, rusqlite::Error> { + conn.query_row( + "SELECT state FROM sync_state_jmap WHERE source_id = ?1 AND type_name = ?2", + params![source_id, ty.jmap_name()], + |row| row.get(0), + ) + .optional() +} + +pub fn upsert( + conn: &Connection, + source_id: i64, + ty: ObjectType, + state: &str, +) -> Result<(), rusqlite::Error> { + conn.execute( + "INSERT INTO sync_state_jmap (source_id, type_name, state) + VALUES (?1, ?2, ?3) + ON CONFLICT (source_id, type_name) + DO UPDATE SET state = excluded.state", + params![source_id, ty.jmap_name(), state], + )?; + Ok(()) +} + +#[cfg(test)] +mod tests { + use super::*; + use crate::db::init; + use crate::db::sources::{SourceKey, upsert_source}; + + fn setup() -> (Connection, i64) { + let c = Connection::open_in_memory().unwrap(); + init::apply_schema(&c).unwrap(); + let sid = upsert_source( + &c, + &SourceKey { + kind: "jmap".to_owned(), + session_url: "https://host".to_owned(), + account_id: "alice".to_owned(), + }, + None, + "alice", + ) + .unwrap(); + (c, sid) + } + + #[test] + fn upsert_then_get_returns_latest_per_type() { + let (c, sid) = setup(); + assert!(get(&c, sid, ObjectType::Email).unwrap().is_none()); + upsert(&c, sid, ObjectType::Email, "s1").unwrap(); + upsert(&c, sid, ObjectType::Mailbox, "m1").unwrap(); + assert_eq!( + get(&c, sid, ObjectType::Email).unwrap().as_deref(), + Some("s1") + ); + assert_eq!( + get(&c, sid, ObjectType::Mailbox).unwrap().as_deref(), + Some("m1") + ); + upsert(&c, sid, ObjectType::Email, "s2").unwrap(); + assert_eq!( + get(&c, sid, ObjectType::Email).unwrap().as_deref(), + Some("s2") + ); + assert_eq!( + get(&c, sid, ObjectType::Mailbox).unwrap().as_deref(), + Some("m1") + ); + } +} diff --git a/src/imap/transport.rs b/src/imap/transport.rs index a538f8e..8ede26a 100644 --- a/src/imap/transport.rs +++ b/src/imap/transport.rs @@ -384,7 +384,6 @@ mod tests { #[test] fn deflate_reads_decompressed_payload_from_inner() { - let mut peer = PeerStream::new(); let mut compressor = flate2::Compress::new(flate2::Compression::default(), false); let plaintext = b"* OK [CAPABILITY IMAP4rev2 COMPRESS=DEFLATE] hi\r\n"; diff --git a/src/jmap/blob.rs b/src/jmap/blob.rs index 19714cd..ae4d2ac 100644 --- a/src/jmap/blob.rs +++ b/src/jmap/blob.rs @@ -21,7 +21,6 @@ where { match value { Value::Object(map) => { - if let Some(Value::String(blob_id)) = map.remove("blobId") { let local = resolve(&blob_id)?; map.insert(SENTINEL_KEY.to_owned(), Value::from(local)); diff --git a/src/jmap/error.rs b/src/jmap/error.rs index f014853..3b63629 100644 --- a/src/jmap/error.rs +++ b/src/jmap/error.rs @@ -33,6 +33,12 @@ pub enum JmapError { #[error("query anchor not found")] AnchorNotFound, + #[error("server cannot calculate changes from the stored state")] + CannotCalculateChanges, + + #[error("server does not implement the requested method")] + UnknownMethod, + #[error("jmap method error in call {call_id}: {error_type}{}", .description.as_deref().map(|d| format!(" ({d})")).unwrap_or_default())] Method { call_id: String, diff --git a/src/jmap/request.rs b/src/jmap/request.rs index 4a63b99..11dcca0 100644 --- a/src/jmap/request.rs +++ b/src/jmap/request.rs @@ -165,6 +165,12 @@ pub fn check_method_error(mr: &MethodCall) -> Result<(), JmapError> { if error_type == "anchorNotFound" { return Err(JmapError::AnchorNotFound); } + if error_type == "cannotCalculateChanges" { + return Err(JmapError::CannotCalculateChanges); + } + if error_type == "unknownMethod" { + return Err(JmapError::UnknownMethod); + } let description = mr .args .get("description") @@ -263,6 +269,7 @@ fn query_pages( pub struct GetResult { pub list: Vec, pub not_found: Vec, + pub state: Option, } impl Default for GetResult { @@ -270,6 +277,7 @@ impl Default for GetResult { GetResult { list: Vec::new(), not_found: Vec::new(), + state: None, } } } @@ -380,6 +388,97 @@ pub fn get_all( Ok(out) } +pub fn get_state( + client: &HttpClient, + api_url: &str, + account_id: &str, + type_name: &str, +) -> Result, JmapError> { + let mut args = Map::new(); + args.insert("accountId".to_owned(), Value::String(account_id.to_owned())); + args.insert("ids".to_owned(), Value::Array(Vec::new())); + let mut req = Request::new(); + req.call(format!("{type_name}/get"), Value::Object(args), "g"); + let resp = req.send(client, api_url)?; + let mr = resp.first()?; + check_method_error(mr)?; + let mut out: GetResult = GetResult::default(); + decode_get(mr, &mut out)?; + Ok(out.state) +} + +#[derive(Debug, Default)] +pub struct ChangesResult { + pub created: Vec, + pub updated: Vec, + pub destroyed: Vec, + pub new_state: String, +} + +pub fn get_changes( + client: &HttpClient, + api_url: &str, + account_id: &str, + type_name: &str, + since_state: &str, + limits: &Limits, +) -> Result { + let mut out = ChangesResult::default(); + let mut since = since_state.to_owned(); + loop { + let mut args = Map::new(); + args.insert("accountId".to_owned(), Value::String(account_id.to_owned())); + args.insert("sinceState".to_owned(), Value::String(since.clone())); + args.insert( + "maxChanges".to_owned(), + Value::from(limits.max_objects_in_get.max(1)), + ); + let mut req = Request::new(); + req.call(format!("{type_name}/changes"), Value::Object(args), "c"); + let resp = req.send(client, api_url)?; + let mr = resp.first()?; + check_method_error(mr)?; + let new_state = mr + .args + .get("newState") + .and_then(Value::as_str) + .ok_or_else(|| JmapError::malformed("changes response has no newState"))? + .to_owned(); + append_id_array(&mr.args, "created", &mut out.created); + append_id_array(&mr.args, "updated", &mut out.updated); + append_id_array(&mr.args, "destroyed", &mut out.destroyed); + let has_more = mr + .args + .get("hasMoreChanges") + .and_then(Value::as_bool) + .unwrap_or(false); + out.new_state = new_state.clone(); + if !has_more || since == new_state { + break; + } + since = new_state; + } + dedup_ids(&mut out.created); + dedup_ids(&mut out.updated); + dedup_ids(&mut out.destroyed); + Ok(out) +} + +fn append_id_array(args: &Value, key: &str, out: &mut Vec) { + if let Some(arr) = args.get(key).and_then(Value::as_array) { + for v in arr { + if let Some(s) = v.as_str() { + out.push(JmapId(s.to_owned())); + } + } + } +} + +fn dedup_ids(ids: &mut Vec) { + let mut seen = IndexSet::new(); + ids.retain(|id| seen.insert(id.0.clone())); +} + fn decode_get( mr: &MethodCall, out: &mut GetResult, @@ -399,6 +498,9 @@ fn decode_get( } } } + if let Some(s) = mr.args.get("state").and_then(Value::as_str) { + out.state = Some(s.to_owned()); + } Ok(()) } diff --git a/src/jmap/wire/calendar.rs b/src/jmap/wire/calendar.rs index 326ff44..d7cd7dc 100644 --- a/src/jmap/wire/calendar.rs +++ b/src/jmap/wire/calendar.rs @@ -13,7 +13,6 @@ use super::JmapId; #[derive(Debug, Clone, Serialize, Deserialize)] #[serde(rename_all = "camelCase")] pub struct Calendar { - #[serde(default, skip_serializing_if = "Option::is_none")] pub id: Option, diff --git a/src/jmap/wire/calendar_event.rs b/src/jmap/wire/calendar_event.rs index 44aa062..638066f 100644 --- a/src/jmap/wire/calendar_event.rs +++ b/src/jmap/wire/calendar_event.rs @@ -13,7 +13,6 @@ use super::JmapId; #[derive(Debug, Clone, Serialize, Deserialize)] #[serde(rename_all = "camelCase")] pub struct CalendarEvent { - #[serde(default, skip_serializing_if = "Option::is_none")] pub id: Option, diff --git a/src/jmap/wire/contact_card.rs b/src/jmap/wire/contact_card.rs index 4eb72cc..332903b 100644 --- a/src/jmap/wire/contact_card.rs +++ b/src/jmap/wire/contact_card.rs @@ -13,7 +13,6 @@ use super::JmapId; #[derive(Debug, Clone, Serialize, Deserialize)] #[serde(rename_all = "camelCase")] pub struct ContactCard { - #[serde(default, skip_serializing_if = "Option::is_none")] pub id: Option, diff --git a/src/jmap/wire/email.rs b/src/jmap/wire/email.rs index bf3df50..90b8af6 100644 --- a/src/jmap/wire/email.rs +++ b/src/jmap/wire/email.rs @@ -12,7 +12,6 @@ use super::{JmapId, UtcDate}; #[derive(Debug, Clone, Serialize, Deserialize)] #[serde(rename_all = "camelCase")] pub struct Email { - #[serde(default, skip_serializing_if = "Option::is_none")] pub id: Option, @@ -29,7 +28,6 @@ pub struct Email { #[derive(Debug, Clone, Serialize, Deserialize)] #[serde(rename_all = "camelCase")] pub struct EmailImport { - pub blob_id: JmapId, pub mailbox_ids: IndexMap, diff --git a/src/jmap/wire/file_node.rs b/src/jmap/wire/file_node.rs index d5cd229..8505e3b 100644 --- a/src/jmap/wire/file_node.rs +++ b/src/jmap/wire/file_node.rs @@ -12,7 +12,6 @@ use super::{JmapId, UtcDate}; #[derive(Debug, Clone, Serialize, Deserialize)] #[serde(rename_all = "camelCase")] pub struct FileNode { - #[serde(default, skip_serializing_if = "Option::is_none")] pub id: Option, diff --git a/src/jmap/wire/identity.rs b/src/jmap/wire/identity.rs index b88144e..4958442 100644 --- a/src/jmap/wire/identity.rs +++ b/src/jmap/wire/identity.rs @@ -11,7 +11,6 @@ use super::{EmailAddress, JmapId}; #[derive(Debug, Clone, Serialize, Deserialize)] #[serde(rename_all = "camelCase")] pub struct Identity { - #[serde(default, skip_serializing_if = "Option::is_none")] pub id: Option, diff --git a/src/jmap/wire/mailbox.rs b/src/jmap/wire/mailbox.rs index 0b4a46f..9699fdb 100644 --- a/src/jmap/wire/mailbox.rs +++ b/src/jmap/wire/mailbox.rs @@ -11,7 +11,6 @@ use super::JmapId; #[derive(Debug, Clone, Serialize, Deserialize)] #[serde(rename_all = "camelCase")] pub struct Mailbox { - #[serde(default, skip_serializing_if = "Option::is_none")] pub id: Option, diff --git a/src/jmap/wire/participant_identity.rs b/src/jmap/wire/participant_identity.rs index ed69d04..deccab0 100644 --- a/src/jmap/wire/participant_identity.rs +++ b/src/jmap/wire/participant_identity.rs @@ -11,7 +11,6 @@ use super::JmapId; #[derive(Debug, Clone, Serialize, Deserialize)] #[serde(rename_all = "camelCase")] pub struct ParticipantIdentity { - #[serde(default, skip_serializing_if = "Option::is_none")] pub id: Option, diff --git a/src/jmap/wire/sieve_script.rs b/src/jmap/wire/sieve_script.rs index 95c5d9c..cdabf7c 100644 --- a/src/jmap/wire/sieve_script.rs +++ b/src/jmap/wire/sieve_script.rs @@ -11,7 +11,6 @@ use super::JmapId; #[derive(Debug, Clone, Serialize, Deserialize)] #[serde(rename_all = "camelCase")] pub struct SieveScript { - #[serde(default, skip_serializing_if = "Option::is_none")] pub id: Option, diff --git a/src/main.rs b/src/main.rs index c6d092c..b088cbf 100644 --- a/src/main.rs +++ b/src/main.rs @@ -99,7 +99,12 @@ fn report(summary: &Summary) { for (type_name, counts) in &summary.per_type { println!( "{type_name}: created={} fetched={} updated={} deleted={} skipped={} failed={}", - counts.created, counts.fetched, counts.updated, counts.deleted, counts.skipped, counts.failed + counts.created, + counts.fetched, + counts.updated, + counts.deleted, + counts.skipped, + counts.failed ); } } diff --git a/src/managesieve/error.rs b/src/managesieve/error.rs index 52570e8..f282cfc 100644 --- a/src/managesieve/error.rs +++ b/src/managesieve/error.rs @@ -141,7 +141,6 @@ mod tests { #[test] fn permanent_codes_override_text_match() { - assert!(!NoError::new("try again later", Some("QUOTA".into())).is_transient()); assert!(!NoError::new("temporarily", Some("AUTH-TOO-WEAK".into())).is_transient()); assert!(!NoError::new("rate limit", Some("TRANSITION-NEEDED".into())).is_transient()); diff --git a/src/sync/import_jmap.rs b/src/sync/import_jmap.rs index 818e322..6a39eca 100644 --- a/src/sync/import_jmap.rs +++ b/src/sync/import_jmap.rs @@ -26,7 +26,7 @@ use crate::jmap::blobxfer; use crate::jmap::connect::{self, Connected}; use crate::jmap::error::JmapError; use crate::jmap::http::{Auth, HttpClient}; -use crate::jmap::request::{get_all, get_objects, query_all_ids}; +use crate::jmap::request::{get_all, get_changes, get_objects, get_state, query_all_ids}; use crate::jmap::session::{Limits, Session}; use crate::jmap::wire::JmapId; use crate::jmap::wire::email::Email; @@ -80,7 +80,7 @@ fn get_props(ty: ObjectType) -> Option<&'static [&'static str]> { } } -type GetMsg = Result<(Vec, usize), JmapError>; +type GetMsg = Result, JmapError>; type BlobRef = (String, String, String); type DlMsg = (String, Result, JmapError>); @@ -271,30 +271,31 @@ fn reconcile_type( }; let local_ids: HashSet = local_map.keys().cloned().collect(); - let (server_ids, preloaded): (Vec, Option>) = if is_queryless(ty) { - let got = get_all::(&net.client, &net.api, &net.account, ty.jmap_name()) + let (server_ids, preloaded, enum_state): (Vec, Option>, Option) = + if is_queryless(ty) { + let got = get_all::(&net.client, &net.api, &net.account, ty.jmap_name()) + .map_err(Error::from)?; + let ids = got + .list + .iter() + .filter_map(|v| { + v.get("id") + .and_then(Value::as_str) + .map(|s| JmapId(s.to_owned())) + }) + .collect(); + (ids, Some(got.list), got.state) + } else { + let ids = query_all_ids( + &net.client, + &net.api, + &net.account, + ty.jmap_name(), + &net.limits, + ) .map_err(Error::from)?; - let ids = got - .list - .iter() - .filter_map(|v| { - v.get("id") - .and_then(Value::as_str) - .map(|s| JmapId(s.to_owned())) - }) - .collect(); - (ids, Some(got.list)) - } else { - let ids = query_all_ids( - &net.client, - &net.api, - &net.account, - ty.jmap_name(), - &net.limits, - ) - .map_err(Error::from)?; - (ids, None) - }; + (ids, None, None) + }; let d = diff(&server_ids, &local_ids); @@ -310,6 +311,18 @@ fn reconcile_type( let source_id = source_id.expect("source_id present outside dry-run"); + let cursor = db::sync_state_jmap::get(&ctx.conn, source_id, ty) + .map_err(|e| Error::Partial(e.to_string()))?; + let run_state = if !supports_changes(ty) { + None + } else if is_queryless(ty) { + enum_state + } else if cursor.is_none() { + get_state(&net.client, &net.api, &net.account, ty.jmap_name()).map_err(Error::from)? + } else { + None + }; + if !d.new.is_empty() { let objects = match preloaded { Some(list) => { @@ -323,7 +336,7 @@ fn reconcile_type( }) .collect() } - None => fetch_objects(net, ty, &d.new, threads, logger, counts), + None => fetch_objects(net, ty, &d.new, get_props(ty), threads, logger, counts), }; insert_objects(ctx, net, ty, source_id, objects, threads, logger, counts)?; } @@ -333,6 +346,10 @@ fn reconcile_type( counts.deleted += d.vanished.len() as u64; } + reconcile_updates( + ctx, net, ty, source_id, &d.present, &local_map, cursor, run_state, threads, logger, counts, + )?; + if logger.enabled(LEVEL_DEFAULT) { eprintln!( "import: {} done (fetched={} deleted={} skipped={} failed={})", @@ -350,18 +367,17 @@ fn fetch_objects( net: &Net, ty: ObjectType, new_ids: &[JmapId], + props: Option<&'static [&'static str]>, threads: usize, logger: &Logger, counts: &mut TypeCounts, ) -> Vec { let chunk = net.limits.max_objects_in_get.max(1) as usize; - let props = get_props(ty); let workers = effective_workers(threads, &net.limits, false); let net = Arc::new(net.clone()); let pool: Pool, GetMsg> = Pool::new(workers, { let net = net.clone(); move |ids: Vec| { - let n = ids.len(); get_objects::( &net.client, &net.api, @@ -371,7 +387,7 @@ fn fetch_objects( props, &net.limits, ) - .map(|r| (r.list, n)) + .map(|r| r.list) } }); @@ -384,7 +400,7 @@ fn fetch_objects( let mut done = 0u64; for res in pool.finish() { match res { - Ok((list, _)) => out.extend(list), + Ok(list) => out.extend(list), Err(e) => { logger.warn(&format!("{} /get chunk failed: {e}", ty.jmap_name())); counts.failed += 1; @@ -689,6 +705,248 @@ fn insert_one( } } +fn supports_changes(ty: ObjectType) -> bool { + !matches!(ty, ObjectType::SieveScript) +} + +fn update_props(ty: ObjectType) -> Option<&'static [&'static str]> { + match ty { + ObjectType::Email => Some(&["mailboxIds", "keywords"]), + other => get_props(other), + } +} + +#[allow(clippy::too_many_arguments)] +fn reconcile_updates( + ctx: &Context, + net: &Net, + ty: ObjectType, + source_id: i64, + present: &[JmapId], + local_map: &HashMap, + cursor: Option, + run_state: Option, + threads: usize, + logger: &Logger, + counts: &mut TypeCounts, +) -> Result<(), Error> { + if supports_changes(ty) + && let Some(since) = cursor + { + match get_changes( + &net.client, + &net.api, + &net.account, + ty.jmap_name(), + &since, + &net.limits, + ) { + Ok(ch) => { + let present_set: HashSet<&str> = present.iter().map(|j| j.0.as_str()).collect(); + let updated: Vec = ch + .updated + .into_iter() + .filter(|j| present_set.contains(j.0.as_str())) + .collect(); + let clean = fetch_and_update( + ctx, net, ty, source_id, &updated, local_map, threads, logger, counts, + )?; + advance_state(ctx, ty, source_id, &ch.new_state, clean, logger)?; + } + Err(err) + if matches!( + err, + JmapError::CannotCalculateChanges | JmapError::UnknownMethod + ) => + { + let reason = if matches!(err, JmapError::UnknownMethod) { + "server does not implement /changes" + } else { + "server cannot calculate changes from stored state" + }; + logger.warn(&format!( + "{}: {reason}; refreshing all present objects", + ty.jmap_name() + )); + let captured = get_state(&net.client, &net.api, &net.account, ty.jmap_name()) + .map_err(Error::from)?; + let clean = fetch_and_update( + ctx, net, ty, source_id, present, local_map, threads, logger, counts, + )?; + if let Some(s) = captured { + advance_state(ctx, ty, source_id, &s, clean, logger)?; + } + } + Err(e) => return Err(Error::from(e)), + } + return Ok(()); + } + + let clean = fetch_and_update( + ctx, net, ty, source_id, present, local_map, threads, logger, counts, + )?; + if let Some(s) = run_state { + advance_state(ctx, ty, source_id, &s, clean, logger)?; + } + Ok(()) +} + +fn advance_state( + ctx: &Context, + ty: ObjectType, + source_id: i64, + state: &str, + clean: bool, + logger: &Logger, +) -> Result<(), Error> { + if !clean { + logger.warn(&format!( + "{}: holding sync state; some updates failed and will be retried on the next run", + ty.jmap_name() + )); + return Ok(()); + } + db::sync_state_jmap::upsert(&ctx.conn, source_id, ty, state) + .map_err(|e| Error::Partial(e.to_string())) +} + +#[allow(clippy::too_many_arguments)] +fn fetch_and_update( + ctx: &Context, + net: &Net, + ty: ObjectType, + source_id: i64, + ids: &[JmapId], + local_map: &HashMap, + threads: usize, + logger: &Logger, + counts: &mut TypeCounts, +) -> Result { + if ids.is_empty() { + return Ok(true); + } + let failed_before = counts.failed; + let objects = fetch_objects(net, ty, ids, update_props(ty), threads, logger, counts); + update_objects( + ctx, net, ty, source_id, objects, local_map, threads, logger, counts, + )?; + Ok(counts.failed == failed_before) +} + +#[allow(clippy::too_many_arguments)] +fn update_objects( + ctx: &Context, + net: &Net, + ty: ObjectType, + source_id: i64, + objects: Vec, + local_map: &HashMap, + threads: usize, + logger: &Logger, + counts: &mut TypeCounts, +) -> Result<(), Error> { + let blob_refs = blob_references(ty, &objects); + let blobs = if blob_refs.is_empty() { + HashMap::new() + } else { + download_blobs(net, blob_refs, threads, logger, counts) + }; + + for batch in objects.chunks(200) { + let tx = ctx + .conn + .unchecked_transaction() + .map_err(|e| Error::Partial(e.to_string()))?; + for obj in batch { + let jmap_id = match obj.get("id").and_then(Value::as_str) { + Some(s) => s.to_owned(), + None => { + counts.failed += 1; + continue; + } + }; + let Some(&local_id) = local_map.get(&jmap_id) else { + continue; + }; + match update_one(&tx, ty, source_id, local_id, obj, &blobs) { + Ok(true) => counts.updated += 1, + Ok(false) => {} + Err(e) => { + logger.warn(&format!("{} {jmap_id} update skipped: {e}", ty.jmap_name())); + counts.failed += 1; + } + } + } + tx.commit().map_err(|e| Error::Partial(e.to_string()))?; + } + Ok(()) +} + +fn update_one( + conn: &Connection, + ty: ObjectType, + source_id: i64, + local_id: i64, + obj: &Value, + blobs: &HashMap>, +) -> Result { + let resolver = DbResolver { conn, source_id }; + match ty { + ObjectType::Mailbox => { + let w = serde_json::from_value(obj.clone())?; + mapping::update_mailbox(conn, local_id, &w, &resolver) + } + ObjectType::Identity => { + let w = serde_json::from_value(obj.clone())?; + mapping::update_identity(conn, local_id, &w) + } + ObjectType::AddressBook => { + let w = serde_json::from_value(obj.clone())?; + mapping::update_address_book(conn, local_id, &w) + } + ObjectType::Calendar => { + let w = serde_json::from_value(obj.clone())?; + mapping::update_calendar(conn, local_id, &w) + } + ObjectType::ParticipantIdentity => { + let w = serde_json::from_value(obj.clone())?; + mapping::update_participant_identity(conn, local_id, &w) + } + ObjectType::Email => mapping::update_email(conn, local_id, obj, &resolver), + ObjectType::SieveScript => { + let w: SieveScript = serde_json::from_value(obj.clone())?; + let data = blobs + .get(&w.blob_id.0) + .ok_or_else(|| JmapError::malformed("sieve blob missing"))?; + let blob_local = db::blobs::intern_blob(conn, data)?; + mapping::update_sieve_script(conn, local_id, &w, blob_local) + } + ObjectType::FileNode => { + let w: FileNode = serde_json::from_value(obj.clone())?; + let blob_local = match (&w.node_type, &w.blob_id) { + (NodeType::File, Some(b)) => { + let data = blobs + .get(&b.0) + .ok_or_else(|| JmapError::malformed("file blob missing"))?; + Some(db::blobs::intern_blob(conn, data)?) + } + _ => None, + }; + mapping::update_file_node(conn, local_id, &w, blob_local, &resolver) + } + ObjectType::ContactCard => { + let w = serde_json::from_value(obj.clone())?; + let mut bi = PrefetchedBlobs { conn, bytes: blobs }; + mapping::update_contact_card(conn, local_id, &w, &resolver, &mut bi) + } + ObjectType::CalendarEvent => { + let w = serde_json::from_value(obj.clone())?; + let mut bi = PrefetchedBlobs { conn, bytes: blobs }; + mapping::update_calendar_event(conn, local_id, &w, &resolver, &mut bi) + } + } +} + fn delete_vanished( conn: &Connection, ty: ObjectType, diff --git a/src/sync/import_jmap/mapping.rs b/src/sync/import_jmap/mapping.rs index 4902d33..9507fef 100644 --- a/src/sync/import_jmap/mapping.rs +++ b/src/sync/import_jmap/mapping.rs @@ -394,6 +394,267 @@ pub fn insert_calendar_event( pub const CALENDAR_EVENT_SELECT: &str = "SELECT id, calendar_ids, is_draft, use_default_alerts, data FROM calendar_events \ WHERE data_type = 'Event'"; +pub fn update_mailbox( + conn: &Connection, + local_id: i64, + wire: &Mailbox, + resolver: &impl LocalResolver, +) -> Result { + let parent = opt_parent(resolver, ObjectType::Mailbox, &wire.parent_id); + let n = conn.execute( + "UPDATE mailboxes SET name = ?1, parent_id = ?2, role = ?3, sort_order = ?4, + is_subscribed = ?5 + WHERE id = ?6 AND (name IS NOT ?1 OR parent_id IS NOT ?2 OR role IS NOT ?3 + OR sort_order IS NOT ?4 OR is_subscribed IS NOT ?5)", + params![ + wire.name, + parent, + wire.role, + wire.sort_order, + wire.is_subscribed as i64, + local_id + ], + )?; + Ok(n > 0) +} + +pub fn update_email( + conn: &Connection, + local_id: i64, + obj: &Value, + resolver: &impl LocalResolver, +) -> Result { + let mailbox_ids: IndexMap = match obj.get("mailboxIds") { + Some(v) => serde_json::from_value(v.clone())?, + None => IndexMap::new(), + }; + let mailbox_locals = translate_in(&mailbox_ids, ObjectType::Mailbox, resolver)?; + if mailbox_locals.is_empty() { + return Err(JmapError::malformed( + "email update has no resolvable mailbox", + )); + } + let keywords: IndexMap = match obj.get("keywords") { + Some(v) => serde_json::from_value(v.clone())?, + None => IndexMap::new(), + }; + let kw: Vec = keywords.keys().cloned().collect(); + let n = conn.execute( + "UPDATE emails SET mailbox_ids = ?1, keywords = ?2 + WHERE id = ?3 AND (mailbox_ids IS NOT ?1 OR keywords IS NOT ?2)", + params![ + id_array_json(&mailbox_locals), + Value::Array(kw.iter().map(|k| Value::from(k.as_str())).collect()).to_string(), + local_id + ], + )?; + Ok(n > 0) +} + +pub fn update_identity( + conn: &Connection, + local_id: i64, + wire: &Identity, +) -> Result { + let n = conn.execute( + "UPDATE identities SET name = ?1, email = ?2, reply_to = ?3, bcc = ?4, + text_signature = ?5, html_signature = ?6 + WHERE id = ?7 AND (name IS NOT ?1 OR email IS NOT ?2 OR reply_to IS NOT ?3 + OR bcc IS NOT ?4 OR text_signature IS NOT ?5 OR html_signature IS NOT ?6)", + params![ + wire.name, + wire.email, + opt_json(&wire.reply_to)?, + opt_json(&wire.bcc)?, + wire.text_signature, + wire.html_signature, + local_id + ], + )?; + Ok(n > 0) +} + +pub fn update_sieve_script( + conn: &Connection, + local_id: i64, + wire: &SieveScript, + blob_local_id: i64, +) -> Result { + let n = conn.execute( + "UPDATE sieve_scripts SET name = ?1, is_active = ?2, blob_id = ?3 + WHERE id = ?4 AND (name IS NOT ?1 OR is_active IS NOT ?2 OR blob_id IS NOT ?3)", + params![wire.name, wire.is_active as i64, blob_local_id, local_id], + )?; + Ok(n > 0) +} + +pub fn update_address_book( + conn: &Connection, + local_id: i64, + wire: &AddressBook, +) -> Result { + let n = conn.execute( + "UPDATE address_books SET name = ?1, description = ?2, sort_order = ?3, is_default = ?4, + is_subscribed = ?5 + WHERE id = ?6 AND (name IS NOT ?1 OR description IS NOT ?2 OR sort_order IS NOT ?3 + OR is_default IS NOT ?4 OR is_subscribed IS NOT ?5)", + params![ + wire.name, + wire.description, + wire.sort_order, + wire.is_default as i64, + wire.is_subscribed as i64, + local_id + ], + )?; + Ok(n > 0) +} + +pub fn update_calendar( + conn: &Connection, + local_id: i64, + wire: &Calendar, +) -> Result { + let n = conn.execute( + "UPDATE calendars SET name = ?1, description = ?2, color = ?3, sort_order = ?4, + is_subscribed = ?5, is_visible = ?6, is_default = ?7, include_in_availability = ?8, + default_alerts_with_time = ?9, default_alerts_without_time = ?10, time_zone = ?11 + WHERE id = ?12 AND (name IS NOT ?1 OR description IS NOT ?2 OR color IS NOT ?3 + OR sort_order IS NOT ?4 OR is_subscribed IS NOT ?5 OR is_visible IS NOT ?6 + OR is_default IS NOT ?7 OR include_in_availability IS NOT ?8 + OR default_alerts_with_time IS NOT ?9 OR default_alerts_without_time IS NOT ?10 + OR time_zone IS NOT ?11)", + params![ + wire.name, + wire.description, + wire.color, + wire.sort_order, + wire.is_subscribed as i64, + wire.is_visible as i64, + wire.is_default as i64, + wire.include_in_availability, + opt_json(&wire.default_alerts_with_time)?, + opt_json(&wire.default_alerts_without_time)?, + wire.time_zone, + local_id + ], + )?; + Ok(n > 0) +} + +pub fn update_participant_identity( + conn: &Connection, + local_id: i64, + wire: &ParticipantIdentity, +) -> Result { + let n = conn.execute( + "UPDATE participant_identities SET name = ?1, calendar_address = ?2, is_default = ?3 + WHERE id = ?4 AND (name IS NOT ?1 OR calendar_address IS NOT ?2 OR is_default IS NOT ?3)", + params![ + wire.name, + wire.calendar_address, + wire.is_default as i64, + local_id + ], + )?; + Ok(n > 0) +} + +pub fn update_file_node( + conn: &Connection, + local_id: i64, + wire: &FileNode, + blob_local_id: Option, + resolver: &impl LocalResolver, +) -> Result { + let parent = opt_parent(resolver, ObjectType::FileNode, &wire.parent_id); + let node_type = serde_json::to_value(wire.node_type)? + .as_str() + .unwrap_or("file") + .to_owned(); + let target = match &wire.target { + Some(t) => Some(serde_json::to_string(t)?), + None => None, + }; + let n = conn.execute( + "UPDATE file_nodes SET parent_id = ?1, node_type = ?2, blob_id = ?3, target = ?4, + name = ?5, media_type = ?6, created = ?7, modified = ?8, is_subscribed = ?9, role = ?10 + WHERE id = ?11 AND (parent_id IS NOT ?1 OR node_type IS NOT ?2 OR blob_id IS NOT ?3 + OR target IS NOT ?4 OR name IS NOT ?5 OR media_type IS NOT ?6 OR created IS NOT ?7 + OR modified IS NOT ?8 OR is_subscribed IS NOT ?9 OR role IS NOT ?10)", + params![ + parent, + node_type, + blob_local_id, + target, + wire.name, + wire.media_type, + format_utc(&wire.created)?, + opt_format_utc(&wire.modified)?, + wire.is_subscribed as i64, + wire.role, + local_id + ], + )?; + Ok(n > 0) +} + +pub fn update_contact_card( + conn: &Connection, + local_id: i64, + wire: &ContactCard, + resolver: &impl LocalResolver, + blobs: &mut impl BlobIntern, +) -> Result { + let address_books = translate_in(&wire.address_book_ids, ObjectType::AddressBook, resolver)?; + let mut data = Value::Object(value_map(&wire.rest)); + let uid = take_string(&mut data, "uid") + .ok_or_else(|| JmapError::malformed("ContactCard has no uid"))?; + rewrite_blobs_in(&mut data, blobs)?; + let n = conn.execute( + "UPDATE contact_cards SET uid = ?1, address_book_ids = ?2, data = ?3 + WHERE id = ?4 AND (uid IS NOT ?1 OR address_book_ids IS NOT ?2 OR data IS NOT ?3)", + params![ + uid, + id_array_json(&address_books), + data.to_string(), + local_id + ], + )?; + Ok(n > 0) +} + +pub fn update_calendar_event( + conn: &Connection, + local_id: i64, + wire: &CalendarEvent, + resolver: &impl LocalResolver, + blobs: &mut impl BlobIntern, +) -> Result { + let calendars = translate_in(&wire.calendar_ids, ObjectType::Calendar, resolver)?; + let mut data = Value::Object(value_map(&wire.rest)); + for drop_key in ["method", "utcStart", "utcEnd", "isOrigin", "baseEventId"] { + if let Value::Object(m) = &mut data { + m.remove(drop_key); + } + } + rewrite_blobs_in(&mut data, blobs)?; + let n = conn.execute( + "UPDATE calendar_events SET calendar_ids = ?1, is_draft = ?2, use_default_alerts = ?3, + data = ?4 + WHERE id = ?5 AND (calendar_ids IS NOT ?1 OR is_draft IS NOT ?2 + OR use_default_alerts IS NOT ?3 OR data IS NOT ?4)", + params![ + id_array_json(&calendars), + wire.is_draft as i64, + wire.use_default_alerts as i64, + data.to_string(), + local_id + ], + )?; + Ok(n > 0) +} + pub fn contact_card_to_wire( uid: &str, address_book_ids: &str, @@ -1017,4 +1278,308 @@ mod tests { assert_eq!(got.wire.parent_id, Some(JmapId("TGTD".to_owned()))); assert_eq!(got.blob_local_id, Some(blob)); } + + fn one_row(c: &Connection, table: &str) -> i64 { + c.query_row(&format!("SELECT count(*) FROM {table}"), [], |r| r.get(0)) + .unwrap() + } + + #[test] + fn delta_update_mailbox_in_place_preserves_id() { + let c = mem(); + let empty = MapResolver { + to_local: HashMap::new(), + to_target: HashMap::new(), + }; + let m: Mailbox = + serde_json::from_value(serde_json::json!({"id":"M1","name":"Personal"})).unwrap(); + let local = insert_mailbox(&c, &m, &empty).unwrap(); + let changed: Mailbox = serde_json::from_value( + serde_json::json!({"id":"M1","name":"PersonalRenamed","role":"archive","sortOrder":4}), + ) + .unwrap(); + assert!( + update_mailbox(&c, local, &changed, &empty).unwrap(), + "a real change reports changed=true" + ); + assert!( + !update_mailbox(&c, local, &changed, &empty).unwrap(), + "re-applying identical values is a no-op (changed=false), so re-runs converge" + ); + assert_eq!(one_row(&c, "mailboxes"), 1); + let (name, role, sort): (String, Option, i64) = c + .query_row( + "SELECT name, role, sort_order FROM mailboxes WHERE id=?1", + params![local], + |r| Ok((r.get(0)?, r.get(1)?, r.get(2)?)), + ) + .unwrap(); + assert_eq!(name, "PersonalRenamed"); + assert_eq!(role.as_deref(), Some("archive")); + assert_eq!(sort, 4); + } + + #[test] + fn delta_update_email_changes_keywords_and_mailboxes_keeps_blob() { + let c = mem(); + let mut res = MapResolver { + to_local: HashMap::new(), + to_target: HashMap::new(), + }; + res.to_local + .insert((ObjectType::Mailbox, "MB".to_owned()), 7); + res.to_local + .insert((ObjectType::Mailbox, "MB2".to_owned()), 8); + let email: Email = serde_json::from_value(serde_json::json!({ + "id":"E1","blobId":"B1","receivedAt":"2021-05-04T10:00:00Z", + "mailboxIds":{"MB":true},"keywords":{"$seen":true} + })) + .unwrap(); + let blob = crate::db::blobs::intern_blob(&c, b"rfc5322").unwrap(); + let local = insert_email(&c, &email, blob, "{}", &res).unwrap(); + + let changed = serde_json::json!({ + "id":"E1","mailboxIds":{"MB2":true},"keywords":{"$seen":true,"$flagged":true} + }); + update_email(&c, local, &changed, &res).unwrap(); + assert_eq!(one_row(&c, "emails"), 1); + let got = c + .query_row( + &format!("{EMAIL_SELECT} WHERE id=?1"), + params![local], + |row| Ok(row_to_email(row)), + ) + .unwrap() + .unwrap(); + assert_eq!(got.blob_local_id, blob, "immutable body blob is untouched"); + assert_eq!(got.mailbox_locals, vec![8], "mailbox membership moved"); + assert!(got.keywords.contains(&"$flagged".to_owned())); + assert!(got.keywords.contains(&"$seen".to_owned())); + } + + #[test] + fn delta_update_identity_in_place() { + let c = mem(); + let id: Identity = + serde_json::from_value(serde_json::json!({"id":"I1","name":"Old","email":"a@x.test"})) + .unwrap(); + let local = insert_identity(&c, &id).unwrap(); + let changed: Identity = serde_json::from_value( + serde_json::json!({"id":"I1","name":"New Name","email":"a@x.test","textSignature":"sig"}), + ) + .unwrap(); + update_identity(&c, local, &changed).unwrap(); + assert_eq!(one_row(&c, "identities"), 1); + let got = c + .query_row( + &format!("{IDENTITY_SELECT} WHERE id=?1"), + params![local], + |row| Ok(row_to_identity(row)), + ) + .unwrap() + .unwrap(); + assert_eq!(got.name, "New Name"); + assert_eq!(got.text_signature, "sig"); + } + + #[test] + fn delta_update_address_book_in_place() { + let c = mem(); + let ab: AddressBook = + serde_json::from_value(serde_json::json!({"id":"A1","name":"Old"})).unwrap(); + let local = insert_address_book(&c, &ab).unwrap(); + let changed: AddressBook = serde_json::from_value( + serde_json::json!({"id":"A1","name":"Renamed","description":"d2","sortOrder":3}), + ) + .unwrap(); + update_address_book(&c, local, &changed).unwrap(); + assert_eq!(one_row(&c, "address_books"), 1); + let got = c + .query_row( + &format!("{ADDRESS_BOOK_SELECT} WHERE id=?1"), + params![local], + |r| Ok(row_to_address_book(r)), + ) + .unwrap() + .unwrap(); + assert_eq!(got.name, "Renamed"); + assert_eq!(got.description.as_deref(), Some("d2")); + } + + #[test] + fn delta_update_calendar_in_place() { + let c = mem(); + let cal: Calendar = + serde_json::from_value(serde_json::json!({"id":"C1","name":"Old","color":"#000"})) + .unwrap(); + let local = insert_calendar(&c, &cal).unwrap(); + let changed: Calendar = serde_json::from_value( + serde_json::json!({"id":"C1","name":"Renamed","color":"#abcdef","timeZone":"Europe/Rome"}), + ) + .unwrap(); + update_calendar(&c, local, &changed).unwrap(); + assert_eq!(one_row(&c, "calendars"), 1); + let got = c + .query_row( + &format!("{CALENDAR_SELECT} WHERE id=?1"), + params![local], + |r| Ok(row_to_calendar(r)), + ) + .unwrap() + .unwrap(); + assert_eq!(got.name, "Renamed"); + assert_eq!(got.color.as_deref(), Some("#abcdef")); + assert_eq!(got.time_zone.as_deref(), Some("Europe/Rome")); + } + + #[test] + fn delta_update_participant_identity_in_place() { + let c = mem(); + let pi: ParticipantIdentity = serde_json::from_value( + serde_json::json!({"id":"P1","name":"Old","calendarAddress":"mailto:me@x.test"}), + ) + .unwrap(); + let local = insert_participant_identity(&c, &pi).unwrap(); + let changed: ParticipantIdentity = serde_json::from_value( + serde_json::json!({"id":"P1","name":"New","calendarAddress":"mailto:me@x.test"}), + ) + .unwrap(); + update_participant_identity(&c, local, &changed).unwrap(); + assert_eq!(one_row(&c, "participant_identities"), 1); + let got = c + .query_row( + &format!("{PARTICIPANT_IDENTITY_SELECT} WHERE id=?1"), + params![local], + |r| Ok(row_to_participant_identity(r)), + ) + .unwrap() + .unwrap(); + assert_eq!(got.name, "New"); + } + + #[test] + fn delta_update_sieve_script_swaps_blob_and_active() { + let c = mem(); + let ss: SieveScript = serde_json::from_value( + serde_json::json!({"id":"S1","name":"main","isActive":false,"blobId":"B1"}), + ) + .unwrap(); + let blob1 = crate::db::blobs::intern_blob(&c, b"keep;").unwrap(); + let local = insert_sieve_script(&c, &ss, blob1).unwrap(); + let blob2 = crate::db::blobs::intern_blob(&c, b"discard;").unwrap(); + let changed: SieveScript = serde_json::from_value( + serde_json::json!({"id":"S1","name":"main","isActive":true,"blobId":"B2"}), + ) + .unwrap(); + update_sieve_script(&c, local, &changed, blob2).unwrap(); + assert_eq!(one_row(&c, "sieve_scripts"), 1); + let got = c + .query_row( + &format!("{SIEVE_SELECT} WHERE id=?1"), + params![local], + |r| Ok(row_to_sieve_script(r)), + ) + .unwrap() + .unwrap(); + assert!(got.is_active); + assert_eq!(got.blob_local_id, blob2, "content blob swapped"); + } + + #[test] + fn delta_update_file_node_changes_blob_and_name() { + let c = mem(); + let res = MapResolver { + to_local: HashMap::new(), + to_target: HashMap::new(), + }; + let file: FileNode = serde_json::from_value(serde_json::json!({ + "id":"F1","nodeType":"file","name":"a.bin", + "type":"application/octet-stream","created":"2022-02-02T02:02:02Z" + })) + .unwrap(); + let blob1 = crate::db::blobs::intern_blob(&c, b"v1").unwrap(); + let local = insert_file_node(&c, &file, Some(blob1), &res).unwrap(); + let blob2 = crate::db::blobs::intern_blob(&c, b"v2").unwrap(); + let changed: FileNode = serde_json::from_value(serde_json::json!({ + "id":"F1","nodeType":"file","name":"renamed.bin", + "type":"text/plain","created":"2022-02-02T02:02:02Z" + })) + .unwrap(); + update_file_node(&c, local, &changed, Some(blob2), &res).unwrap(); + assert_eq!(one_row(&c, "file_nodes"), 1); + let (name, mt, b): (String, Option, Option) = c + .query_row( + "SELECT name, media_type, blob_id FROM file_nodes WHERE id=?1", + params![local], + |r| Ok((r.get(0)?, r.get(1)?, r.get(2)?)), + ) + .unwrap(); + assert_eq!(name, "renamed.bin"); + assert_eq!(mt.as_deref(), Some("text/plain")); + assert_eq!(b, Some(blob2)); + } + + #[test] + fn delta_update_contact_card_rewrites_data_and_columns() { + let c = mem(); + let res = MapResolver { + to_local: HashMap::from([((ObjectType::AddressBook, "AB".to_owned()), 1)]), + to_target: HashMap::new(), + }; + let card: ContactCard = serde_json::from_value(serde_json::json!({ + "id":"C1","addressBookIds":{"AB":true},"uid":"u-1","name":{"full":"Old"} + })) + .unwrap(); + let mut blobs = FakeBlobs; + let local = insert_contact_card(&c, &card, &res, &mut blobs).unwrap(); + let changed: ContactCard = serde_json::from_value(serde_json::json!({ + "id":"C1","addressBookIds":{"AB":true},"uid":"u-1","name":{"full":"New Name"} + })) + .unwrap(); + update_contact_card(&c, local, &changed, &res, &mut blobs).unwrap(); + assert_eq!(one_row(&c, "contact_cards"), 1); + let (uid, data): (String, String) = c + .query_row( + &format!("{CONTACT_CARD_SELECT} WHERE id=?1"), + params![local], + |r| Ok((r.get(1)?, r.get(3)?)), + ) + .unwrap(); + assert_eq!(uid, "u-1"); + let stored: Value = serde_json::from_str(&data).unwrap(); + assert_eq!(stored["name"]["full"], Value::from("New Name")); + } + + #[test] + fn delta_update_calendar_event_rewrites_data_and_columns() { + let c = mem(); + let res = MapResolver { + to_local: HashMap::from([((ObjectType::Calendar, "CAL".to_owned()), 5)]), + to_target: HashMap::new(), + }; + let ev: CalendarEvent = serde_json::from_value(serde_json::json!({ + "id":"EV1","calendarIds":{"CAL":true},"isDraft":false, + "useDefaultAlerts":false,"title":"Old","@type":"Event" + })) + .unwrap(); + let mut blobs = FakeBlobs; + let local = insert_calendar_event(&c, &ev, &res, &mut blobs).unwrap(); + let changed: CalendarEvent = serde_json::from_value(serde_json::json!({ + "id":"EV1","calendarIds":{"CAL":true},"isDraft":true, + "useDefaultAlerts":false,"title":"Rescheduled","@type":"Event" + })) + .unwrap(); + update_calendar_event(&c, local, &changed, &res, &mut blobs).unwrap(); + assert_eq!(one_row(&c, "calendar_events"), 1); + let (dr, data): (i64, String) = c + .query_row( + &format!("{CALENDAR_EVENT_SELECT} AND id=?1"), + params![local], + |r| Ok((r.get(2)?, r.get(4)?)), + ) + .unwrap(); + assert_eq!(dr, 1, "isDraft column updated"); + let stored: Value = serde_json::from_str(&data).unwrap(); + assert_eq!(stored["title"], Value::from("Rescheduled")); + } } diff --git a/src/sync/import_maildir/keywords.rs b/src/sync/import_maildir/keywords.rs index c0e74be..5afb158 100644 --- a/src/sync/import_maildir/keywords.rs +++ b/src/sync/import_maildir/keywords.rs @@ -182,7 +182,6 @@ mod tests { #[test] fn flags_from_filename_dovecot_extension_metadata_passes_through() { - let mut flags = flags_from_filename("uid,S=1234,W=1300:2,RS"); flags.sort(); let mut want = vec![Flag::Replied, Flag::Seen]; @@ -192,7 +191,6 @@ mod tests { #[test] fn flags_from_filename_unknown_alpha_chars_skipped() { - let mut flags = flags_from_filename("uid:2,SaT"); flags.sort(); let mut want = vec![Flag::Seen, Flag::Trashed]; diff --git a/src/sync/import_maildir/messages.rs b/src/sync/import_maildir/messages.rs index e381169..fe17f5d 100644 --- a/src/sync/import_maildir/messages.rs +++ b/src/sync/import_maildir/messages.rs @@ -76,7 +76,6 @@ pub fn list_folder(folder_path: &Path) -> std::io::Result { continue; } if seen.contains_key(&unique_id) { - continue; } let flags = flags_from_filename(&filename); @@ -173,7 +172,6 @@ pub fn insert_new( #[derive(Debug, Clone, Copy, PartialEq, Eq)] pub enum PresentOutcome { - Unchanged, KeywordsUpdated, diff --git a/src/sync/import_maildir/tree.rs b/src/sync/import_maildir/tree.rs index 68c3d51..d64515f 100644 --- a/src/sync/import_maildir/tree.rs +++ b/src/sync/import_maildir/tree.rs @@ -331,7 +331,6 @@ mod tests { #[test] fn discover_rejects_dovecot_layout_fs_tree() { - let td = tempfile::tempdir().unwrap(); make_maildir(td.path(), &["Sent"]); let err = discover(td.path(), true).unwrap_err(); diff --git a/src/sync/import_managesieve/coordinator.rs b/src/sync/import_managesieve/coordinator.rs index fa683cb..ac11cb3 100644 --- a/src/sync/import_managesieve/coordinator.rs +++ b/src/sync/import_managesieve/coordinator.rs @@ -420,7 +420,6 @@ fn account_id_for(auth: &ManageSieveAuth) -> String { #[derive(Debug)] enum SieveAuthError { - TerminallyRefused(String), NoUsableMechanism(String), diff --git a/src/sync/keys.rs b/src/sync/keys.rs index 046f667..6d80604 100644 --- a/src/sync/keys.rs +++ b/src/sync/keys.rs @@ -48,7 +48,6 @@ pub enum EmailKey { #[derive(Debug, Clone, PartialEq, Eq)] pub struct EmailIndex { - pub mids: Vec, pub fb: [u8; 32], diff --git a/tests/mock_exchange_ews.rs b/tests/mock_exchange_ews.rs index 153595d..a1363bd 100644 --- a/tests/mock_exchange_ews.rs +++ b/tests/mock_exchange_ews.rs @@ -454,7 +454,6 @@ fn mailbox_kinds_are_three_separate_sources() { #[test] fn get_item_batches_chunk_the_id_list() { - let ids: Vec = (0..7).map(|i| ItemId::new(format!("I{i}"), "K")).collect(); let chunks: Vec<&[ItemId]> = ids.chunks(3).collect(); assert_eq!(chunks.len(), 3); diff --git a/tests/mock_imap.rs b/tests/mock_imap.rs index 3617284..9a3593f 100644 --- a/tests/mock_imap.rs +++ b/tests/mock_imap.rs @@ -807,7 +807,11 @@ fn coordinator_present_run_is_convergent() { assert_eq!(email.1.deleted, 0, "convergent run deletes nothing"); assert_eq!(email.1.updated, 0, "unchanged flags update nothing"); let dbc = Connection::open(&archive).unwrap(); - assert_eq!(count(&dbc, "blobs"), 1, "no body re-fetched on a present-only run"); + assert_eq!( + count(&dbc, "blobs"), + 1, + "no body re-fetched on a present-only run" + ); } #[test] @@ -830,7 +834,10 @@ fn coordinator_present_flag_change_updates_keywords() { .iter() .find(|(k, _)| *k == "email") .unwrap(); - assert_eq!(email.1.updated, 1, "a changed flag set is counted as updated"); + assert_eq!( + email.1.updated, 1, + "a changed flag set is counted as updated" + ); assert_eq!(email.1.created, 0, "present message is not re-created"); assert_eq!(email.1.fetched, 0, "no body fetched on a present-only run"); let dbc = Connection::open(&archive).unwrap(); @@ -869,7 +876,11 @@ fn coordinator_present_newly_deleted_is_left_intact() { "a present message that newly gained \\Deleted is skipped, not updated" ); let dbc = Connection::open(&archive).unwrap(); - assert_eq!(count(&dbc, "emails"), 1, "the archived message is preserved"); + assert_eq!( + count(&dbc, "emails"), + 1, + "the archived message is preserved" + ); let kw: String = dbc .query_row("SELECT keywords FROM emails LIMIT 1", [], |r| r.get(0)) .unwrap(); @@ -1080,7 +1091,6 @@ fn coordinator_dispatches_to_multiple_worker_connections() { continue; } if cmd.starts_with("UID FETCH") { - let after = cmd.strip_prefix("UID FETCH ").unwrap_or(""); let set = after.split_whitespace().next().unwrap_or(""); let uids = parse_uid_set(set); diff --git a/tests/mock_jmap.rs b/tests/mock_jmap.rs index 1833e22..4fd0f6f 100644 --- a/tests/mock_jmap.rs +++ b/tests/mock_jmap.rs @@ -10,7 +10,9 @@ use serde_json::json; use vandelay::jmap::account::{self, AccountSelector}; use vandelay::jmap::error::JmapError; use vandelay::jmap::http::{Auth, HttpClient, RetryPolicy}; -use vandelay::jmap::request::{self, SetRequest, get_all, get_objects, set_call}; +use vandelay::jmap::request::{ + self, SetRequest, get_all, get_changes, get_objects, get_state, set_call, +}; use vandelay::jmap::session::{Limits, Session}; use vandelay::jmap::wire::JmapId; use vandelay::jmap::wire::identity::Identity; @@ -1040,3 +1042,130 @@ fn shared_throttle_level_grows_across_concurrent_workers() { "both workers saw a 429 (or more) before recovery" ); } + +#[test] +fn get_changes_paginates_until_no_more() { + let mut server = mockito::Server::new(); + let api = "/jmap/api"; + let _p1 = server + .mock("POST", api) + .match_body(mockito::Matcher::Regex("\"sinceState\":\"s1\"".into())) + .with_body( + json!({"methodResponses":[["Mailbox/changes",{"accountId":"w","oldState":"s1", + "newState":"s2","hasMoreChanges":true,"created":["A"],"updated":["U1"], + "destroyed":["D1"]},"c"]]}) + .to_string(), + ) + .expect(1) + .create(); + let _p2 = server + .mock("POST", api) + .match_body(mockito::Matcher::Regex("\"sinceState\":\"s2\"".into())) + .with_body( + json!({"methodResponses":[["Mailbox/changes",{"accountId":"w","oldState":"s2", + "newState":"s3","hasMoreChanges":false,"created":[],"updated":["U2"], + "destroyed":[]},"c"]]}) + .to_string(), + ) + .expect(1) + .create(); + + let url = format!("{}{}", server.url(), api); + let r = get_changes(&client(0), &url, "w", "Mailbox", "s1", &limits(500)).expect("changes"); + let updated: Vec = r.updated.iter().map(|i| i.0.clone()).collect(); + assert_eq!(updated, vec!["U1".to_owned(), "U2".to_owned()]); + let created: Vec = r.created.iter().map(|i| i.0.clone()).collect(); + assert_eq!(created, vec!["A".to_owned()]); + let destroyed: Vec = r.destroyed.iter().map(|i| i.0.clone()).collect(); + assert_eq!(destroyed, vec!["D1".to_owned()]); + assert_eq!(r.new_state, "s3", "cursor advances to the final newState"); +} + +#[test] +fn get_changes_cannot_calculate_changes_is_typed_error() { + let mut server = mockito::Server::new(); + let api = "/jmap/api"; + let _m = server + .mock("POST", api) + .with_body( + json!({"methodResponses":[["error",{"type":"cannotCalculateChanges"},"c"]]}) + .to_string(), + ) + .create(); + let url = format!("{}{}", server.url(), api); + let err = + get_changes(&client(0), &url, "w", "Email", "stale", &limits(500)).expect_err("must error"); + assert!( + matches!(err, JmapError::CannotCalculateChanges), + "got {err:?}" + ); +} + +#[test] +fn get_changes_dedups_repeated_ids_across_pages() { + let mut server = mockito::Server::new(); + let api = "/jmap/api"; + let _p1 = server + .mock("POST", api) + .match_body(mockito::Matcher::Regex("\"sinceState\":\"s1\"".into())) + .with_body( + json!({"methodResponses":[["Mailbox/changes",{"accountId":"w","oldState":"s1", + "newState":"s2","hasMoreChanges":true,"created":[],"updated":["U1","U2"], + "destroyed":[]},"c"]]}) + .to_string(), + ) + .expect(1) + .create(); + let _p2 = server + .mock("POST", api) + .match_body(mockito::Matcher::Regex("\"sinceState\":\"s2\"".into())) + .with_body( + json!({"methodResponses":[["Mailbox/changes",{"accountId":"w","oldState":"s2", + "newState":"s3","hasMoreChanges":false,"created":[],"updated":["U1","U3"], + "destroyed":[]},"c"]]}) + .to_string(), + ) + .expect(1) + .create(); + + let url = format!("{}{}", server.url(), api); + let r = get_changes(&client(0), &url, "w", "Mailbox", "s1", &limits(500)).expect("changes"); + let updated: Vec = r.updated.iter().map(|i| i.0.clone()).collect(); + assert_eq!( + updated, + vec!["U1".to_owned(), "U2".to_owned(), "U3".to_owned()], + "an id repeated across pages is fetched once, first-seen order preserved" + ); +} + +#[test] +fn get_changes_unknown_method_is_typed_error() { + let mut server = mockito::Server::new(); + let api = "/jmap/api"; + let _m = server + .mock("POST", api) + .with_body(json!({"methodResponses":[["error",{"type":"unknownMethod"},"c"]]}).to_string()) + .create(); + let url = format!("{}{}", server.url(), api); + let err = + get_changes(&client(0), &url, "w", "Email", "s1", &limits(500)).expect_err("must error"); + assert!(matches!(err, JmapError::UnknownMethod), "got {err:?}"); +} + +#[test] +fn get_state_reads_state_from_empty_get() { + let mut server = mockito::Server::new(); + let api = "/jmap/api"; + let _m = server + .mock("POST", api) + .match_body(mockito::Matcher::Regex("Email/get".into())) + .with_body( + json!({"methodResponses":[["Email/get",{"accountId":"w","state":"snap-1", + "list":[],"notFound":[]},"g"]]}) + .to_string(), + ) + .create(); + let url = format!("{}{}", server.url(), api); + let st = get_state(&client(0), &url, "w", "Email").expect("get_state"); + assert_eq!(st.as_deref(), Some("snap-1")); +} diff --git a/tests/mock_maildir.rs b/tests/mock_maildir.rs index 10c39ce..4b89d70 100644 --- a/tests/mock_maildir.rs +++ b/tests/mock_maildir.rs @@ -379,7 +379,6 @@ fn qmail_root_only_maildir_imports_inbox_only() { } fn cyrus_export_fixture(root: &Path) { - ensure_maildir(root); write_message( root, @@ -537,7 +536,6 @@ fn rejects_path_without_cur_subdir() { #[test] fn rejects_dovecot_layout_fs_tree() { - let td = TempDir::new().unwrap(); ensure_maildir(td.path()); for s in ["cur", "new", "tmp"] { @@ -720,7 +718,6 @@ fn unreadable_file_is_counted_and_warned_not_aborted() { #[test] fn pointing_at_a_dot_subfolder_warns_but_imports() { - let parent = TempDir::new().unwrap(); ensure_maildir(parent.path()); let sub = ensure_subfolder(parent.path(), ".Sent"); @@ -738,7 +735,6 @@ fn pointing_at_a_dot_subfolder_warns_but_imports() { #[test] fn malformed_message_yields_zero_message_match_but_imports() { - let td = TempDir::new().unwrap(); ensure_maildir(td.path()); write_message(td.path(), "cur", "1.M0.host:2,S", b""); diff --git a/tests/mock_managesieve.rs b/tests/mock_managesieve.rs index cf71fff..aa34e71 100644 --- a/tests/mock_managesieve.rs +++ b/tests/mock_managesieve.rs @@ -533,7 +533,6 @@ fn vanished_script_is_deleted_locally() { #[test] fn active_flag_flip_does_not_violate_partial_unique_index() { - let seed: Script = Box::new(|conn| { auth_then(conn, |c| { let _ = c.read_line()?; @@ -580,7 +579,6 @@ fn active_flag_flip_does_not_violate_partial_unique_index() { #[test] fn resume_after_partial_run_completes_remainder() { - let seed: Script = Box::new(|conn| { auth_then(conn, |c| { let _ = c.read_line()?; @@ -694,7 +692,6 @@ fn transient_no_on_getscript_retries_then_succeeds() { #[test] fn bye_mid_listscripts_reconnects_and_succeeds() { - let first: Script = Box::new(|conn| { conn.write_capability("PLAIN", false)?; conn.write_line("OK")?; @@ -749,7 +746,6 @@ fn bye_mid_listscripts_reconnects_and_succeeds() { #[test] fn bye_mid_getscript_reconnects_and_completes_remaining_scripts() { - let first: Script = Box::new(|conn| { conn.write_capability("PLAIN", false)?; conn.write_line("OK")?; @@ -812,7 +808,6 @@ fn bye_mid_getscript_reconnects_and_completes_remaining_scripts() { #[test] fn post_auth_unsolicited_capability_is_consumed() { - let server = MockSieveServer::start(|conn| { conn.write_capability("PLAIN", false)?; conn.write_line("OK")?; @@ -994,7 +989,6 @@ fn three_scripts_with_middle_active_lands_the_middle_active() { #[test] fn transient_no_exhausting_max_retries_skips_script_and_followup_picks_it_up() { - let archive = tempfile("retry_exhaustion"); let first = MockSieveServer::start(|conn| { auth_then(conn, |c| { diff --git a/tests/mock_sync.rs b/tests/mock_sync.rs index cc74349..63bc2ed 100644 --- a/tests/mock_sync.rs +++ b/tests/mock_sync.rs @@ -433,7 +433,7 @@ fn import_removes_vanished_mailbox_from_archive_on_second_pass() { .mock("POST", api) .match_body(Matcher::Regex("Mailbox/get".into())) .with_body( - json!({"methodResponses":[["Mailbox/get",{"accountId":"w","list":[ + json!({"methodResponses":[["Mailbox/get",{"accountId":"w","state":"s1","list":[ {"id":"A","name":"alpha","parentId":null,"role":null,"sortOrder":0,"isSubscribed":true}, {"id":"B","name":"bravo","parentId":null,"role":null,"sortOrder":0,"isSubscribed":true}, {"id":"C","name":"charlie","parentId":null,"role":null,"sortOrder":0,"isSubscribed":true} @@ -473,6 +473,16 @@ fn import_removes_vanished_mailbox_from_archive_on_second_pass() { ) .expect(1) .create(); + let _ch2 = server + .mock("POST", api) + .match_body(Matcher::Regex("Mailbox/changes".into())) + .with_body( + json!({"methodResponses":[["Mailbox/changes",{"accountId":"w","oldState":"s1", + "newState":"s2","hasMoreChanges":false,"created":[],"updated":[],"destroyed":[]},"c"]]}) + .to_string(), + ) + .expect(1) + .create(); let s2 = sync::import_jmap::run( common(&archive), @@ -486,6 +496,7 @@ fn import_removes_vanished_mailbox_from_archive_on_second_pass() { .map(|(_, c)| c.clone()) .expect("mailbox counts"); assert_eq!(mb2.fetched, 0, "second pass fetches nothing"); + assert_eq!(mb2.updated, 0, "no changed mailboxes reported"); assert_eq!(mb2.deleted, 1, "vanished mailbox B is deleted"); { let conn = rusqlite::Connection::open(&archive).unwrap(); @@ -502,7 +513,7 @@ fn import_removes_vanished_mailbox_from_archive_on_second_pass() { } #[test] -fn import_present_item_change_on_server_is_not_propagated() { +fn import_present_item_change_is_propagated_via_changes() { let mut server = mockito::Server::new(); let base = server.url(); let api = "/jmap/api"; @@ -530,7 +541,7 @@ fn import_present_item_change_on_server_is_not_propagated() { .mock("POST", api) .match_body(Matcher::Regex("Mailbox/get".into())) .with_body( - json!({"methodResponses":[["Mailbox/get",{"accountId":"w","list":[ + json!({"methodResponses":[["Mailbox/get",{"accountId":"w","state":"s1","list":[ {"id":"A","name":"OriginalName","parentId":null,"role":null,"sortOrder":0,"isSubscribed":true} ],"notFound":[]},"g"]]}) .to_string(), @@ -543,6 +554,17 @@ fn import_present_item_change_on_server_is_not_propagated() { import_cfg_objects(&base, vec![ObjectType::Mailbox]), ) .expect("first import"); + { + let conn = rusqlite::Connection::open(&archive).unwrap(); + let cursor: String = conn + .query_row( + "SELECT state FROM sync_state_jmap WHERE type_name='Mailbox'", + [], + |r| r.get(0), + ) + .expect("first import records the state cursor"); + assert_eq!(cursor, "s1"); + } let _q2 = server .mock("POST", api) @@ -554,10 +576,29 @@ fn import_present_item_change_on_server_is_not_propagated() { ) .expect(1) .create(); - let nope_get = server + let changes = server + .mock("POST", api) + .match_body(Matcher::AllOf(vec![ + Matcher::Regex("Mailbox/changes".into()), + Matcher::Regex("\"sinceState\":\"s1\"".into()), + ])) + .with_body( + json!({"methodResponses":[["Mailbox/changes",{"accountId":"w","oldState":"s1", + "newState":"s2","hasMoreChanges":false,"created":[],"updated":["A"],"destroyed":[]},"c"]]}) + .to_string(), + ) + .expect(1) + .create(); + let _g2 = server .mock("POST", api) .match_body(Matcher::Regex("Mailbox/get".into())) - .expect(0) + .with_body( + json!({"methodResponses":[["Mailbox/get",{"accountId":"w","state":"s2","list":[ + {"id":"A","name":"UpdatedName","parentId":null,"role":null,"sortOrder":0,"isSubscribed":true} + ],"notFound":[]},"g"]]}) + .to_string(), + ) + .expect(1) .create(); let s2 = sync::import_jmap::run( @@ -565,22 +606,34 @@ fn import_present_item_change_on_server_is_not_propagated() { import_cfg_objects(&base, vec![ObjectType::Mailbox]), ) .expect("second import"); - nope_get.assert(); + changes.assert(); let mb2 = s2 .per_type .iter() .find(|(t, _)| *t == "Mailbox") .map(|(_, c)| c.clone()) .expect("mailbox counts"); + assert_eq!(mb2.fetched, 0, "no new objects on the second pass"); assert_eq!( - mb2.fetched, 0, - "present items must not be re-fetched: changes on server are intentionally ignored" + mb2.updated, 1, + "the changed mailbox is detected via /changes and refreshed in place" ); - let name: String = rusqlite::Connection::open(&archive) - .unwrap() + let conn = rusqlite::Connection::open(&archive).unwrap(); + let name: String = conn .query_row("SELECT name FROM mailboxes WHERE id=1", [], |r| r.get(0)) .unwrap(); - assert_eq!(name, "OriginalName", "archive name was not overwritten"); + assert_eq!( + name, "UpdatedName", + "a server-side property change is propagated into the archive" + ); + let cursor: String = conn + .query_row( + "SELECT state FROM sync_state_jmap WHERE type_name='Mailbox'", + [], + |r| r.get(0), + ) + .unwrap(); + assert_eq!(cursor, "s2", "cursor advances to the changes newState"); let _ = std::fs::remove_file(&archive); } @@ -614,7 +667,7 @@ fn import_removes_vanished_email_and_drops_cross_ref() { .mock("POST", api) .match_body(Matcher::Regex("Mailbox/get".into())) .with_body( - json!({"methodResponses":[["Mailbox/get",{"accountId":"w","list":[ + json!({"methodResponses":[["Mailbox/get",{"accountId":"w","state":"sm1","list":[ {"id":"MX","name":"Inbox","parentId":null,"role":"inbox","sortOrder":0,"isSubscribed":true} ],"notFound":[]},"g"]]}) .to_string(), @@ -635,7 +688,7 @@ fn import_removes_vanished_email_and_drops_cross_ref() { .mock("POST", api) .match_body(Matcher::Regex("Email/get".into())) .with_body( - json!({"methodResponses":[["Email/get",{"accountId":"w","list":[ + json!({"methodResponses":[["Email/get",{"accountId":"w","state":"se1","list":[ {"id":"E1","blobId":"BLB1","receivedAt":"2020-01-01T00:00:00Z","mailboxIds":{"MX":true},"keywords":{"$seen":true}}, {"id":"E2","blobId":"BLB2","receivedAt":"2020-01-02T00:00:00Z","mailboxIds":{"MX":true},"keywords":{}} ],"notFound":[]},"g"]]}) @@ -688,6 +741,26 @@ fn import_removes_vanished_email_and_drops_cross_ref() { ) .expect(1) .create(); + let _mbch2 = server + .mock("POST", api) + .match_body(Matcher::Regex("Mailbox/changes".into())) + .with_body( + json!({"methodResponses":[["Mailbox/changes",{"accountId":"w","oldState":"sm1", + "newState":"sm2","hasMoreChanges":false,"created":[],"updated":[],"destroyed":[]},"c"]]}) + .to_string(), + ) + .expect(1) + .create(); + let _ech2 = server + .mock("POST", api) + .match_body(Matcher::Regex("Email/changes".into())) + .with_body( + json!({"methodResponses":[["Email/changes",{"accountId":"w","oldState":"se1", + "newState":"se2","hasMoreChanges":false,"created":[],"updated":[],"destroyed":[]},"c"]]}) + .to_string(), + ) + .expect(1) + .create(); let s2 = sync::import_jmap::run( common(&archive), @@ -1297,7 +1370,10 @@ fn export_sieve_scripts_identical_content_different_names_both_created() { counts.created, 2, "two scripts with identical content but distinct names must both reach the target" ); - assert_eq!(counts.skipped, 0, "neither distinct name collapses onto the other"); + assert_eq!( + counts.skipped, 0, + "neither distinct name collapses onto the other" + ); assert_eq!(counts.failed, 0); let _ = std::fs::remove_file(&archive); } @@ -1539,3 +1615,735 @@ fn import_dry_run_does_not_write_archive_or_download_blobs() { assert_eq!(source_rows, 0, "dry-run must not record the JMAP source"); let _ = std::fs::remove_file(&archive); } + +#[test] +fn import_email_keyword_change_is_propagated_without_blob_refetch() { + let mut server = mockito::Server::new(); + let base = server.url(); + let api = "/jmap/api"; + let archive = tmp(); + + let _root = server.mock("GET", "/").with_status(404).create(); + let _wk = server + .mock("GET", "/.well-known/jmap") + .with_body(session_body_full(&base)) + .expect_at_least(1) + .create(); + let _mbterm = anchor_terminator(&mut server, api, "Mailbox"); + let _emterm = anchor_terminator(&mut server, api, "Email"); + + let _mbq1 = server + .mock("POST", api) + .match_body(Matcher::Regex("Mailbox/query".into())) + .with_body( + json!({"methodResponses":[["Mailbox/query",{"accountId":"w","ids":["MX"]},"q"]]}) + .to_string(), + ) + .expect(1) + .create(); + let _mbg1 = server + .mock("POST", api) + .match_body(Matcher::Regex("Mailbox/get".into())) + .with_body(json!({"methodResponses":[["Mailbox/get",{"accountId":"w","state":"sm1","list":[ + {"id":"MX","name":"Inbox","parentId":null,"role":"inbox","sortOrder":0,"isSubscribed":true} + ],"notFound":[]},"g"]]}).to_string()) + .expect(1) + .create(); + let _eq1 = server + .mock("POST", api) + .match_body(Matcher::Regex("Email/query".into())) + .with_body( + json!({"methodResponses":[["Email/query",{"accountId":"w","ids":["E1"]},"q"]]}) + .to_string(), + ) + .expect(1) + .create(); + let _eg1 = server + .mock("POST", api) + .match_body(Matcher::Regex("Email/get".into())) + .with_body(json!({"methodResponses":[["Email/get",{"accountId":"w","state":"se1","list":[ + {"id":"E1","blobId":"BLB1","receivedAt":"2020-01-01T00:00:00Z","mailboxIds":{"MX":true},"keywords":{"$seen":true}} + ],"notFound":[]},"g"]]}).to_string()) + .expect(1) + .create(); + let _dl1 = server + .mock("GET", Matcher::Regex("/jmap/dl/w/BLB1/.*".into())) + .with_body("From: a@x\r\nMessage-ID: <1@h>\r\n\r\nbody-one") + .expect(1) + .create(); + + sync::import_jmap::run( + common(&archive), + import_cfg_objects(&base, vec![ObjectType::Mailbox, ObjectType::Email]), + ) + .expect("first import"); + + let _mbq2 = server + .mock("POST", api) + .match_body(Matcher::Regex("Mailbox/query".into())) + .with_body( + json!({"methodResponses":[["Mailbox/query",{"accountId":"w","ids":["MX"]},"q"]]}) + .to_string(), + ) + .expect(1) + .create(); + let _mbch2 = server + .mock("POST", api) + .match_body(Matcher::Regex("Mailbox/changes".into())) + .with_body(json!({"methodResponses":[["Mailbox/changes",{"accountId":"w","oldState":"sm1","newState":"sm2","hasMoreChanges":false,"created":[],"updated":[],"destroyed":[]},"c"]]}).to_string()) + .expect(1) + .create(); + let _eq2 = server + .mock("POST", api) + .match_body(Matcher::Regex("Email/query".into())) + .with_body( + json!({"methodResponses":[["Email/query",{"accountId":"w","ids":["E1"]},"q"]]}) + .to_string(), + ) + .expect(1) + .create(); + let echanges = server + .mock("POST", api) + .match_body(Matcher::Regex("Email/changes".into())) + .with_body(json!({"methodResponses":[["Email/changes",{"accountId":"w","oldState":"se1","newState":"se2","hasMoreChanges":false,"created":[],"updated":["E1"],"destroyed":[]},"c"]]}).to_string()) + .expect(1) + .create(); + let _eg2 = server + .mock("POST", api) + .match_body(Matcher::Regex("Email/get".into())) + .with_body( + json!({"methodResponses":[["Email/get",{"accountId":"w","state":"se2","list":[ + {"id":"E1","mailboxIds":{"MX":true},"keywords":{"$seen":true,"$flagged":true}} + ],"notFound":[]},"g"]]}) + .to_string(), + ) + .expect(1) + .create(); + let no_blob_refetch = server + .mock("GET", Matcher::Regex("/jmap/dl/w/BLB1/.*".into())) + .with_body("should-not-be-fetched") + .expect(0) + .create(); + + let s2 = sync::import_jmap::run( + common(&archive), + import_cfg_objects(&base, vec![ObjectType::Mailbox, ObjectType::Email]), + ) + .expect("second import"); + echanges.assert(); + no_blob_refetch.assert(); + let em2 = s2 + .per_type + .iter() + .find(|(t, _)| *t == "Email") + .map(|(_, c)| c.clone()) + .expect("email counts"); + assert_eq!(em2.updated, 1, "the changed email is refreshed"); + assert_eq!(em2.fetched, 0, "no new emails"); + let conn = rusqlite::Connection::open(&archive).unwrap(); + let kw: String = conn + .query_row("SELECT keywords FROM emails LIMIT 1", [], |r| r.get(0)) + .unwrap(); + assert!( + kw.contains("$seen") && kw.contains("$flagged"), + "keyword change propagated into the archive: {kw}" + ); + let blobs: i64 = conn + .query_row("SELECT count(*) FROM blobs", [], |r| r.get(0)) + .unwrap(); + assert_eq!( + blobs, 1, + "the immutable body blob is not re-downloaded on update" + ); + let _ = std::fs::remove_file(&archive); +} + +#[test] +fn import_cannot_calculate_changes_falls_back_to_full_refresh() { + let mut server = mockito::Server::new(); + let base = server.url(); + let api = "/jmap/api"; + let archive = tmp(); + + let _root = server.mock("GET", "/").with_status(404).create(); + let _wk = server + .mock("GET", "/.well-known/jmap") + .with_body(session_body_full(&base)) + .expect_at_least(1) + .create(); + let _term = anchor_terminator(&mut server, api, "Mailbox"); + + let _q1 = server + .mock("POST", api) + .match_body(Matcher::Regex("Mailbox/query".into())) + .with_body( + json!({"methodResponses":[["Mailbox/query",{"accountId":"w","ids":["A"]},"q"]]}) + .to_string(), + ) + .expect(1) + .create(); + let _g1 = server + .mock("POST", api) + .match_body(Matcher::Regex("Mailbox/get".into())) + .with_body(json!({"methodResponses":[["Mailbox/get",{"accountId":"w","state":"s1","list":[ + {"id":"A","name":"OriginalName","parentId":null,"role":null,"sortOrder":0,"isSubscribed":true} + ],"notFound":[]},"g"]]}).to_string()) + .expect(1) + .create(); + + sync::import_jmap::run( + common(&archive), + import_cfg_objects(&base, vec![ObjectType::Mailbox]), + ) + .expect("first import"); + + let _q2 = server + .mock("POST", api) + .match_body(Matcher::Regex("Mailbox/query".into())) + .with_body( + json!({"methodResponses":[["Mailbox/query",{"accountId":"w","ids":["A"]},"q"]]}) + .to_string(), + ) + .expect(1) + .create(); + let cannot = server + .mock("POST", api) + .match_body(Matcher::Regex("Mailbox/changes".into())) + .with_body( + json!({"methodResponses":[["error",{"type":"cannotCalculateChanges"},"c"]]}) + .to_string(), + ) + .expect(1) + .create(); + let capture_state = server + .mock("POST", api) + .match_body(Matcher::AllOf(vec![ + Matcher::Regex("Mailbox/get".into()), + Matcher::Regex("\"ids\":\\[\\]".into()), + ])) + .with_body(json!({"methodResponses":[["Mailbox/get",{"accountId":"w","state":"s9","list":[],"notFound":[]},"g"]]}).to_string()) + .expect(1) + .create(); + let refresh_get = server + .mock("POST", api) + .match_body(Matcher::AllOf(vec![ + Matcher::Regex("Mailbox/get".into()), + Matcher::Regex("\"A\"".into()), + ])) + .with_body(json!({"methodResponses":[["Mailbox/get",{"accountId":"w","state":"s9","list":[ + {"id":"A","name":"RefreshedName","parentId":null,"role":null,"sortOrder":0,"isSubscribed":true} + ],"notFound":[]},"g"]]}).to_string()) + .expect(1) + .create(); + + let s2 = sync::import_jmap::run( + common(&archive), + import_cfg_objects(&base, vec![ObjectType::Mailbox]), + ) + .expect("second import"); + cannot.assert(); + capture_state.assert(); + refresh_get.assert(); + let mb2 = s2 + .per_type + .iter() + .find(|(t, _)| *t == "Mailbox") + .map(|(_, c)| c.clone()) + .expect("mailbox counts"); + assert_eq!(mb2.updated, 1, "fallback refreshes the present object"); + let conn = rusqlite::Connection::open(&archive).unwrap(); + let name: String = conn + .query_row("SELECT name FROM mailboxes WHERE id=1", [], |r| r.get(0)) + .unwrap(); + assert_eq!(name, "RefreshedName", "A-fallback propagated the change"); + let cursor: String = conn + .query_row( + "SELECT state FROM sync_state_jmap WHERE type_name='Mailbox'", + [], + |r| r.get(0), + ) + .unwrap(); + assert_eq!( + cursor, "s9", + "fallback captured a fresh cursor for the next run" + ); + let _ = std::fs::remove_file(&archive); +} + +#[test] +fn import_failed_update_holds_cursor_for_retry() { + let mut server = mockito::Server::new(); + let base = server.url(); + let api = "/jmap/api"; + let archive = tmp(); + + let _root = server.mock("GET", "/").with_status(404).create(); + let _wk = server + .mock("GET", "/.well-known/jmap") + .with_body(session_body_full(&base)) + .expect_at_least(1) + .create(); + let _mbterm = anchor_terminator(&mut server, api, "Mailbox"); + let _emterm = anchor_terminator(&mut server, api, "Email"); + + let _mbq1 = server + .mock("POST", api) + .match_body(Matcher::Regex("Mailbox/query".into())) + .with_body( + json!({"methodResponses":[["Mailbox/query",{"accountId":"w","ids":["MX"]},"q"]]}) + .to_string(), + ) + .expect(1) + .create(); + let _mbg1 = server + .mock("POST", api) + .match_body(Matcher::Regex("Mailbox/get".into())) + .with_body(json!({"methodResponses":[["Mailbox/get",{"accountId":"w","state":"sm1","list":[ + {"id":"MX","name":"Inbox","parentId":null,"role":"inbox","sortOrder":0,"isSubscribed":true} + ],"notFound":[]},"g"]]}).to_string()) + .expect(1) + .create(); + let _eq1 = server + .mock("POST", api) + .match_body(Matcher::Regex("Email/query".into())) + .with_body( + json!({"methodResponses":[["Email/query",{"accountId":"w","ids":["E1"]},"q"]]}) + .to_string(), + ) + .expect(1) + .create(); + let _eg1 = server + .mock("POST", api) + .match_body(Matcher::Regex("Email/get".into())) + .with_body(json!({"methodResponses":[["Email/get",{"accountId":"w","state":"se1","list":[ + {"id":"E1","blobId":"BLB1","receivedAt":"2020-01-01T00:00:00Z","mailboxIds":{"MX":true},"keywords":{"$seen":true}} + ],"notFound":[]},"g"]]}).to_string()) + .expect(1) + .create(); + let _dl1 = server + .mock("GET", Matcher::Regex("/jmap/dl/w/BLB1/.*".into())) + .with_body("From: a@x\r\nMessage-ID: <1@h>\r\n\r\nbody-one") + .expect(1) + .create(); + + sync::import_jmap::run( + common(&archive), + import_cfg_objects(&base, vec![ObjectType::Mailbox, ObjectType::Email]), + ) + .expect("first import"); + + let _mbq2 = server + .mock("POST", api) + .match_body(Matcher::Regex("Mailbox/query".into())) + .with_body( + json!({"methodResponses":[["Mailbox/query",{"accountId":"w","ids":["MX"]},"q"]]}) + .to_string(), + ) + .expect(1) + .create(); + let _mbch2 = server + .mock("POST", api) + .match_body(Matcher::Regex("Mailbox/changes".into())) + .with_body(json!({"methodResponses":[["Mailbox/changes",{"accountId":"w","oldState":"sm1","newState":"sm2","hasMoreChanges":false,"created":[],"updated":[],"destroyed":[]},"c"]]}).to_string()) + .expect(1) + .create(); + let _eq2 = server + .mock("POST", api) + .match_body(Matcher::Regex("Email/query".into())) + .with_body( + json!({"methodResponses":[["Email/query",{"accountId":"w","ids":["E1"]},"q"]]}) + .to_string(), + ) + .expect(1) + .create(); + let _ech2 = server + .mock("POST", api) + .match_body(Matcher::Regex("Email/changes".into())) + .with_body(json!({"methodResponses":[["Email/changes",{"accountId":"w","oldState":"se1","newState":"se2","hasMoreChanges":false,"created":[],"updated":["E1"],"destroyed":[]},"c"]]}).to_string()) + .expect(1) + .create(); + let bad_update = server + .mock("POST", api) + .match_body(Matcher::Regex("Email/get".into())) + .with_body( + json!({"methodResponses":[["Email/get",{"accountId":"w","state":"se2","list":[ + {"id":"E1","mailboxIds":{},"keywords":{"$seen":true,"$flagged":true}} + ],"notFound":[]},"g"]]}) + .to_string(), + ) + .expect(1) + .create(); + + let s2 = sync::import_jmap::run( + common(&archive), + import_cfg_objects(&base, vec![ObjectType::Mailbox, ObjectType::Email]), + ) + .expect("second import"); + bad_update.assert(); + let em2 = s2 + .per_type + .iter() + .find(|(t, _)| *t == "Email") + .map(|(_, c)| c.clone()) + .expect("email counts"); + assert_eq!( + em2.updated, 0, + "the failed update is not counted as applied" + ); + assert!(em2.failed >= 1, "the unresolvable update is counted failed"); + let conn = rusqlite::Connection::open(&archive).unwrap(); + let cursor: String = conn + .query_row( + "SELECT state FROM sync_state_jmap WHERE type_name='Email'", + [], + |r| r.get(0), + ) + .unwrap(); + assert_eq!( + cursor, "se1", + "cursor is held at the pre-change state so the failed update retries next run" + ); + let _ = std::fs::remove_file(&archive); +} + +#[test] +fn import_unknown_method_changes_falls_back_to_full_refresh() { + let mut server = mockito::Server::new(); + let base = server.url(); + let api = "/jmap/api"; + let archive = tmp(); + + let _root = server.mock("GET", "/").with_status(404).create(); + let _wk = server + .mock("GET", "/.well-known/jmap") + .with_body(session_body_full(&base)) + .expect_at_least(1) + .create(); + let _term = anchor_terminator(&mut server, api, "Mailbox"); + + let _q1 = server + .mock("POST", api) + .match_body(Matcher::Regex("Mailbox/query".into())) + .with_body( + json!({"methodResponses":[["Mailbox/query",{"accountId":"w","ids":["A"]},"q"]]}) + .to_string(), + ) + .expect(1) + .create(); + let _g1 = server + .mock("POST", api) + .match_body(Matcher::Regex("Mailbox/get".into())) + .with_body(json!({"methodResponses":[["Mailbox/get",{"accountId":"w","state":"s1","list":[ + {"id":"A","name":"OriginalName","parentId":null,"role":null,"sortOrder":0,"isSubscribed":true} + ],"notFound":[]},"g"]]}).to_string()) + .expect(1) + .create(); + + sync::import_jmap::run( + common(&archive), + import_cfg_objects(&base, vec![ObjectType::Mailbox]), + ) + .expect("first import"); + + let _q2 = server + .mock("POST", api) + .match_body(Matcher::Regex("Mailbox/query".into())) + .with_body( + json!({"methodResponses":[["Mailbox/query",{"accountId":"w","ids":["A"]},"q"]]}) + .to_string(), + ) + .expect(1) + .create(); + let unknown = server + .mock("POST", api) + .match_body(Matcher::Regex("Mailbox/changes".into())) + .with_body(json!({"methodResponses":[["error",{"type":"unknownMethod"},"c"]]}).to_string()) + .expect(1) + .create(); + let _capture_state = server + .mock("POST", api) + .match_body(Matcher::AllOf(vec![ + Matcher::Regex("Mailbox/get".into()), + Matcher::Regex("\"ids\":\\[\\]".into()), + ])) + .with_body(json!({"methodResponses":[["Mailbox/get",{"accountId":"w","state":"s9","list":[],"notFound":[]},"g"]]}).to_string()) + .expect(1) + .create(); + let refresh_get = server + .mock("POST", api) + .match_body(Matcher::AllOf(vec![ + Matcher::Regex("Mailbox/get".into()), + Matcher::Regex("\"A\"".into()), + ])) + .with_body(json!({"methodResponses":[["Mailbox/get",{"accountId":"w","state":"s9","list":[ + {"id":"A","name":"RefreshedName","parentId":null,"role":null,"sortOrder":0,"isSubscribed":true} + ],"notFound":[]},"g"]]}).to_string()) + .expect(1) + .create(); + + let s2 = sync::import_jmap::run( + common(&archive), + import_cfg_objects(&base, vec![ObjectType::Mailbox]), + ) + .expect("second import"); + unknown.assert(); + refresh_get.assert(); + let mb2 = s2 + .per_type + .iter() + .find(|(t, _)| *t == "Mailbox") + .map(|(_, c)| c.clone()) + .expect("mailbox counts"); + assert_eq!( + mb2.updated, 1, + "a server without Mailbox/changes degrades to a full refresh instead of aborting" + ); + let conn = rusqlite::Connection::open(&archive).unwrap(); + let name: String = conn + .query_row("SELECT name FROM mailboxes WHERE id=1", [], |r| r.get(0)) + .unwrap(); + assert_eq!(name, "RefreshedName"); + let _ = std::fs::remove_file(&archive); +} + +fn sieve_get_body(name: &str, blob: &str) -> String { + json!({"methodResponses":[["SieveScript/get",{"accountId":"w","state":"x","list":[ + {"id":"S1","name":name,"isActive":true,"blobId":blob} + ],"notFound":[]},"g"]]}) + .to_string() +} + +#[test] +fn import_sieve_script_reimport_unchanged_is_convergent() { + let mut server = mockito::Server::new(); + let base = server.url(); + let api = "/jmap/api"; + let archive = tmp(); + + let _root = server.mock("GET", "/").with_status(404).create(); + let _wk = server + .mock("GET", "/.well-known/jmap") + .with_body(session_body_full(&base)) + .expect_at_least(1) + .create(); + let _term = anchor_terminator(&mut server, api, "SieveScript"); + let _dl = server + .mock("GET", Matcher::Regex("/jmap/dl/w/B1/.*".into())) + .with_body("keep;\n") + .create(); + + let _q1 = server + .mock("POST", api) + .match_body(Matcher::Regex("SieveScript/query".into())) + .with_body( + json!({"methodResponses":[["SieveScript/query",{"accountId":"w","ids":["S1"]},"q"]]}) + .to_string(), + ) + .expect(1) + .create(); + let _g1 = server + .mock("POST", api) + .match_body(Matcher::Regex("SieveScript/get".into())) + .with_body(sieve_get_body("main", "B1")) + .expect(1) + .create(); + + sync::import_jmap::run( + common(&archive), + import_cfg_objects(&base, vec![ObjectType::SieveScript]), + ) + .expect("first import"); + + let _q2 = server + .mock("POST", api) + .match_body(Matcher::Regex("SieveScript/query".into())) + .with_body( + json!({"methodResponses":[["SieveScript/query",{"accountId":"w","ids":["S1"]},"q"]]}) + .to_string(), + ) + .expect(1) + .create(); + let _g2 = server + .mock("POST", api) + .match_body(Matcher::Regex("SieveScript/get".into())) + .with_body(sieve_get_body("main", "B1")) + .expect(1) + .create(); + + let s2 = sync::import_jmap::run( + common(&archive), + import_cfg_objects(&base, vec![ObjectType::SieveScript]), + ) + .expect("second import"); + let ss = s2 + .per_type + .iter() + .find(|(t, _)| *t == "SieveScript") + .map(|(_, c)| c.clone()) + .expect("sieve counts"); + assert_eq!(ss.created, 0, "no new scripts"); + assert_eq!( + ss.updated, 0, + "an unchanged SieveScript must not be counted as updated on re-import (convergent)" + ); + let _ = std::fs::remove_file(&archive); +} + +#[test] +fn import_sieve_script_content_change_is_propagated() { + let mut server = mockito::Server::new(); + let base = server.url(); + let api = "/jmap/api"; + let archive = tmp(); + + let _root = server.mock("GET", "/").with_status(404).create(); + let _wk = server + .mock("GET", "/.well-known/jmap") + .with_body(session_body_full(&base)) + .expect_at_least(1) + .create(); + let _term = anchor_terminator(&mut server, api, "SieveScript"); + let _dl1 = server + .mock("GET", Matcher::Regex("/jmap/dl/w/B1/.*".into())) + .with_body("keep;\n") + .create(); + let _dl2 = server + .mock("GET", Matcher::Regex("/jmap/dl/w/B2/.*".into())) + .with_body("discard;\n") + .create(); + + let _q1 = server + .mock("POST", api) + .match_body(Matcher::Regex("SieveScript/query".into())) + .with_body( + json!({"methodResponses":[["SieveScript/query",{"accountId":"w","ids":["S1"]},"q"]]}) + .to_string(), + ) + .expect(1) + .create(); + let _g1 = server + .mock("POST", api) + .match_body(Matcher::Regex("SieveScript/get".into())) + .with_body(sieve_get_body("main", "B1")) + .expect(1) + .create(); + + sync::import_jmap::run( + common(&archive), + import_cfg_objects(&base, vec![ObjectType::SieveScript]), + ) + .expect("first import"); + + let _q2 = server + .mock("POST", api) + .match_body(Matcher::Regex("SieveScript/query".into())) + .with_body( + json!({"methodResponses":[["SieveScript/query",{"accountId":"w","ids":["S1"]},"q"]]}) + .to_string(), + ) + .expect(1) + .create(); + let _g2 = server + .mock("POST", api) + .match_body(Matcher::Regex("SieveScript/get".into())) + .with_body(sieve_get_body("main", "B2")) + .expect(1) + .create(); + + let s2 = sync::import_jmap::run( + common(&archive), + import_cfg_objects(&base, vec![ObjectType::SieveScript]), + ) + .expect("second import"); + let ss = s2 + .per_type + .iter() + .find(|(t, _)| *t == "SieveScript") + .map(|(_, c)| c.clone()) + .expect("sieve counts"); + assert_eq!( + ss.updated, 1, + "the changed script content is re-fetched and updated" + ); + let conn = rusqlite::Connection::open(&archive).unwrap(); + let body: Vec = conn + .query_row( + "SELECT b.data FROM blobs b JOIN sieve_scripts s ON s.blob_id = b.id LIMIT 1", + [], + |r| r.get(0), + ) + .unwrap(); + assert_eq!( + String::from_utf8(body).unwrap(), + "discard;\n", + "new script content propagated into the archive blob" + ); + let _ = std::fs::remove_file(&archive); +} + +#[test] +fn first_run_cursor_is_captured_up_front_not_from_the_fetch() { + let mut server = mockito::Server::new(); + let base = server.url(); + let api = "/jmap/api"; + let archive = tmp(); + + let _root = server.mock("GET", "/").with_status(404).create(); + let _wk = server + .mock("GET", "/.well-known/jmap") + .with_body(session_body_full(&base)) + .expect_at_least(1) + .create(); + let _term = anchor_terminator(&mut server, api, "Mailbox"); + let _q = server + .mock("POST", api) + .match_body(Matcher::Regex("Mailbox/query".into())) + .with_body( + json!({"methodResponses":[["Mailbox/query",{"accountId":"w","ids":["A"]},"q"]]}) + .to_string(), + ) + .expect(1) + .create(); + // Up-front state snapshot (ids:[]) reports an EARLIER state than the new-fetch. + let _state = server + .mock("POST", api) + .match_body(Matcher::AllOf(vec![ + Matcher::Regex("Mailbox/get".into()), + Matcher::Regex("\"ids\":\\[\\]".into()), + ])) + .with_body(json!({"methodResponses":[["Mailbox/get",{"accountId":"w","state":"before","list":[],"notFound":[]},"g"]]}).to_string()) + .expect(1) + .create(); + // The new-fetch reports a LATER state; if we (incorrectly) captured from here, the cursor + // would be "after" and an object changed mid-run could be missed next run. + let _g = server + .mock("POST", api) + .match_body(Matcher::AllOf(vec![ + Matcher::Regex("Mailbox/get".into()), + Matcher::Regex("\"A\"".into()), + ])) + .with_body(json!({"methodResponses":[["Mailbox/get",{"accountId":"w","state":"after","list":[ + {"id":"A","name":"Personal","parentId":null,"role":null,"sortOrder":0,"isSubscribed":true} + ],"notFound":[]},"g"]]}).to_string()) + .expect(1) + .create(); + + sync::import_jmap::run( + common(&archive), + import_cfg_objects(&base, vec![ObjectType::Mailbox]), + ) + .expect("import"); + let cursor: String = rusqlite::Connection::open(&archive) + .unwrap() + .query_row( + "SELECT state FROM sync_state_jmap WHERE type_name='Mailbox'", + [], + |r| r.get(0), + ) + .unwrap(); + assert_eq!( + cursor, "before", + "cursor must be the pre-fetch snapshot (lower bound), not the post-fetch state" + ); + let _ = std::fs::remove_file(&archive); +} diff --git a/tests/sync_jmap.rs b/tests/sync_jmap.rs index b23c290..59c6dfa 100644 --- a/tests/sync_jmap.rs +++ b/tests/sync_jmap.rs @@ -11,8 +11,11 @@ use std::path::{Path, PathBuf}; use integration::stalwart::shared as shared_stalwart; use rusqlite::Connection; -use vandelay::jmap::account::AccountSelector; -use vandelay::jmap::http::Auth; +use serde_json::{Map, Value, json}; +use vandelay::jmap::account::{self, AccountSelector}; +use vandelay::jmap::http::{Auth, HttpClient, RetryPolicy}; +use vandelay::jmap::request::Request; +use vandelay::jmap::session::Session; use vandelay::logging::Logger; use vandelay::sync::{self, CommonConfig, ConnectConfig, ExportConfig, ImportConfig}; @@ -1148,3 +1151,90 @@ fn live_blob_quota_429_triggers_retry_after_then_succeeds() { drop(_ttl_guard); seeder::teardown(base_url()).expect("teardown"); } + +#[test] +#[ignore = "requires Docker"] +fn import_delta_propagates_email_keyword_change_via_changes() { + let fx = seeder::provision(base_url()).expect("provision"); + let acc = fx.account("test1").expect("test1"); + let archive = tmp_archive("delta-email"); + + sync::import_jmap::run( + common(&archive, false), + import_cfg(AccountSelector::Id(acc.account_id.clone())), + ) + .expect("first import"); + + let (jmap_id, local_id): (String, i64) = { + let conn = Connection::open(&archive).unwrap(); + conn.query_row( + "SELECT s.jmap_id, s.local_id FROM sync_id_jmap s + JOIN emails e ON e.id = s.local_id + WHERE s.type_name = 'Email' AND e.keywords NOT LIKE '%$flagged%' + LIMIT 1", + [], + |r| Ok((r.get(0)?, r.get(1)?)), + ) + .expect("an unflagged email exists in the archive") + }; + + let client = HttpClient::new(basic("test1"), RetryPolicy::new(5), true); + let session = Session::discover(&client, base_url()).expect("session discovered"); + let account = account::resolve( + &AccountSelector::Id(acc.account_id.clone()), + &session, + &client, + ) + .expect("account resolved"); + let mut update = Map::new(); + update.insert(jmap_id.clone(), json!({ "keywords/$flagged": true })); + let mut req = Request::new(); + req.call( + "Email/set", + json!({ "accountId": account, "update": Value::Object(update) }), + "s", + ); + let resp = req.send(&client, &session.api_url).expect("Email/set sent"); + let mr = resp.first().expect("a method response"); + assert!( + mr.args + .get("updated") + .and_then(|u| u.get(&jmap_id)) + .is_some(), + "server accepted the keyword update: {:?}", + mr.args + ); + + let s2 = sync::import_jmap::run( + common(&archive, false), + import_cfg(AccountSelector::Id(acc.account_id.clone())), + ) + .expect("second import"); + + let conn = Connection::open(&archive).unwrap(); + let kw: String = conn + .query_row( + "SELECT keywords FROM emails WHERE id = ?1", + [local_id], + |r| r.get(0), + ) + .unwrap(); + assert!( + kw.contains("$flagged"), + "delta re-import propagated the new flag into the archive: {kw}" + ); + let em = s2 + .per_type + .iter() + .find(|(t, _)| *t == "Email") + .map(|(_, c)| c.clone()) + .expect("email counts"); + assert!( + em.updated >= 1, + "the changed email was detected via Email/changes and refreshed (updated={})", + em.updated + ); + drop(conn); + seeder::teardown(base_url()).expect("teardown"); + let _ = std::fs::remove_file(&archive); +} diff --git a/tests/sync_maildir.rs b/tests/sync_maildir.rs index 82caa34..ad570e0 100644 --- a/tests/sync_maildir.rs +++ b/tests/sync_maildir.rs @@ -428,7 +428,6 @@ fn dry_run_then_real_run_produces_same_counts_for_new() { #[test] fn trashed_flag_added_between_runs_deletes_present_row() { - let td = tempfile::TempDir::new().unwrap(); ensure_maildir(td.path()); let path = write( @@ -459,7 +458,6 @@ fn trashed_flag_added_between_runs_deletes_present_row() { #[cfg(unix)] #[test] fn symlinked_subfolder_is_followed_and_appears_as_its_own_folder() { - let td = tempfile::TempDir::new().unwrap(); ensure_maildir(td.path()); let real = ensure_subfolder(td.path(), ".Real"); @@ -602,7 +600,6 @@ fn blob_hashes(conn: &Connection) -> std::collections::HashSet { #[test] #[ignore = "requires Docker"] fn maildir_message_count_matches_jmap_for_same_corpus() { - let fx = seeder::provision(base_url()).expect("provision"); let acc = fx.account("test1").expect("test1"); let corpus = seeder::data::load_mbox(30).expect("mbox corpus"); diff --git a/tests/sync_managesieve.rs b/tests/sync_managesieve.rs index 72585e3..4050740 100644 --- a/tests/sync_managesieve.rs +++ b/tests/sync_managesieve.rs @@ -209,7 +209,6 @@ fn managesieve_second_run_is_convergent() { #[test] #[ignore = "requires Docker"] fn managesieve_and_jmap_imports_share_blob_bytes() { - let fx = seeder::provision(base_url()).expect("provision"); let acc = fx.account("test1").expect("test1"); let msieve_archive = tmp_archive("parity-msieve"); @@ -310,7 +309,6 @@ fn managesieve_dry_run_reports_diff_without_writing() { #[test] #[ignore = "requires Docker"] fn managesieve_source_change_protection_refuses_second_account() { - let fx = seeder::provision(base_url()).expect("provision"); let acc1 = fx.account("test1").expect("test1"); let acc2 = fx.account("test2").expect("test2"); @@ -330,7 +328,6 @@ fn managesieve_source_change_protection_refuses_second_account() { #[test] #[ignore = "requires Docker"] fn managesieve_implicit_tls_path_succeeds_when_offered() { - let fx = seeder::provision(base_url()).expect("provision"); let acc = fx.account("test1").expect("test1"); let archive = tmp_archive("implicit_tls"); @@ -350,7 +347,6 @@ fn managesieve_implicit_tls_path_succeeds_when_offered() { #[test] #[ignore = "requires Docker"] fn managesieve_round_trip_via_jmap_export_converges() { - let fx = seeder::provision(base_url()).expect("provision"); let src = fx.account("test1").expect("test1"); let dst = fx.account("test4").expect("test4");