From ceb8e1815c776f7d3a65921a9a0617f9d0cdd052 Mon Sep 17 00:00:00 2001 From: Maurus Decimus <11444311+mdecimus@users.noreply.github.com> Date: Sun, 28 Jun 2026 11:11:23 +0200 Subject: [PATCH] Self heal on `blobNotFound` errors when exporting data (closes #13) --- CHANGELOG.md | 9 ++ Cargo.lock | 2 +- Cargo.toml | 2 +- src/sync/export.rs | 41 +++++ src/sync/export/email.rs | 278 ++++++++++++++++++++-------------- src/sync/export/sieve.rs | 30 ++-- src/sync/export/tree.rs | 29 +++- src/sync/export/uidtype.rs | 17 ++- tests/mock_sync.rs | 296 +++++++++++++++++++++++++++++++++++++ 9 files changed, 574 insertions(+), 130 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 5bd1c26..752d9a1 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -2,6 +2,15 @@ All notable changes to this project will be documented in this file. This project adheres to [Semantic Versioning](http://semver.org/). +## [1.0.6] - 2026-06-XX + +### Added + +### Changed + +### Fixed +- Self heal on `blobNotFound` errors when exporting data (#13). + ## [1.0.5] - 2026-06-27 ### Added diff --git a/Cargo.lock b/Cargo.lock index 196b56b..e74e8d5 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2892,7 +2892,7 @@ dependencies = [ [[package]] name = "vandelay" -version = "1.0.5" +version = "1.0.6" dependencies = [ "base64", "blake3", diff --git a/Cargo.toml b/Cargo.toml index 56e984f..f554878 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -1,7 +1,7 @@ [package] name = "vandelay" description = "JMAP account migration utility" -version = "1.0.5" +version = "1.0.6" authors = ["Stalwart Labs LLC "] license = "Apache-2.0 OR MIT" repository = "https://github.com/stalwartlabs/vandelay" diff --git a/src/sync/export.rs b/src/sync/export.rs index 51e3cc5..0efa30b 100644 --- a/src/sync/export.rs +++ b/src/sync/export.rs @@ -60,6 +60,7 @@ struct Uploader<'a> { net: &'a Net, conn: &'a Connection, cache: HashMap, + touched: Vec, } impl<'a> Uploader<'a> { @@ -68,10 +69,12 @@ impl<'a> Uploader<'a> { net, conn, cache: HashMap::new(), + touched: Vec::new(), } } fn upload_with(&mut self, local_id: i64, content_type: &str) -> Result { + self.touched.push(local_id); if let Some(id) = self.cache.get(&local_id) { return Ok(id.clone()); } @@ -93,6 +96,14 @@ impl<'a> Uploader<'a> { self.cache.insert(local_id, id.clone()); Ok(id) } + + fn invalidate(&mut self, local_id: i64) { + self.cache.remove(&local_id); + } + + fn take_touched(&mut self) -> Vec { + std::mem::take(&mut self.touched) + } } impl BlobUpload for Uploader<'_> { @@ -485,6 +496,36 @@ mod common { ) } + 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") + }) + } + + pub fn retry_if_blob_missing( + net: &Net, + ty: ObjectType, + cid: &str, + uploader: &mut Uploader<'_>, + touched: Vec, + outcome: crate::jmap::request::SetOutcome, + mut rebuild: F, + ) -> Result + where + F: FnMut(&mut Uploader<'_>) -> Result, + { + if !blob_not_found(&outcome, cid) { + return Ok(outcome); + } + for id in &touched { + uploader.invalidate(*id); + } + let _ = uploader.take_touched(); + let wire = rebuild(uploader)?; + let _ = uploader.take_touched(); + create_batch(net, ty, vec![(cid.to_owned(), wire)]).map_err(Error::from) + } + fn synthesize_dry_run_outcome( ty: ObjectType, creates: &[(String, Value)], diff --git a/src/sync/export/email.rs b/src/sync/export/email.rs index f1d75ec..b8a7135 100644 --- a/src/sync/export/email.rs +++ b/src/sync/export/email.rs @@ -15,7 +15,7 @@ use crate::jmap::error::JmapError; use crate::jmap::request::{Request, check_method_error, get_objects}; use crate::jmap::wire::JmapId; use crate::logging::Logger; -use crate::sync::import_jmap::mapping::{EMAIL_SELECT, TargetResolver, row_to_email}; +use crate::sync::import_jmap::mapping::{EMAIL_SELECT, EmailRow, TargetResolver, row_to_email}; use crate::sync::keys::{EmailIndex, EmailKey, email_index, email_keys, index_from_json}; use crate::sync::{Context, TypeCounts}; use crate::types::ObjectType; @@ -91,7 +91,7 @@ pub fn reconcile( } let target_keys: HashSet = email_keys(&indices).into_iter().collect(); - let local: Vec<(i64, crate::sync::import_jmap::mapping::EmailRow)> = { + let local: Vec<(i64, EmailRow)> = { let mut stmt = ctx .conn .prepare(EMAIL_SELECT) @@ -120,134 +120,188 @@ pub fn reconcile( continue; } let (local_id, row) = &local[i]; - let mut mids = Map::new(); - let mut all_resolved = true; - for ml in &row.mailbox_locals { - match maps.target(ObjectType::Mailbox, *ml) { - Some(t) => { - mids.insert(t.0, Value::Bool(true)); - } - None => { - all_resolved = false; - break; - } - } - } - if !all_resolved { - logger.warn(&format!( - "email local {local_id} skipped: mailbox not on target" - )); - counts.failed += 1; - continue; - } - let blob = match uploader.upload_with(row.blob_local_id, "message/rfc822") { - Ok(b) => b.0, - Err(e) => { - logger.warn(&format!("email blob upload failed: {e}")); - counts.failed += 1; - continue; - } - }; - let mut kw = Map::new(); - for k in &row.keywords { - kw.insert(k.clone(), Value::Bool(true)); - } - let item = json!({ - "blobId": blob, - "mailboxIds": Value::Object(mids), - "keywords": Value::Object(kw), - "receivedAt": row.received_at, - }); - send_import_chunk(net, &[(format!("e{local_id}"), item)], counts, logger); + export_one(net, &mut uploader, maps, *local_id, row, counts, logger); } Ok(Plan::default()) } -fn send_import_chunk( +fn build_mailbox_ids(row: &EmailRow, maps: &Maps) -> Option> { + let mut mids = Map::new(); + for ml in &row.mailbox_locals { + let t = maps.target(ObjectType::Mailbox, *ml)?; + mids.insert(t.0, Value::Bool(true)); + } + Some(mids) +} + +fn build_keywords(row: &EmailRow) -> Map { + let mut kw = Map::new(); + for k in &row.keywords { + kw.insert(k.clone(), Value::Bool(true)); + } + kw +} + +fn import_item( + blob: String, + mids: Map, + kw: Map, + received_at: &str, +) -> Value { + json!({ + "blobId": blob, + "mailboxIds": Value::Object(mids), + "keywords": Value::Object(kw), + "receivedAt": received_at, + }) +} + +fn export_one( net: &Net, - items: &[(String, Value)], + uploader: &mut Uploader, + maps: &Maps, + local_id: i64, + row: &EmailRow, counts: &mut TypeCounts, logger: &Logger, ) { - if items.is_empty() { - return; - } + let cid = format!("e{local_id}"); + let mids = match build_mailbox_ids(row, maps) { + Some(m) => m, + None => { + logger.warn(&format!( + "email local {local_id} skipped: mailbox not on target" + )); + counts.failed += 1; + return; + } + }; + let blob = match uploader.upload_with(row.blob_local_id, "message/rfc822") { + Ok(b) => b.0, + Err(e) => { + logger.warn(&format!("email blob upload failed: {e}")); + counts.failed += 1; + return; + } + }; if net.dry_run { - counts.created += items.len() as u64; + counts.created += 1; return; } - let mut map = Map::new(); - for (k, v) in items { - map.insert(k.clone(), v.clone()); + let item = import_item(blob, mids, build_keywords(row), &row.received_at); + match send_single_import(net, &cid, item) { + Ok(SingleImport::Created) => counts.created += 1, + Ok(SingleImport::Skipped) => counts.skipped += 1, + Ok(SingleImport::NotCreated { error_type, .. }) if error_type == "blobNotFound" => { + retry_after_reupload(net, uploader, maps, &cid, row, counts, logger); + } + Ok(SingleImport::NotCreated { detail, .. }) => { + logger.warn(&format!("Email/import {cid} failed: {detail}")); + counts.failed += 1; + } + Err(e) => { + logger.warn(&format!("Email/import {cid} send failed: {e}")); + counts.failed += 1; + } } +} + +fn retry_after_reupload( + net: &Net, + uploader: &mut Uploader, + maps: &Maps, + cid: &str, + row: &EmailRow, + counts: &mut TypeCounts, + logger: &Logger, +) { + uploader.invalidate(row.blob_local_id); + let blob = match uploader.upload_with(row.blob_local_id, "message/rfc822") { + Ok(b) => b.0, + Err(e) => { + logger.warn(&format!("Email/import {cid}: blob re-upload failed: {e}")); + counts.failed += 1; + return; + } + }; + let mids = match build_mailbox_ids(row, maps) { + Some(m) => m, + None => { + logger.warn(&format!( + "Email/import {cid} skipped: mailbox not on target" + )); + counts.failed += 1; + return; + } + }; + let item = import_item(blob, mids, build_keywords(row), &row.received_at); + match send_single_import(net, cid, item) { + Ok(SingleImport::Created) => counts.created += 1, + Ok(SingleImport::Skipped) => counts.skipped += 1, + Ok(SingleImport::NotCreated { detail, .. }) => { + logger.warn(&format!( + "Email/import {cid} failed after blob re-upload: {detail}" + )); + counts.failed += 1; + } + Err(e) => { + logger.warn(&format!( + "Email/import {cid} send failed after blob re-upload: {e}" + )); + counts.failed += 1; + } + } +} + +enum SingleImport { + Created, + Skipped, + NotCreated { error_type: String, detail: String }, +} + +fn send_single_import(net: &Net, cid: &str, item: Value) -> Result { + let mut emails = Map::new(); + emails.insert(cid.to_owned(), item); let mut req = Request::new(); req.call( "Email/import", - json!({ "accountId": net.account, "emails": Value::Object(map) }), + json!({ "accountId": net.account, "emails": Value::Object(emails) }), "i", ); - if req.fits(&net.limits).is_err() { - resplit_or_fail(net, items, counts, logger); - return; - } - match req.send(&net.client, &net.api) { - Ok(resp) => match resp.first().and_then(|mr| { - check_method_error(mr)?; - Ok(mr) - }) { - Ok(mr) => absorb_import(mr, counts, logger), - Err(JmapError::RequestTooLarge) => resplit_or_fail(net, items, counts, logger), - Err(JmapError::Method { error_type, .. }) if error_type == "requestTooLarge" => { - resplit_or_fail(net, items, counts, logger); - } - Err(e) => { - logger.warn(&format!( - "Email/import method error ({} items): {e}", - items.len() - )); - counts.failed += items.len() as u64; - } - }, - Err(JmapError::RequestTooLarge) => resplit_or_fail(net, items, counts, logger), - Err(e) => { - logger.warn(&format!( - "Email/import send failed ({} items): {e}", - items.len() - )); - counts.failed += items.len() as u64; - } - } -} - -fn resplit_or_fail(net: &Net, items: &[(String, Value)], counts: &mut TypeCounts, logger: &Logger) { - if items.len() <= 1 { - for (cid, _) in items { - logger.warn(&format!( - "Email/import {cid} exceeds maxSizeRequest alone; skipped" - )); - } - counts.failed += items.len() as u64; - return; - } - let mid = items.len() / 2; - send_import_chunk(net, &items[..mid], counts, logger); - send_import_chunk(net, &items[mid..], counts, logger); -} - -fn absorb_import(mr: &crate::jmap::request::MethodCall, counts: &mut TypeCounts, logger: &Logger) { - if let Some(created) = mr.args.get("created").and_then(Value::as_object) { - counts.created += created.len() as u64; - } - if let Some(nc) = mr.args.get("notCreated").and_then(Value::as_object) { - for (cid, err) in nc { - let etype = err.get("type").and_then(Value::as_str).unwrap_or(""); - if etype == "alreadyExists" { - counts.skipped += 1; - } else { - logger.warn(&format!("Email/import {cid} failed: {err}")); - counts.failed += 1; - } + req.fits(&net.limits)?; + let resp = req.send(&net.client, &net.api)?; + let mr = resp.first()?; + check_method_error(mr)?; + if let Some(err) = mr + .args + .get("notCreated") + .and_then(Value::as_object) + .and_then(|nc| nc.get(cid)) + { + let error_type = err + .get("type") + .and_then(Value::as_str) + .unwrap_or("") + .to_owned(); + if error_type == "alreadyExists" { + return Ok(SingleImport::Skipped); } + return Ok(SingleImport::NotCreated { + error_type, + detail: err.to_string(), + }); } + if mr + .args + .get("created") + .and_then(Value::as_object) + .is_some_and(|c| !c.is_empty()) + { + return Ok(SingleImport::Created); + } + Ok(SingleImport::NotCreated { + error_type: String::new(), + detail: format!("Email/import returned neither created nor notCreated for {cid}"), + }) } diff --git a/src/sync/export/sieve.rs b/src/sync/export/sieve.rs index cbb88dc..8928202 100644 --- a/src/sync/export/sieve.rs +++ b/src/sync/export/sieve.rs @@ -8,7 +8,7 @@ use std::collections::{HashMap, HashSet}; use serde_json::{Value, json}; -use super::common::{create_batch, jid, target_get_all}; +use super::common::{create_batch, jid, retry_if_blob_missing, target_get_all}; use super::{Maps, Net, Plan, Uploader}; use crate::error::Error; use crate::jmap::request::Request; @@ -64,16 +64,24 @@ pub fn reconcile( counts.skipped += 1; id } else { - let blob_id = uploader - .upload_with(*blob_local, "application/sieve") - .map_err(Error::from)?; - let mut obj = serde_json::Map::new(); - if let Some(n) = name { - obj.insert("name".to_owned(), Value::String(n.clone())); - } - obj.insert("blobId".to_owned(), Value::String(blob_id.0)); - let outcome = create_batch(net, ty, vec![(format!("c{local}"), Value::Object(obj))]) - .map_err(Error::from)?; + let cid = format!("c{local}"); + let build = |up: &mut Uploader<'_>| -> Result { + let blob_id = up + .upload_with(*blob_local, "application/sieve") + .map_err(Error::from)?; + let mut obj = serde_json::Map::new(); + if let Some(n) = name { + obj.insert("name".to_owned(), Value::String(n.clone())); + } + obj.insert("blobId".to_owned(), Value::String(blob_id.0)); + Ok(Value::Object(obj)) + }; + let _ = uploader.take_touched(); + let wire = build(&mut uploader)?; + let touched = uploader.take_touched(); + let outcome = create_batch(net, ty, vec![(cid.clone(), wire)]).map_err(Error::from)?; + let outcome = + retry_if_blob_missing(net, ty, &cid, &mut uploader, touched, outcome, build)?; match outcome.created.first().and_then(|(_, v)| jid(v)) { Some(id) => { counts.created += 1; diff --git a/src/sync/export/tree.rs b/src/sync/export/tree.rs index 96ea498..ffac84d 100644 --- a/src/sync/export/tree.rs +++ b/src/sync/export/tree.rs @@ -9,7 +9,7 @@ use std::collections::{HashMap, HashSet}; use rusqlite::params; use serde_json::Value; -use super::common::{create_batch, jid, target_query_get}; +use super::common::{create_batch, jid, retry_if_blob_missing, target_query_get}; use super::{Net, Plan, Uploader}; use crate::error::Error; use crate::jmap::wire::JmapId; @@ -182,6 +182,8 @@ pub fn reconcile( counts.failed += 1; continue; } + let cid = format!("c{}", n.local); + let _ = uploader.take_touched(); let obj = match build_create(ctx, ty, n.local, maps, &mut uploader) { Ok(o) => o, Err(e) => { @@ -194,8 +196,29 @@ pub fn reconcile( continue; } }; - let outcome = create_batch(net, ty, vec![(format!("c{}", n.local), obj)]) - .map_err(Error::from)?; + let touched = uploader.take_touched(); + let outcome = + create_batch(net, ty, vec![(cid.clone(), obj)]).map_err(Error::from)?; + let outcome = match retry_if_blob_missing( + net, + ty, + &cid, + &mut uploader, + touched, + outcome, + |up| build_create(ctx, ty, n.local, maps, up), + ) { + Ok(o) => o, + Err(e) => { + logger.warn(&format!( + "{} local {} skipped: {e}", + ty.jmap_name(), + n.local + )); + counts.failed += 1; + continue; + } + }; for (cid, v) in &outcome.created { if let Some(local) = cid.strip_prefix('c').and_then(|s| s.parse::().ok()) && let Some(id) = jid(v) diff --git a/src/sync/export/uidtype.rs b/src/sync/export/uidtype.rs index 61734c2..90da298 100644 --- a/src/sync/export/uidtype.rs +++ b/src/sync/export/uidtype.rs @@ -8,7 +8,7 @@ use std::collections::HashSet; use serde_json::Value; -use super::common::{create_batch, jid, target_query_get}; +use super::common::{create_batch, jid, retry_if_blob_missing, target_query_get}; use super::{Maps, Net, Plan, Uploader}; use crate::error::Error; use crate::logging::Logger; @@ -74,6 +74,8 @@ pub fn reconcile( counts.skipped += 1; continue; } + let cid = format!("c{local}"); + let _ = uploader.take_touched(); let wire = match build_wire(ctx, ty, *local, maps, &mut uploader) { Ok(w) => w, Err(e) => { @@ -82,8 +84,19 @@ pub fn reconcile( continue; } }; + let touched = uploader.take_touched(); + let outcome = create_batch(net, ty, vec![(cid.clone(), wire)]).map_err(Error::from)?; let outcome = - create_batch(net, ty, vec![(format!("c{local}"), wire)]).map_err(Error::from)?; + match retry_if_blob_missing(net, ty, &cid, &mut uploader, touched, outcome, |up| { + build_wire(ctx, ty, *local, maps, up) + }) { + Ok(o) => o, + Err(e) => { + logger.warn(&format!("{} local {local} skipped: {e}", ty.jmap_name())); + counts.failed += 1; + continue; + } + }; for (cid, v) in &outcome.created { if let Some(parsed) = cid.strip_prefix('c').and_then(|s| s.parse::().ok()) && let Some(id) = jid(v) diff --git a/tests/mock_sync.rs b/tests/mock_sync.rs index bb3d876..58ce1b0 100644 --- a/tests/mock_sync.rs +++ b/tests/mock_sync.rs @@ -317,6 +317,177 @@ fn email_export_sends_one_email_per_import_call() { let _ = std::fs::remove_file(&archive); } +#[test] +fn export_email_blob_not_found_reuploads_and_retries() { + let mut server = mockito::Server::new(); + let base = server.url(); + let api = "/jmap/api"; + + let archive = tmp(); + { + let conn = db::init::open(&archive).unwrap(); + conn.execute( + "INSERT INTO mailboxes (id,name,parent_id,role) VALUES (1,'Inbox',NULL,'inbox')", + [], + ) + .unwrap(); + let raw = b"From: a@x\r\nSubject: dup\r\nMessage-ID: \r\n\r\nbody"; + let blob = db::blobs::intern_blob(&conn, raw).unwrap(); + for _ in 0..2 { + conn.execute( + "INSERT INTO emails (blob_id,received_at,mailbox_ids,keywords) + VALUES (?1,'2020-01-01T00:00:00Z','[1]','[]')", + rusqlite::params![blob], + ) + .unwrap(); + } + } + + let _root = server.mock("GET", "/").with_status(404).create(); + let _wk = server + .mock("GET", "/.well-known/jmap") + .with_body(session_body(&base)) + .create(); + + let _mq1 = server + .mock("POST", api) + .match_body(Matcher::Regex("Mailbox/query".into())) + .with_body( + json!({"methodResponses":[["Mailbox/query", + {"accountId":"w","ids":["t1"]},"q"]]}) + .to_string(), + ) + .expect(1) + .create(); + let _mq2 = server + .mock("POST", api) + .match_body(Matcher::Regex("Mailbox/query".into())) + .with_body( + json!({"methodResponses":[["Mailbox/query", + {"accountId":"w","ids":[]},"q"]]}) + .to_string(), + ) + .expect(1) + .create(); + let _mg = server + .mock("POST", api) + .match_body(Matcher::Regex("Mailbox/get".into())) + .with_body( + json!({"methodResponses":[["Mailbox/get",{"accountId":"w","list":[ + {"id":"t1","name":"Inbox","role":"inbox","parentId":null, + "myRights":{"mayDelete":true}}],"notFound":[]},"g"]]}) + .to_string(), + ) + .expect(1) + .create(); + let _eq = server + .mock("POST", api) + .match_body(Matcher::Regex("Email/query".into())) + .with_body( + json!({"methodResponses":[["Email/query", + {"accountId":"w","ids":[]},"q"]]}) + .to_string(), + ) + .expect(1) + .create(); + + let up1 = server + .mock("POST", Matcher::Regex("/jmap/upload/".into())) + .with_body(json!({"blobId":"UP1"}).to_string()) + .expect(1) + .create(); + let up2 = server + .mock("POST", Matcher::Regex("/jmap/upload/".into())) + .with_body(json!({"blobId":"UP2"}).to_string()) + .expect(1) + .create(); + + let imp_e1 = server + .mock("POST", api) + .match_body(Matcher::AllOf(vec![ + Matcher::Regex("Email/import".into()), + Matcher::Regex("e1".into()), + ])) + .with_body( + json!({"methodResponses":[["Email/import",{"accountId":"w", + "created":{"e1":{"id":"x1","blobId":"UP1","threadId":"t","size":10}}},"i"]]}) + .to_string(), + ) + .expect(1) + .create(); + let imp_e2_stale = server + .mock("POST", api) + .match_body(Matcher::AllOf(vec![ + Matcher::Regex("Email/import".into()), + Matcher::Regex("e2".into()), + Matcher::Regex("UP1".into()), + ])) + .with_body( + json!({"methodResponses":[["Email/import",{"accountId":"w", + "notCreated":{"e2":{"type":"blobNotFound"}}},"i"]]}) + .to_string(), + ) + .expect(1) + .create(); + let imp_e2_fresh = server + .mock("POST", api) + .match_body(Matcher::AllOf(vec![ + Matcher::Regex("Email/import".into()), + Matcher::Regex("e2".into()), + Matcher::Regex("UP2".into()), + ])) + .with_body( + json!({"methodResponses":[["Email/import",{"accountId":"w", + "created":{"e2":{"id":"x2","blobId":"UP2","threadId":"t","size":10}}},"i"]]}) + .to_string(), + ) + .expect(1) + .create(); + + let summary = sync::export::run( + CommonConfig { + archive: archive.clone(), + threads: 1, + dry_run: false, + max_retries: 1, + allow_invalid_certs: false, + logger: Logger::from_flags(true, 0), + }, + ExportConfig { + connect: ConnectConfig { + url: base.clone(), + auth: Auth::Basic { + user: "u".into(), + password: "p".into(), + }, + account: AccountSelector::Id("w".into()), + }, + objects: None, + prune: false, + yes: true, + }, + ) + .expect("export run"); + + let email = summary + .per_type + .iter() + .find(|(t, _)| *t == "Email") + .map(|(_, c)| c.clone()) + .expect("email counts"); + assert_eq!(email.created, 2, "both emails end up created"); + assert_eq!(email.failed, 0, "blobNotFound self-heals, not a failure"); + assert_eq!(email.skipped, 0); + assert!(!summary.any_failed()); + + up1.assert(); + up2.assert(); + imp_e1.assert(); + imp_e2_stale.assert(); + imp_e2_fresh.assert(); + let _ = std::fs::remove_file(&archive); +} + fn session_body_full(base: &str) -> String { json!({ "apiUrl": format!("{base}/jmap/api"), @@ -1485,6 +1656,131 @@ fn export_sieve_scripts_identical_content_different_names_both_created() { let _ = std::fs::remove_file(&archive); } +#[test] +fn export_sieve_blob_not_found_on_dedup_reuse_reuploads_and_retries() { + let mut server = mockito::Server::new(); + let base = server.url(); + let api = "/jmap/api"; + let archive = tmp(); + let shared = b"require [\"fileinto\"];\nfileinto \"Archive\";\n"; + { + let conn = db::init::open(&archive).unwrap(); + let blob = db::blobs::intern_blob(&conn, shared).unwrap(); + conn.execute( + "INSERT INTO sieve_scripts (id,name,is_active,blob_id) VALUES (1,'dup-A',0,?1)", + rusqlite::params![blob], + ) + .unwrap(); + conn.execute( + "INSERT INTO sieve_scripts (id,name,is_active,blob_id) VALUES (2,'dup-B',0,?1)", + rusqlite::params![blob], + ) + .unwrap(); + } + + 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 _g = server + .mock("POST", api) + .match_body(Matcher::Regex("SieveScript/get".into())) + .with_body( + json!({"methodResponses":[["SieveScript/get", + {"accountId":"w","list":[],"notFound":[]},"g"]]}) + .to_string(), + ) + .expect(1) + .create(); + let up1 = server + .mock("POST", Matcher::Regex("/jmap/upload/".into())) + .with_body(json!({"blobId":"UP1"}).to_string()) + .expect(1) + .create(); + let up2 = server + .mock("POST", Matcher::Regex("/jmap/upload/".into())) + .with_body(json!({"blobId":"UP2"}).to_string()) + .expect(1) + .create(); + let create_a = server + .mock("POST", api) + .match_body(Matcher::AllOf(vec![ + Matcher::Regex("SieveScript/set".into()), + Matcher::Regex("dup-A".into()), + ])) + .with_body( + json!({"methodResponses":[["SieveScript/set",{"accountId":"w", + "created":{"c1":{"id":"S1"}}},"s"]]}) + .to_string(), + ) + .expect(1) + .create(); + let create_b_stale = server + .mock("POST", api) + .match_body(Matcher::AllOf(vec![ + Matcher::Regex("SieveScript/set".into()), + Matcher::Regex("dup-B".into()), + Matcher::Regex("UP1".into()), + ])) + .with_body( + json!({"methodResponses":[["SieveScript/set",{"accountId":"w", + "notCreated":{"c2":{"type":"blobNotFound"}}},"s"]]}) + .to_string(), + ) + .expect(1) + .create(); + let create_b_fresh = server + .mock("POST", api) + .match_body(Matcher::AllOf(vec![ + Matcher::Regex("SieveScript/set".into()), + Matcher::Regex("dup-B".into()), + Matcher::Regex("UP2".into()), + ])) + .with_body( + json!({"methodResponses":[["SieveScript/set",{"accountId":"w", + "created":{"c2":{"id":"S2"}}},"s"]]}) + .to_string(), + ) + .expect(1) + .create(); + let _deactivate = server + .mock("POST", api) + .match_body(Matcher::Regex("onSuccessDeactivateScript".into())) + .with_body( + json!({"methodResponses":[["SieveScript/set",{"accountId":"w"},"a"]]}).to_string(), + ) + .expect(1) + .create(); + + let summary = sync::export::run( + common(&archive), + export_cfg_objects(&base, vec![ObjectType::SieveScript]), + ) + .expect("export"); + up1.assert(); + up2.assert(); + create_a.assert(); + create_b_stale.assert(); + create_b_fresh.assert(); + let counts = summary + .per_type + .iter() + .find(|(t, _)| *t == "SieveScript") + .map(|(_, c)| c.clone()) + .expect("sieve counts"); + assert_eq!( + counts.created, 2, + "the dedup-reused blobId that came back blobNotFound self-heals via re-upload" + ); + assert_eq!( + counts.failed, 0, + "blobNotFound on a reused blob is not a failure" + ); + let _ = std::fs::remove_file(&archive); +} + #[test] fn import_deeply_nested_mailbox_tree_orders_correctly() { let mut server = mockito::Server::new();