Self heal on blobNotFound errors when exporting data (closes #13)

This commit is contained in:
Maurus Decimus
2026-06-28 11:11:23 +02:00
parent 6ca561fff5
commit ceb8e1815c
9 changed files with 574 additions and 130 deletions
+41
View File
@@ -60,6 +60,7 @@ struct Uploader<'a> {
net: &'a Net,
conn: &'a Connection,
cache: HashMap<i64, JmapId>,
touched: Vec<i64>,
}
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<JmapId, JmapError> {
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<i64> {
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<F>(
net: &Net,
ty: ObjectType,
cid: &str,
uploader: &mut Uploader<'_>,
touched: Vec<i64>,
outcome: crate::jmap::request::SetOutcome,
mut rebuild: F,
) -> Result<crate::jmap::request::SetOutcome, Error>
where
F: FnMut(&mut Uploader<'_>) -> Result<Value, Error>,
{
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)],
+166 -112
View File
@@ -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<EmailKey> = 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<Map<String, Value>> {
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<String, Value> {
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<String, Value>,
kw: Map<String, Value>,
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<SingleImport, JmapError> {
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}"),
})
}
+19 -11
View File
@@ -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<Value, Error> {
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;
+26 -3
View File
@@ -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::<i64>().ok())
&& let Some(id) = jid(v)
+15 -2
View File
@@ -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::<i64>().ok())
&& let Some(id) = jid(v)