v1.0.0
This commit is contained in:
+288
-30
@@ -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<Value>, usize), JmapError>;
|
||||
type GetMsg = Result<Vec<Value>, JmapError>;
|
||||
type BlobRef = (String, String, String);
|
||||
type DlMsg = (String, Result<Vec<u8>, JmapError>);
|
||||
|
||||
@@ -271,30 +271,31 @@ fn reconcile_type(
|
||||
};
|
||||
let local_ids: HashSet<String> = local_map.keys().cloned().collect();
|
||||
|
||||
let (server_ids, preloaded): (Vec<JmapId>, Option<Vec<Value>>) = if is_queryless(ty) {
|
||||
let got = get_all::<Value>(&net.client, &net.api, &net.account, ty.jmap_name())
|
||||
let (server_ids, preloaded, enum_state): (Vec<JmapId>, Option<Vec<Value>>, Option<String>) =
|
||||
if is_queryless(ty) {
|
||||
let got = get_all::<Value>(&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<Value> {
|
||||
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<Vec<JmapId>, GetMsg> = Pool::new(workers, {
|
||||
let net = net.clone();
|
||||
move |ids: Vec<JmapId>| {
|
||||
let n = ids.len();
|
||||
get_objects::<Value>(
|
||||
&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<String, i64>,
|
||||
cursor: Option<String>,
|
||||
run_state: Option<String>,
|
||||
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<JmapId> = 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<String, i64>,
|
||||
threads: usize,
|
||||
logger: &Logger,
|
||||
counts: &mut TypeCounts,
|
||||
) -> Result<bool, Error> {
|
||||
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<Value>,
|
||||
local_map: &HashMap<String, i64>,
|
||||
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<String, Vec<u8>>,
|
||||
) -> Result<bool, JmapError> {
|
||||
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,
|
||||
|
||||
@@ -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<bool, JmapError> {
|
||||
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<bool, JmapError> {
|
||||
let mailbox_ids: IndexMap<JmapId, bool> = 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<String, bool> = match obj.get("keywords") {
|
||||
Some(v) => serde_json::from_value(v.clone())?,
|
||||
None => IndexMap::new(),
|
||||
};
|
||||
let kw: Vec<String> = 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<bool, JmapError> {
|
||||
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<bool, JmapError> {
|
||||
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<bool, JmapError> {
|
||||
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<bool, JmapError> {
|
||||
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<bool, JmapError> {
|
||||
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<i64>,
|
||||
resolver: &impl LocalResolver,
|
||||
) -> Result<bool, JmapError> {
|
||||
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<bool, JmapError> {
|
||||
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<bool, JmapError> {
|
||||
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<String>, 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":"[email protected]"}))
|
||||
.unwrap();
|
||||
let local = insert_identity(&c, &id).unwrap();
|
||||
let changed: Identity = serde_json::from_value(
|
||||
serde_json::json!({"id":"I1","name":"New Name","email":"[email protected]","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:[email protected]"}),
|
||||
)
|
||||
.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:[email protected]"}),
|
||||
)
|
||||
.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<String>, Option<i64>) = 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"));
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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];
|
||||
|
||||
@@ -76,7 +76,6 @@ pub fn list_folder(folder_path: &Path) -> std::io::Result<DiskListing> {
|
||||
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,
|
||||
|
||||
@@ -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();
|
||||
|
||||
@@ -420,7 +420,6 @@ fn account_id_for(auth: &ManageSieveAuth) -> String {
|
||||
|
||||
#[derive(Debug)]
|
||||
enum SieveAuthError {
|
||||
|
||||
TerminallyRefused(String),
|
||||
|
||||
NoUsableMechanism(String),
|
||||
|
||||
@@ -48,7 +48,6 @@ pub enum EmailKey {
|
||||
|
||||
#[derive(Debug, Clone, PartialEq, Eq)]
|
||||
pub struct EmailIndex {
|
||||
|
||||
pub mids: Vec<String>,
|
||||
|
||||
pub fb: [u8; 32],
|
||||
|
||||
Reference in New Issue
Block a user