+6
-4
@@ -20,7 +20,7 @@ use crate::jmap::request::{Request, SetRequest, get_all, get_objects, query_all_
|
||||
use crate::jmap::session::{Limits, Session};
|
||||
use crate::jmap::wire::JmapId;
|
||||
use crate::logging::{LEVEL_DEFAULT, Logger};
|
||||
use crate::sync::import_jmap::mapping::{BlobUpload, TargetResolver};
|
||||
use crate::sync::import_jmap::mapping::{BlobBytes, TargetResolver};
|
||||
use crate::sync::{CommonConfig, Context, ExportConfig, Summary, TypeCounts};
|
||||
use crate::types::ObjectType;
|
||||
|
||||
@@ -110,9 +110,10 @@ impl<'a> Uploader<'a> {
|
||||
}
|
||||
}
|
||||
|
||||
impl BlobUpload for Uploader<'_> {
|
||||
fn upload(&mut self, local_id: i64) -> Result<JmapId, JmapError> {
|
||||
self.upload_with(local_id, "application/octet-stream")
|
||||
impl BlobBytes for Uploader<'_> {
|
||||
fn bytes(&self, local_id: i64) -> Result<Vec<u8>, JmapError> {
|
||||
db::blobs::blob_bytes(self.conn, local_id)?
|
||||
.ok_or_else(|| JmapError::malformed(format!("blob local id {local_id} missing")))
|
||||
}
|
||||
}
|
||||
|
||||
@@ -172,6 +173,7 @@ pub fn run(common: CommonConfig, config: ExportConfig) -> Result<Summary, Error>
|
||||
);
|
||||
let plan = match res {
|
||||
Ok(p) => p,
|
||||
Err(e) if e.aborts_run() => return Err(e),
|
||||
Err(e) => {
|
||||
logger.warn(&format!("type {} aborted: {e}", ty.jmap_name()));
|
||||
counts.failed += 1;
|
||||
|
||||
@@ -12,7 +12,10 @@ use super::common::{jid, target_query_get};
|
||||
use super::{Maps, Net, Plan, Uploader};
|
||||
use crate::error::Error;
|
||||
use crate::jmap::error::JmapError;
|
||||
use crate::jmap::request::{Request, check_method_error, get_objects};
|
||||
use crate::jmap::request::{
|
||||
MethodCall, Request, check_method_error, get_objects, retry_method_call,
|
||||
};
|
||||
use crate::jmap::retry::MethodCallKind;
|
||||
use crate::jmap::wire::JmapId;
|
||||
use crate::logging::Logger;
|
||||
use crate::sync::import_jmap::mapping::{EMAIL_SELECT, EmailRow, TargetResolver, row_to_email};
|
||||
@@ -219,7 +222,7 @@ fn export_one(
|
||||
return;
|
||||
}
|
||||
let item = import_item(blob, mids, build_keywords(row), &row.received_at);
|
||||
match send_single_import(net, &cid, item) {
|
||||
match send_single_import(net, &cid, item, logger) {
|
||||
Ok(SingleImport::Created) => counts.created += 1,
|
||||
Ok(SingleImport::Skipped) => counts.skipped += 1,
|
||||
Ok(SingleImport::NotCreated { error_type, .. }) if error_type == "blobNotFound" => {
|
||||
@@ -277,7 +280,7 @@ fn retry_after_reupload(
|
||||
}
|
||||
};
|
||||
let item = import_item(blob, mids, build_keywords(row), &row.received_at);
|
||||
match send_single_import(net, cid, item) {
|
||||
match send_single_import(net, cid, item, logger) {
|
||||
Ok(SingleImport::Created) => counts.created += 1,
|
||||
Ok(SingleImport::Skipped) => counts.skipped += 1,
|
||||
Ok(SingleImport::NotCreated { detail, .. }) => {
|
||||
@@ -304,7 +307,12 @@ enum SingleImport {
|
||||
NotCreated { error_type: String, detail: String },
|
||||
}
|
||||
|
||||
fn send_single_import(net: &Net, cid: &str, item: Value) -> Result<SingleImport, JmapError> {
|
||||
fn send_single_import(
|
||||
net: &Net,
|
||||
cid: &str,
|
||||
item: Value,
|
||||
logger: &Logger,
|
||||
) -> Result<SingleImport, JmapError> {
|
||||
let mut emails = Map::new();
|
||||
emails.insert(cid.to_owned(), item);
|
||||
let mut req = Request::new();
|
||||
@@ -314,8 +322,15 @@ fn send_single_import(net: &Net, cid: &str, item: Value) -> Result<SingleImport,
|
||||
"i",
|
||||
);
|
||||
req.fits(&net.limits)?;
|
||||
let resp = req.send(&net.client, &net.api)?;
|
||||
let mr = resp.first()?;
|
||||
retry_method_call(
|
||||
&net.client,
|
||||
MethodCallKind::SingleObjectWrite,
|
||||
logger,
|
||||
|| interpret_import(req.send(&net.client, &net.api)?.first()?, cid),
|
||||
)
|
||||
}
|
||||
|
||||
fn interpret_import(mr: &MethodCall, cid: &str) -> Result<SingleImport, JmapError> {
|
||||
check_method_error(mr)?;
|
||||
if let Some(err) = mr
|
||||
.args
|
||||
|
||||
+32
-20
@@ -5,10 +5,11 @@
|
||||
*/
|
||||
|
||||
use std::collections::HashSet;
|
||||
use std::fmt::Write as _;
|
||||
|
||||
use serde_json::Value;
|
||||
|
||||
use super::common::{create_batch, jid, retry_if_blob_missing, target_query_get};
|
||||
use super::common::{create_batch, jid, target_query_get};
|
||||
use super::{Maps, Net, Plan, Uploader};
|
||||
use crate::error::Error;
|
||||
use crate::logging::Logger;
|
||||
@@ -23,6 +24,15 @@ fn target_uid(v: &Value) -> Option<String> {
|
||||
v.get("uid").and_then(Value::as_str).map(str::to_owned)
|
||||
}
|
||||
|
||||
fn describe(ty: ObjectType, local: i64, uid: &str) -> String {
|
||||
let mut out = String::new();
|
||||
let _ = write!(out, "{} local {local}", ty.jmap_name());
|
||||
if !uid.is_empty() {
|
||||
let _ = write!(out, " (uid {uid})");
|
||||
}
|
||||
out
|
||||
}
|
||||
|
||||
pub fn reconcile(
|
||||
ctx: &Context,
|
||||
net: &Net,
|
||||
@@ -66,7 +76,7 @@ pub fn reconcile(
|
||||
};
|
||||
|
||||
let mut matched_uids: HashSet<String> = HashSet::new();
|
||||
let mut uploader = Uploader::new(net, &ctx.conn);
|
||||
let blobs = Uploader::new(net, &ctx.conn);
|
||||
for (local, uid) in &rows {
|
||||
if let Some(tid) = by_uid.get(uid) {
|
||||
maps.insert(ty, *local, crate::jmap::wire::JmapId(tid.clone()));
|
||||
@@ -75,28 +85,30 @@ pub fn reconcile(
|
||||
continue;
|
||||
}
|
||||
let cid = format!("c{local}");
|
||||
let _ = uploader.take_touched();
|
||||
let wire = match build_wire(ctx, ty, *local, maps, &mut uploader) {
|
||||
let wire = match build_wire(ctx, ty, *local, maps, &blobs) {
|
||||
Ok(w) => w,
|
||||
Err(e) if e.aborts_run() => return Err(e),
|
||||
Err(e) => {
|
||||
logger.warn(&format!("{} local {local} skipped: {e}", ty.jmap_name()));
|
||||
logger.warn(&format!("{} skipped: {e}", describe(ty, *local, uid)));
|
||||
counts.failed += 1;
|
||||
continue;
|
||||
}
|
||||
};
|
||||
let touched = uploader.take_touched();
|
||||
let outcome = create_batch(net, ty, vec![(cid.clone(), wire)]).map_err(Error::from)?;
|
||||
let outcome =
|
||||
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;
|
||||
let outcome = match create_batch(net, ty, vec![(cid.clone(), wire)]) {
|
||||
Ok(o) => o,
|
||||
Err(e) => {
|
||||
let mapped = Error::from(e);
|
||||
if mapped.aborts_run() {
|
||||
return Err(mapped);
|
||||
}
|
||||
};
|
||||
logger.warn(&format!(
|
||||
"{} not created: {mapped}",
|
||||
describe(ty, *local, uid)
|
||||
));
|
||||
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)
|
||||
@@ -136,7 +148,7 @@ fn build_wire(
|
||||
ty: ObjectType,
|
||||
local: i64,
|
||||
maps: &Maps,
|
||||
up: &mut Uploader<'_>,
|
||||
blobs: &Uploader<'_>,
|
||||
) -> Result<Value, Error> {
|
||||
if ty == ObjectType::ContactCard {
|
||||
let (uid, abids, data): (String, String, String) = ctx
|
||||
@@ -147,7 +159,7 @@ fn build_wire(
|
||||
|r| Ok((r.get(1)?, r.get(2)?, r.get(3)?)),
|
||||
)
|
||||
.map_err(|e| Error::Partial(e.to_string()))?;
|
||||
contact_card_to_wire(&uid, &abids, &data, maps, up).map_err(Error::from)
|
||||
contact_card_to_wire(&uid, &abids, &data, maps, blobs).map_err(Error::from)
|
||||
} else {
|
||||
let (cal, dr, ud, data): (String, i64, i64, String) = ctx
|
||||
.conn
|
||||
@@ -157,6 +169,6 @@ fn build_wire(
|
||||
|r| Ok((r.get(1)?, r.get(2)?, r.get(3)?, r.get(4)?)),
|
||||
)
|
||||
.map_err(|e| Error::Partial(e.to_string()))?;
|
||||
calendar_event_to_wire(&cal, dr != 0, ud != 0, &data, maps, up).map_err(Error::from)
|
||||
calendar_event_to_wire(&cal, dr != 0, ud != 0, &data, maps, blobs).map_err(Error::from)
|
||||
}
|
||||
}
|
||||
|
||||
+45
-1
@@ -10,4 +10,48 @@ pub mod coordinator;
|
||||
pub mod items;
|
||||
pub mod tree;
|
||||
|
||||
pub use coordinator::{DavAuth, DavImportConfig, DavKindArg, run};
|
||||
pub use coordinator::{DavAuth, DavImportConfig, DavKindArg, run, run_reporting};
|
||||
|
||||
use crate::error::Error;
|
||||
use crate::jmap::error::JmapError;
|
||||
|
||||
pub(crate) fn per_collection_failure(err: JmapError) -> Error {
|
||||
match err {
|
||||
JmapError::Sqlite(e) => Error::Db(crate::db::init::OpenError::Sqlite(e)),
|
||||
scoped => Error::Partial(scoped.to_string()),
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
#[test]
|
||||
fn a_forbidden_collection_is_a_per_unit_failure() {
|
||||
for e in [
|
||||
JmapError::Auth("server returned 403: forbidden".to_owned()),
|
||||
JmapError::HttpStatus {
|
||||
status: 405,
|
||||
body: "method not allowed".to_owned(),
|
||||
},
|
||||
JmapError::RetriesExhausted("PROPFIND kept returning 503".to_owned()),
|
||||
JmapError::Transport("io: connection reset".to_owned()),
|
||||
JmapError::Malformed("multistatus parse".to_owned()),
|
||||
] {
|
||||
let mapped = per_collection_failure(e);
|
||||
assert!(
|
||||
!mapped.aborts_run(),
|
||||
"{mapped} is scoped to one collection and must not abort the run"
|
||||
);
|
||||
assert_eq!(mapped.exit_code(), 5);
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn an_archive_failure_still_aborts() {
|
||||
let mapped =
|
||||
per_collection_failure(JmapError::Sqlite(rusqlite::Error::QueryReturnedNoRows));
|
||||
assert!(mapped.aborts_run());
|
||||
assert_eq!(mapped.exit_code(), 7);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -13,7 +13,7 @@ use crate::db::sources::SourceKey;
|
||||
use crate::error::Error;
|
||||
use crate::jmap::http::{Auth, RetryPolicy};
|
||||
use crate::logging::{LEVEL_DEFAULT, LEVEL_PROGRESS, Logger};
|
||||
use crate::sync::{CommonConfig, Summary, TypeCounts};
|
||||
use crate::sync::{CommonConfig, RunOutcome, Summary, TypeCounts};
|
||||
|
||||
use super::collections;
|
||||
use super::items;
|
||||
@@ -83,6 +83,20 @@ pub struct DavImportConfig {
|
||||
}
|
||||
|
||||
pub fn run(common: CommonConfig, config: DavImportConfig) -> Result<Summary, Error> {
|
||||
run_reporting(common, config).into_result()
|
||||
}
|
||||
|
||||
pub fn run_reporting(common: CommonConfig, config: DavImportConfig) -> RunOutcome {
|
||||
let mut summary = Summary::default();
|
||||
let error = run_into(common, config, &mut summary).err();
|
||||
RunOutcome { summary, error }
|
||||
}
|
||||
|
||||
fn run_into(
|
||||
common: CommonConfig,
|
||||
config: DavImportConfig,
|
||||
summary: &mut Summary,
|
||||
) -> Result<(), Error> {
|
||||
let logger = common.logger;
|
||||
enforce_tls_policy(&config.url, config.allow_cleartext)?;
|
||||
|
||||
@@ -138,38 +152,42 @@ pub fn run(common: CommonConfig, config: DavImportConfig) -> Result<Summary, Err
|
||||
}
|
||||
|
||||
if common.dry_run {
|
||||
let mut summary = run_dry_diff(&conn, &client, &discovery, &config, logger)?;
|
||||
let result = run_dry_diff(&conn, &client, &discovery, &config, logger, summary);
|
||||
summary.retries_observed = client.retries_observed();
|
||||
summary.retry_after_sleeps = client.retry_after_sleeps();
|
||||
return Ok(summary);
|
||||
return result;
|
||||
}
|
||||
|
||||
let username = config.auth.username();
|
||||
let source_id = db::sources::upsert_source(&conn, &key, Some(&account_id), &username)
|
||||
.map_err(|e| Error::Partial(e.to_string()))?;
|
||||
|
||||
let summary = match config.kind {
|
||||
let phase = match config.kind {
|
||||
DavKindArg::Caldav => {
|
||||
run_caldav(&mut conn, &client, source_id, &discovery, &config, logger)?
|
||||
run_caldav(&mut conn, &client, source_id, &discovery, &config, logger)
|
||||
}
|
||||
DavKindArg::Carddav => {
|
||||
run_carddav(&mut conn, &client, source_id, &discovery, &config, logger)?
|
||||
run_carddav(&mut conn, &client, source_id, &discovery, &config, logger)
|
||||
}
|
||||
DavKindArg::Webdav => {
|
||||
run_webdav(&mut conn, &client, source_id, &discovery, &config, logger)?
|
||||
run_webdav(&mut conn, &client, source_id, &discovery, &config, logger)
|
||||
}
|
||||
};
|
||||
|
||||
*summary = phase.summary;
|
||||
summary.retries_observed = client.retries_observed();
|
||||
summary.retry_after_sleeps = client.retry_after_sleeps();
|
||||
if let Some(e) = phase.error {
|
||||
return Err(e);
|
||||
}
|
||||
|
||||
if !summary.any_failed()
|
||||
&& let Err(e) = run_gc(&conn)
|
||||
{
|
||||
logger.warn(&format!("blob GC skipped: {e}"));
|
||||
}
|
||||
|
||||
let mut summary = summary;
|
||||
summary.retries_observed = client.retries_observed();
|
||||
summary.retry_after_sleeps = client.retry_after_sleeps();
|
||||
Ok(summary)
|
||||
Ok(())
|
||||
}
|
||||
|
||||
type ReconcileCollections = fn(
|
||||
@@ -198,25 +216,71 @@ fn run_collection_phase(
|
||||
config: &DavImportConfig,
|
||||
logger: Logger,
|
||||
phase: ItemPhase,
|
||||
) -> Result<Summary, Error> {
|
||||
) -> RunOutcome {
|
||||
let mut summary = Summary::default();
|
||||
let mut container_counts = TypeCounts::default();
|
||||
let mut item_counts = TypeCounts::default();
|
||||
let mut counts = PhaseCounts::default();
|
||||
let error = collection_phase_into(
|
||||
conn,
|
||||
client,
|
||||
source_id,
|
||||
PhaseInput {
|
||||
discovery,
|
||||
config,
|
||||
logger,
|
||||
phase: &phase,
|
||||
},
|
||||
&mut counts,
|
||||
)
|
||||
.err();
|
||||
|
||||
summary
|
||||
.per_type
|
||||
.push((phase.container_label, counts.container));
|
||||
summary.per_type.push((phase.item_label, counts.items));
|
||||
RunOutcome { summary, error }
|
||||
}
|
||||
|
||||
#[derive(Default)]
|
||||
struct PhaseCounts {
|
||||
container: TypeCounts,
|
||||
items: TypeCounts,
|
||||
}
|
||||
|
||||
struct PhaseInput<'a> {
|
||||
discovery: &'a Discovery,
|
||||
config: &'a DavImportConfig,
|
||||
logger: Logger,
|
||||
phase: &'a ItemPhase,
|
||||
}
|
||||
|
||||
fn collection_phase_into(
|
||||
conn: &mut Connection,
|
||||
client: &DavClient,
|
||||
source_id: i64,
|
||||
input: PhaseInput<'_>,
|
||||
counts: &mut PhaseCounts,
|
||||
) -> Result<(), Error> {
|
||||
let PhaseInput {
|
||||
discovery,
|
||||
config,
|
||||
logger,
|
||||
phase,
|
||||
} = input;
|
||||
|
||||
let upserted = (phase.reconcile_collections)(
|
||||
conn,
|
||||
source_id,
|
||||
&discovery.collections,
|
||||
&mut container_counts,
|
||||
&mut counts.container,
|
||||
logger,
|
||||
)?;
|
||||
if logger.enabled(LEVEL_DEFAULT) {
|
||||
eprintln!(
|
||||
"import: {} done (upserted={} deleted={} failed={})",
|
||||
phase.container_label,
|
||||
container_counts.created + container_counts.fetched,
|
||||
container_counts.deleted,
|
||||
container_counts.failed
|
||||
counts.container.created + counts.container.fetched,
|
||||
counts.container.deleted,
|
||||
counts.container.failed
|
||||
);
|
||||
}
|
||||
|
||||
@@ -229,23 +293,19 @@ fn run_collection_phase(
|
||||
logger,
|
||||
};
|
||||
for (collection_href, local_id) in &upserted {
|
||||
match (phase.reconcile_items)(conn, &ctx, collection_href, *local_id, &mut item_counts) {
|
||||
match (phase.reconcile_items)(conn, &ctx, collection_href, *local_id, &mut counts.items) {
|
||||
Ok(()) => {}
|
||||
Err(e) if e.aborts_run() => return Err(e),
|
||||
Err(e) => {
|
||||
logger.warn(&format!(
|
||||
"{} {collection_href:?}: items failed: {e}",
|
||||
phase.container_label
|
||||
));
|
||||
item_counts.failed += 1;
|
||||
counts.items.failed += 1;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
summary
|
||||
.per_type
|
||||
.push((phase.container_label, container_counts));
|
||||
summary.per_type.push((phase.item_label, item_counts));
|
||||
Ok(summary)
|
||||
Ok(())
|
||||
}
|
||||
|
||||
fn run_dry_diff(
|
||||
@@ -254,11 +314,11 @@ fn run_dry_diff(
|
||||
discovery: &Discovery,
|
||||
config: &DavImportConfig,
|
||||
logger: Logger,
|
||||
) -> Result<Summary, Error> {
|
||||
summary: &mut Summary,
|
||||
) -> Result<(), Error> {
|
||||
use crate::dav::href::join_absolute;
|
||||
use crate::dav::xml;
|
||||
use crate::db::dav_ids;
|
||||
let mut summary = Summary::default();
|
||||
let (container_label, item_label, container_type, item_type) = match config.kind {
|
||||
DavKindArg::Caldav => (
|
||||
"calendar",
|
||||
@@ -295,33 +355,57 @@ fn run_dry_diff(
|
||||
for coll in &discovery.collections {
|
||||
let url = join_absolute(&discovery.home_set_url, coll.href.as_str())
|
||||
.map_err(|e| Error::Partial(e.to_string()))?;
|
||||
let ms = client
|
||||
match client
|
||||
.propfind_responses(&url, 1, &xml::propfind_dav_items(), &url)
|
||||
.map_err(Error::from)?;
|
||||
let new_count = ms
|
||||
.responses
|
||||
.iter()
|
||||
.filter(|r| !r.props.is_collection)
|
||||
.count();
|
||||
item_counts.created += new_count as u64;
|
||||
if logger.enabled(LEVEL_DEFAULT) {
|
||||
eprintln!(" {} items: {new_count}", coll.href.as_str());
|
||||
.map_err(super::per_collection_failure)
|
||||
{
|
||||
Ok(ms) => {
|
||||
let new_count = ms
|
||||
.responses
|
||||
.iter()
|
||||
.filter(|r| !r.props.is_collection)
|
||||
.count();
|
||||
item_counts.created += new_count as u64;
|
||||
if logger.enabled(LEVEL_DEFAULT) {
|
||||
eprintln!(" {} items: {new_count}", coll.href.as_str());
|
||||
}
|
||||
}
|
||||
Err(e) if e.aborts_run() => return Err(e),
|
||||
Err(e) => {
|
||||
logger.warn(&format!(
|
||||
"{container_label} {:?}: enumeration failed: {e}",
|
||||
coll.href.as_str()
|
||||
));
|
||||
item_counts.failed += 1;
|
||||
}
|
||||
}
|
||||
}
|
||||
} else if let Some(root) = discovery.collections.first() {
|
||||
let url = join_absolute(&discovery.home_set_url, root.href.as_str())
|
||||
.map_err(|e| Error::Partial(e.to_string()))?;
|
||||
let ms = client
|
||||
match client
|
||||
.propfind_responses(&url, 1, &xml::propfind_webdav_listing(), &url)
|
||||
.map_err(Error::from)?;
|
||||
let new_count = ms
|
||||
.responses
|
||||
.iter()
|
||||
.filter(|r| !r.props.is_collection)
|
||||
.count();
|
||||
item_counts.created += new_count as u64;
|
||||
if logger.enabled(LEVEL_DEFAULT) {
|
||||
eprintln!(" root {} files: {new_count}", root.href.as_str());
|
||||
.map_err(super::per_collection_failure)
|
||||
{
|
||||
Ok(ms) => {
|
||||
let new_count = ms
|
||||
.responses
|
||||
.iter()
|
||||
.filter(|r| !r.props.is_collection)
|
||||
.count();
|
||||
item_counts.created += new_count as u64;
|
||||
if logger.enabled(LEVEL_DEFAULT) {
|
||||
eprintln!(" root {} files: {new_count}", root.href.as_str());
|
||||
}
|
||||
}
|
||||
Err(e) if e.aborts_run() => return Err(e),
|
||||
Err(e) => {
|
||||
logger.warn(&format!(
|
||||
"{container_label} {:?}: enumeration failed: {e}",
|
||||
root.href.as_str()
|
||||
));
|
||||
container_counts.failed += 1;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -332,7 +416,7 @@ fn run_dry_diff(
|
||||
if !matches!(config.kind, DavKindArg::Webdav) {
|
||||
summary.per_type.push((item_label, item_counts));
|
||||
}
|
||||
Ok(summary)
|
||||
Ok(())
|
||||
}
|
||||
|
||||
fn run_caldav(
|
||||
@@ -342,7 +426,7 @@ fn run_caldav(
|
||||
discovery: &Discovery,
|
||||
config: &DavImportConfig,
|
||||
logger: Logger,
|
||||
) -> Result<Summary, Error> {
|
||||
) -> RunOutcome {
|
||||
run_collection_phase(
|
||||
conn,
|
||||
client,
|
||||
@@ -366,7 +450,7 @@ fn run_carddav(
|
||||
discovery: &Discovery,
|
||||
config: &DavImportConfig,
|
||||
logger: Logger,
|
||||
) -> Result<Summary, Error> {
|
||||
) -> RunOutcome {
|
||||
run_collection_phase(
|
||||
conn,
|
||||
client,
|
||||
@@ -390,13 +474,16 @@ fn run_webdav(
|
||||
discovery: &Discovery,
|
||||
config: &DavImportConfig,
|
||||
logger: Logger,
|
||||
) -> Result<Summary, Error> {
|
||||
) -> RunOutcome {
|
||||
let mut summary = Summary::default();
|
||||
let mut file_counts = TypeCounts::default();
|
||||
|
||||
if discovery.collections.is_empty() {
|
||||
summary.per_type.push(("filenode", file_counts));
|
||||
return Ok(summary);
|
||||
return RunOutcome {
|
||||
summary,
|
||||
error: None,
|
||||
};
|
||||
}
|
||||
let root = &discovery.collections[0];
|
||||
let ctx = tree::WebDavCtx {
|
||||
@@ -406,10 +493,10 @@ fn run_webdav(
|
||||
dav_connections: config.dav_connections,
|
||||
logger,
|
||||
};
|
||||
tree::reconcile_filenodes(conn, &ctx, root, &mut file_counts)?;
|
||||
let error = tree::reconcile_filenodes(conn, &ctx, root, &mut file_counts).err();
|
||||
|
||||
summary.per_type.push(("filenode", file_counts));
|
||||
Ok(summary)
|
||||
RunOutcome { summary, error }
|
||||
}
|
||||
|
||||
fn map_discovery_error(err: DiscoveryError) -> Error {
|
||||
|
||||
@@ -113,7 +113,7 @@ fn enumerate_items(client: &DavClient, url: &str) -> Result<Vec<ServerItem>, Err
|
||||
let body = xml::propfind_dav_items();
|
||||
let ms = client
|
||||
.propfind_responses(url, 1, &body, url)
|
||||
.map_err(Error::from)?;
|
||||
.map_err(super::per_collection_failure)?;
|
||||
if ms.status >= 400 {
|
||||
return Err(Error::Partial(format!(
|
||||
"enumerate {url}: http {}",
|
||||
|
||||
@@ -115,6 +115,7 @@ pub fn reconcile_filenodes(
|
||||
}
|
||||
}
|
||||
}
|
||||
Err(e) if e.aborts_run() => return Err(e),
|
||||
Err(e) => {
|
||||
logger.warn(&format!("PROPFIND {url}: {e}"));
|
||||
counts.failed += 1;
|
||||
@@ -394,7 +395,7 @@ fn walk_one(
|
||||
let body = xml::propfind_webdav_listing();
|
||||
let ms = client
|
||||
.propfind_responses(url, 1, &body, url)
|
||||
.map_err(Error::from)?;
|
||||
.map_err(super::per_collection_failure)?;
|
||||
if ms.status >= 400 {
|
||||
return Err(Error::Partial(format!("http {}", ms.status)));
|
||||
}
|
||||
|
||||
@@ -11,6 +11,7 @@ use serde_json::{Value, json};
|
||||
|
||||
use crate::db::exchange_graph_ids;
|
||||
use crate::error::Error;
|
||||
use crate::exchange::jscalendar::override_patch_from_event;
|
||||
use crate::exchange_graph::api::{self, PREFER_BODY_HTML, PREFER_BODY_TEXT, PREFER_TIMEZONE_UTC};
|
||||
use crate::exchange_graph::calendar_map::{
|
||||
ConvertedEvent, EventType, classify_event_type, convert_event,
|
||||
@@ -185,7 +186,7 @@ fn merge_exception_into(master: &mut ConvertedEvent, ex: &ConvertedEvent) {
|
||||
let Value::Object(overrides) = overrides else {
|
||||
return;
|
||||
};
|
||||
overrides.insert(key, ex.data.clone());
|
||||
overrides.insert(key, override_patch_from_event(&ex.data));
|
||||
}
|
||||
|
||||
fn merge_exception_into_existing(
|
||||
@@ -253,7 +254,7 @@ fn merge_persisted_master(
|
||||
.entry("recurrenceOverrides".to_owned())
|
||||
.or_insert_with(|| Value::Object(serde_json::Map::new()));
|
||||
if let Value::Object(overrides) = entry {
|
||||
overrides.insert(key, ex_data.clone());
|
||||
overrides.insert(key, override_patch_from_event(ex_data));
|
||||
}
|
||||
}
|
||||
tx.execute(
|
||||
|
||||
@@ -16,11 +16,10 @@ use crate::exchange_graph::error::GraphError;
|
||||
use crate::exchange_graph::oauth::{
|
||||
AcquiredToken, OAuthFlow, acquire, default_authority, refresh_access_token,
|
||||
};
|
||||
use crate::exchange_graph::types::{EventBodyFormat, MailboxKind, synthetic_account_id};
|
||||
use crate::exchange_graph::types::{EventBodyFormat, MailboxKind, Surfaces, synthetic_account_id};
|
||||
use crate::jmap::http::RetryPolicy;
|
||||
use crate::logging::LEVEL_DEFAULT;
|
||||
use crate::sync::{CommonConfig, Summary, TypeCounts};
|
||||
use crate::types::ObjectType;
|
||||
|
||||
#[derive(Debug, Clone)]
|
||||
pub enum GraphAuth {
|
||||
@@ -39,7 +38,7 @@ pub struct GraphImportConfig {
|
||||
pub api_base: String,
|
||||
pub user_target: Option<String>,
|
||||
pub mailbox_kind: MailboxKind,
|
||||
pub objects: Option<Vec<ObjectType>>,
|
||||
pub surfaces: Surfaces,
|
||||
pub event_body_format: EventBodyFormat,
|
||||
pub graph_connections: usize,
|
||||
pub top: usize,
|
||||
@@ -141,32 +140,10 @@ pub fn run(common: CommonConfig, config: GraphImportConfig) -> Result<Summary, E
|
||||
event_body_format: config.event_body_format,
|
||||
};
|
||||
|
||||
let want_mail = config
|
||||
.objects
|
||||
.as_ref()
|
||||
.map(|set| {
|
||||
set.iter()
|
||||
.any(|o| matches!(o, ObjectType::Mailbox | ObjectType::Email))
|
||||
})
|
||||
.unwrap_or(true);
|
||||
let want_calendar = config
|
||||
.objects
|
||||
.as_ref()
|
||||
.map(|set| {
|
||||
set.iter()
|
||||
.any(|o| matches!(o, ObjectType::Calendar | ObjectType::CalendarEvent))
|
||||
})
|
||||
.unwrap_or(true)
|
||||
&& !matches!(config.mailbox_kind, MailboxKind::Archive);
|
||||
let want_contacts = config
|
||||
.objects
|
||||
.as_ref()
|
||||
.map(|set| {
|
||||
set.iter()
|
||||
.any(|o| matches!(o, ObjectType::AddressBook | ObjectType::ContactCard))
|
||||
})
|
||||
.unwrap_or(true)
|
||||
&& !matches!(config.mailbox_kind, MailboxKind::Archive);
|
||||
let primary = !matches!(config.mailbox_kind, MailboxKind::Archive);
|
||||
let want_mail = config.surfaces.mail;
|
||||
let want_calendar = config.surfaces.calendar && primary;
|
||||
let want_contacts = config.surfaces.contacts && primary;
|
||||
|
||||
if want_mail {
|
||||
let folders = super::folders::reconcile_mail(
|
||||
@@ -509,7 +486,7 @@ mod tests {
|
||||
api_base: "https://graph.microsoft.com/v1.0".to_owned(),
|
||||
user_target: None,
|
||||
mailbox_kind: MailboxKind::Primary,
|
||||
objects: None,
|
||||
surfaces: Surfaces::ALL,
|
||||
event_body_format: EventBodyFormat::Text,
|
||||
graph_connections: 4,
|
||||
top: 100,
|
||||
|
||||
@@ -13,4 +13,4 @@ pub mod pool;
|
||||
|
||||
pub mod coordinator;
|
||||
|
||||
pub use coordinator::{ImapAuth, ImapImportConfig, run};
|
||||
pub use coordinator::{ImapAuth, ImapImportConfig, run, run_reporting};
|
||||
|
||||
@@ -29,7 +29,7 @@ use crate::imap::transport::Connector;
|
||||
use crate::logging::{LEVEL_DEFAULT, LEVEL_PROGRESS, Logger};
|
||||
use crate::sync::emailmeta::email_index_from_blob;
|
||||
use crate::sync::keys::index_to_json;
|
||||
use crate::sync::{CommonConfig, Summary, TypeCounts};
|
||||
use crate::sync::{CommonConfig, RunOutcome, Summary, TypeCounts};
|
||||
|
||||
use super::fetch;
|
||||
use super::folders::{
|
||||
@@ -199,6 +199,20 @@ pub enum ImapAuth {
|
||||
}
|
||||
|
||||
pub fn run(common: CommonConfig, config: ImapImportConfig) -> Result<Summary, Error> {
|
||||
run_reporting(common, config).into_result()
|
||||
}
|
||||
|
||||
pub fn run_reporting(common: CommonConfig, config: ImapImportConfig) -> RunOutcome {
|
||||
let mut summary = Summary::default();
|
||||
let error = run_into(common, config, &mut summary).err();
|
||||
RunOutcome { summary, error }
|
||||
}
|
||||
|
||||
fn run_into(
|
||||
common: CommonConfig,
|
||||
config: ImapImportConfig,
|
||||
summary: &mut Summary,
|
||||
) -> Result<(), Error> {
|
||||
let logger = common.logger;
|
||||
let mut conn = db::init::open(&common.archive)?;
|
||||
|
||||
@@ -341,14 +355,15 @@ pub fn run(common: CommonConfig, config: ImapImportConfig) -> Result<Summary, Er
|
||||
|
||||
if common.dry_run {
|
||||
let existing_source = db::sources::find_source(&conn, &source_key)?;
|
||||
return dry_run_summary(
|
||||
*summary = dry_run_summary(
|
||||
&conn,
|
||||
existing_source,
|
||||
&mut client,
|
||||
&control_ctx,
|
||||
&resolved,
|
||||
logger,
|
||||
);
|
||||
)?;
|
||||
return Ok(());
|
||||
}
|
||||
|
||||
let mut mailbox_counts = TypeCounts::default();
|
||||
@@ -410,6 +425,16 @@ pub fn run(common: CommonConfig, config: ImapImportConfig) -> Result<Summary, Er
|
||||
&mut email_counts,
|
||||
) {
|
||||
Ok(()) => {}
|
||||
Err(e) if e.aborts_run() => {
|
||||
pool.shutdown();
|
||||
let _ = client.logout();
|
||||
*summary = Summary {
|
||||
per_type: vec![("mailbox", mailbox_counts), ("email", email_counts)],
|
||||
retries_observed: backoff.total_retries(),
|
||||
retry_after_sleeps: backoff.transient_retries() as u64,
|
||||
};
|
||||
return Err(e);
|
||||
}
|
||||
Err(e) => {
|
||||
log_at(
|
||||
logger,
|
||||
@@ -424,11 +449,12 @@ pub fn run(common: CommonConfig, config: ImapImportConfig) -> Result<Summary, Er
|
||||
pool.shutdown();
|
||||
let _ = client.logout();
|
||||
|
||||
Ok(Summary {
|
||||
*summary = Summary {
|
||||
per_type: vec![("mailbox", mailbox_counts), ("email", email_counts)],
|
||||
retries_observed: backoff.total_retries(),
|
||||
retry_after_sleeps: backoff.transient_retries() as u64,
|
||||
})
|
||||
};
|
||||
Ok(())
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone)]
|
||||
|
||||
+19
-4
@@ -33,7 +33,7 @@ use crate::jmap::wire::email::Email;
|
||||
use crate::jmap::wire::file_node::{FileNode, NodeType};
|
||||
use crate::jmap::wire::sieve_script::SieveScript;
|
||||
use crate::logging::{LEVEL_DEFAULT, LEVEL_PROGRESS, Logger};
|
||||
use crate::sync::{CommonConfig, Context, ImportConfig, Summary, TypeCounts};
|
||||
use crate::sync::{CommonConfig, Context, ImportConfig, RunOutcome, Summary, TypeCounts};
|
||||
use crate::types::ObjectType;
|
||||
|
||||
const IMPORT_ORDER: [ObjectType; 10] = [
|
||||
@@ -147,6 +147,20 @@ fn username_of(auth: &Auth) -> String {
|
||||
}
|
||||
|
||||
pub fn run(common: CommonConfig, config: ImportConfig) -> Result<Summary, Error> {
|
||||
run_reporting(common, config).into_result()
|
||||
}
|
||||
|
||||
pub fn run_reporting(common: CommonConfig, config: ImportConfig) -> RunOutcome {
|
||||
let mut summary = Summary::default();
|
||||
let error = run_into(common, config, &mut summary).err();
|
||||
RunOutcome { summary, error }
|
||||
}
|
||||
|
||||
fn run_into(
|
||||
common: CommonConfig,
|
||||
config: ImportConfig,
|
||||
summary: &mut Summary,
|
||||
) -> Result<(), Error> {
|
||||
let logger = common.logger;
|
||||
let ctx = Context::open(common, &config.connect)?;
|
||||
let connected = connect::prepare(&ctx, &config.connect)?;
|
||||
@@ -196,7 +210,6 @@ pub fn run(common: CommonConfig, config: ImportConfig) -> Result<Summary, Error>
|
||||
session: connected.session.clone(),
|
||||
};
|
||||
|
||||
let mut summary = Summary::default();
|
||||
let mut dry_rows: Vec<(&'static str, u64, u64, u64)> = Vec::new();
|
||||
let threads = ctx.common.threads;
|
||||
|
||||
@@ -216,6 +229,7 @@ pub fn run(common: CommonConfig, config: ImportConfig) -> Result<Summary, Error>
|
||||
&mut dry_rows,
|
||||
) {
|
||||
Ok(()) => {}
|
||||
Err(e) if e.aborts_run() => return Err(e),
|
||||
Err(e) => {
|
||||
logger.warn(&format!(
|
||||
"type {} aborted: {e}; continuing (run will exit 5)",
|
||||
@@ -229,7 +243,8 @@ pub fn run(common: CommonConfig, config: ImportConfig) -> Result<Summary, Error>
|
||||
|
||||
if ctx.dry_run() {
|
||||
print_dry_run(&dry_rows);
|
||||
return Ok(Summary::default());
|
||||
*summary = Summary::default();
|
||||
return Ok(());
|
||||
}
|
||||
|
||||
if !summary.any_failed()
|
||||
@@ -240,7 +255,7 @@ pub fn run(common: CommonConfig, config: ImportConfig) -> Result<Summary, Error>
|
||||
|
||||
summary.retries_observed = ctx.client.retries_observed();
|
||||
summary.retry_after_sleeps = ctx.client.retry_after_sleeps();
|
||||
Ok(summary)
|
||||
Ok(())
|
||||
}
|
||||
|
||||
fn work_list(config: &ImportConfig, connected: &Connected) -> Vec<ObjectType> {
|
||||
|
||||
+171
-30
@@ -8,7 +8,7 @@ use indexmap::IndexMap;
|
||||
use rusqlite::{Connection, Row, params};
|
||||
use serde_json::{Map, Value};
|
||||
|
||||
use crate::jmap::blob::{export_blob_ids, import_blob_ids};
|
||||
use crate::jmap::blob::{BlobWalkError, InlineShape, import_blob_ids, inline_blob_data_uris};
|
||||
use crate::jmap::error::JmapError;
|
||||
use crate::jmap::wire::JmapId;
|
||||
use crate::jmap::wire::address_book::AddressBook;
|
||||
@@ -35,8 +35,8 @@ pub trait BlobIntern {
|
||||
fn intern(&mut self, jmap_blob_id: &str) -> Result<i64, JmapError>;
|
||||
}
|
||||
|
||||
pub trait BlobUpload {
|
||||
fn upload(&mut self, local_id: i64) -> Result<JmapId, JmapError>;
|
||||
pub trait BlobBytes {
|
||||
fn bytes(&self, local_id: i64) -> Result<Vec<u8>, JmapError>;
|
||||
}
|
||||
|
||||
fn translate_in(
|
||||
@@ -662,10 +662,10 @@ pub fn contact_card_to_wire(
|
||||
address_book_ids: &str,
|
||||
data: &str,
|
||||
resolver: &impl TargetResolver,
|
||||
blobs: &mut impl BlobUpload,
|
||||
blobs: &impl BlobBytes,
|
||||
) -> Result<Value, JmapError> {
|
||||
let mut value: Value = serde_json::from_str(data)?;
|
||||
restore_blobs_in(&mut value, blobs)?;
|
||||
inline_blobs_in(&mut value, InlineShape::JsContactResource, blobs)?;
|
||||
let locals = parse_local_id_array(address_book_ids)?;
|
||||
let abids = translate_out(&locals, ObjectType::AddressBook, resolver)?;
|
||||
if let Value::Object(map) = &mut value {
|
||||
@@ -681,10 +681,10 @@ pub fn calendar_event_to_wire(
|
||||
use_default_alerts: bool,
|
||||
data: &str,
|
||||
resolver: &impl TargetResolver,
|
||||
blobs: &mut impl BlobUpload,
|
||||
blobs: &impl BlobBytes,
|
||||
) -> Result<Value, JmapError> {
|
||||
let mut value: Value = serde_json::from_str(data)?;
|
||||
restore_blobs_in(&mut value, blobs)?;
|
||||
inline_blobs_in(&mut value, InlineShape::JsCalendarLink, blobs)?;
|
||||
let locals = parse_local_id_array(calendar_ids)?;
|
||||
let calids = translate_out(&locals, ObjectType::Calendar, resolver)?;
|
||||
if let Value::Object(map) = &mut value {
|
||||
@@ -725,21 +725,20 @@ fn take_string(value: &mut Value, key: &str) -> Option<String> {
|
||||
|
||||
fn rewrite_blobs_in(data: &mut Value, blobs: &mut impl BlobIntern) -> Result<(), JmapError> {
|
||||
import_blob_ids(data, |jmap_blob_id| {
|
||||
blobs
|
||||
.intern(jmap_blob_id)
|
||||
.map_err(|e| crate::jmap::blob::BlobWalkError::Resolver(e.to_string()))
|
||||
blobs.intern(jmap_blob_id).map_err(BlobWalkError::resolver)
|
||||
})
|
||||
.map_err(JmapError::from)
|
||||
.map_err(BlobWalkError::into_source)
|
||||
}
|
||||
|
||||
fn restore_blobs_in(data: &mut Value, blobs: &mut impl BlobUpload) -> Result<(), JmapError> {
|
||||
export_blob_ids(data, |local_id| {
|
||||
blobs
|
||||
.upload(local_id)
|
||||
.map(|id| id.0)
|
||||
.map_err(|e| crate::jmap::blob::BlobWalkError::Resolver(e.to_string()))
|
||||
fn inline_blobs_in(
|
||||
data: &mut Value,
|
||||
shape: InlineShape,
|
||||
blobs: &impl BlobBytes,
|
||||
) -> Result<(), JmapError> {
|
||||
inline_blob_data_uris(data, shape, |local_id| {
|
||||
blobs.bytes(local_id).map_err(BlobWalkError::resolver)
|
||||
})
|
||||
.map_err(JmapError::from)
|
||||
.map_err(BlobWalkError::into_source)
|
||||
}
|
||||
|
||||
fn opt_json<T: serde::Serialize>(value: &Option<T>) -> Result<Option<String>, JmapError> {
|
||||
@@ -904,9 +903,35 @@ mod tests {
|
||||
Ok(42)
|
||||
}
|
||||
}
|
||||
impl BlobUpload for FakeBlobs {
|
||||
fn upload(&mut self, local_id: i64) -> Result<JmapId, JmapError> {
|
||||
Ok(JmapId(format!("T{local_id}")))
|
||||
impl BlobBytes for FakeBlobs {
|
||||
fn bytes(&self, local_id: i64) -> Result<Vec<u8>, JmapError> {
|
||||
Ok(format!("payload-{local_id}").into_bytes())
|
||||
}
|
||||
}
|
||||
|
||||
fn decode_data_uri(value: &Value, media_type: &str) -> Vec<u8> {
|
||||
use base64::Engine;
|
||||
let uri = value.as_str().unwrap_or_else(|| panic!("{value} is a URI"));
|
||||
let prefix = format!("data:{media_type};base64,");
|
||||
let payload = uri
|
||||
.strip_prefix(&prefix)
|
||||
.unwrap_or_else(|| panic!("{uri} does not start with {prefix}"));
|
||||
base64::engine::general_purpose::STANDARD
|
||||
.decode(payload)
|
||||
.expect("base64 payload")
|
||||
}
|
||||
|
||||
fn assert_no_blob_id(value: &Value) {
|
||||
match value {
|
||||
Value::Object(map) => {
|
||||
assert!(map.get("blobId").is_none(), "blobId present in {value}");
|
||||
assert!(map.get("@blob").is_none(), "@blob present in {value}");
|
||||
for child in map.values() {
|
||||
assert_no_blob_id(child);
|
||||
}
|
||||
}
|
||||
Value::Array(items) => items.iter().for_each(assert_no_blob_id),
|
||||
_ => {}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1032,7 +1057,7 @@ mod tests {
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn contact_card_strips_uid_and_rewrites_blob() {
|
||||
fn contact_card_strips_uid_and_inlines_media_as_data_uri() {
|
||||
let c = mem();
|
||||
let res = MapResolver {
|
||||
to_local: HashMap::from([((ObjectType::AddressBook, "AB".to_owned()), 1)]),
|
||||
@@ -1043,7 +1068,10 @@ mod tests {
|
||||
"addressBookIds": { "AB": true },
|
||||
"uid": "urn:uuid:42",
|
||||
"name": { "full": "Jane" },
|
||||
"photos": { "p1": { "blobId": "PB", "mediaType": "image/png" } }
|
||||
"media": { "photo": {
|
||||
"@type": "Media", "kind": "photo",
|
||||
"blobId": "PB", "mediaType": "image/png"
|
||||
} }
|
||||
}))
|
||||
.unwrap();
|
||||
let mut blobs = FakeBlobs;
|
||||
@@ -1058,14 +1086,127 @@ mod tests {
|
||||
assert_eq!(uid, "urn:uuid:42");
|
||||
let stored: Value = serde_json::from_str(&data).unwrap();
|
||||
assert!(stored.get("uid").is_none());
|
||||
assert_eq!(stored["photos"]["p1"]["@blob"], Value::from(42));
|
||||
assert_eq!(stored["media"]["photo"]["@blob"], Value::from(42));
|
||||
|
||||
let mut up = FakeBlobs;
|
||||
let wire = contact_card_to_wire(&uid, &abids, &data, &res, &mut up).unwrap();
|
||||
let up = FakeBlobs;
|
||||
let wire = contact_card_to_wire(&uid, &abids, &data, &res, &up).unwrap();
|
||||
assert_eq!(wire["uid"], Value::from("urn:uuid:42"));
|
||||
assert_eq!(wire["addressBookIds"]["TAB"], Value::Bool(true));
|
||||
assert_eq!(wire["photos"]["p1"]["blobId"], Value::from("T42"));
|
||||
assert!(wire["photos"]["p1"].get("@blob").is_none());
|
||||
assert_eq!(
|
||||
decode_data_uri(&wire["media"]["photo"]["uri"], "image/png"),
|
||||
b"payload-42"
|
||||
);
|
||||
assert_eq!(
|
||||
wire["media"]["photo"]["mediaType"],
|
||||
Value::from("image/png")
|
||||
);
|
||||
assert!(wire["media"]["photo"].get("href").is_none());
|
||||
assert_no_blob_id(&wire);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn calendar_event_inlines_link_enclosure_as_data_uri() {
|
||||
let c = mem();
|
||||
let res = MapResolver {
|
||||
to_local: HashMap::from([((ObjectType::Calendar, "CAL".to_owned()), 5)]),
|
||||
to_target: HashMap::from([((ObjectType::Calendar, 5), "TCAL".to_owned())]),
|
||||
};
|
||||
let ev: CalendarEvent = serde_json::from_value(serde_json::json!({
|
||||
"id": "EV1",
|
||||
"calendarIds": { "CAL": true },
|
||||
"uid": "ev-with-enclosure",
|
||||
"title": "Review",
|
||||
"@type": "Event",
|
||||
"links": { "1": {
|
||||
"@type": "Link", "rel": "enclosure",
|
||||
"blobId": "AB", "contentType": "text/plain", "title": "agenda.txt"
|
||||
} }
|
||||
}))
|
||||
.unwrap();
|
||||
let mut blobs = FakeBlobs;
|
||||
let local = insert_calendar_event(&c, &ev, &res, &mut blobs).unwrap();
|
||||
let (cal, dr, ud, data): (String, i64, i64, String) = c
|
||||
.query_row(
|
||||
&format!("{CALENDAR_EVENT_SELECT} AND id = ?1"),
|
||||
params![local],
|
||||
|r| Ok((r.get(1)?, r.get(2)?, r.get(3)?, r.get(4)?)),
|
||||
)
|
||||
.unwrap();
|
||||
let stored: Value = serde_json::from_str(&data).unwrap();
|
||||
assert_eq!(stored["links"]["1"]["@blob"], Value::from(42));
|
||||
|
||||
let up = FakeBlobs;
|
||||
let wire = calendar_event_to_wire(&cal, dr != 0, ud != 0, &data, &res, &up).unwrap();
|
||||
assert_eq!(
|
||||
decode_data_uri(&wire["links"]["1"]["href"], "text/plain"),
|
||||
b"payload-42"
|
||||
);
|
||||
assert_eq!(wire["links"]["1"]["contentType"], Value::from("text/plain"));
|
||||
assert_eq!(wire["links"]["1"]["rel"], Value::from("enclosure"));
|
||||
assert!(wire["links"]["1"].get("uri").is_none());
|
||||
assert_no_blob_id(&wire);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn contact_card_media_without_media_type_defaults_to_octet_stream() {
|
||||
let c = mem();
|
||||
let res = MapResolver {
|
||||
to_local: HashMap::from([((ObjectType::AddressBook, "AB".to_owned()), 1)]),
|
||||
to_target: HashMap::from([((ObjectType::AddressBook, 1), "TAB".to_owned())]),
|
||||
};
|
||||
c.execute(
|
||||
"INSERT INTO contact_cards (id,uid,address_book_ids,data)
|
||||
VALUES (1,'u-1','[1]',?1)",
|
||||
params![
|
||||
serde_json::json!({ "@type": "Card", "media": { "photo": { "@blob": 9 } } })
|
||||
.to_string()
|
||||
],
|
||||
)
|
||||
.unwrap();
|
||||
let (uid, abids, data): (String, String, String) = c
|
||||
.query_row(&format!("{CONTACT_CARD_SELECT} WHERE id = 1"), [], |r| {
|
||||
Ok((r.get(1)?, r.get(2)?, r.get(3)?))
|
||||
})
|
||||
.unwrap();
|
||||
let up = FakeBlobs;
|
||||
let wire = contact_card_to_wire(&uid, &abids, &data, &res, &up).unwrap();
|
||||
assert_eq!(
|
||||
decode_data_uri(&wire["media"]["photo"]["uri"], "application/octet-stream"),
|
||||
b"payload-9"
|
||||
);
|
||||
assert_no_blob_id(&wire);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn an_archive_read_failure_while_inlining_is_an_archive_error_not_a_unit_failure() {
|
||||
struct BrokenBlobs;
|
||||
impl BlobBytes for BrokenBlobs {
|
||||
fn bytes(&self, _local_id: i64) -> Result<Vec<u8>, JmapError> {
|
||||
Err(JmapError::Sqlite(rusqlite::Error::QueryReturnedNoRows))
|
||||
}
|
||||
}
|
||||
let res = MapResolver {
|
||||
to_local: HashMap::new(),
|
||||
to_target: HashMap::from([
|
||||
((ObjectType::AddressBook, 1), "TAB".to_owned()),
|
||||
((ObjectType::Calendar, 1), "TCAL".to_owned()),
|
||||
]),
|
||||
};
|
||||
let card = serde_json::json!({ "@type": "Card", "media": { "photo": { "@blob": 9 } } })
|
||||
.to_string();
|
||||
let event =
|
||||
serde_json::json!({ "@type": "Event", "links": { "1": { "@blob": 9 } } }).to_string();
|
||||
let failures = [
|
||||
contact_card_to_wire("u-1", "[1]", &card, &res, &BrokenBlobs).expect_err("card"),
|
||||
calendar_event_to_wire("[1]", false, false, &event, &res, &BrokenBlobs)
|
||||
.expect_err("event"),
|
||||
];
|
||||
for err in failures {
|
||||
assert!(matches!(err, JmapError::Sqlite(_)), "{err:?}");
|
||||
let mapped = crate::error::Error::from(err);
|
||||
assert!(mapped.aborts_run(), "{mapped} must abort the run");
|
||||
assert_eq!(mapped.exit_code(), 7);
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
@@ -1103,8 +1244,8 @@ mod tests {
|
||||
assert!(stored.get("calendarIds").is_none());
|
||||
assert_eq!(stored["title"], Value::from("Sprint"));
|
||||
|
||||
let mut up = FakeBlobs;
|
||||
let wire = calendar_event_to_wire(&cal, dr != 0, ud != 0, &data, &res, &mut up).unwrap();
|
||||
let up = FakeBlobs;
|
||||
let wire = calendar_event_to_wire(&cal, dr != 0, ud != 0, &data, &res, &up).unwrap();
|
||||
assert_eq!(wire["calendarIds"]["TCAL"], Value::Bool(true));
|
||||
assert_eq!(wire["isDraft"], Value::Bool(true));
|
||||
assert_eq!(wire["title"], Value::from("Sprint"));
|
||||
|
||||
@@ -9,4 +9,4 @@ pub mod keywords;
|
||||
pub mod messages;
|
||||
pub mod tree;
|
||||
|
||||
pub use coordinator::{MaildirImportConfig, run};
|
||||
pub use coordinator::{MaildirImportConfig, run, run_reporting};
|
||||
|
||||
@@ -15,7 +15,7 @@ use crate::db;
|
||||
use crate::db::sources::SourceKey;
|
||||
use crate::error::Error;
|
||||
use crate::logging::{LEVEL_DEFAULT, LEVEL_PROGRESS, Logger};
|
||||
use crate::sync::{CommonConfig, Summary, TypeCounts};
|
||||
use crate::sync::{CommonConfig, RunOutcome, Summary, TypeCounts};
|
||||
|
||||
use super::messages;
|
||||
use super::tree;
|
||||
@@ -32,6 +32,20 @@ pub struct MaildirImportConfig {
|
||||
}
|
||||
|
||||
pub fn run(common: CommonConfig, config: MaildirImportConfig) -> Result<Summary, Error> {
|
||||
run_reporting(common, config).into_result()
|
||||
}
|
||||
|
||||
pub fn run_reporting(common: CommonConfig, config: MaildirImportConfig) -> RunOutcome {
|
||||
let mut summary = Summary::default();
|
||||
let error = run_into(common, config, &mut summary).err();
|
||||
RunOutcome { summary, error }
|
||||
}
|
||||
|
||||
fn run_into(
|
||||
common: CommonConfig,
|
||||
config: MaildirImportConfig,
|
||||
summary: &mut Summary,
|
||||
) -> Result<(), Error> {
|
||||
let logger = common.logger;
|
||||
if common.threads > 1 {
|
||||
log_at(
|
||||
@@ -121,8 +135,8 @@ pub fn run(common: CommonConfig, config: MaildirImportConfig) -> Result<Summary,
|
||||
);
|
||||
|
||||
if common.dry_run {
|
||||
let summary = build_dry_run_summary(&conn, source_id, &resolved, logger)?;
|
||||
return Ok(summary);
|
||||
*summary = build_dry_run_summary(&conn, source_id, &resolved, logger)?;
|
||||
return Ok(());
|
||||
}
|
||||
|
||||
let mut mailbox_counts = TypeCounts::default();
|
||||
@@ -174,6 +188,14 @@ pub fn run(common: CommonConfig, config: MaildirImportConfig) -> Result<Summary,
|
||||
&mut email_counts,
|
||||
logger,
|
||||
) {
|
||||
if e.aborts_run() {
|
||||
*summary = Summary {
|
||||
per_type: vec![("mailbox", mailbox_counts), ("email", email_counts)],
|
||||
retries_observed: 0,
|
||||
retry_after_sleeps: 0,
|
||||
};
|
||||
return Err(e);
|
||||
}
|
||||
log_at(
|
||||
logger,
|
||||
LEVEL_DEFAULT,
|
||||
@@ -207,11 +229,12 @@ pub fn run(common: CommonConfig, config: MaildirImportConfig) -> Result<Summary,
|
||||
),
|
||||
);
|
||||
|
||||
Ok(Summary {
|
||||
*summary = Summary {
|
||||
per_type: vec![("mailbox", mailbox_counts), ("email", email_counts)],
|
||||
retries_observed: 0,
|
||||
retry_after_sleeps: 0,
|
||||
})
|
||||
};
|
||||
Ok(())
|
||||
}
|
||||
|
||||
struct RunFlags {
|
||||
|
||||
@@ -13,4 +13,4 @@ pub mod mbox;
|
||||
pub mod tree;
|
||||
pub mod walk;
|
||||
|
||||
pub use coordinator::{TakeoutImportConfig, run};
|
||||
pub use coordinator::{TakeoutImportConfig, run, run_reporting};
|
||||
|
||||
@@ -15,7 +15,7 @@ use crate::db::sources::SourceKey;
|
||||
use crate::db::takeout_ids;
|
||||
use crate::error::Error;
|
||||
use crate::logging::{LEVEL_DEFAULT, LEVEL_PROGRESS, Logger};
|
||||
use crate::sync::{CommonConfig, Summary, TypeCounts};
|
||||
use crate::sync::{CommonConfig, RunOutcome, Summary, TypeCounts};
|
||||
|
||||
use super::calendar;
|
||||
use super::contacts;
|
||||
@@ -31,6 +31,20 @@ pub struct TakeoutImportConfig {
|
||||
}
|
||||
|
||||
pub fn run(common: CommonConfig, config: TakeoutImportConfig) -> Result<Summary, Error> {
|
||||
run_reporting(common, config).into_result()
|
||||
}
|
||||
|
||||
pub fn run_reporting(common: CommonConfig, config: TakeoutImportConfig) -> RunOutcome {
|
||||
let mut summary = Summary::default();
|
||||
let error = run_into(common, config, &mut summary).err();
|
||||
RunOutcome { summary, error }
|
||||
}
|
||||
|
||||
fn run_into(
|
||||
common: CommonConfig,
|
||||
config: TakeoutImportConfig,
|
||||
summary: &mut Summary,
|
||||
) -> Result<(), Error> {
|
||||
let logger = common.logger;
|
||||
if common.threads > 1 && logger.enabled(LEVEL_PROGRESS) {
|
||||
eprintln!("takeout importer is single-threaded; --threads value will be ignored");
|
||||
@@ -90,7 +104,8 @@ pub fn run(common: CommonConfig, config: TakeoutImportConfig) -> Result<Summary,
|
||||
};
|
||||
|
||||
if common.dry_run {
|
||||
return Ok(dry_run_summary(&conn, source_id, &walk_result, logger));
|
||||
*summary = dry_run_summary(&conn, source_id, &walk_result, logger);
|
||||
return Ok(());
|
||||
}
|
||||
|
||||
let options = MappingOptions {
|
||||
@@ -107,7 +122,8 @@ pub fn run(common: CommonConfig, config: TakeoutImportConfig) -> Result<Summary,
|
||||
let mut mailbox_cache: HashMap<String, i64> =
|
||||
takeout_ids::all_for_type(&conn, source_id, takeout_ids::MAILBOX)?;
|
||||
|
||||
process_mbox_files(
|
||||
let mut aborted: Option<Error> = None;
|
||||
if let Err(e) = process_mbox_files(
|
||||
&mut conn,
|
||||
source_id,
|
||||
&walk_result,
|
||||
@@ -116,27 +132,38 @@ pub fn run(common: CommonConfig, config: TakeoutImportConfig) -> Result<Summary,
|
||||
&mut mailbox_counts,
|
||||
&mut email_counts,
|
||||
logger,
|
||||
)?;
|
||||
) {
|
||||
aborted = Some(e);
|
||||
}
|
||||
|
||||
process_ics_files(
|
||||
&mut conn,
|
||||
source_id,
|
||||
&walk_result,
|
||||
&mut calendar_counts,
|
||||
&mut calendar_event_counts,
|
||||
logger,
|
||||
)?;
|
||||
if aborted.is_none()
|
||||
&& let Err(e) = process_ics_files(
|
||||
&mut conn,
|
||||
source_id,
|
||||
&walk_result,
|
||||
&mut calendar_counts,
|
||||
&mut calendar_event_counts,
|
||||
logger,
|
||||
)
|
||||
{
|
||||
aborted = Some(e);
|
||||
}
|
||||
|
||||
process_vcf_files(
|
||||
&mut conn,
|
||||
source_id,
|
||||
&walk_result,
|
||||
&mut book_counts,
|
||||
&mut card_counts,
|
||||
logger,
|
||||
)?;
|
||||
if aborted.is_none()
|
||||
&& let Err(e) = process_vcf_files(
|
||||
&mut conn,
|
||||
source_id,
|
||||
&walk_result,
|
||||
&mut book_counts,
|
||||
&mut card_counts,
|
||||
logger,
|
||||
)
|
||||
{
|
||||
aborted = Some(e);
|
||||
}
|
||||
|
||||
let no_failures = mailbox_counts.failed == 0
|
||||
let no_failures = aborted.is_none()
|
||||
&& mailbox_counts.failed == 0
|
||||
&& email_counts.failed == 0
|
||||
&& calendar_counts.failed == 0
|
||||
&& calendar_event_counts.failed == 0
|
||||
@@ -171,7 +198,7 @@ pub fn run(common: CommonConfig, config: TakeoutImportConfig) -> Result<Summary,
|
||||
);
|
||||
}
|
||||
|
||||
Ok(Summary {
|
||||
*summary = Summary {
|
||||
per_type: vec![
|
||||
("mailbox", mailbox_counts),
|
||||
("email", email_counts),
|
||||
@@ -182,7 +209,11 @@ pub fn run(common: CommonConfig, config: TakeoutImportConfig) -> Result<Summary,
|
||||
],
|
||||
retries_observed: 0,
|
||||
retry_after_sleeps: 0,
|
||||
})
|
||||
};
|
||||
match aborted {
|
||||
Some(e) => Err(e),
|
||||
None => Ok(()),
|
||||
}
|
||||
}
|
||||
|
||||
#[allow(clippy::too_many_arguments)]
|
||||
@@ -212,6 +243,7 @@ fn process_mbox_files(
|
||||
};
|
||||
match mail::process_file(conn, &file.path, ctx, mailbox_counts, email_counts, logger) {
|
||||
Ok(()) => {}
|
||||
Err(e) if e.aborts_run() => return Err(e),
|
||||
Err(e) => {
|
||||
logger.warn(&format!("mbox {:?}: {e}", file.path));
|
||||
email_counts.failed += 1;
|
||||
@@ -244,6 +276,9 @@ fn process_ics_files(
|
||||
event_counts,
|
||||
logger,
|
||||
) {
|
||||
if e.aborts_run() {
|
||||
return Err(e);
|
||||
}
|
||||
logger.warn(&format!("ics {:?}: {e}", file.path));
|
||||
event_counts.failed += 1;
|
||||
}
|
||||
@@ -274,6 +309,9 @@ fn process_vcf_files(
|
||||
card_counts,
|
||||
logger,
|
||||
) {
|
||||
if e.aborts_run() {
|
||||
return Err(e);
|
||||
}
|
||||
logger.warn(&format!("vcf {:?}: {e}", file.path));
|
||||
card_counts.failed += 1;
|
||||
}
|
||||
|
||||
@@ -94,6 +94,34 @@ impl Summary {
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Debug)]
|
||||
pub struct RunOutcome {
|
||||
pub summary: Summary,
|
||||
pub error: Option<Error>,
|
||||
}
|
||||
|
||||
impl RunOutcome {
|
||||
pub fn from_result(result: Result<Summary, Error>) -> RunOutcome {
|
||||
match result {
|
||||
Ok(summary) => RunOutcome {
|
||||
summary,
|
||||
error: None,
|
||||
},
|
||||
Err(error) => RunOutcome {
|
||||
summary: Summary::default(),
|
||||
error: Some(error),
|
||||
},
|
||||
}
|
||||
}
|
||||
|
||||
pub fn into_result(self) -> Result<Summary, Error> {
|
||||
match self.error {
|
||||
Some(error) => Err(error),
|
||||
None => Ok(self.summary),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
pub struct Context {
|
||||
pub conn: Connection,
|
||||
pub client: HttpClient,
|
||||
@@ -119,3 +147,28 @@ impl Context {
|
||||
self.common.dry_run
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
#[test]
|
||||
fn an_aborted_run_keeps_its_error_and_its_partial_counts() {
|
||||
let mut summary = Summary::default();
|
||||
summary.per_type.push(("Mailbox", TypeCounts::default()));
|
||||
let outcome = RunOutcome {
|
||||
summary,
|
||||
error: Some(Error::Connection("http status 404".to_owned())),
|
||||
};
|
||||
assert_eq!(outcome.summary.per_type.len(), 1);
|
||||
let err = outcome.into_result().expect_err("aborted");
|
||||
assert_eq!(err.exit_code(), 2);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn a_result_without_a_native_outcome_reports_nothing_on_abort() {
|
||||
let outcome = RunOutcome::from_result(Err(Error::Partial("one object".to_owned())));
|
||||
assert!(outcome.summary.per_type.is_empty());
|
||||
assert!(outcome.error.is_some());
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user